mirror of
https://github.com/vectorize-io/hindsight.git
synced 2026-09-14 19:31:49 +08:00
perf/profiling-env
8 Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
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.
|
||
|
|
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 |
||
|
|
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). |
||
|
|
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. |
||
|
|
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 |
||
|
|
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. |
||
|
|
84330d0453 | docs(api): repoint Worker Configuration link to #distributed-workers (#1869) | ||
|
|
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.
|