Commit Graph

8 Commits

Author SHA1 Message Date
Nicolò Boschi b1de1b9418 fix(operations): allow cancelling in-flight operations (#4131)
`DELETE /v1/default/banks/{bank}/operations/{id}` only accepted `pending`
operations, so an operation stranded in `processing` — orphaned when a worker
was killed before it could write a terminal status — could only be cleared by
hand-editing `async_operations` and restarting the container.

Cancel now accepts `processing` too. It stays cooperative and is never
immediate: the row is flipped to `cancelled` and the worker running it stops at
its next `_check_op_alive` checkpoint (between retain sub-batches/documents,
between consolidation LLM batches). For the orphaned case nothing is running, so
the flip is the whole fix. No heartbeat and no per-batch bookkeeping is added.

Making the flip stick required guarding the worker writes that had none, and
would otherwise overwrite it:

- `_schedule_retry` — the one that actually resurrected cancelled work: a task
  failing after cancellation went back to `pending` and was re-claimed.
- `_mark_failed` (poller and engine), `_defer_operation`.

`_mark_completed` already guarded on `status='processing'`; tests now pin it.

The sibling rollup counted only `completed`/`failed` as done, so a cancelled
child stranded its `batch_retain` parent in `processing` forever — the same
wedge one level up. Both rollup copies now treat `cancelled` as done and settle
the parent on `cancelled` (a real failure still outranks it), cancel performs
the rollup itself so cancelling the last outstanding child terminalizes the
parent, and a cancelled parent is never flipped back by a child finishing later.

Control plane: the Cancel button was gated on `pending`, hiding the fix from the
UI an operator would reach for. It now shows for `processing` rows too.
2026-09-07 11:21:06 +02:00
Nicolò Boschi b75e941916 fix(worker): rotate slots across banks so bulk ingest stops starving them (#3861) (#3980)
* fix(worker): rotate slots across banks so bulk ingest stops starving them (#3861)

Claiming was a strict global FIFO on created_at, so a bank under sustained bulk
ingest held every worker slot for as long as its queue lasted. Measured on a
six-bank instance: one bulk bank owned 17 of 17 retain slots in 90.6% of daily
samples, and a write to a bank with an empty queue timed out at 300s behind the
backlog. The queue drains correctly once ingest stops — this is fairness, not
correctness.

Deficit round robin, with a quantum of one slot. Every operation costs exactly
one slot, so DRR's deficit counter is always zero and drops out; what is left is
the round-robin walk. The poller keeps a cursor over the bank id space per
schema — one level below the tenant rotation it already had — and the claim
takes one row for the first bank sorting after it.

Both tiers are one statement: a `rot` CTE (bounded index seek past the cursor)
unioned with the `fifo` claim that was always there, joined back and locked with
FOR UPDATE OF ... SKIP LOCKED. Same query count as before. 0.10ms against 0.02ms
for the bare FIFO claim, at 50k pending rows.

The cursor is a *range*, not a set of known banks: the starved bank is by
definition one this worker has never claimed for, so only a range can discover
it. And `fifo` is what makes the rotation safe at both ends — it is the wrap
(once the cursor passes the last bank, `rot` matches nothing and `fifo` claims
the whole pool, so the end of a round costs neither a query nor an empty claim,
and a single-bank deployment sits in that state permanently), and it is what
keeps this work-conserving: a bank alone with work still takes every slot, so
nothing is throttled and no slot is held open for an idle bank.

Rejected along the way, both measured: an ordering predicate that re-scanned the
bank's queue per candidate row (11s at 10k pending), and a separate seek before
the claim (correct, but a round trip on every claim, for every schema, on every
poll). No new index — `idx_async_operations_bank_status` and
`idx_async_operations_bank_created_desc` already serve both branches.

Oracle keeps the plain FIFO claim, deliberately: its ROWNUM rewrite for
FOR UPDATE + LIMIT applies before ORDER BY, so a rotation there would land on
whichever bank is scanned first — under bulk ingest, the one it exists to rotate
away from.

claim_tasks now returns ClaimedOperations (rows + the rotation's next cursor)
rather than a bare row list, so learning where the rotation got to costs no
second statement; callers updated.

Claude-Session: https://claude.ai/code/session_017ufCz6qrNxn36Stug7ek8A

* chore(docs): regenerate the docs skill for the worker-tuning paragraph

hindsight-docs/docs/developer/api/operations.mdx is the source; the skill
reference is generated from it by scripts/generate-docs-skill.sh, and
verify-generated-files fails on the drift.

Claude-Session: https://claude.ai/code/session_017ufCz6qrNxn36Stug7ek8A
2026-09-01 12:08:09 +02:00
Nicolò Boschi 681a79e6b4 fix(graph-maintenance): drive the entity prune off a queue, not a bank-wide sweep (#3222) (#3409)
* fix(graph-maintenance): drive the entity prune off a queue, not a bank-wide sweep (#3222)

The graph_maintenance job's Pass 2/3 were two bank-wide statements re-evaluated
on every invocation, whether or not anything had changed: the orphan-entity
prune probed once per entity in the bank, and the stale-cooccurrence prune
evaluated an INTERSECT per cooccurrence row in the bank. Their cost tracked the
size of the bank rather than the size of the delete, so past a few million rows
neither could finish inside asyncpg's 60s command timeout. The job then failed
on every run with a bare TimeoutError, forever, on exactly the banks that most
needed it — and, because db_utils treats a timeout as transient, re-ran the
doomed statement nine times per attempt, holding a worker slot for ~10 minutes
each time.

Both prunes are now driven by `entity_maintenance_queue`, filled inside the
deleting transaction the way `graph_maintenance_queue` already is for the relink
pass. A run claims a bounded batch of candidate entities, prunes what is
genuinely dead, and commits — so the cost is O(delta), and the work already done
survives whatever stops the run.

Measured on a dense fixture (100k entities, 1.5M unit_entities, 2.86M
cooccurrences, hub entities holding 150-400 postings):

  bank-wide orphan prune          9.6s      → batch of 50:   15ms
  bank-wide cooccurrence prune    >11 min   → batch of 50:   2.0s
                                  (cancelled; ~1.7ms per pair over 2.86M pairs)

Also:

* A wall-clock budget for the whole job. Both passes commit per batch, so
  exhausting it is not a failure — the run reports `queues_drained: false`,
  logs it, and chains a follow-up (under a real queue; a synchronous backend
  would recurse instead of schedule). Large backlogs converge over runs.
* The scoping predicate is a UNION of the two endpoint columns, not
  `entity_id_1 = ANY(...) OR entity_id_2 = ANY(...)` — that OR is the #3387
  shape and cannot be driven from either index.
* Every site that removes units or replaces entity postings now enqueues
  candidates: document delete, single and bulk memory delete, curation
  edit/invalidate, document re-ingest, and the delta-retain chunk cascade.
* The migration seeds the queue with every existing entity, so garbage a bank
  accumulated while its sweep was failing is still reclaimed — incrementally,
  a bounded batch per run, instead of in one statement that cannot finish.

* fix(graph-maintenance): compose the queue-scoped prune with the set-based staleness check

Rebase reconciliation with #3408, which landed the same statement while this
was in review.

staleness against a set of live pairs built once, instead of a correlated
INTERSECT re-scanned per row — removing the (rows judged) x (hub degree)
product. That is the better predicate, and it composes with the queue scoping
rather than competing with it: `live` is now seeded from the *claimed
candidates'* units instead of the whole bank's. Correctness holds because every
pair being judged has a candidate as an endpoint, so any unit still grounding
one of those pairs references a candidate and is in the seeded set.

Measured on the dense fixture (100k entities, 1.5M unit_entities, 2.86M
cooccurrences):

  bank-wide, set-based (#3408 as merged)   did not finish in 10 min
  batch of 50, per-pair INTERSECT (mine)   2.0 s
  batch of 50, composed                    16-65 ms

The batch-size rationale is updated to the new numbers; 50 still holds, now
with three orders of magnitude of margin instead of one. #3367's hub/bank-
scoping regression test is kept, adapted to seed candidates.

Also re-chains the migration onto d9c1a7b4e2f6, which took the same parent on
main and would otherwise leave two alembic heads.

* fix(graph-maintenance): restore the review fixes without the birth-time enqueue

Drops the "queue every entity at creation" change and keeps the rest of the
review round (dataclass pass results, the Oracle IN-list chunking on the by-unit
enqueue, the per-site enqueue tests, the migration-seed test, the budget's
follow-up-chain test, the stale-comment sweep).

The birth-time enqueue existed to reclaim an entity created in retain's Phase 1
whose Phase-2 link never landed. It is not worth what it costs: such a row is a
single entry in the registry with no postings and no cooccurrences, and #2662
exists because the retry is *supposed* to adopt it — Phase 2 reasserts resolved
parents under FOR KEY SHARE precisely so a pruner cannot delete one out from
under it. Pointing the pruner at every freshly created entity leans on that race
for a leak that is one row wide. #3408 landing the set-based predicate is what
made the trade obviously bad: the expensive half of this job was never those
rows.

Entities created but never linked are therefore no longer proactively reclaimed.
The migration's one-time seed still clears the population a bank has already
accumulated.

* fix(graph-maintenance): don't backfill the entity queue on upgrade

The migration seeded one queue row per existing entity so a bank could reclaim
what it stranded while its bank-wide sweep was failing. That is the wrong trade:
the INSERT runs inside a migration at API startup, so a large deployment pays a
slow upgrade writing a row per entity, and then a prune check for every one of
them — a self-inflicted backlog to collect rows that cost the bank nothing.

The queue now starts empty and fills from real deletes. Historical strays stay
until something touches them; they are single registry rows with no postings and
no cooccurrences.

The migration test pins the two properties that are easy to lose later: the
upgrade enqueues nothing, and the composite key collapses overlapping deletes
into one row (which is also what the #3034 locking upsert conflicts on).
2026-08-12 16:14:19 +02:00
Nicolò Boschi d61e8ffff6 fix(worker): discover expired-operation schemas in one query, default retention off (#2819)
Follow-up to #2708, which bounded terminal `async_operations` history. Two
issues with what landed:

1. The cleanup worker did not use a cross-tenant routine. It opened a connection
   and a prune transaction against *every* tenant schema on every cleanup cycle,
   paying the full per-tenant cost even when nothing was prunable — the query
   storm the server-side maintenance routines exist to avoid.
2. It shipped as a breaking change, silently switching deployments from
   unbounded operation history to a 30-day TTL on upgrade.

Adds `schemas_with_expired_operations(p_days int) RETURNS SETOF text` — the
`async_operations` counterpart to `schemas_with_expired_rows`. One round-trip
returns just the schemas holding expired terminal rows; the worker then acquires
a connection and prunes only there. It needs its own routine rather than reusing
`schemas_with_expired_rows` because eligibility here isn't "row older than N
days" — pending and processing rows are never prunable, so the status filter has
to be part of the predicate.

Install policy follows b6d2f8a4c1e7 (#2638/#2824): the routine is database-global
(it enumerates pg_class across every schema), so exactly one copy is installed —
into the schema this deployment is configured to use, which is the one the worker
calls via fq_routine. Gating on the literal "public" instead of the configured
schema is what left non-public deployments without the sibling routines (#2638);
installing into every schema would leave a dead duplicate per tenant.

Exactly one run satisfies that predicate, so concurrent per-schema runs never
issue competing CREATE OR REPLACE against the same pg_proc row and cannot hit
`tuple concurrently updated`. No cross-process coordination, and in particular no
advisory lock, which is unusable behind connection poolers and managed PG
(#2817). Runs targeting any other schema drop the routine there instead.

The worker calls the routine through schema.fq_routine() (added in #2824) rather
than a hardcoded public. qualifier — duplicating that qualifier across callers is
precisely how #2638 recurs.

Vanishing schemas are skipped rather than fatal (c7e9f1a3b5d2), and an absent
routine degrades cost, not correctness — Oracle and un-migrated PostgreSQL fall
back to the previous full sweep with a warning.

DEFAULT_OPERATION_RETENTION_DAYS 30 -> 0. Operation history is a user-visible
audit trail, so bounding it is an opt-in policy decision rather than something an
upgrade applies silently. Set HINDSIGHT_API_OPERATION_RETENTION_DAYS to a
positive number of days to enable pruning. Docs, .env.example and the bundled
embed template updated to match.

- test_schemas_with_expired_operations — drives the real routine against pg0 in a
  throwaway schema: old pending/processing rows alone don't make a schema
  eligible, a terminal row does, a too-old cutoff doesn't, p_days <= 0 is empty.
- test_expired_operations_routine_installs_in_the_configured_schema —
  parametrized over base / default public / non-public single-tenant; guards
  against reintroducing the #2638 literal gate or an advisory lock.
- test_expired_operations_tenant_runs_install_nothing — tenant runs emit no
  CREATE and drop any copy in their own schema.
- test_discovery_targets_the_configured_non_public_schema — the worker calls the
  copy in its configured schema, not a hardcoded public one.
- TestWorkerOperationCleanupSchemaNarrowing — only reported schemas are pruned,
  nothing expired means no pruning, unclaimed schemas are skipped, a missing
  routine falls back to the full sweep, Oracle never calls the routine.
2026-07-20 18:35:54 +02:00
Voscko 86ff344c93 fix(worker): bound terminal operation history (#2708)
Add configurable TTL (default 30 days, 0=keep-forever) for terminal async_operations rows. Expired completed/failed/cancelled rows are pruned in bounded batches (1000/cycle) by a background task that never touches pending/processing work. Batch children are protected until their parent is pruned; cancelled-child cleanup atomically cancels a pending parent first. PG uses FOR UPDATE SKIP LOCKED; both PG and Oracle re-check eligibility under the row lock before deleting. Includes indexes, docs, and regenerated SDKs.

184 retention/worker/operation-status tests pass locally. All CI green.

Fixes #2705
2026-07-14 10:04:26 -04:00
Evo c1089698b5 docs(api): document the progress snapshot + include_payload on the operation status endpoint (#2037)
PR #2013 added a durable progress snapshot (OperationProgress: stage/at/
processed/total/detail) plus an updated_at heartbeat and an include_payload
query param yielding task_payload to GET .../operations/{operation_id}, but
the 'Get operation status' docs had no response-field prose for any of them
(the example even passes include_payload without explaining it). Added a
response-fields subsection sourced from http.py. Regenerated the skills mirror.
2026-06-08 10:35:52 +02:00
Evo 84330d0453 docs(api): repoint Worker Configuration link to #distributed-workers (#1869) 2026-06-01 09:51:13 +02:00
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