Area: mq — external NATS shard ownership
ResetOrphaned (internal/mq/external.go:795) decides whether a shard's unacked rows can be safely redelivered to a new owner by checking: no one currently holds the server-side pin, rows are actually pending ack, and the durable has not delivered or acked anything within an activity window (recentlyActive/activeWindow, external.go:849-865, which defaults to max(PinnedTTL, minPinnedTTL), about 10s). That per-consumer check is sound on its own, but the decision to call it for a given unit is made by the membership layer in internal/ingest/claims.go, based on whether that unit's previous owner still looks alive — and that membership view can be wrong for a process that is alive but has just lost its lease.
claimLoop.reap (claims.go:394) drops a process's own membership term once its lease context ends, then joinMembers (claims.go:409) rejoins at the lowest free slot — not necessarily the same slot j the process held before. Between the lease loss and the rejoin (and the giveUp/take reconciliation that follows), other processes' view of "who owns unit U" can point at a slot that is no longer live, even though the process that held it is still running and may still be actively delivering that unit's rows.
Failure scenario: a network blip or a slow membership tick causes process A's lease for slot j to lapse (reap) while A is still up and mid-drain on some units — for example handing units over gracefully after detecting the loss, or simply slow to renew. If the same disruption also prevents A from renewing its NATS-side pin on a unit it still holds, the pin lapses (PinnedTTL, ~10s) before A finishes. Another process, seeing slot j no longer live in its own membership read, calls take() → ResetOrphaned for that unit. If A has not delivered or acked anything on it within the activity window (plausible if A itself is affected by whatever caused the network blip), the reset succeeds and A's in-flight, unacked rows are redelivered to the new owner — a duplicate at-least-once insert if A later also writes the same rows once it recovers.
- Evidence: inferred, from reading
claims.go's reap/joinMembers/take and external.go's ResetOrphaned; not reproduced against an actual lease loss under load.
- This is distinct from — and not fixed by — the orderly handover path (
claims.go's giveUp/handOver), which keeps a unit's pin held through its own drain specifically so ResetOrphaned sees it as still held. The gap is for a process that loses its membership lease without a graceful handover, where the pin's own fate is decoupled from the membership view.
Constraints: duplicate inserts from a reset are the accepted design contract (decided for #624: WaveHouse may RESET/UNPIN an operator-owned durable on takeover, and order is not promised across a change), and are at-least-once by design — dedupe (where enabled) covers this. The question this issue raises is whether the specific reap-then-rejoin timing widens that window beyond what the pin-based check already accounts for.
Reproduce/confirm: in an integration test, hold a unit's pin while forcing the owning process's membership lease to lapse (e.g. pause its coordinator calls) without going through the graceful giveUp/handOver path, then have a second process attempt take()/ResetOrphaned on that unit and check whether it succeeds while the first process is still mid-delivery.
Found in review of #624.
Related: #624, #613.
Area: mq — external NATS shard ownership
ResetOrphaned(internal/mq/external.go:795) decides whether a shard's unacked rows can be safely redelivered to a new owner by checking: no one currently holds the server-side pin, rows are actually pending ack, and the durable has not delivered or acked anything within an activity window (recentlyActive/activeWindow,external.go:849-865, which defaults tomax(PinnedTTL, minPinnedTTL), about 10s). That per-consumer check is sound on its own, but the decision to call it for a given unit is made by the membership layer ininternal/ingest/claims.go, based on whether that unit's previous owner still looks alive — and that membership view can be wrong for a process that is alive but has just lost its lease.claimLoop.reap(claims.go:394) drops a process's own membership term once its lease context ends, thenjoinMembers(claims.go:409) rejoins at the lowest free slot — not necessarily the same slotjthe process held before. Between the lease loss and the rejoin (and thegiveUp/takereconciliation that follows), other processes' view of "who owns unit U" can point at a slot that is no longer live, even though the process that held it is still running and may still be actively delivering that unit's rows.Failure scenario: a network blip or a slow membership tick causes process A's lease for slot
jto lapse (reap) while A is still up and mid-drain on some units — for example handing units over gracefully after detecting the loss, or simply slow to renew. If the same disruption also prevents A from renewing its NATS-side pin on a unit it still holds, the pin lapses (PinnedTTL, ~10s) before A finishes. Another process, seeing slotjno longer live in its own membership read, callstake()→ResetOrphanedfor that unit. If A has not delivered or acked anything on it within the activity window (plausible if A itself is affected by whatever caused the network blip), the reset succeeds and A's in-flight, unacked rows are redelivered to the new owner — a duplicate at-least-once insert if A later also writes the same rows once it recovers.claims.go'sreap/joinMembers/takeandexternal.go'sResetOrphaned; not reproduced against an actual lease loss under load.claims.go'sgiveUp/handOver), which keeps a unit's pin held through its own drain specifically soResetOrphanedsees it as still held. The gap is for a process that loses its membership lease without a graceful handover, where the pin's own fate is decoupled from the membership view.Constraints: duplicate inserts from a reset are the accepted design contract (decided for #624: WaveHouse may RESET/UNPIN an operator-owned durable on takeover, and order is not promised across a change), and are at-least-once by design — dedupe (where enabled) covers this. The question this issue raises is whether the specific reap-then-rejoin timing widens that window beyond what the pin-based check already accounts for.
Reproduce/confirm: in an integration test, hold a unit's pin while forcing the owning process's membership lease to lapse (e.g. pause its coordinator calls) without going through the graceful
giveUp/handOverpath, then have a second process attempttake()/ResetOrphanedon that unit and check whether it succeeds while the first process is still mid-delivery.Found in review of #624.
Related: #624, #613.