Repository navigation
docs(mq): external NATS durability rests on replicas, not fsync - #654
Merged
Merged
Conversation
The embedded broker fsyncs every event before the 200. Under mq.backend: nats the partition stream acks after a Raft quorum has stored the event, so WaveHouse does not require sync_always there. Say so in the deployment guide and Durability & Storage, and scope the fsync paragraphs to the embedded broker. The topology check now reports the history and dead-letter streams' replica count as it did the partitions', and at one replica the recommended finding says an ack then rests on one server's disk and its sync_interval. A one-server development cluster still boots. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
A stream in persist_mode: async flushes in the background even under sync_always, so an ack precedes the write. Refuse it on an ingest partition and recommend against it on the dead-letter stream, and say so in the one-replica durability docs. Check the shipped Helm values for a sync option by reading the file: the test server renders only config.merge, so asserting on its SyncAlways could not catch one. Scope durability.md's "no event is dropped" to limits, and stop calling embedded the only mq backend. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
The generated history kept 2h. The settings seed's gap window is 15 minutes, so default the history's maxAge to 15m, and say in the generated header and the deployment guide that it must be at least the longest stream.gap_window_minutes among the tenants served. The sweeper's warning for a longer window is unchanged. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
PurgeAcked compared the sweeper's cutoff, taken before the call, with its own time.Now() minus max_age, so a window exactly as long as the history read as longer. At the new 15m default that warned for every tenant on the seed's 15-minute window. Allow a second's slack, and pin the equal case. Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
Part of #613. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL
|
Important Review skippedAuto reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Advanced Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
7 tasks done
EricAndrechek
added a commit
that referenced
this pull request
Sep 29, 2026
…at/mq-nats-wiring Collapses #654, #646 and #644 into #639. Conflicts in CHANGELOG.md and architecture.md kept both sides. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01FyrXjhR7iDg33paioLQHFq
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Part of #613.
Stacking: the base is
feat/coord-nats-kv(#646). This PR's own change is the four commits after5ba2bc13.Why
The embedded broker is a single JetStream node, so it
fsyncs every event before the200(SyncAlways). Undermq.backend: nats, a publish to a stream with 3 or more replicas is acked only once a Raft quorum has stored it. That quorum ack is the durability WaveHouse relies on, so it must not requiresync_alwayson the operator's servers. This PR states that contract and has the topology check report the case where it does not hold.Measured before the change (at
5ba2bc13)Nothing in nats mode requires or configures fsync:
deployments/nats/values.yamlhas nosync,sync_intervalorsync_alwayskey anywhere. A new unit test now pins this.nats_topology.go) had no sync check.num_replicas < 3was a recommended finding on the partitions and the lease bucket, but not on the history or dead-letter streams.nats_manifests.go,wavehouse mq manifests) already defaulted to 3 replicas, and it sets no sync or persist option.ExternalNATS.publishcallsjs.PublishMsgsynchronously withWithExpectStream, so it returns after the stream's PubAck. For R>1 that is the quorum ack.defaultSyncInterval = 2 * time.Minute(server/filestore.go:333). A stream store and its Raft log inherit the server'sSyncAlways/SyncInterval. A stream withpersist_mode: async, which is allowed only at R1, forcesSyncAlways=falseand asynchronous flushing, even when the server setssync_always(stream.go:995-1003,1844-1850).SyncAlwaysin WaveHouse ismq.EmbeddedSyncAlways(embedded only).What
internal/mq/nats_topology.go):num_replicas < 3is now a recommended finding on the ingest partitions, the history and the DLQ, through one helper. At R1 the text says an ack then rests on one server's disk, and a crash loses what it stored since its last sync (sync_interval). At R2 it recommends 3 across failure domains. Because it is only recommended, a one-server dev cluster still boots.persist_mode: asyncis a required finding on an ingest partition (the ack precedes the write) and a recommended one on the DLQ. It came up in review:sync_alwaysdoes not reach such a stream.--replicas 1and gets the recommended findings.maxAgechanges from 2h to 15m, matching the settings seed'sstream.gap_window_minutesof 15. The generated header and deployment.md say to set it to at least the longeststream.gap_window_minutesamong the tenants served.deployments/nats/jetstream.yamlis regenerated.max_age. It compared the sweeper's cutoff with a latertime.Now(), so at the new default it would have warned for every tenant on the seed's window. It now allows one second of slack.200; (b) nats acks after the stream's quorum stores it, with no fsync required, so durability comes from replicas spread across failure domains; (c) at R1 the server'ssync_intervalgoverns, and events since the last sync can be lost on a crash. The Persistent Storage fsync paragraph (~line 178) is now scoped to the embedded broker.embeddedis the only backend.Tests
TestReplicasProblem(unit): R0 and R1 namesync_interval, R2 recommends 3 without it, and R3 and R5 give no finding.TestVerifyNATSTopology_Findings: new cases for one replica on a partition, the history and the DLQ, and forpersist_mode: asyncon a partition (required) and the DLQ (recommended). The shipped-manifest finding counts are updated (N+2, and N+3 with the lease bucket).TestShippedValues_SetNoSync(unit): walksvalues.yamlfor any key containingsync. The fixture'sServerConfigrenders onlyconfig.merge, so a server-side assertion alone could not catch aconfig.jetstreamoverride.TestExternalNATS_PublishesWithoutSyncAlways(integration): against a server withSyncAlwaysoff, a publish to an R1 partition is acked and stored, and the verifier reports only recommended findings, including that partition'snum_replicas.TestExternalNATS_PurgeAckedWarnsOnAShortHistorygains an equal-window tenant that must not warn. Measured: it fails against the previousexternal.goand passes with the fix over-count=3.Evidence
make ci(shared queue,GOTOOLCHAIN=go1.26.6) atd563fa9a: exit 0, "All CI checks passed".internal/mqintegration: 65 tests, 3 skipped.make build-docs: all internal links valid.Review
The reviewer markers are keyed to the main checkout (#454), so the verdicts are recorded here. No marker was hand-written.
pre-push-reviewer(opus), four rounds:SyncAlwaysassertion could not see the values file. Fixed withTestShippedValues_SetNoSync.persist_mode: asyncwas not covered. Fixed with the verifier findings and docs.docs-reviewer(opus), five rounds:Left to later PRs / not changed
mq.sync_intervalknob for the embedded broker (mq: expose JetStream sync_interval as a config knob (throughput vs durability tradeoff) #139 still tracks it).PurgeAckedcomparison and tests.🤖 Generated with Claude Code
https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL