Skip to content

bug(mq): a durable's last ack can land after Close, or never if the process exits #665

Description

@EricAndrechek

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.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.

Evidence

  • Measured (tests): the same late write is what made t.TempDir cleanup fail with directory not empty under obs/<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.
  • 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.

Related: #383, #442.

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions