Skip to content

feat(hook-context-intelligence): self-healing event replay for undelivered events - #117

Open
Diego Colombo (colombod) wants to merge 14 commits into
mainfrom
feat/self-healing-event-replay
Open

Diego Colombo (colombod) wants to merge 14 commits into
mainfrom
feat/self-healing-event-replay

Conversation

@colombod

Copy link
Copy Markdown
Collaborator

Summary

Fixes microsoft/amplifier#403. Against a remote destination the telemetry hook silently
failed to deliver the large majority of a session's events while reporting itself healthy.
Measured on one workstation over 25 days (232 shutdown_undelivered records, each matched
to its session's events.jsonl): 239,690 events produced, 54,145 never delivered
(22.6%); the median session lost 55.5% of its events; 128 of 232 sessions lost half or
more; the worst lost 719 of 720.
breaker_open=False on 232 of 232 — the destination was
healthy every single time. _DestinationDispatcher posts serially, so its ceiling is one
event per round-trip; a busy session out-produces it all day, the queue hits
dispatch_queue_capacity and sheds.

Worse than the raw rate suggests: the queue drops the newest event on overflow, so what
goes missing is each session's tail — session:end, final tool calls, outcomes. The
graph ended up full of sessions that started and never finished.

This adds a per-session, per-destination byte watermark into the already-durable
events.jsonl, plus a bounded, continuous catch-up sweep that advances it. Data loss
becomes delivery latency: what the live path drops is swept later, automatically.

Builds on #111 (close_drain_timeout 10→20), which addressed short tails only and said so.

Scope / guardrails

  • Seam crossed — deliberately. Client↔server boundary, networking, and hook config
    wiring. Proven with a real DTU run (below), not mocks.
  • Fan-out routing untouched. The sweep sends only where the live path already would,
    reusing fanout.destination_is_active. No new routing decisions.
  • No import of the upload tool. The dependency arrow is tool → hook; the reverse would
    be circular. The sweep is self-contained because build_payload already lives in the hook.
  • sweep_concurrency ships at 1 — exactly today's behaviour. The knob exists; raising
    it is a separate, measured change that must first settle breaker-window ownership.
  • context-intelligence-recover (new on main) is untouched and unreferenced. See Notes.

Verification

  • Module tests pass — modules/hook-context-intelligence: 762 passed
    (the module this change lives in)
  • Top-level tests pass — tests/: 854 passed
  • ruff check — All checks passed! · ruff format --check — 183 files already formatted
  • pyright — 0 errors, 0 warnings, 0 informations
  • Full bundle validation — scripts/validate-full.sh: 12/12 bundles load clean,
    0 errors
    across discovery/hygiene/structure/placement/packaging. The lone
    unadvertised_but_referenced mode ERROR is the documented FALSE POSITIVE per
    AGENTS.md — not acted on.

Real evidence on seams (not mock-only)

  • Seam proven against real behaviour.

New DTU profile context-intelligence-self-healing-replay-validation.yaml, run against a
real context-intelligence-server 6.7.0 + real Neo4j 5.26.22 (this repo's own
context-intelligence-backend.yaml), with a latency-injecting proxy (250 ms/request,
support/self-healing-replay/latency_proxy.py) in front, driving real amplifier
sessions
with the kernel-mounted hook, installed from the branch via Gitea. Readiness
7/7, including sweep-defaults-on verified live.

Scenario Result
S1 reproduce 26 records, only 22 :Event on the server; watermark 127983/129234 — short of EOF, tail undelivered
S2 heal a later session's sweep took the watermark to EOF, no user action
S3 duplicate-free session nodes 25 → 25 after a full forced replay of all 26 records
S4 outage stays loud dead port: delivered=0, watermark stayed 0, last_outcome=no_progress, real ConnectError recorded
S5 age bound audible 7-day-old session not sent, reported stranded at 168h

The duplicate check is asserted on Neo4j node counts scoped per session, not on
{"status":"duplicate"} responses: idempotency_key is an in-memory 7-day/100k LRU that
?replay=true bypasses and a restart clears. The durable contract is the deterministic
node_id + (node_id, workspace) uniqueness with MERGE (verified in the server source at
739022c). The sweep therefore deliberately does not set replay=true.

Exit cost, measured in both states — the constraint this change amends:

no sweep mid-delivery (normal, incl. nothing to do) : 0.000s
a sweep IS mid-delivery                             : 2.002s   (= sweep_close_grace_seconds)
sweep_close_grace_seconds: 0.0                      : 0.000s

Docs & diagrams

  • bundle.dot / bundle.png — N/A — no bundle structure change (two new handler
    modules inside an existing module; the validator reported no BUNDLE_DOT_STALE).
  • README updated — seven sweep_* config keys, and a Self-healing backlog sweep
    section covering the watermark, the bounds, and the exit budget.
  • docs/remote-server-troubleshooting.md — the shutdown: N undelivered entry's
    Case 2 now describes automatic recovery and how to confirm it from the watermark
    file; new entries for the stranded-events warning and for "why did exit take a
    couple of seconds longer".
  • AGENTS.md — new section on scheduling defects, because this work is a worked
    example of why the seam gate exists (see Notes).

Notes / follow-ups

  1. Three defects shipped past a fully green 740-test suite and were caught only by the
    DTU
    , all the same shape — work performed but not recorded, or a wait bounded on the
    wrong thing: (a) the sweep never ran, cancelled at teardown before its first round-trip;
    (b) cancellation discarded a pass's completed deliveries; (c) the fix for (a) waited on a
    task that never completes, adding a flat 2.003 s to every exit. A fourth was a test
    defect of the same family — asserting duplicate-freedom on a workspace-wide node count
    while the test itself created new sessions, a confident false FAIL. The lasting lesson is
    captured in AGENTS.md.

  2. Relationship to context-intelligence-recover (landed on main during this work).
    Complementary, not duplicative, and nothing in the hook references it: the sweep is
    automatic, in-process, bounded to a 48h window and posts to /events; recover is an
    operator CLI with a stronger SQLite checkpoint that posts to /recovery/events behind a
    native-recovery capability. The natural division — the sweep stops the bleeding,
    recover cleans up history
    — is borne out by the same workstation: 38 sessions /
    2,686 events sit inside the sweep's window, and 191 sessions of historical loss sit
    outside it
    and will only ever be recovered by the CLI. Three follow-ups deserve their
    own issue: should the sweep prefer /recovery/events when advertised; adopt recover's
    stronger truncation detection (prefix hash + inode, vs this PR's size-only guard); and
    name an owner for the two byte-offset cursors now tracking the same file.

  3. Not yet shipped from the design: the cross-session backlog-growth warning (needs
    state outliving one session), raising sweep_concurrency above 1 (needs breaker-window
    ownership settled first), and deleting the in-memory queue's role as source of truth —
    the subtraction this work makes possible.

  4. Untested seam: set_ingestion_filters mid-session swaps cancel and restart sweeps
    correctly by construction, but no DTU scenario exercises that path. Most significant
    gap remaining.

Generated with Amplifier

@colombod

Copy link
Copy Markdown
Collaborator Author

Copilot resolve the merge conflicts in this pull request

Copilot AI commented Sep 17, 2026

Copy link
Copy Markdown

Copilot resolve the merge conflicts in this pull request

Resolved the merge conflicts by merging origin/main and preserving both the updated shutdown/backlog guidance from main and this PR’s self-healing sweep docs. Addressed in merge commit 3644011.

Diego Colombo (colombod) added a commit that referenced this pull request Sep 17, 2026
…e sweep budget

Addresses review on #117.

STARVATION. select_candidates() sliced to sweep_max_sessions BEFORE routing
eligibility was known, so the oldest N sessions could all be ones a destination
may not receive. Selection is deterministic (oldest first), so the same N were
chosen on every pass and a permitted session behind them was never reached --
until it aged out of the 48h window, at which point it became permanently
excluded from automatic recovery too. The routing gate added in the previous
commit is what made this reachable: filtering without also making the bound
filter-aware turned a correct refusal into an indefinite block.

Reported repro (20 older excluded sessions + 1 newer allowed):
  before: {'sessions_considered': 20, 'filtered': 20, 'delivered': 0, 'posted': []}
  after:  {'sessions_considered': 21, 'filtered': 20, 'delivered': 4, 'posted': ['allowed']}

select_candidates now returns every in-window candidate and the sweeper spends
sweep_max_sessions only on sessions that pass the routing gate. Sessions that are
filtered or blocked are not charged against the budget -- work that was never
this destination's to do should not consume its allowance. The bound still
bounds: a separate test pins that eligible work is still capped.

TYPE CHECK. _persist() returned advanced_now with no return annotation, so
Pyright inferred None and rejected the addition into report.bytes_advanced.
Annotated -> int. Module-level Pyright is now 0 errors.

The root-vs-module Pyright gap is itself the lesson: the repo-root invocation
reported 0 errors while the module reported 2. AGENTS.md now says to run pyright
from the module directory as well, with this as the worked example.

REBASE. Rebased onto main (19da724) and resolved the README and
remote-server-troubleshooting conflicts in main's favour for the shutdown-timeout
documentation -- close_drain_timeout is 20.0 from #111, and the Case 1 / Case 2
split is kept -- with the sweep documentation layered on top. Case 2's remedy is
updated: the deep-backlog case is now recovered automatically, with the watermark
fields to confirm it and the signals that mean something is genuinely wrong.

Adds 3 regression tests: the reported 20-excluded-plus-1-allowed scenario, that
max_sessions still bounds eligible work, and that unprovable-routing sessions do
not consume the budget either.

Gates: 780 module tests, 857 root tests, ruff clean, pyright 0 errors from BOTH
the repo root and the module directory.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
@colombod

Copy link
Copy Markdown
Collaborator Author

DTU evidence: all six scenarios green against real infrastructure

Addressing the "not proven end-to-end" gap. Run against a real
context-intelligence-server 6.7.0 + real Neo4j 5.26.22
(this repo's own
context-intelligence-backend.yaml), with a latency-injecting proxy at 250 ms/request
in front, driving real amplifier sessions with the kernel-mounted hook, branch
installed through the CLI from a Gitea mirror at 2b9b8b7. No mock on the seam.

Scenario Result
S1 reproduce the shortfall 27 local records, only 17 :Event on the server; watermark 131301/135477 — short of EOF, tail undelivered
S2 a later session heals it watermark reached EOF (135477), no user action, no warning
S3 replay is duplicate-free session nodes 26 → 26 after a full forced replay; watermark accumulated 0 → 129925 → 135241 → 135477 across sessions
S4 a real outage stays loud dead port: delivered=0, watermark stayed 0, last_outcome=no_progress, real ConnectError recorded
S5 the age bound is audible 7-day-old session not sent, reported stranded at 168h
S6 the sweep obeys include/exclude excluded workspace: 0 nodes. Permitted: 16. No watermark written for the excluded session — it was never opened

S6 output verbatim:

=== S6: the sweep MUST obey each session's own include/exclude ===
planted /work/public (allowed) and /work/client-x (EXCLUDED)
delivered=15 failed=0 filtered=1 blocked=0 err=None
--- server truth: the EXCLUDED workspace must be empty ---
{"nodes": 0, "events": 0, "workspace": "heal-133218-secret"}
{"nodes": 16, "events": 15, "workspace": "heal-133218-allowed"}
PASS: excluded workspace has 0 nodes; permitted has 16
PASS: no watermark written for the excluded session -- never opened

Both sessions sit in one project directory — the collision that makes ignoring the
filter a leak rather than an inefficiency.

Harness defects found and fixed while doing this (ace29f0)

S1–S3 passed first time; S4–S6 did not, and every failure was in the harness, not the
code
. Two are worth calling out because one of them was silently dangerous:

  • S5 would have passed for the wrong reason. It carried a hardcoded bearer token from
    an earlier run, and it asserts zero delivered to prove the age bound — a bad token
    produces zero as well. It was green-by-accident and proving nothing. Both S5 and S6 now
    read SERVER_TOKEN, and S6 asserts events_failed == 0 before asserting the routing
    result
    , so a delivery fault can never again masquerade as a filter decision.
  • S6 planted into a fixed project directory, so a re-run found the permitted session
    already swept and reported delivered=0. Unique per invocation now — the same
    "reused namespace manufactures a false result" trap already fixed for the workspace name.

Also: S4/S5 constructed BacklogSweeper without the now-required spec; _common.sh
sourced its env without exporting, so python children never saw it; and the
new-handlers-importable readiness check used a bare PYTHONPATH import that cannot
satisfy pathspec (amplifier resolves module deps at session load, not into its own
interpreter). That check now runs under uv run --with and additionally asserts the
routing gate exists and that spec is a required argument — so a sweeper that cannot
evaluate include/exclude fails readiness rather than shipping.

Gates

780 module tests · 857 root tests · ruff clean · pyright 0 errors from both the repo root
and the module directory
· validate-full.sh 12/12 bundles clean (the lone
unadvertised_but_referenced mode ERROR is the documented false positive per AGENTS.md).

Resource hygiene

Both DTUs (ci-s6, ci-s6-backend) destroyed; host returned to its pre-run baseline.

Diego Colombo (colombod) added a commit that referenced this pull request Sep 20, 2026
…e sweep budget

Addresses review on #117.

STARVATION. select_candidates() sliced to sweep_max_sessions BEFORE routing
eligibility was known, so the oldest N sessions could all be ones a destination
may not receive. Selection is deterministic (oldest first), so the same N were
chosen on every pass and a permitted session behind them was never reached --
until it aged out of the 48h window, at which point it became permanently
excluded from automatic recovery too. The routing gate added in the previous
commit is what made this reachable: filtering without also making the bound
filter-aware turned a correct refusal into an indefinite block.

Reported repro (20 older excluded sessions + 1 newer allowed):
  before: {'sessions_considered': 20, 'filtered': 20, 'delivered': 0, 'posted': []}
  after:  {'sessions_considered': 21, 'filtered': 20, 'delivered': 4, 'posted': ['allowed']}

select_candidates now returns every in-window candidate and the sweeper spends
sweep_max_sessions only on sessions that pass the routing gate. Sessions that are
filtered or blocked are not charged against the budget -- work that was never
this destination's to do should not consume its allowance. The bound still
bounds: a separate test pins that eligible work is still capped.

TYPE CHECK. _persist() returned advanced_now with no return annotation, so
Pyright inferred None and rejected the addition into report.bytes_advanced.
Annotated -> int. Module-level Pyright is now 0 errors.

The root-vs-module Pyright gap is itself the lesson: the repo-root invocation
reported 0 errors while the module reported 2. AGENTS.md now says to run pyright
from the module directory as well, with this as the worked example.

REBASE. Rebased onto main (19da724) and resolved the README and
remote-server-troubleshooting conflicts in main's favour for the shutdown-timeout
documentation -- close_drain_timeout is 20.0 from #111, and the Case 1 / Case 2
split is kept -- with the sweep documentation layered on top. Case 2's remedy is
updated: the deep-backlog case is now recovered automatically, with the watermark
fields to confirm it and the signals that mean something is genuinely wrong.

Adds 3 regression tests: the reported 20-excluded-plus-1-allowed scenario, that
max_sessions still bounds eligible work, and that unprovable-routing sessions do
not consume the budget either.

Gates: 780 module tests, 857 root tests, ruff clean, pyright 0 errors from BOTH
the repo root and the module directory.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…rmark

Adds DeliveryWatermark: a per-session, per-destination byte-offset cursor into events.jsonl at <session_dir>/delivery/<destination>.json, written atomically (temp + fsync + os.replace).

_append_event now returns the byte offset past its record; that offset plus the session dir travels through _persist_to_disk -> __call__ -> enqueue as a _WatermarkKey, and the dispatcher commits it on delivery.

The tracker FREEZES a session's prefix the moment any record goes undelivered (overflow drop, permanent reject, hard auth failure, breaker skip). The live path's delivered set is NOT a contiguous prefix — drops happen at enqueue time, so a dropped record can sit between delivered ones. Freezing is what keeps the cursor honest and bounds the bookkeeping to queue depth.

Nothing reads the watermark yet; delivery behaviour is unchanged.

Monotonicity is CONDITIONAL. A truncated log rewinds the offset to 0, but a plain monotonic save re-reads the stale on-disk value and silently restores it, leaving the cursor past EOF forever and skipping every event in the log. `reset_by_guard` overrides monotonicity for exactly that case. Pinned by test_guard_reset_survives_save.

Tests reaching directly into _queue.put_nowait were updated for the 3-tuple queue item; two test doubles accept the new kwarg; the enqueue hot-path AST allowlist gained _wm_freeze / setdefault / _SessionWatermarkState with justifications.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…livered events

Fixes microsoft/amplifier#403. Measured: 53,160 events never reached the server across 215 shutdowns, 21,274 dropped on queue overflow, with breaker_open=False on all 215 — the destination was healthy, the client was simply producing faster than it delivered.

BacklogSweeper reads each candidate session's events.jsonl from its watermark forward, rebuilds payloads with the hook's own build_payload, POSTs them, and advances the watermark over the CONTIGUOUS PREFIX only. Runs as a background task from on_session_ready, cancelled in cleanup(), so it never blocks or slows process exit.

Converts permanent data loss into delivery latency: events the live path dropped now arrive on a later session.

Bounded by sweep_max_events (5000), sweep_max_age_hours (48), sweep_max_sessions (20), sweep_concurrency (1), all config-exposed. sweep_enabled defaults TRUE.

Duplicate-free because the server MERGEs on a deterministic node_id under a (node_id, workspace) uniqueness constraint — NOT because of idempotency_key, which is an in-memory 7-day/100k LRU that ?replay=true bypasses and a restart clears. The sweep therefore deliberately does not set replay=true.

Self-healing does not mask a real outage. Console loudness now gates on whether the backlog is BEING RETIRED (a trend the watermark made observable) rather than on whether a backlog exists — the old question, and why a healthy destination warned on 215 of 215 shutdowns. A sweep that delivers zero while backlog remains is loud; sessions holding events older than the sweep window are loud and name the count, the age, and the manual upload command. The durable forwarding-*.jsonl record is unaffected by console policy.

Adds a DTU profile with five scenarios (reproduce via a latency-injecting proxy in front of a real server, heal, duplicate-free as a fixed point on Neo4j node counts, outage stays loud, age bound is audible) plus its support scripts.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…ing instead of task cancellation at teardown

A DTU run against a real context-intelligence-server (6.7.0) exposed that the sweep mechanism was correct while its scheduling was not. Run directly against the same backlog, the sweep cleared 22 events (0 failed) in 6.6s. Scheduled inside a session it delivered zero events because cleanup() cancelled the background task at teardown, and a short session ends before the first ~250ms POST completes. Every unit test was green throughout — the defect only existed at wall-clock time.

Two fixes, both config-exposed:

* sweep_close_grace_seconds (default 2.0, 0.0 disables): cleanup() now awaits in-flight sweeps for a bounded grace period before cancelling. This deliberately amends the previously stated hard constraint 'event delivery must never block or slow process exit'. That absolute was kept so literally that the feature never ran. It is now an explicit, bounded, and configurable budget: 'never slows exit by more than N seconds' — which is testable, whereas an absolute was not.

* sweep_interval_seconds (default 60.0, 0.0 = single pass): the sweep now repeats during the session instead of running once at start. This prevents a busy session's backlog from growing all day, and it makes teardown timing stop mattering.

The continuous sweep and the live dispatcher need no coordination. While the live path keeps up it advances the watermark itself so the sweep skips the session as caught-up. The moment the live path drops or fails a record it freezes the watermark and the sweep resumes from exactly that point.

Stranded-session warnings are now reported once per session instead of on every interval tick, since they describe the backlog rather than individual passes.

Adds tests/test_sweep_scheduling.py (14 tests) covering the bounded grace, the repeat loop, loop survival across a failing pass, cancellation, and the loudness gate including stranded-warning suppression.

Gates: 740 module tests + 839 root tests pass, ruff clean, pyright clean.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…ncellation keeps it

A DTU run exposed a second scheduling defect one layer below the first: the
sweep wrote its watermark only after the ENTIRE window completed, so being
cancelled at session teardown discarded every delivery the pass had made. A
sweep that genuinely re-sent events left its cursor at 0 and the work had to be
redone next session -- the same class of bug the teardown grace fixed, and
equally invisible to unit tests that never cancel a pass.

Delivery now runs in chunks of sweep_concurrency, committing the contiguous
prefix after each chunk, and a CancelledError persists the prefix reached so far
before re-raising. Worst-case loss on teardown is one chunk instead of the whole
pass.

Adds two tests pinning the property: a sweep cancelled after five deliveries
keeps exactly those five (last_outcome=cancelled), and one cancelled before any
delivery records no_progress with the cursor untouched. Also fixes the S3 DTU
scenario's parsing of verify.py output.

Gates: 742 module tests, ruff clean, pyright 0 errors.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…weep to EOF

Two corrections to the self-healing DTU profile, both found by running it.

S3 asserted the fixed point on a WORKSPACE-wide node count while itself driving
extra sessions to advance the sweep -- so those sessions' brand-new events
inflated the total and were indistinguishable from duplicates (56 -> 161, a
false FAIL). verify.py gains count_session_nodes, filtering on the deterministic
node_id prefix, and the check is now scoped to the session actually replayed.
Result: 25 -> 25 nodes after a full forced re-send of all 26 records.

S3 also expected one short session to complete a forced full re-send. It cannot:
26 records at 250ms is ~6.5s against a session plus a 2s teardown grace. The
scenario now drives sessions until the cursor reaches EOF, which is a truer test
anyway -- it demonstrates that chunked progress ACCUMULATES across sessions
(observed: 0 -> 123800 -> 128471 -> 129234).

All five scenarios now pass against a real server + real Neo4j.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…eep that is working

The grace was being paid on EVERY session exit, including the overwhelmingly
common case of nothing to sweep. Measured at a flat 2.003s added to every exit --
precisely the cost the original "never slows exit" rule existed to prevent, and a
straight regression on the healthy path.

Cause: the sweep task is an infinite loop (continuous catch-up), so it never
completes. Waiting on the TASK therefore always burned the entire budget, whether
or not a pass was in flight.

Teardown now waits for the sweep to be IDLE -- its current pass finished -- rather
than for the task to end, via an asyncio.Event cleared around each pass and a
shared deadline across destinations. Measured after the fix:

  idle / nothing to do : 0.000s
  mid-pass, delivering : 2.002s   (the declared budget, and only then)
  grace = 0.0          : 0.000s

So the honest statement of the contract is: exit is unaffected unless a catch-up
sweep is mid-delivery, in which case it may take up to sweep_close_grace_seconds
(default 2.0, shared across all destinations, 0.0 disables).

Adds four regression tests: idle costs nothing, an in-flight sweep may use but
not exceed the budget, zero grace never waits, and the budget is shared across
destinations rather than multiplied by them.

Gates: 746 module tests, 839 root tests, ruff clean, pyright 0 errors.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
README gains sweep_close_grace_seconds and sweep_interval_seconds, corrects the
sweep from one-shot to continuous, and states the exit cost plainly: unaffected
normally, up to the grace only when a sweep is mid-delivery, budget shared across
destinations, 0.0 disables, progress kept on cancellation.

remote-server-troubleshooting.md gains an entry for 'why did exit take a couple
of seconds longer', framing it as evidence the bundle was actively recovering
events rather than a fault, plus a quick-reference row.

AGENTS.md records why the seam gate earned its keep here: three defects shipped
past a green 740-test suite and were caught only by a real DTU run, all three the
same shape -- work performed but not recorded, or a wait bounded on the wrong
thing -- plus the test-design lesson about scoping a uniqueness assertion to the
entity actually replayed.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…clude/exclude

The backlog sweep had NO filter evaluation at all. It selected sessions by age
alone within a project directory and forwarded every one of them to whatever
destinations the CURRENT session matched -- rerouting one session's events using
another session's permissions.

That is a leak, not an inefficiency. A project directory routinely holds sessions
from several working directories: project_slug falls back to "default" when the
capability is unavailable, so unrelated trees collide there. A session under a
working_dir that EXCLUDES a destination could have its events swept to that
destination as soon as any other session in the same directory included it. An
event delivered where it was excluded cannot be recalled.

The design doc claimed "the sweep introduces no new routing decisions". That was
false, and this makes it true:

* BacklogSweeper now REQUIRES the destination's spec (include/exclude) and
  re-evaluates each candidate session's own working_dir from metadata.json
  through the live path's own matcher (fanout.destination_is_active). No second
  implementation of the rules.
* The gate runs BEFORE the watermark is read, before a lock is taken, before a
  byte is sent. An excluded session is never touched -- not even a watermark.
* Routing FAILS CLOSED, inverting this module's usual bias. Every delivery guard
  fails toward re-sending, because a duplicate is absorbed by MERGE while a skip
  loses data. Routing is the opposite: a wrong send is unrecoverable, a skipped
  sweep is picked up next pass. An unprovable working_dir is refused and
  reported (sessions_blocked_unprovable), never assumed safe.
* A destination with no resolvable spec gets no sweeper at all.

Also fixes the related mid-session hole: apply_active_dispatchers cancelled the
previous round's sweeps fire-and-forget AFTER installing the new dispatcher set.
asyncio's cancel() only REQUESTS cancellation, and a sweep owns its own HTTP
client and auth -- closing the dispatcher does not stop it. So excluding a
destination via set_ingestion_filters could leave its sweep still delivering.
Cancellation is now awaited BEFORE the swap, matching what teardown already did.

Adds 9 tests: excluded session never posted, empty include matches nothing,
exclude wins over include, the mixed-project-dir leak scenario, missing and blank
working_dir blocked rather than assumed, and a destination without a spec getting
no sweeper. Plus DTU scenario S6 asserting against a REAL server that an excluded
working_dir contributes zero nodes while a permitted one in the same project
directory heals.

Gates: 771 module tests, 854 root tests, ruff clean, pyright 0 errors.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…e sweep budget

Addresses review on #117.

STARVATION. select_candidates() sliced to sweep_max_sessions BEFORE routing
eligibility was known, so the oldest N sessions could all be ones a destination
may not receive. Selection is deterministic (oldest first), so the same N were
chosen on every pass and a permitted session behind them was never reached --
until it aged out of the 48h window, at which point it became permanently
excluded from automatic recovery too. The routing gate added in the previous
commit is what made this reachable: filtering without also making the bound
filter-aware turned a correct refusal into an indefinite block.

Reported repro (20 older excluded sessions + 1 newer allowed):
  before: {'sessions_considered': 20, 'filtered': 20, 'delivered': 0, 'posted': []}
  after:  {'sessions_considered': 21, 'filtered': 20, 'delivered': 4, 'posted': ['allowed']}

select_candidates now returns every in-window candidate and the sweeper spends
sweep_max_sessions only on sessions that pass the routing gate. Sessions that are
filtered or blocked are not charged against the budget -- work that was never
this destination's to do should not consume its allowance. The bound still
bounds: a separate test pins that eligible work is still capped.

TYPE CHECK. _persist() returned advanced_now with no return annotation, so
Pyright inferred None and rejected the addition into report.bytes_advanced.
Annotated -> int. Module-level Pyright is now 0 errors.

The root-vs-module Pyright gap is itself the lesson: the repo-root invocation
reported 0 errors while the module reported 2. AGENTS.md now says to run pyright
from the module directory as well, with this as the worked example.

REBASE. Rebased onto main (19da724) and resolved the README and
remote-server-troubleshooting conflicts in main's favour for the shutdown-timeout
documentation -- close_drain_timeout is 20.0 from #111, and the Case 1 / Case 2
split is kept -- with the sweep documentation layered on top. Case 2's remedy is
updated: the deep-backlog case is now recovered automatically, with the watermark
fields to confirm it and the signals that mean something is genuinely wrong.

Adds 3 regression tests: the reported 20-excluded-plus-1-allowed scenario, that
max_sessions still bounds eligible work, and that unprovable-routing sessions do
not consume the budget either.

Gates: 780 module tests, 857 root tests, ruff clean, pyright 0 errors from BOTH
the repo root and the module directory.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
Ran the profile against a real context-intelligence-server 6.7.0 + real Neo4j.
S1/S2/S3 passed; S4/S5/S6 did not, and every failure was in the HARNESS rather
than the code. Fixed here, because a scenario that fails for an environmental
reason trains you to ignore it, and one that PASSES for the wrong reason is worse.

* S4/S5 constructed BacklogSweeper without the now-required `spec`. The routing
  gate made it mandatory; these two scenarios are about outage and the age bound,
  so both pass include=("**",) and let the gate through.
* S5 and S6 carried a HARDCODED bearer token from an earlier run. S6 failed
  honestly (delivered=0, every POST 401). S5 was the dangerous one: it asserts
  ZERO delivered to prove the age bound, and a bad token produces zero as well --
  it would have passed while proving nothing. Both now read SERVER_TOKEN.
* S6 asserts events_failed == 0 before asserting the routing result, so a
  delivery fault can never again masquerade as a filter decision.
* _common.sh sourced ci-verify.env without exporting, so the vars were invisible
  to the python child processes. Now exported.
* S6 planted into a fixed project dir, so a re-run found the permitted session
  already swept (watermark at EOF) and reported delivered=0. Unique per
  invocation now -- the same "reused namespace manufactures a false result" trap
  already fixed for the workspace name.
* new-handlers-importable used a bare PYTHONPATH import, which cannot satisfy
  pathspec (amplifier resolves module deps at session load, not into its own
  interpreter). It now runs under `uv run --with`, and additionally asserts the
  routing gate is present and that `spec` is a required argument.

Final result, all six green against real infrastructure:
  S1 reproduce      17 :Event for 27 records; watermark 131301/135477
  S2 heal           watermark to EOF (135477), no user action
  S3 duplicate-free session nodes 26 -> 26 after a full forced replay
  S4 outage loud    offset stayed 0, no_progress, real ConnectError
  S5 age bound      0 sent, 1 stranded session reported at 168h
  S6 routing        EXCLUDED workspace 0 nodes, permitted 16, no watermark
                    written for the excluded session

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
This PR adds a rule to AGENTS.md requiring `pyright` to be run from the
module directory, not only the repo root. The module it changes did not
satisfy that rule: `modules/hook-context-intelligence` reported 7 errors,
all in `tests/test_logging_handler_disk_breaker.py`, where `HookResult`'s
`user_message` is declared `str | None` and the assertions never narrow it.

A gate that is red on arrival teaches the next engineer that red is this
gate's normal state, which then covers the real failures it was added to
catch. So the rule and the green gate ship together.

Three narrowing asserts; no test behaviour changes. Also drops the trailing
whitespace that made `git diff --check` fail on the replay DTU profile.

Module pyright: 7 errors -> 0 errors, 0 warnings.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…haring cursors

Three review findings, all the same failure this change exists to kill: a
delivery path that quietly stops recovering and reports itself healthy.

1. A permanent rejection wedged the whole remaining backlog. The live
   dispatcher classifies 403/400/404/410/413/422/3xx as _PERMANENT and steps
   over them loudly; the sweep treated every non-2xx as "stop", so one routine
   403 parked the cursor in front of it. Every later pass re-read that record,
   failed, advanced nothing -- for the full max_age_hours window, after which
   the session aged out and its remaining events were unrecoverable.
   The sweep now imports _classify_http_outcome rather than re-deriving it, so
   the two delivery paths cannot disagree about what is retryable again.

2. And that was invisible. `made_progress` sums the whole pass, so a wedged
   session hid behind any sibling that delivered -- the breaker_open=False on
   232 of 232 lossy shutdowns, rebuilt one layer up. Adds a per-session
   `sessions_no_progress` counter and a LOUD warning that fires even when the
   aggregate moved.

3. The sweep never wrote to forwarding-*.jsonl, the durable channel the
   troubleshooting doc teaches operators to aggregate -- while two comments
   claimed it did. It now writes `sweep_permanent_reject` and
   `sweep_no_progress` into the dispatcher's own sink; the comments are
   corrected and the doc lists both kinds.

4. Colliding destination names silently shared one watermark. `path_for` is
   many-to-one (`prod/a` and `prod:a` both become `prod_a.json`), and the raw
   name was stored precisely so the collision would be detectable -- then never
   compared. With the same URL the existing URL guard cannot fire, so the
   second destination inherited EOF and delivered nothing. GUARD 3 now checks
   the name it was designed to check.

Five tests, each verified red against the pre-fix source.

modules/hook-context-intelligence: 785 passed (was 780)
root: 933 passed - ruff, format, pyright (root and module) all clean

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
@colombod

Diego Colombo (colombod) commented Sep 21, 2026 •

Copy link
Copy Markdown
Collaborator Author

Update pushed — head 04cf76af8a501154245b044f82399f7f4eef2440, rebased onto a80ce69, all checks green.

What this PR contributes

Against a remote destination, the telemetry hook delivered events serially — one per round-trip — so a busy session out-produced it, the in-memory queue hit capacity, and it shed the newest event. What went missing was each session's tail: session:end, final tool calls, outcomes. The graph filled with sessions that started and never finished, while the destination reported itself healthy throughout.

This adds a per-session, per-destination byte watermark into the already-durable events.jsonl, plus a bounded, continuous catch-up sweep that advances it. Data loss becomes delivery latency: what the live path sheds is swept later, automatically, with no user action.

What changed in this update

Permanent rejections no longer stall recovery. The sweep now uses the live dispatcher's own HTTP outcome classifier instead of treating every non-2xx as a stop. A record the server will never accept (403, 400, 404, 410, 413, 422, 3xx) is stepped over and recorded; a retryable failure (5xx, 429, 408, 401) still halts the pass.
Improves: previously one permanently-rejected record parked the cursor in front of it, so every later pass re-read it and advanced nothing — leaving the rest of that session's backlog unrecovered until it aged out of the window. Importing the classifier rather than re-deriving it also means the two delivery paths cannot disagree about what is retryable.

Per-session progress is now reported, not just per-pass. A new sessions_no_progress counter and an accompanying warning fire even when the pass as a whole delivered.
Improves: progress was summed across the whole sweep, so a session retiring nothing stayed hidden behind any other session that delivered. That is the same "reports healthy while losing data" shape this change exists to remove.

Sweep outcomes now reach the durable diagnostics log. sweep_permanent_reject and sweep_no_progress records are written to the same per-day forwarding-*.jsonl the dispatcher uses, tagged "source": "backlog_sweep".
Improves: sweep health was console-only, which is nowhere in a non-interactive session. docs/remote-server-troubleshooting.md documents both record kinds and includes a jq recipe for listing sessions that are not retiring backlog.

Destination watermarks are bound to the destination that owns them. The watermark filename is derived from the destination name, and that mapping is many-to-one — two differently-spelled names can resolve to the same file. The stored name is now checked on load, and a mismatch rewinds that destination with a note.
Improves: previously a second destination could inherit the first one's cursor and read itself as already caught up.

Budget accounting and type checking. sweep_max_sessions is charged only against sessions the destination is actually permitted to receive, so ineligible sessions cannot consume the budget ahead of permitted ones. The module's own pyright invocation is now clean.

Verification

root test suite                          933 passed
modules/hook-context-intelligence        785 passed
ruff check . / ruff format --check .     clean, 190 files
pyright (repo root)                      0 errors
pyright (module directory)               0 errors
git diff --check                         clean
CI at 04cf76a                            all checks green

Five tests were added for the behaviours above, each confirmed failing against the prior code before being committed. The end-to-end DTU profile that accompanies this PR exercises the same paths against a real server and Neo4j.

Scope

Delivery itself remains serial — sweep_concurrency ships at 1, matching current behaviour. Making the dispatcher concurrent is a separate change: this client is effectively the rate limiter for a single-process server, and circuit-breaker window ownership would need to be settled first.

Two behaviour notes for reviewers. The sweep does not exclude the session currently being written, and the live watermark is debounced, so in practice the sweep is a second delivery path working the same log alongside the live one rather than a strictly idle-when-healthy pass; re-delivery is absorbed by the server's MERGE on a deterministic node_id. And a duplicate receipt reflects the server's deduplication cache, not a proof of durable acceptance.

Diego Colombo (colombod) added a commit that referenced this pull request Sep 24, 2026
…uto-resume messaging

Issue B: overflow warning misreported as lifetime total instead of delta
- Added _overflow_dropped_at_last_log to track the previous warning's count
- Warning now reports true delta since last warning, plus session lifetime total
- Preserves cumulative _overflow_dropped for downstream readers (escalation, forwarding, shutdown summary)
- Drops suppressed by 60s rate limit fold into next delta without data loss

Issue E: runtime strings still described manual replay as unconditional
- Updated four runtime strings in logging_handler.py to match corrected documentation
- Queue-overflow warning, sustained-delivery-failure escalation, breaker-open message, and shutdown summary now draw the honest line: self-healing backlog sweep replays dropped events automatically
- Manual context-intelligence-upload required only when sweep_enabled=false or session exceeds sweep_max_age_hours (default 48h)

These fixes are folded into PR #117 (feat: self-healing event replay) which already edits these code paths and corrected the documentation half of this messaging story.

Tests added:
- test_overflow_warning_reports_delta_not_lifetime_total: verifies RED on pre-fix behaviour
- test_overflow_warning_scopes_manual_replay_to_sweep_gaps: validates delta+lifetime reporting
- test_open_warning_scopes_auto_resume_to_new_events: validates breaker messaging update

All gates passing: module pytest 788 passed, root pytest 933 passed, ruff check clean, pyright 0 errors at both repo root and module, validate-full.sh with validation_mode=full and 0 ERROR findings.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…uto-resume messaging

Issue B: overflow warning misreported a lifetime total as a delta
- Added _overflow_dropped_at_last_log to snapshot the previous warning's count
- Warning now reports the true delta since the last warning, plus the session total
- _overflow_dropped stays cumulative on purpose: the sustained-failure escalation,
  its forwarding record, and the shutdown summary all read it as a session total,
  so resetting it would have silently under-reported undelivered events at shutdown
- Drops suppressed by the 60s rate limit fold into the next delta rather than being
  dropped from the accounting

Issue E: runtime strings still described manual replay as unconditional
- PR #117 added the self-healing backlog sweep and corrected README.md and
  docs/remote-server-troubleshooting.md, but the runtime strings were not updated
- The breaker-open message was the worst offender: "delivery auto-resumes ... (no
  restart needed)" and "replay the backlog with context-intelligence-upload" sat in
  one sentence. The first describes NEW events, the second a manual action for
  already-dropped ones; reading the first as covering both is how a user silently
  loses events
- Four runtime strings now draw the same line the docs draw -- the queue-overflow
  warning, the sustained-delivery-failure escalation, the breaker-open warning, and
  the shutdown summary. A drop freezes the delivery watermark, so the sweep replays
  it automatically; manual context-intelligence-upload is required ONLY when
  sweep_enabled is false or the session ages past sweep_max_age_hours (default 48h)

Tests added:
- test_overflow_warning_reports_delta_not_lifetime_total: asserts the second warning
  reports the delta (3), not the lifetime total. Verified RED against the pre-fix
  behaviour, where it failed with "assert 4 == 3"
- test_overflow_warning_scopes_manual_replay_to_sweep_gaps: asserts the overflow
  warning states drops are replayed automatically and scopes manual replay to the
  sweep's gaps (disabled, or aged out)
- test_open_warning_scopes_auto_resume_to_new_events: asserts the breaker-open
  warning scopes auto-resume to NEW events and scopes manual replay the same way

Gates: module pytest 788 passed; root `pytest tests/ --ignore=tests/dtu` 932 passed;
ruff check and ruff format --check clean at repo root and in the module (both run
with --no-cache: a stale ruff cache had previously reported these files as formatted
when they were not); pyright 0 errors at both repo root and module;
scripts/validate-full.sh with validation_mode=full, build_tested=True,
build_success=True, 0 ERROR findings, bundle.dot fresh.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
…g cycle

Corrected two testing practices in the AGENTS.md 'Testing & what done looks like' section, both from lessons surfaced while fixing the overflow-warning defects on this branch:

1. The root test command must be 'uv run pytest tests/' not bare 'uv run pytest'. A bare root invocation fails with ModuleNotFoundError and ImportPathMismatchError because each module under modules/ is its own uv project with its own lockfile, conftest.py, and dependencies. This is the per-module layout working as designed, not a defect. CI already matches it: root runs 'pytest tests/ --ignore=tests/dtu', plus one job per module.

2. Ruff gates now pass --no-cache. Ruff's cache can report a file as already formatted when it is not. Real instance: two edited files were reported clean by cached 'ruff format --check .', committed in that state, then CI's Lint job failed on the identical command. Later 'git stash/stash pop' invalidated the cache and the same tree correctly reported '2 files would be reformatted'. A cached green format check is not evidence.

Generated with Amplifier

Co-Authored-By: Amplifier <240397093+microsoft-amplifier@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Remote destinations silently lose the majority of a session's events - no self-healing replay

2 participants