Skip to content

Bound negentropy phase deadlines and Parquet shutdown replay - #64

Merged
erskingardner merged 4 commits into
masterfrom
codex/negentropy-phase-progress
Oct 10, 2026
Merged

erskingardner merged 4 commits into
masterfrom
codex/negentropy-phase-progress

Conversation

@erskingardner

@erskingardner erskingardner commented Oct 10, 2026 •

Copy link
Copy Markdown
Contributor

Problem

The bounded negentropy canary timed out during legitimate comparison work because its idle timer only saw admitted events. Follow-up attempts also exposed a replay-floor configuration mismatch and an ingester shutdown that waited for the entire historical Parquet replay until systemd killed it. These failures interrupt reconciliation testing and make routine ingester restarts unsafe.

Fix

Track inventory, connection, comparison, fetching, and draining separately, extending phase deadlines only for validated progress while retaining the overall wall-clock cap. Allow the Parquet shadow to finish its current immutable segment and defer remaining work to startup replay during shutdown; deployments without startup replay still drain their live queue. Preserve replay identity, improve mismatch diagnostics, and document reuse of the recorded floor and verification of graceful shutdown. Grant the isolated socket directory in the checked-in systemd sandbox.

Validation and limits

  • Regression coverage includes slow comparison, repeated/chatter messages, stalled connection reporting, fixed overall deadline, deferred Parquet replay, no-replay queue draining, and restoring the original replay floor without resetting inventory.
  • Full just precommit and Linux release builds are required before the bounded live canary.
  • Shutdown does not interrupt an in-flight conversion/upload or bound ClickHouse draining. Historical Parquet startup still verifies existing segments; optimizing that scan is not part of this change, and the durable replay floor must not be advanced casually to skip it.
  • Live trial remains one relay, three fixed 15-minute windows, six attempts, 20-minute intake and 25-minute maintenance. No rolling activation. Prior failure ledgers remain intact.

@erskingardner erskingardner left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review metadata Value
Reviewed at (UTC) 2026-10-10T13:18:03Z
Commit reviewed 46226d66e8fa94c04a81b18e77a8c9c98b8e8b01
Model claude-opus-5-5
Reasoning level Not exposed by runtime
Recommended action Fix blocking issues before merge

(Posted as a comment because GitHub does not allow requesting changes on your own PR. Treat it as request-changes.)

The premise holds. The canary failed because a slow, empty comparison could not refresh a budget that only admissions refreshed, so phase-aware deadlines are the right fix. The core mechanism is correct: the watch-driven select! re-reads the deadline when it changes and re-checks it before failing, the 9-minute wall still caps everything, and the new worker tests cover the slow-empty, wall-cap and chatter cases.

Blocking (both small):

  1. Phase::Connect is given exactly the same 20s as try_connect's own timeout, and the phase deadline is armed first. An unreachable relay now usually ends as a silent Deadline with no AttemptFailed{Relay} report, where it used to end as a reported Incomplete. That loses diagnostics, and diagnostics are what this PR is meant to improve. See the inline comment.
  2. The doc now says the ingester sandbox must include ReadWritePaths=/run/pensieve, but ops/systemd/pensieve-ingest.service (which the project treats as authoritative) still has ReadWritePaths=/data /archive /var/lib/pensieve. Fix the unit in this PR, or move the deployment-permissions paragraph into a separate ops PR that does.

Non-blocking:

  • The reply-digest set adds little protection and no test exercises it (inline).
  • Duplicated timeout logging and a redundant progress_until (inline).

CI: clippy and rustfmt pass; the test shards and doctests were still pending when I reviewed.

