You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
bug(mq): a durable's last ack can land after Close, or never if the process exits #665
A consumer's last ack state can reach disk afterEmbeddedNATS.Close() returns, or never, if the process exits first. The durable then keeps its previous ack floor on disk, and rows that were already inserted and double-acked may be redelivered on the next boot, producing duplicate inserts on tables without dedupe or a ReplacingMergeTree.
This is a different cause from #383 (shutdown returning before the drain completes) and would survive its fix.
Mechanism (read from nats-server v2.14.6)
consumerFileStore.flushLoop is started with a bare go statement (server/filestore.go), so neither Server.Shutdown nor WaitForShutdown joins it.
writeState clears o.dirty before its unlocked atomic write (o.dat.tmp, then a rename onto o.dat).
consumerFileStore.Stop writes state, and waits for the flusher (at most 100 ms), only while dirty is set. A write already under way therefore outlives Stop, Shutdown and WaitForShutdown.
Measured (probe, SyncAlways on as in production): two messages, the second acked 150 ms after the first, then Close() at once. Under in-process CPU contention, o.dat changed after Close returned in 20 of 400 runs, and o.dat.tmp was present at the moment Close returned in 19 of 400. Without contention, 0 of 400 in two runs.
Inferred, not reproduced end to end: that a process exiting right after Close loses that write, and that the next boot redelivers the acked rows.
Also affected (inferred, not observed)
Restart tests that reopen the same store directory, for example TestEmbeddedNATS_ADurableOnDiskIsReusedAcrossARestart, can reopen it before the first broker's late write lands, and then see a redelivery the test does not expect.
Possible directions
Upstream: consumerFileStore.Stop should always wait for an in-flight flushLoop write, not only while dirty is set.
Here: after WaitForShutdown, wait until each durable's o.dat.tmp is gone (bounded), or write each consumer's state synchronously on the shutdown path.
Summary
A consumer's last ack state can reach disk after
EmbeddedNATS.Close()returns, or never, if the process exits first. The durable then keeps its previous ack floor on disk, and rows that were already inserted and double-acked may be redelivered on the next boot, producing duplicate inserts on tables without dedupe or a ReplacingMergeTree.This is a different cause from #383 (shutdown returning before the drain completes) and would survive its fix.
Mechanism (read from nats-server v2.14.6)
consumerFileStore.flushLoopis started with a baregostatement (server/filestore.go), so neitherServer.ShutdownnorWaitForShutdownjoins it.writeStateclearso.dirtybefore its unlocked atomic write (o.dat.tmp, then a rename ontoo.dat).consumerFileStore.Stopwrites state, and waits for the flusher (at most 100 ms), only whiledirtyis set. A write already under way therefore outlivesStop,ShutdownandWaitForShutdown.Evidence
t.TempDircleanup fail withdirectory not emptyunderobs/<consumer>/in fix(test): ingest dispatch-loop test races t.TempDir cleanup, reddens main #442. Under six parallel race-enabled test processes it hit about 4% of teardowns.SyncAlwayson as in production): two messages, the second acked 150 ms after the first, thenClose()at once. Under in-process CPU contention,o.datchanged afterClosereturned in 20 of 400 runs, ando.dat.tmpwas present at the momentClosereturned in 19 of 400. Without contention, 0 of 400 in two runs.Closeloses that write, and that the next boot redelivers the acked rows.Also affected (inferred, not observed)
Restart tests that reopen the same store directory, for example
TestEmbeddedNATS_ADurableOnDiskIsReusedAcrossARestart, can reopen it before the first broker's late write lands, and then see a redelivery the test does not expect.Possible directions
consumerFileStore.Stopshould always wait for an in-flightflushLoopwrite, not only whiledirtyis set.WaitForShutdown, wait until each durable'so.dat.tmpis gone (bounded), or write each consumer's state synchronously on the shutdown path.Related: #383, #442.