Files
Nicolò Boschi cc3ba4a37c feat(api): async link recompute to fix outgoing-link staleness after deletes (#1772)
* feat(api): async link recompute to fix outgoing-link staleness after deletes

When a memory_unit is deleted (via delete_document, delete_memory_unit, or
document re-ingest via handle_document_tracking), the FK cascade removes its
incoming temporal/semantic links. Other units that had this unit in their
top-K neighbours therefore lose links and stay permanently under-capped —
retain only generates links for newly-inserted units, never re-evaluates
surviving ones.

This adds a reactive top-up:

* Inside the delete transaction, capture from_unit_ids that pointed at the
  doomed units and write them to a new link_recompute_queue table (PG: ON
  CONFLICT DO NOTHING, Oracle: IGNORE_ROW_ON_DUPKEY_INDEX hint for dedup).
* After commit, submit_async_link_recompute schedules a new task type
  ("link_recompute"), deduplicating per bank.
* Worker drains the queue in batches of 50; for each victim it counts
  current outgoing temporal/semantic links and, if below cap, runs the
  same probes used at retain time (fetch_temporal_neighbours,
  compute_semantic_links_ann) to find replacements. bulk_insert_links has
  ON CONFLICT DO NOTHING, so re-probing freely is safe.

submit_async_link_recompute is also called after every retain, where it
short-circuits with no_work=True when the queue is empty — that lets the
upsert path (handle_document_tracking) enqueue victims inline without
needing a return-value plumbing change.

Worker slot is opt-in (default 0) via HINDSIGHT_API_WORKER_LINK_RECOMPUTE_MAX_SLOTS.

Tests cover enqueue correctness (cross-doc, self-exclude, entity-link
skip, dedup), worker behaviour (empty drain, missing-victim no-op,
top-up to cap, no-op at cap), and a cap-parity guard against retain-side
constants drifting.

* docs: revamp /developer/api/operations with all 6 operation types

The page previously listed only batch_retain + consolidate. Rewritten to
cover every async task type Hindsight runs: retain, file_convert_retain,
consolidation, refresh_mental_model, link_recompute (new), and
webhook_delivery — with triggers, lifecycle states, bank-dedup notes, and
the full list/status/cancel/retry endpoint surface.

Also adds HINDSIGHT_API_WORKER_LINK_RECOMPUTE_MAX_SLOTS to the worker
configuration table.

* refactor(api): rename link_recompute → graph_maintenance + kind discriminator

Generalize the queue and worker so future post-mutation cleanups (orphan
entity pruning, stale cooccurrence removal, etc.) can ride on the same
async surface without spawning their own task types.

Schema (alembic b5a4c3e2f1d8): table renamed to graph_maintenance_queue
with shape (bank_id, kind, target_id, enqueued_at) and PK on
(bank_id, kind, target_id). Today the only kind is 'relink_unit', which
holds the same payload as the previous link_recompute_queue.

Renames (mechanical):
* task_type and operation_type: link_recompute → graph_maintenance
* env var: HINDSIGHT_API_WORKER_LINK_RECOMPUTE_MAX_SLOTS →
           HINDSIGHT_API_WORKER_GRAPH_MAINTENANCE_MAX_SLOTS
* module hindsight_api/engine/link_recompute.py →
         hindsight_api/engine/graph_maintenance.py
* engine helpers: enqueue_link_recompute_victims → enqueue_relink_victims;
                  run_link_recompute_job → run_graph_maintenance_job;
                  submit_async_link_recompute → submit_async_graph_maintenance;
                  _handle_link_recompute → _handle_graph_maintenance
* ops methods: enqueue_link_recompute_victims → enqueue_graph_maintenance
               (now takes kind + target_ids);
               claim_link_recompute_batch → claim_graph_maintenance_batch
               (now returns (kind, target_id) tuples)
* worker job result keys: victims_processed → targets_processed,
                          links_added → relink_links_added

Worker now groups each claimed batch by kind and dispatches to a per-kind
handler; unknown kinds are dequeued and logged without crashing (added
test_skips_unknown_kind_without_failing). The 'relink_unit' handler is
the same code that previously lived inline in run_link_recompute_job.

Docs updated: operations.md reframes the section around graph_maintenance
as a framework with kinds, with relink_unit documented as the first one;
configuration.md gets the new env var name.

Revision ID bumped from d8f1e2c3a4b5 to b5a4c3e2f1d8 since the table
schema changed shape — dev/staging DBs that already applied the previous
revision get a fresh migration instead of a silent no-op.

* docs(operations): rework per review — trim, link out, multi-language tabs

- Drop the unsupported Kafka note and the type-summary table; the
  per-section headings carry the same info without duplication.
- Add a parent-op section for retain_batch explaining how Hindsight splits
  large submissions into a parent + N children and how exclude_parents
  hides the parent rows.
- file_convert_retain: point at Configuration → File Processing for which
  converter runs (markitdown / Docling / LlamaParse).
- consolidation: shorten to a one-liner pointing at the Observations page
  instead of restating it.
- refresh_mental_model: mention the auto-refresh trigger and drop the
  LLM-provider gate caveat (the model-level check covers it).
- graph_maintenance: shorter why/what framing without the algorithm walk,
  drop the PG/Oracle asymmetry note (matches retain-time semantic behaviour
  and isn't operations-doc material).
- Convert curl examples to <Tabs>/<CodeSnippet> with Python, Node.js, CLI,
  and Go variants, matching the pattern used by recall/retain/documents.
  Added examples/api/operations.{py,mjs,sh,go} with sections wired into
  the Tabs blocks.

Page renamed .md → .mdx so the Tabs/CodeSnippet imports work.

* docs(operations): correct file-parser list

Hindsight ships three parsers: markitdown (default), iris (Vectorize Iris
cloud), and llama_parse. Docling was never wired up — drop it from the
file_convert_retain note and name the actual options + the
HINDSIGHT_API_FILE_PARSER env var that selects between them.

* refactor(api): drop kind discriminator; add entity + cooccurrence prune passes

graph_maintenance is one job now, not a dispatcher of subtypes. Every
invocation runs three passes:

1. Link top-up — drains graph_maintenance_queue (the only queued work) and
   tops up each victim unit's outgoing temporal/semantic links via the same
   probes retain uses.
2. Orphan entity prune (NEW) — deletes entities in the bank that no longer
   have any unit_entities references. FK ON DELETE CASCADE on
   entity_cooccurrences cleans up cooccurrences pointing at pruned entities
   automatically.
3. Stale cooccurrence prune (NEW) — defensive sweep for cooccurrence rows
   where both endpoints still exist but no current memory_unit references
   both of them (the cooccurrence was real when recorded, but every unit
   witnessing it has since been deleted).

Schema change: graph_maintenance_queue loses the kind column. It's now just
(bank_id, unit_id, enqueued_at) with PK (bank_id, unit_id). Renamed
target_id → unit_id to make intent obvious. The bank-wide sweeps in passes
2 and 3 don't need per-target queueing — they're backed by entities(bank_id)
and unit_entities(entity_id) indexes.

Ops surface: enqueue_graph_maintenance / claim_graph_maintenance_batch lose
the kind parameter and return unit-id-only payloads. Added
prune_orphan_entities and prune_stale_cooccurrences as ops methods with PG
and Oracle implementations.

Triggers: delete_document and delete_memory_unit now submit
graph_maintenance whenever any unit is removed (not gated on whether relink
victims were enqueued), so the entity/cooccurrence sweeps fire even when a
deleted unit had no incoming links.

Test surface: dropped the unknown-kind test and the cross-kind enqueue
test. Added TestOrphanEntityPrune (scoped sweep, doesn't cross banks) and
TestStaleCooccurrencePrune (prunes when no shared unit, keeps when shared).
All 14 tests in tests/test_graph_maintenance.py pass.

Docs: operations.mdx graph_maintenance section drops the kinds framing and
describes the three passes directly.

* docs(ops_oracle): correct misleading rowcount comment

The Oracle DatabaseConnection wrapper reshapes cursor.rowcount into a
PG-compatible "DELETE N" status string before returning, so the shared
parsing in prune_orphan_entities works on both dialects. The previous
comment claimed the opposite.

* fix(ci): test/example bugs surfaced by CI run

* test_graph_maintenance: _insert_cooccurrence now sorts the two entity
  IDs before insert. entity_cooccurrences has a CHECK constraint
  entity_id_1 < entity_id_2 (canonical ordering to dedupe (A,B) vs (B,A))
  which my helper ignored. asyncpg surfaced this as a CheckViolationError
  in test_keeps_cooccurrence_with_shared_unit.

* examples/api/operations.py: collapsed two top-level asyncio.run() calls
  into a single asyncio.run(main()). Multiple event loops on the same
  Hindsight client broke the SDK's async HTTP context ("Timeout context
  manager should be used inside a task"). The doc snippets also use a
  real operation_id pulled from list_operations rather than a hardcoded
  one that doesn't exist.

* examples/api/operations.sh: was using a hardcoded UUID, so cancel/retry
  returned 404 against the live API. Now creates a real pending op via
  --async retain, exercises get/cancel on it, then creates a second op
  and cancels it so retry has something to re-queue.

* operations.mdx: added the CLI tab to the async-retain Tabs block —
  code-parity check requires all four language tabs and was rejecting
  the build.

* fix(ci): cooccurrence assertions + python example loop reuse

* tests/test_graph_maintenance.py: both stale-cooccurrence assertions
  now query (entity_id_1, entity_id_2) with the same canonical sort the
  insert helper applies. The test_keeps_cooccurrence_with_shared_unit
  failure ("None == 5") was caused by inserting (sorted_a, sorted_b)
  but reading (ent_a, ent_b) — the SELECT just missed the row.

* examples/api/operations.py: dropped the sync client.retain() seed call
  in favour of aretain_batch inside the async main(). Mixing sync
  (client.retain → _run_async → its own event loop) with the async
  operations API (asyncio.run(main) → fresh loop) left the underlying
  HTTP client bound to a dead loop, surfacing as
  "Timeout context manager should be used inside a task".

* skills/hindsight-docs/references/developer/api/operations.md: regenerated
  to match the .mdx — verify-generated-files caught the drift from the
  previous CLI-tab edit.
2026-05-27 14:41:05 +02:00

69 lines
2.0 KiB
Bash

#!/bin/bash
# Operations API examples for Hindsight CLI
# Run: bash examples/api/operations.sh
set -e
BANK_ID="my-bank"
# =============================================================================
# Setup (not shown in docs)
# =============================================================================
# Submit an async retain so we have a real pending operation_id to exercise
# get / cancel below.
OPERATION_ID=$(
hindsight memory retain "$BANK_ID" "setup-1 for ops examples" --async -o json \
| jq -r '.operation_id'
)
# =============================================================================
# Doc Examples
# =============================================================================
# [docs:operations-list]
hindsight operation list my-bank
# [/docs:operations-list]
# [docs:operations-get]
hindsight operation get my-bank "$OPERATION_ID"
# [/docs:operations-get]
# [docs:operations-cancel]
hindsight operation cancel my-bank "$OPERATION_ID"
# [/docs:operations-cancel]
# Setup for retry: create another op and cancel it so it's in a retryable state.
OPERATION_ID=$(
hindsight memory retain "$BANK_ID" "setup-2 for retry example" --async -o json \
| jq -r '.operation_id'
)
hindsight operation cancel "$BANK_ID" "$OPERATION_ID" >/dev/null
# [docs:operations-retry]
hindsight operation retry my-bank "$OPERATION_ID"
# [/docs:operations-retry]
# [docs:operations-async-retain]
# Submit an async retain and capture the operation_id from the JSON response.
OPERATION_ID=$(
hindsight memory retain my-bank "Alice joined Google in 2023" --async -o json \
| jq -r '.operation_id'
)
# Poll until the worker finishes — completed/failed/cancelled are all terminal.
while true; do
STATUS=$(hindsight operation get my-bank "$OPERATION_ID" -o json | jq -r '.status')
if [ "$STATUS" = "completed" ] || [ "$STATUS" = "failed" ] || [ "$STATUS" = "cancelled" ]; then
echo "finished: $STATUS"
break
fi
sleep 2
done
# [/docs:operations-async-retain]