s.phase = phase;
s.deadline = Instant::now()
+ match phase {
Phase::Connect => super::SETUP_TIME,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocking: the connect budget races try_connect's own 20s timeout.

progress.phase(Phase::Connect) arms now + SETUP_TIME (20s) just before relay.try_connect(Duration::from_secs(20)) starts its own timer (_try_connect(timeout, ..) in nostr-relay-pool 0.44.3). For an unreachable or slow-TLS relay, the phase deadline expires first or in the same timer tick. If select! handles the sleep arm first, attempt returns WorkerError::Deadline and drops the SDK future. That path skips report_failure, so the parent gets EOF instead of AttemptFailed{Relay}.

Before this PR, the 2-minute idle budget comfortably covered connect, so this case was always a reported Incomplete. Now the classification depends on timer and poll ordering.

The simplest fix is to drop the special case and let try_connect bound itself:

Suggested change
Phase::Connect => super::SETUP_TIME,
Phase::Compare => self.compare,

(then delete the now-redundant Phase::Compare line below). If you want a distinct connect bound, make it strictly larger than the try_connect timeout.

inherits that IPC group. The worker unit has `pensieve-ipc` only as a
supplementary group; it must not acquire the ingester's archive/data group.
The ingester's `ProtectSystem=strict` sandbox must explicitly include
`ReadWritePaths=/run/pensieve` in addition to its existing writable paths.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocking: the doc and the checked-in unit disagree.

This says the ingester's ProtectSystem=strict sandbox must include ReadWritePaths=/run/pensieve, but ops/systemd/pensieve-ingest.service:63 is still ReadWritePaths=/data /archive /var/lib/pensieve. Under strict, /run is read-only, so an operator who installs the unit from ops/ will fail to bind the socket on the next canary.

The fix belongs in the unit file (ReadWritePaths=/data /archive /var/lib/pensieve /run/pensieve), not only in prose. This paragraph and the /data/negentropy provisioning notes are also unrelated to phase-progress budgets. Either fix the unit here, or split both into a small ops PR.

return Err(WorkerError::Incomplete);
}
let bytes = hex::decode(message.as_ref()).map_err(|_| WorkerError::Incomplete)?;
if replies.len() >= 4096 || !replies.insert(Sha256::digest(&bytes)) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking: the digest set looks speculative, and nothing tests it.

The negentropy initiator holds no per-round state. It answers whatever ranges arrive, so a hostile relay can send endless distinct valid messages and pass this check. That leaves the 4096 cap and the 9-minute wall as the real bounds, and the wall already caps how long any relay can hold the worker. A benign relay cannot repeat a reply without looping deterministically, which the wall also ends.

I'd drop the HashSet/Sha256 (and the sha2 import). Keep only a plain round counter if you want an early stop. Whichever you keep is a new terminal failure mode, so give it a test (for example, the fake relay replays its last NEG-MSG). The doc line "repeated reply digests ... do not buy time" also undersells what happens: they end the attempt as Incomplete.

