diff --git a/.beads/issues.jsonl b/.beads/issues.jsonl index f34564477d..f425706af7 100644 --- a/.beads/issues.jsonl +++ b/.beads/issues.jsonl @@ -58,6 +58,8 @@ {"_type":"issue","id":"polylogue-tf2.1","title":"Rerun forensics on current archive; price origin_reported providers","description":"Rerun scripts/agent_forensics.py against the current archive (v23+); price origin_reported providers via the vendored LiteLLM catalog (match last path segment); all-provider headline or explicitly-labeled per-provenance figures that cannot be misread; record deltas vs 06-27; verify chart SVGs render. Cache-inclusion must be disambiguated (Codex input INCLUDES cached ~96%; see bd memories). Also blocked on logical-session token attribution — the headline must not be double-counted.","notes":"Correction to close_reason monetary values: stored/provider-priced subset was $239,453.14; catalog API-equivalent was $318,650.88; origin_reported catalog estimate was $79,197.74. The original close_reason text lost dollar-prefixed digits due shell expansion, not measurement drift.","status":"closed","priority":0,"issue_type":"task","assignee":"Sinity","owner":"ezo.dev@gmail.com","created_at":"2026-07-03T04:31:33Z","created_by":"Sinity","updated_at":"2026-07-03T09:59:13Z","started_at":"2026-07-03T09:28:10Z","closed_at":"2026-07-03T09:59:02Z","close_reason":"Completed with blocker caveat preserved: scripts/agent_forensics.py now prices origin_reported rows through the shared vendored LiteLLM pricing catalog while preserving stored provenance; report separates stored/provider-priced cost from catalog API-equivalent estimates and carries logical-session/cache caveats instead of claiming final billing reconciliation. Regenerated current artifact at .agent/demos/agent-forensics against /home/sinity/.local/share/polylogue schema v23: 16,498 physical sessions, 4,142,175 messages, 356.5B tokens, ,453.14 stored/provider-priced subset, ,650.88 catalog API-equivalent, and ,197.74 origin_reported catalog estimate. SVG parse check passed for 9 charts; devtools test tests/unit/scripts/test_agent_forensics.py passed; devtools verify --quick passed run 20260703T095718Z-quick-753466-96559776; devloop-review clean. Remaining final-reconciliation blocker stays open as polylogue-4ts.2.","labels":["area:usage","campaign"],"dependencies":[{"issue_id":"polylogue-tf2.1","depends_on_id":"polylogue-4ts.2","type":"blocks","created_at":"2026-07-03T06:32:45Z","created_by":"Sinity","metadata":"{}"},{"issue_id":"polylogue-tf2.1","depends_on_id":"polylogue-sru.7","type":"blocks","created_at":"2026-07-03T06:31:33Z","created_by":"Sinity","metadata":"{}"},{"issue_id":"polylogue-tf2.1","depends_on_id":"polylogue-tf2","type":"parent-child","created_at":"2026-07-03T06:31:33Z","created_by":"Sinity","metadata":"{}"}],"dependency_count":2,"dependent_count":2,"comment_count":0} {"_type":"issue","id":"polylogue-tf2","title":"Campaign: agent-forensics regeneration + all-provider repricing","description":"Regenerate the agent-forensics packet on the current archive with an honest all-provider headline. The 2026-06-27 report (546.6B tokens, $89,368 API-list equivalent, 216x cache amplification) is the most stranger-legible artifact on any shelf, but its numbers are pre-dedup stale and the headline prices only the priced-provenance subset (Claude Code cost_usd rows); Codex/ChatGPT/Gemini are origin_reported token counts with no dollar value (operator estimate ~$150K all-provider). Sequenced after claim-vs-evidence per operator direction 2026-07-02.","design":"Current slice design: turn the existing agent-forensics/cost headline into a product-backed all-provider repricing artifact. First inspect devtools/scripts and polylogue analyze surfaces for agent_forensics/cost code. Use active archive usage headline (detail=headline) for authoritative physical_session and logical_session_model_high_water token totals. Keep priced-provenance dollars and origin-reported token estimates separate: do not multiply every token by one blended price without a labeled lane. Add or reuse a shared pricing/projection helper so the demo artifact is regenerated from Polylogue product code, not ad hoc SQL. Acceptance for this slice: the generated agent-forensics artifact names archive root/schema, includes physical vs logical token grain, separates priced subset from origin-reported estimate lanes, gives reproduction commands, and has focused tests for any new repricing helper/surface.","acceptance_criteria":"Terminal state: regenerated forensics packet on the current archive with an honest all-provider headline (priced subset AND origin-reported estimate lanes separated), agent_forensics.py folded into polylogue analyze (tf2.2), artifact on the demo shelf with reproduction commands, cold-reader gate passed. Epic closes only when that artifact is recorded.","status":"closed","priority":0,"issue_type":"epic","assignee":"Sinity","owner":"ezo.dev@gmail.com","created_at":"2026-07-03T04:31:32Z","created_by":"Sinity","updated_at":"2026-07-03T19:06:44Z","started_at":"2026-07-03T18:47:23Z","closed_at":"2026-07-03T19:06:44Z","close_reason":"Completed: provider usage headline now exposes product-backed pricing lanes in polylogue analyze usage --detail headline, separating stored/provider-priced cost from catalog API-equivalent estimates for origin_reported rows. Regenerated the current .agent/demos/agent-forensics artifact against /home/sinity/.local/share/polylogue schema v23: physical-session tokens 395,320,980,423; logical high-water tokens 288,741,229,728; stored/provider-priced USD 243,392.189328; catalog API-equivalent USD 337,565.031618; priced lane 13,889 rows / 12,331 sessions / 12,650 matched rows; origin_reported lane 2,308 rows / 2,270 sessions / 2,302 matched rows. Verification: live polylogue --plain analyze usage --detail headline --format json --limit 0 wrote /realm/tmp/polylogue-usage-headline-pricing-current.json; devtools test tests/unit/storage/test_provider_usage_report.py tests/unit/cli/test_diagnostics.py passed 23 tests; devtools verify --quick passed run 20260703T190553Z-quick-2226137-d91d4e8f; devtools workspace demo-shelf --json reported ok. Non-claim preserved: this is not final billing reconciliation and physical/logical token grains stay explicitly separated.","labels":["area:usage","campaign","size:M","spine"],"dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"polylogue-sru","title":"Campaign: claim-vs-evidence report to finding-grade","description":"Terminal state: an externally publishable finding ('how often do coding agents proceed past failed tool calls, by model/tool') with stated sample frame, calibrated markers, benign/consequential split, seeded stranger-runnable reproduction, and a passed cold-reader gate. Slice closure is NOT campaign closure; this epic stays top-of-frame until its terminal state is recorded.\\n\\nState as of 2026-07-03 after calibrated active-archive regeneration: archive root /home/sinity/.local/share/polylogue, index schema v23, 41,886 structured failures total, 5,000 origin-stratified failures inspected (3,746 claude-code-session, 1,247 codex-session, 7 claude-ai-export), 100 unpaired structured failures. Marker vocabulary was tightened to avoid broad issue/fix/block/gitignored false positives. Immediate next-turn totals: acknowledged=420, silent_proceed=1,205, ambiguous=3,375 (2,624 wordless tool continuations; 751 prose without marker). Lower-bound silent rate is 24.1%; among classified immediate next turns, silent rate is 74.2%. Next-3 sensitivity window, stopping before the next user message, finds 302 acknowledgments that appear only after the next turn; window3 silent lower bound is 37.0%. Calibration: 50 hand-labeled immediate-next-turn rows, acknowledged-marker precision=1.0, recall=0.8421052631578947, invalid rows=0. Artifact: .agent/demos/claim-vs-evidence/claim-vs-evidence.report.json.","notes":"2026-07-03 update: methodology package is now cold-read gated. .agent/demos/claim-vs-evidence contains aggregate live evidence, public-summary.json, PUBLIC_REPRODUCTION.md, COLD_READER_GATE.md, and COLD_READ_RESULT.md. Seeded reproduction is meaningful, not empty: 4 structured failures, 2 acknowledged follow-ups, 2 silent-proceed follow-ups, 0 unpaired. Cold-reader subagent PASS recovered claim/non-claim, sample frame, rates, calibration, caveats, and reproduction commands from the artifact directory only. Remaining campaign child: polylogue-sru.1 productizes action-unit outcome/followup_class capability.","status":"closed","priority":0,"issue_type":"epic","owner":"ezo.dev@gmail.com","created_at":"2026-07-03T04:31:26Z","created_by":"Sinity","updated_at":"2026-07-03T09:28:09Z","closed_at":"2026-07-03T09:28:09Z","close_reason":"Completed: all seven campaign children are closed. The claim-vs-evidence finding now has bounded sample-frame reporting, calibrated marker precision/recall, handler-class and next-3 sensitivity splits, meaningful seeded reproduction, cold-reader PASS, and productized action-unit followup_class/followup_message_ref query capability. Current artifact lives under .agent/demos/claim-vs-evidence and was regenerated against /home/sinity/.local/share/polylogue schema v23.","labels":["area:substrate","campaign"],"dependency_count":0,"dependent_count":0,"comment_count":0} +{"_type":"issue","id":"polylogue-xb4i","title":"Parse prefetch/cache admission must bound parsed-tree bytes, not raw payload bytes","description":"Two earlyoom kills of the v42 rebuild driver (19.3G and 20.2G peaks, 2026-07-20) on a whale-dense page: the DaemonParseStage admission budget (#3195, RAM/16 clamp) and the prefetch cache both account raw PAYLOAD bytes, but parsed ParsedSession trees inflate ~10x+ payload, so a 2GiB payload admission can resident tens of GB of trees; clamping inflight to 256MiB did not help because the CACHE retains the whole page of parsed trees regardless. Fix: account estimated in-memory tree size in both admission and cache retention, with eviction/spill for whales. Interim mitigation in the live walk: 500-raw pages + MemoryHigh=14G on the unit.","status":"closed","priority":1,"issue_type":"task","owner":"ezo.dev@gmail.com","created_at":"2026-07-20T13:38:46Z","created_by":"Sinity","updated_at":"2026-07-20T14:03:58Z","closed_at":"2026-07-20T14:03:58Z","close_reason":"Shipped in PR #3209 (merged): estimate_parsed_tree_bytes structural estimator (two-term fit calibrated against deep-walk measurements, constants rounded up so misestimation biases to eviction/reparse), adaptive cached-tree budget RAM/8 clamped [256MiB,4GiB] with POLYLOGUE_DAEMON_PARSE_STAGE_MAX_CACHED_TREE_BYTES override, side-ledger tracking with largest-first eviction and whale-never-retained. Root cause of the two 2026-07-20 earlyoom kills (19.3/20.2G peaks). 12 tests, mypy --strict, verify --quick green.","labels":["area:daemon","area:perf","horizon:frontier"],"dependency_count":0,"dependent_count":0,"comment_count":0} +{"_type":"issue","id":"polylogue-623q","title":"Import performance envelope: full-corpus rebuild must land in well under an hour","description":"Operator mandate 2026-07-20: the v42 full-archive rebuild (101K raws) taking multiple hours-to-days is unacceptable — import at this scale, including all remaining derived steps, must complete in well under an hour. Measured state: parallel warm parse is ~30-60 raws/s (fine); the serial engine pass was ~3 raws/s with \u003e50% of it census/spill cache overhead + per-unit fsync commits. Landed levers: #3208 (stage telemetry, no-fsync NVMe spill + decoded RAM layer, commit_batch_size=200 in the rebuild path). Remaining candidate levers, to be driven by per-page stage timings: census receipt cost, index full_replace batching, model_usage_seed, census-parse vs warm-cache dedup, writer-thread pipelining of serialization vs SQLite execution. Exit criterion: a measured full-corpus rebuild receipt under 60 min on this machine, recorded on this bead.","status":"open","priority":1,"issue_type":"epic","owner":"ezo.dev@gmail.com","created_at":"2026-07-20T13:38:45Z","created_by":"Sinity","updated_at":"2026-07-20T13:38:45Z","labels":["area:ingest","area:perf","horizon:frontier"],"dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"polylogue-q88p","title":"Content-address the embeddings tier: vectors keyed by identity-free input hash","description":"Operator ruling 2026-07-20: reindexing must not lose embeddings - vectors are about content, not transient index identity. Current defect: message_embeddings_meta binds vectors to messages.content_hash, and _message_content_hash INCLUDES session_id/position/variant_index (identity-contaminated), so rebuilds/lineage shifts invalidate vectors whose text never changed - hence the 777K-vector rescue (04kl). Fix: key the vector store by embedding_input_hash = H(model, normalized embedder input text) - identity-free, same philosophy as the svfj block evidence hash which deliberately excludes identity. Index side keeps a rebuildable message_id -\u003e input_hash mapping; freshness = input_hash lacks a vector; dedup free (fork-replayed identical messages embed once - real API savings in a lineage-heavy archive). End state: rebuilds CANNOT lose embeddings by construction; the rescue concept is retired (automagic doctrine). ORDERING: design this first, then execute the one-time 04kl rescue directly INTO the content-addressed layout (avoid double migration). Embeddings tier schema bump = derived-tier regime (edit canonical DDL + rebuild plan = the rescue itself).","status":"closed","priority":1,"issue_type":"task","owner":"ezo.dev@gmail.com","created_at":"2026-07-20T00:56:55Z","created_by":"Sinity","updated_at":"2026-07-20T05:06:31Z","closed_at":"2026-07-20T05:06:31Z","close_reason":"PR #3192 merged: embeddings tier v4 content-addressed — vectors keyed by identity-free embedding_input_hash(model, NFC input text); rebuildable message_id→hash refs in embeddings tier (index version untouched); all consumers retargeted both twins; 04kl rescue lands into v4; rebuild-survival/dedup/property tests. Reindexing can no longer lose embeddings by construction. Follow-up debt noted in PR: reconcile-path vector GC deferred.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"polylogue-fbte","title":"Rebuild resume re-walks entire corpus: replay phase never populates its cursor","description":"Observed on operation ab5bad1f (2026-07-20): the rebuild transaction record has last_raw_id=None, processed_raw_count=0 even after committing 31,882 sessions - the replay phase never writes its positional cursor, so every resume re-walks the ENTIRE raw corpus relying on per-raw skip fast-paths (byte-proven supersedence #3146, content-hash match). Measured cost: ~2.25h of pure re-verification walk per resume on the ~50K-raw archive, proportional to corpus size instead of remaining work. Fix directions: (a) populate last_raw_id/processed_raw_count during replay batches (fields already exist in the transaction schema), resume seeks past them; or (b) resume-time cheap skip via indexed committed-membership lookup (raw_id already classified in the generation) instead of parse+hash per raw. Either makes recovery O(remaining). Related: polylogue-6mvg (rebuild throughput program).","notes":"2026-07-20 operator-driven rescope: do NOT build the cursor fix on the CLI resume surface - gd6v deletes that command on proven equivalence, so investment there is throwaway. The O(remaining)-resume property is a REQUIREMENT OF THE REPLACEMENT: gd6v daemon bulk path must record replay-phase progress (or skip via committed-membership lookup) so interruption recovery is proportional to remaining work, verified as part of gd6v equivalence gate. This bead stays as the requirement record; implementation lands in gd6v.","status":"open","priority":1,"issue_type":"bug","owner":"ezo.dev@gmail.com","created_at":"2026-07-20T00:51:28Z","created_by":"Sinity","updated_at":"2026-07-20T00:59:01Z","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"polylogue-o7hx","title":"Hook-spool isolation must derive from archive root, not global XDG","description":"Hazard bitten 4+ times (lrou twice, ajmu lane 2026-07-20 drained 6262 real hook events into a scratch archive): hooks_sidecar_dir() in paths/_roots.py resolves data_home()/hooks from pure XDG, independent of the archive root, so ANY daemon - whatever --root / POLYLOGUE_ARCHIVE_ROOT says - drains the one real global spool. The POLYLOGUE_HOOK_SIDECAR_DIR env override is a manual escape hatch agents must remember (and did not, four times) - exactly the pattern the automagic doctrine purges. Fix: the spool path derives from the RESOLVED archive root (default \u003carchive_root\u003e/hooks), which is byte-identical to the current path for the default production root (archive root IS data_home) - zero migration. All consumers (daemon drain, hook_paste_enrichment, cli init, agent_integration installer-rendered writer scripts) use the same resolved path; a scratch-rooted daemon then reads a scratch spool by construction. Fold the env override away per no-compat doctrine (config key if genuine configurability is wanted). Writer scripts get the concrete path baked at install time by the installer.","status":"closed","priority":1,"issue_type":"bug","owner":"ezo.dev@gmail.com","created_at":"2026-07-19T23:37:08Z","created_by":"Sinity","updated_at":"2026-07-20T00:54:45Z","closed_at":"2026-07-20T00:54:45Z","close_reason":"PR #3187 merged: spool path now derives from resolved archive root (byte-identical for prod root, scratch daemons isolated by construction); second instance of the bug fixed in config.py captured default; env override deleted; installer bakes --sidecar-dir into writer scripts; stale hazard docs removed. Deploy step recorded on dcz5 (bake paths when pin bumps). Closes the 4x-bitten scratch-daemon spool drain hazard.","dependency_count":0,"dependent_count":0,"comment_count":0} diff --git a/polylogue/maintenance/replay.py b/polylogue/maintenance/replay.py index 3363dc1765..b004c9fb4c 100644 --- a/polylogue/maintenance/replay.py +++ b/polylogue/maintenance/replay.py @@ -99,6 +99,14 @@ #: never becomes a second archive-wide materialization. _REBUILD_CENSUS_SPILL_CACHE_BYTES: Final[int] = 512 * 1024 * 1024 +#: Rebuild-scale commit batching (polylogue-amg1/oikv machinery; the rebuild +#: caller previously never opted in, paying one fsync'd commit per censused +#: raw and per replayed cohort -- thousands per page). A crash discards at +#: most one open batch and the resume reprocesses it from scratch (contract +#: pinned by test_backfill_resumes_after_replay_batch_crash_discards_whole_ +#: batch_cleanly et al.), so the loss window is bounded and cheap. +_REBUILD_COMMIT_BATCH_UNITS: Final[int] = 200 + # --------------------------------------------------------------------------- # Target dispatch @@ -196,6 +204,8 @@ async def rebuild_index_from_source( owned_inactive_generation=owned_inactive_generation, max_cached_payload_bytes=_REBUILD_CENSUS_SPILL_CACHE_BYTES, ingest_workers=resolved_ingest_workers, + commit_batch_size=_REBUILD_COMMIT_BATCH_UNITS, + replay_commit_batch_size=1, bulk_fts=bulk_fts, bulk_build=bulk_build, prefetch_cache=prefetch_cache, diff --git a/polylogue/sources/revision_backfill.py b/polylogue/sources/revision_backfill.py index 212f584d5f..88a7197514 100644 --- a/polylogue/sources/revision_backfill.py +++ b/polylogue/sources/revision_backfill.py @@ -8,14 +8,16 @@ import sqlite3 import tempfile import threading +import time from collections.abc import Callable, Iterator, Sequence from contextlib import contextmanager from dataclasses import dataclass from io import BytesIO from pathlib import Path from types import TracebackType -from typing import BinaryIO, Literal +from typing import BinaryIO, Final, Literal +from polylogue import logging as _polylogue_logging from polylogue.archive.ingest_flags import ( COMPACT_BROWSER_CAPTURE_INGEST_FLAG, DOM_FALLBACK_INGEST_FLAG, @@ -38,6 +40,8 @@ from polylogue.sources.sqlite_snapshot import looks_like_sqlite_bytes from polylogue.storage.sqlite.archive_tiers.archive import ArchiveStore +_LOGGER = _polylogue_logging.get_logger(__name__) + def _browser_snapshot_fidelity(ingest_flags: Sequence[str]) -> Literal["dom", "native"] | None: """Derive membership-classification browser fidelity from parser ingest flags. @@ -596,6 +600,7 @@ def backfill_historical_revision_evidence( max_cached_payload_bytes: int | None = None, ingest_workers: int = 1, commit_batch_size: int | None = None, + replay_commit_batch_size: int | None = None, bulk_fts: bool = False, bulk_build: bool = False, prefetch_cache: RawParsePrefetchCache | None = None, @@ -663,9 +668,24 @@ def backfill_historical_revision_evidence( """ adoption_deferred = 0 quarantined = 0 + stage_timings: dict[str, float] = {} logical_keys: set[str] = set() - replay_batch_size = commit_batch_size if commit_batch_size is not None and commit_batch_size > 0 else None - replay_batched = replay_batch_size is not None + # The REPLAY phase's batch size is separately tunable + # (``replay_commit_batch_size``; ``None`` inherits ``commit_batch_size``): + # each replayed cohort may flush blob-publication receipts on a SEPARATE + # source.db connection (a deliberate GC-safety design, see + # storage/blob_publication.py), and that connection waits at BEGIN + # IMMEDIATE behind the batch's held write lock -- a long replay batch + # window therefore deadlocks into 'database is locked' once the 30s busy + # timeout expires. Callers that batch aggressively (the full rebuild) + # pass ``replay_commit_batch_size=1`` to keep replay at per-cohort + # commits while still batching the census phase, which has no separate- + # connection writers inside its window. + effective_replay_batch = replay_commit_batch_size if replay_commit_batch_size is not None else commit_batch_size + replay_batch_size = ( + effective_replay_batch if effective_replay_batch is not None and effective_replay_batch > 0 else None + ) + replay_batched = replay_batch_size is not None and replay_batch_size > 1 archive_context = ( ArchiveStore.open_owned_inactive_generation( archive_root, @@ -680,6 +700,7 @@ def backfill_historical_revision_evidence( archive_context as archive, _ParsedSessionSpill(archive_root, max_cached_payload_bytes=spill_cache_bytes) as spill, ): + census_started = time.perf_counter() census = _census_historical_revision_evidence( archive, spill, @@ -689,6 +710,8 @@ def backfill_historical_revision_evidence( commit_batch_size=commit_batch_size, prefetch_cache=prefetch_cache, ) + stage_timings["census"] = time.perf_counter() - census_started + receipt_started = time.perf_counter() censused_raw_ids, _censused_keys = archive.expand_raw_membership_selection(selected_raw_ids) # The direct backfill entry point must publish the same current-parser # receipt as the census-only entry point before it assigns or applies @@ -696,6 +719,7 @@ def backfill_historical_revision_evidence( # durable receipt writer observes one complete source snapshot. archive.commit() _record_raw_authority_parser_census(archive_root, tuple(censused_raw_ids)) + stage_timings["census_receipt"] = time.perf_counter() - receipt_started membership_candidates = census.membership_candidates provisional_full_raw_ids = census.provisional_full_raw_ids @@ -723,7 +747,11 @@ def commit_replay_unit() -> None: # cohort to membership governance and let parsed-content # prefix rules decide it; append chains remain byte-governed. for raw_id in archive.convertible_full_revision_raw_ids(logical_key): + spill_started = time.perf_counter() sessions, _payload_bytes = spill.for_raw(archive, raw_id) + stage_timings["spill_load"] = stage_timings.get("spill_load", 0.0) + ( + time.perf_counter() - spill_started + ) if len(sessions) != 1: raise RuntimeError(f"full revision {raw_id} no longer parses to one session") archive.replace_raw_membership_census( @@ -740,7 +768,11 @@ def commit_replay_unit() -> None: parsed_by_raw_id: dict[str, ParsedSession] = {} retained_bytes = 0 for raw_id in plan.accepted_raw_ids: + spill_started = time.perf_counter() sessions, payload_bytes = spill.for_raw(archive, raw_id) + stage_timings["spill_load"] = stage_timings.get("spill_load", 0.0) + ( + time.perf_counter() - spill_started + ) if len(sessions) != 1: raise RuntimeError(f"classified raw revision {raw_id} no longer parses to one session") parsed_by_raw_id[raw_id] = sessions[0] @@ -761,6 +793,7 @@ def commit_replay_unit() -> None: plan, parsed_by_raw_id, acquired_at_ms=0, + stage_timings_s=stage_timings, manage_transaction=not replay_batched, bulk_fts=bulk_fts, bulk_build=bulk_build, @@ -797,7 +830,11 @@ def commit_replay_unit() -> None: if head_raw_id is not None and archive._raw_revision_authority(head_raw_id) == "quarantined": candidate_raw_ids.add(head_raw_id) for raw_id in sorted(candidate_raw_ids): + spill_started = time.perf_counter() sessions, payload_bytes = spill.for_raw(archive, raw_id) + stage_timings["spill_load"] = stage_timings.get("spill_load", 0.0) + ( + time.perf_counter() - spill_started + ) for session in sessions: session_logical_key = f"{session.source_name.value}:{session.provider_session_id}" if session_logical_key != logical_key: @@ -839,6 +876,7 @@ def commit_replay_unit() -> None: member_sessions, projections, acquired_at_ms=0, + stage_timings_s=stage_timings, manage_transaction=not replay_batched, bulk_fts=bulk_fts, bulk_build=bulk_build, @@ -854,6 +892,12 @@ def commit_replay_unit() -> None: commit_replay_unit() if replay_batched: archive.commit() + if stage_timings: + stage_timings["total"] = time.perf_counter() - census_started + _LOGGER.info( + "backfill stage timings: %s", + " ".join(f"{key}={value:.1f}s" for key, value in sorted(stage_timings.items(), key=lambda kv: -kv[1])), + ) return RevisionBackfillResult( census.scanned, census.classified, @@ -1264,11 +1308,31 @@ class _ParsedSessionSpill: I/O for completeness and makes no raw cohort silently disappear. """ + #: Decoded-session RAM layer budget (payload-equivalent bytes). Replay + #: consumes a cohort almost immediately after census parses it, so a + #: small hot layer turns the common for_raw() into a dict hit instead of + #: a pickle.loads round-trip. Bounded independently of the sqlite layer. + _DECODED_CACHE_PAYLOAD_BYTES: Final[int] = 256 * 1024 * 1024 + def __init__(self, archive_root: Path, *, max_cached_payload_bytes: int | None) -> None: - fd, name = tempfile.mkstemp(prefix=".revision-census-", suffix=".sqlite", dir=archive_root) + # Place the spill beside the RESOLVED index tier, not the archive + # root: on deployments where the .db files are symlinks (e.g. root + # SSD config dir -> NVMe data disk), a spill in archive_root would + # put census churn on the wear-limited disk the symlinks exist to + # protect. + index_path = archive_root / "index.db" + spill_dir = index_path.resolve().parent if index_path.exists() else archive_root + fd, name = tempfile.mkstemp(prefix=".revision-census-", suffix=".sqlite", dir=spill_dir) os.close(fd) self.path = Path(name) self.conn = sqlite3.connect(self.path) + # Disposable single-connection cache: durability is meaningless (the + # fallback is reparsing durable source evidence), so skip the + # journal and every fsync -- the per-add commit previously paid a + # synchronous journal cycle per censused raw. + self.conn.execute("PRAGMA journal_mode=OFF") + self.conn.execute("PRAGMA synchronous=OFF") + self.conn.execute("PRAGMA temp_store=MEMORY") self.conn.execute( """ CREATE TABLE parsed_sessions ( @@ -1283,6 +1347,8 @@ def __init__(self, archive_root: Path, *, max_cached_payload_bytes: int | None) self.conn.execute("CREATE INDEX parsed_sessions_logical ON parsed_sessions(logical_key, raw_id)") self.max_cached_payload_bytes = max_cached_payload_bytes self.cached_payload_bytes = 0 + self._decoded: dict[str, tuple[list[ParsedSession], int]] = {} + self._decoded_payload_bytes = 0 def __enter__(self) -> _ParsedSessionSpill: return self @@ -1316,8 +1382,22 @@ def add(self, raw_id: str, sessions: list[ParsedSession], *, payload_bytes: int) ), ) self.cached_payload_bytes += payload_bytes + self._retain_decoded(raw_id, sessions, payload_bytes=payload_bytes) + + def _retain_decoded(self, raw_id: str, sessions: list[ParsedSession], *, payload_bytes: int) -> None: + if payload_bytes > self._DECODED_CACHE_PAYLOAD_BYTES: + return + while self._decoded and self._decoded_payload_bytes + payload_bytes > self._DECODED_CACHE_PAYLOAD_BYTES: + oldest_raw = next(iter(self._decoded)) + _evicted_sessions, evicted_bytes = self._decoded.pop(oldest_raw) + self._decoded_payload_bytes -= evicted_bytes + self._decoded[raw_id] = (sessions, payload_bytes) + self._decoded_payload_bytes += payload_bytes def for_raw(self, archive: ArchiveStore, raw_id: str) -> tuple[list[ParsedSession], int]: + decoded = self._decoded.get(raw_id) + if decoded is not None: + return decoded rows = self.conn.execute( "SELECT parsed, payload_bytes FROM parsed_sessions WHERE raw_id = ? ORDER BY logical_key", (raw_id,) ).fetchall()