let state = progress.state();
match outcome {
Err(_) => {
tracing::warn!(relay = %relay, job, phase = ?state.phase, deadline = "wall", rounds = state.rounds, missing = state.missing, batches = state.batches, admitted = state.admitted, "isolated deadline exceeded");

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: the two tracing::warn! calls repeat the same eight fields and differ only in deadline. Compute the label once, for example Err(_) => Some("wall"), Ok(Err(WorkerError::Deadline)) => Some("phase_or_parent_ack"), otherwise None, then log once. Also, progress_until (line 137) duplicates Progress::new's initial deadline. Using progress.state().deadline for the inventory timeout keeps a single source of truth.

@erskingardner erskingardner left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review metadata Value
Reviewed at (UTC) 2026-10-10T13:25:31Z
Commit reviewed 46226d66e8fa94c04a81b18e77a8c9c98b8e8b01
Model Grok 4.7
Reasoning level Not exposed by runtime
Recommended action Fix blocking issues before merge

The compare budget is the right fix for the canary that died during a long inventory comparison, and CI is green. The nine-minute wall, the 4,096-reply cap, and the rule that chatter does not buy time all hold.

Fix the connect-phase deadline before the next canary. It is armed for the same 20 seconds as try_connect, and it starts first, so a hung handshake becomes an unreported Deadline. The parent records that EOF as WorkerLost rather than the RelayFailure it recorded when try_connect returned Incomplete.

.await
.map_err(|_| WorkerError::Incomplete)?;
tracing::info!(relay = %relay_url, phase = "connect", "isolated reconciliation");
progress.phase(Phase::Connect);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

phase(Connect) starts a 20-second deadline, then try_connect(Duration::from_secs(20)) starts its own timer a moment later. On a hung handshake the select loop reaches Deadline first and attempt returns without report_failure. The parent sees EOF and records RetryReason::WorkerLost (runtime.rs treats a protocol error as disappearance without a reliable classification). A connect error that try_connect actually returns still becomes Incomplete plus AttemptFailed, which the parent records as RelayFailure.

Previously the progress budget was two minutes, so the SDK's 20-second connect timeout won and the failure was reported.

Keep the SDK timeout strictly shorter than this backstop so a normal connect timeout still produces AttemptFailed. Call try_connect with something like 15 seconds and leave the phase budget at 20, or drop the Connect deadline and let the existing SDK timeout stand inside the remaining inventory budget. Compare and fetch can keep their phase deadlines: those waits have no SDK timeout, and Deadline there matches the old idle behavior.

A test that hangs the handshake and expects an AttemptFailed frame, while the phase budget is still in the future, would lock this in.

inherits that IPC group. The worker unit has `pensieve-ipc` only as a
supplementary group; it must not acquire the ingester's archive/data group.
The ingester's `ProtectSystem=strict` sandbox must explicitly include
`ReadWritePaths=/run/pensieve` in addition to its existing writable paths.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking. ops/systemd/pensieve-ingest.service has ProtectSystem=strict and ReadWritePaths=/data /archive /var/lib/pensieve. Under that sandbox /run is read-only, so the ingester cannot create /run/pensieve/negentropy.sock. Add /run/pensieve to ReadWritePaths in the unit. The tmpfiles entry already creates the directory; systemd requires that path to exist when the service starts.

@erskingardner

Copy link
Copy Markdown
Contributor Author

Addressed both reviewers in bcd34f2: connect backstop is now 30 seconds around the SDK’s 20-second connect timeout, with a real stalled-handshake test proving AttemptFailed(Relay). The checked-in unit permits the provisioned socket directory (optional when isolated mode is disabled). Added duplicate-reply and 4096-reply-cap tests, clarified their failure semantics, and removed the redundant inventory deadline. All 26 focused worker tests and full just precommit pass. The updated-head CI and bounded production canary remain gates; this comment is not a fresh reviewer approval.

@erskingardner

Copy link
Copy Markdown
Contributor Author

Updated head bcd34f2 passed full just precommit, all seven CI checks, and the Linux release build. Both review blockers were addressed. The authorized bounded canary was attempted but stopped before worker activation: shared inventory rejected the newly selected replay floor with replay cursor identity differs; preserve existing state. Three queued jobs and failure evidence are preserved; no state reset or retry performed. Live archive health is HTTP 200 and the worker remains stopped. Separately, the preceding ingester shutdown reached systemd’s stop timeout and was SIGKILLed after final segment indexing; this also requires investigation before another restart. This is not a successful production canary.

@erskingardner erskingardner left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review metadata Value
Reviewed at (UTC) 2026-10-10T14:25:53Z
Commit reviewed 019524bd5c94548d8660a97aad322dafef912d87
Model claude-opus-5-5
Reasoning level Not exposed by runtime
Recommended action Resolve serious concerns before merge

(Posted as a comment because GitHub does not allow requesting changes on your own PR.)

The negentropy part is now in good shape. bcd34f2 fixed both blockers from the earlier round. The 30s connect backstop now sits outside try_connect's 20s timeout, the stalled-handshake test shows the failure is reported as AttemptFailed(Relay), and the unit file now includes -/run/pensieve. The phase and progress mechanism is correct. The 9-minute wall still caps everything. slow_comparison_can_complete_empty_without_event_admission reproduces the canary failure, and it fails on the old single idle budget.

Serious concern: scope. 019524b adds two changes that the PR title, the PR body and the commit message (which has no body) do not describe:

  1. A new ingester-wide shutdown behavior for the Parquet shadow. It applies to every Parquet-enabled deployment, whether or not negentropy is in use.
  2. A replay-cursor diagnostic and test in sync/inventory.rs, with runbook guidance.

Item 2 is small and clearly comes from the failed canary, so it is fine to keep. Item 1 should be its own PR, or at the very least the title and body should explain it so whoever merges knows that ingester shutdown semantics change. It also looks like it treats a symptom (see the inline comment on main.rs). Startup replay re-hashes every segment at or above the floor on every start, and that is the likely reason shutdown waited.

Non-blocking: the stop-channel semantics and flag could be simpler, the duplicated timeout warn! from the last round is still there, and the new handshake test takes about 20s of real time.

CI at this head: rustfmt, clippy and doctests pass. The three test shards were still pending when this was posted.

// The archive is durable now. Do not hold shutdown hostage to a full
// historical Parquet replay; finish its current immutable work unit.
if let Some(ref handle) = parquet_shadow_handle {
handle.request_shutdown();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This changes ingester shutdown for every Parquet-enabled deployment, but the PR title and body only describe negentropy worker deadlines. Please move it to its own PR, or at least describe it in this PR's title and body.

On the premise: run_notepack_work_unit calls sha256_file(input) before it checks whether the work unit is already Published (pensieve-lake/src/work_unit.rs:102-104). As a result, every start re-reads every 256 MB segment from the replay floor up, even segments that were fully published long ago. That is the likely reason the previous shutdown hit the systemd stop timeout. Stopping at a segment boundary bounds the shutdown wait, which is reasonable on its own. But the backlog is re-hashed on every restart and gets larger as the archive grows, and live notifications queue behind it. A cheap pre-check, for example a path+size+mtime index lookup before hashing, or simply moving the replay floor forward once a segment is published, would fix the cause. With that in place, this shutdown change may not be needed.

self.process(&sealed.path, false);
loop {
crossbeam_channel::select_biased! {
recv(stopping) -> _ => return,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Two small things:

  • recv(stopping) -> _ also fires when the channel is disconnected, i.e. when the handle is dropped without join. In the no-replay case, that gives up the drain-to-completion guarantee the doc comment promises, and it is inconsistent with the replay loop, where try_recv().is_ok() ignores a disconnect. Today this probably can't be reached outside process exit, but the two loops should behave the same way.
  • Simpler: decide once in start_parquet_shadow. Pass crossbeam_channel::never() as stopping when replay_dir is None. That removes the replay_enabled field and the branch in request_shutdown, and the no-replay drain semantics then follow from the channel itself.

}
Ok(result) => {
if matches!(result, Err(WorkerError::Deadline)) {
tracing::warn!(relay = %relay, job, phase = ?state.phase, deadline = "phase_or_parent_ack", rounds = state.rounds, missing = state.missing, batches = state.batches, admitted = state.admitted, "isolated deadline exceeded");

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit, still open from the last round: this duplicates the warn! at line 154 except for deadline. Compute let deadline = match &outcome { Err(_) => Some("wall"), Ok(Err(WorkerError::Deadline)) => Some("phase_or_parent_ack"), _ => None };, log once, and then flatten outcome. That is about 8 fewer lines, and the two log lines can't drift apart.

}

#[tokio::test]
async fn stalled_handshake_reports_relay_failure_before_connect_backstop() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: this test always runs for the full real 20s try_connect timeout, and its budget is 28s, which is tight on a slow CI shard. The test is valuable, so keep it. Consider making the 20s connect timeout a named const (it is already duplicated as SETUP_TIME in the Phase::Connect arm and as a literal in attempt) so that a test build can shorten it.

@erskingardner erskingardner left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review metadata Value
Reviewed at (UTC) 2026-10-10T14:32:32Z
Commit reviewed 019524bd5c94548d8660a97aad322dafef912d87
Model Grok 4.7
Reasoning level Not exposed by runtime
Recommended action Merge

(Posted as a comment because GitHub does not allow approving your own pull request. Treat it as an approval.)

The phase budgets are the right fix for the canary that died during a long empty comparison. Both earlier blockers are fixed on this commit: connect is a 30-second backstop around the SDK's 20-second timeout, and the stalled-handshake test still gets AttemptFailed(Relay); the unit file grants the socket directory. CI is green. The nine-minute wall still wins, and chatter still does not buy time.

The parquet change matches the failed restart: with startup replay, shutdown finishes the current segment and leaves the rest for the next scan. The cursor mismatch stays fail-closed, which is what you want; the next canary has to reuse the stored floor.

Residual, not a code defect: shutdown still joins the ClickHouse indexer before it can exit. If the SIGKILL was Waiting for ClickHouse indexer to finish... rather than parquet replay, this does not shorten that wait.

"replay cursor identity differs; preserve existing state".to_owned(),
));
return Err(Error::Config(format!(
"replay cursor identity differs; preserve existing state (stored floor={}, next={}, configured floor={floor})",

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Non-blocking. This condition also fails when the archive path or prefix differs, but the message only prints floors. A path or prefix mismatch then shows the same floor twice and hides the field that actually changed. Include the stored and configured archive and prefix next to the floor and next.

@erskingardner erskingardner changed the title Track negentropy comparison progress with bounded phase deadlines Bound negentropy phase deadlines and Parquet shutdown replay Oct 10, 2026
@erskingardner

Copy link
Copy Markdown
Contributor Author

Addressed this review round in 06af078: the no-replay sink now uses a never-ready stop receiver, disconnected control is not treated as an explicit shutdown, and a regression covers dropped control while draining. Replay diagnostics now show stored/configured archive, prefix and floor. Updated the title/body to explicitly include ingester-wide Parquet shutdown behavior. Full just precommit passed after the edits. Historical startup hashing remains deliberately unchanged; changing durable floors or trusting only path metadata is not part of this fix.

Canary is held before restart: the still-running bcd34f2 production ingester reported a PutObject dispatch failure on segment 29810 at 14:34:23 UTC and a HeadObject error on 29811 at 14:35:41. Segment 29812 subsequently published, and endpoint DNS/TLS checks pass, so this is not evidence of a persistent total outage. Canonical archive and failed-run evidence are preserved. No new ingester deployment or bounded canary has occurred this turn; worker remains stopped. Latest-head CI remains to be validated.

@erskingardner

Copy link
Copy Markdown
Contributor Author

Final head 06af078 passed all seven CI checks and the Linux release build. Six previously failed Parquet work units were already retried on the server; read-only remote HEAD checks confirmed matching sizes and SHA-256 metadata for all six. Controlled old-ingester shutdown completed without timeout/SIGKILL, with its existing replay-identity fault reported after durability cleanup. New ingester SHA ce643dff93cffc7b6fceaa5eea820ee03e5bd73a81cf1723d28e1001531b5373 is running with restored inventory floor 29808; archive health is 200 and isolated readiness succeeded.

Bounded canary FAILED at 2026-10-10 16:11:19 UTC: relay.damus.io reached the 300-second Compare phase deadline with rounds=0, missing=0, batches=0, admitted=0. Worker was stopped before an automatic retry. Controlled worker loss earlier retained job 1, but recovery to durable completion was NOT proven; zero of three jobs completed. Current Damus NIP-11 metadata does not advertise NIP-77, so relay capability needs explicit preflight before another trial; do not infer that a longer timeout is the fix. Evidence and consistent ledger backup retained at /var/lib/pensieve-canary/20261010T160337Z. Ingester remains active with NRestarts=0; worker remains stopped. No rolling runtime or relay expansion activated.

@erskingardner
erskingardner merged commit 4273ba6 into master Oct 10, 2026
7 checks passed
@erskingardner
erskingardner deleted the codex/negentropy-phase-progress branch October 10, 2026 16:55
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.

1 participant