diff --git a/sw-block/design/v3-recovery-algorithm-consensus.md b/sw-block/design/v3-recovery-algorithm-consensus.md new file mode 100644 index 000000000..bfcb41104 --- /dev/null +++ b/sw-block/design/v3-recovery-algorithm-consensus.md @@ -0,0 +1,306 @@ +# V3 / G7-recovery — algorithm consensus (normative) + +**Status**: Normative — **single source of truth** for recover *algorithm* (**第一性原理** + **可追溯推论**). + +**Audience**: Architects, protocol authors, implementers, QA. + +**How to read this document** + +| Layer | Sections | Stability | +|-------|----------|-----------| +| **I — Recover 算法基础 / 第一性原理** | **§I — P1–P7** | **Maximal stability.** Change requires architect + explicit revision log entry. Implementations **must** preserve every principle or document a **controlled exception**. | +| **II — Architecture derived from §I** | **§II — §4–§6** | **Stable**, but refactorable if wording improves without violating **§I**. | +| **III — Coordination / wire / taxonomy** | **§III — §8–§10** | **Operational detail.** Child docs (**§V §15**) deepen bytes and file paths — **cannot contradict §I–§II**. | +| **IV — 未定细节 · 已知风险登记** | **§IV** | **Explicitly unfinished.** Skipping **§IV** risks **logical bugs at protocol→code handoff**. Each row should close into **§III** or a leaf spec + test. | + +**Rule**: If a **detail is not pinned** anywhere in this doc suite, reviewers **still** MUST check it against **§I (P1–P7)** before merge. + +--- + +## I. Recover algorithm foundation & first principles + +These are **not** preferences — they are the **minimal joint axioms**. Everything in **§II+** serves them. + +### 1. Boundary & truth (**P1**) + +- **Recover is a bounded session**, not ambient state (`StartSession … EndSession` / fail). +- **Authoritative convergence** toward **healthy replicated state** is proven only via **explicit barrier handshake** **`AchievedLSN ≥ engineered target`** together with **`baseDone` / layer-1 conjunct as defined by product**. **Inferring “done” from traffic silence, Kind ratios, or sender idle windows is forbidden.** + +### 2. Exactly one causal WAL frontier per peer-session (**P2**) + +- For a given replica during a recover session there is **exactly one** logical sequence of WAL application that must not be forked (`LSN`-monotonic ingest where required). +- **Double-ship**: the Primary **must not** deliver the **same logical WAL mutation twice** down **two uncoordinated live paths**. Resolving ambiguity by “substrate might dedupe” is **non-compliant** (#2 is preventive, not rehabilitative). + +### 3. Homogeneous WAL (**P3**) + +- Any path that carries **(LSN, LBA, data)** to the Replica is **semantically WAL replication** regardless of framing (`MsgShipEntry` vs `frameWALEntry`). **Rebuild does not invent a second Wal tape** — only **scheduling** (**which LSN ships next** off the **same ordered log**) and **framing** differ. + +### 4. Separation of transports vs separation of truths (**P4**) + +- **Multiple sockets / envelopes** ARE allowed (**legacy vs dual-lane** today); **forked causal ownership of WAL shipping** is NOT. +- The **minimal strict “dual line”** is the **extent/base bulk channel** (**no authoring LSN** on wire). WAL toward a replica is **one shipper decision stream** (**cursor + pin**, **§6**) — **not** a second phantom **`core/recovery`** WAL pump. + +### 5. Replica is the semantic hinge (**P5**) + +- **Steady WAL apply** → substrate path (LSN contiguous where enforced). +- **Recover WAL apply** → same substrate **plus** WAL-claims visible to base (`bitmap` / `MarkApplied`). +- **Recover base bytes** → `bitmap` arbitration **against** WAL-claims (**base cannot argue with LWW across LSN** alone). +- **Substrate cross-LSN LWW resolves WAL-vs-WAL** (including hypothetical duplicate ingress); **bitmap resolves WAL-vs-base**. + +### 6. Pin costs the Primary (**P6**) + +- **Recycle / retention MUST respect** **`MinPinAcrossActiveSessions`** and **honest `BaseBatchAck → SetPinFloor`**. Designing for **parallel base ∥ WAL drain** reduces wall-clock pin exposure (**performance / ops invariant**, correctness rests on **§I P1–P5**). +- **Substrate recycle** (`S` advancement, ring reclaim) **consult** **`RecycleFloorSource`** — see **`v3-storage-logical-pin-gate.md`** (**WALStore** / **memorywal** gated; **smartwal** MUST close before dual-lane is default-eligible). + +### 7. Engine owns trigger numbers (**P7**) + +- **`fromLSN` / `pinLSN` / `targetLSN`** are **frozen or engine-authored** facts at session admission; transport **silently overwriting** (`fromLSN := 0` when engine chose `R+1`) violates **parity** with decision logic. + +--- + +## II. Derived architecture (**from §I**) + +### 4. Architectural picture (**two lines, homogeneous WAL**) + +#### 4.1 Steady state + +``` +Primary ──[ WAL exchange · steady path ]──► Replica ──► substrate ApplyEntry(...) +``` + +Conceptually **steady path** — **append → ship** on **legacy bearer** (**today**: SWRP / `MsgShipEntry`). + +#### 4.2 Recover session (**superposition**) + +| Line | Payload | Consensus | +|------|---------|-----------| +| **WAL** | **(LSN, LBA, data)** tuples | **One `WalShipper` / peer: **pin + `cursor` + `send(incoming, debt)` / `send(∅, debt)` + timer** (**§6**). **Dual-lane `frameWALEntry` is still P3 homogeneous WAL.** | +| **Recover / extent** | Snapshot blocks | **Dedicated bulk channel** (**P4** strict second line). | + +**Anti-pattern (**P2/P4**)**: duplicated **recovery-only WAL pumps** hiding in **`core/recovery/sender`** without **unified shipper cursor** — forks ownership and restores **double-ship** risk. Archive such branches unless architect re-opens. + +### 5. Recover / extent lane (**minimal**) + +Purpose: converge faster than pure WAL for cold images. +Uses **`frameBaseBlock`**, **`BaseDone`**, **`BaseBatchAck` → `pin_floor`** (bytes: **`v3-recovery-pin-floor-wire.md`**). +**Bitmap + `ApplyBaseBlock`** encodes **P5** for base. + +### 6. WalShipper (**single owning entity**) — cursor + pin, not a heavy “mode taxonomy” (**primary implementation locus**) + +**Implementable procedure** (pseudo-code, INV names, **`shipMu`/R1 double-check`): **`v3-recovery-wal-shipper-spec.md`**. Implementations **must** match that doc or obtain **architect exception** logged in consensus **§13**. + +Implementations **MAY** avoid a visible **`Realtime | Backlog` enum** provided the **behavior below** holds. Naming **“modes”** in prose is optional; **state** is **`cursor[replica]`**, **`pin` / `fromLSN`**, and **`head`** on the **one Primary Wal order**. + +#### 6.1 Session injection (**recover manager**) + +On recover **Start** for peer **P**, the **manager/coordinator** sets **`pin`** / **`fromLSN`** (lower replay bound) and initializes **`cursor[P] := fromLSN`** (or product-defined watermark). **Engine owns numbers** (**P7**). + +#### 6.2 Debt / backlog (**same tape**) + +**Operational debt** for peer **P** (what the serializer **must** drain) means **`cursor[P] < head`** — **unsent Wal prefix** on the **single Primary log**. This is **not** a second queue of “different Wal”; it is **`head − cursor`** on **one LSN order**. + +**Policy debt (recovery)**: Upper layers **open** operational debt by admitting a recover session with **`fromLSN` / conservative start** anchored to an **extent-persisted** (or otherwise **truthful**) lower bound (**P7**). **Closing recovery** drives this **gap → 0** — i.e. the algorithm **closes the prefix gap**, not “catch up to **`head`’s instantaneous coordinate**”, which moves every append (**moving target**; **§IV `T2`** refines barrier cuts). + +**Steady (no reopened policy hole)**: After **`cursor == head`** for the session’s purpose, **each new Wal append** triggers **normal ship** but **does not open a new conservative debt class** — only **bounded implementation lag** (buffers, batches) remains. **Semantic debt reopening** requires a **new** session / coordinator decision (extent re-anchor, re-seed, failure path), not “every append adds debt.” + +#### 6.3 Emit priority & **strategy abstraction** (**per append + timer**) + +Implementations **must** obey the **scheduler choices** below. Prose **`send(...)`** names **WalShipper’s decision** (“what goes out next **on the Wal tape for P**”), **not** a second protocol frame type unless product maps it to bytes. + +**(A) `send(incoming, debt)` — tail visible this opportunity** + +- **`debt` empty (`cursor == head`)**: **emit `incoming`** (the **new tail** Wal assigned on append, or the next unsent tail under **R1**). **Steady**: **one append → one forward emit** when transport allows — **debt-free send**. +- **`debt` non-empty (`cursor < head`)**: **emit the oldest unsent** (`cursor`) first; the **new tail `incoming` does not jump ahead** — it **extends** the same tape and **waits** behind **`cursor..head−1`**. **Right model**: **priority dequeue on one LSN order**; **wrong model**: “K live + M backlog” as **two mailboxes**. **`ModeBacklog`**: **`MUST`** obey this literally. **`ModeRealtime`**: gap repair via substrate scan is **forbidden** on the **`NotifyAppend` hot path** — see **§13**. + +**(B) `send(∅, debt)` — timer / auxiliary loop (**no newly paired append**)** + +**MUST** exist: **periodic `ShipOpportunity`** (timer or shipper-internal loop **equivalent**) that runs **under the same serializer** and attempts **`emit-from-cursor`** when **`cursor < head`**, so debt **cannot** depend solely on new appends (**Primary idle starvation forbidden**). **Interpretation**: **Wal-side payload paired to *this tick’s append* may be absent** (**`∅`**) yet **debt drain still runs**. **Literal requirement applies in `ModeBacklog`**; **`ModeRealtime`** **carve-out** — **consensus §13 E‑WALSHIPPER‑DUAL‑MODE**. + +**P2**: **Exactly one serializer** chooses **the next LSN emitted to P** — **no parallel** “steady ship” and “recover pump” **both** advancing unsynchronized cursors for the same peer. + +**Formal recovery arc (informative)**: Session proceeds **with-debt sends** (`send(incoming, debt)` while **`cursor < head`**, plus **`send(∅, debt)`**) **until operational debt clears** (**`cursor == head`** under **§6.4** consistency), then **without-debt sends** for steady tail (**`send(incoming, ∅)`**). + +**Dual-mode carve-out (**`WalShipper`**)**: **`ModeBacklog`** vs **`ModeRealtime`** (**steady / no substrate replay**) is **architect‑approved** under **`§13 — E‑WALSHIPPER‑DUAL‑MODE`**. **`§6.3**(A)(B)** and **§6.8**(3)(4)(9) remain **literally normative in Backlog**; Realtime semantics **do not** replace **P2 / §6.8(1)/(5) / §6.4**. + +#### 6.4 Flip / “mode transition” atomicity = **Primary Wal + ship serialization** (**R1**) + +The historical **gap-at-flip** bug is **not** a magical “mode bit” — it is **check-then-act** between **“`cursor` caught `head`”** and **a concurrent append** extending **`head`**. + +**Normative**: **Cursor catch-up decisions** and **append admission** that affect **`head` / `cursor`** **MUST** be reconciled under **one lock or equivalent serial point**, or use a **double-check** after claiming “caught up” so **no LSN is neither backlog-emitted nor append-emitted**. **“Backlog empty” is `cursor == head` under that serializer’s consistent read.** + +#### 6.5 Worst-case bandwidth (**1:1** ship ≤ append) + +If **link capacity ≈ append rate**, the shipper **MAY** apply **only scheduling**, not magic: + +| Policy | Meaning | +|--------|---------| +| **Drain-only** | **Emit only from `cursor`** until budget; stall or slow append (**admission** product choice). | +| **Lockstep** | **Each append admits at most one forward ship op** pacing **`head−cursor`**. | +| **Append-off, drain-only** | **Temporary freeze** frontend Wal append; **emit until** **debt / `cursor sustainable ship rate** for the session, **`head−cursor` diverges** without bound — **timer + priority alone cannot fix math**. + +**Normative**: implementation **MUST** surface **fail-close** (end session, engine retry/escalate) **and/or** **admission throttle** after **product thresholds** — **unbounded wait is non-compliant**. **Ops “enough bandwidth”** is **necessary not sufficient**. + +#### 6.7 Replica + +**One rewind** at session open (**`expected = fromLSN+1` first frame**) then **same monotonic `checkMonotonic`** on recover Wal ingress (**§7**). + +**Note**: **`§IV T1–T3, T7`** refine edge cases; **§6.4–§6.6** already narrow them. + +#### 6.8 WalShipper scheduling — **implementer checklist** (**SW / code review**) + +**Audience**: `seaweedfs` docs + **`seaweed_block`** reviewers — verifies implementation against **§6** without re-deriving prose. Violating any **MUST** is **non-compliant** unless logged in **§13** as an architect-approved exception. + +1. **SINGLE-SERIALIZER (`P2`)** — Exactly **one** WalShipper (or mechanically equivalent owning struct) **per peer `P`** runs **`cursor` / emit decisions**. **MUST NOT** attach a second live WAL pump (`core/recovery` backlog scan vs steady **`Ship`** both advancing **`cursor`** for **`P`** without **one** `shipMu` / equivalent). + +2. **ONE TAPE** — **`cursor < head`** means unsent prefix on the **same** Primary WAL order; **MUST NOT** model as unrelated “steady mailbox vs recover mailbox”. + +3. **PRIORITY (`§6.3`)** — When **`cursor < head`**, **MUST** ship **from `cursor` forward** until policy stops the round (**budget**, **`ctx`**, transport back-pressure) **before** serving **pure new tail** driven only by arrivals that skipped backlog. + +4. **TIMER — SELF-DRIVING BACKLOG DRAIN (`§6.3` MUST)** — **MUST** implement a **periodic `ShipOpportunity`** (timer or shipper-internal loop **equivalent**) that attempts **`emit-from-cursor`** so backlog **cannot** depend solely on new appends (**Primary idle starvation forbidden**). + +5. **FRAMING ≠ SECOND TAPE (`P3–P4`)** — **`MsgShipEntry`** (legacy bearer) vs **`frameWALEntry`** (dual-lane bearer) is **encode/profile + conn** selection at emit time; **does not** authorize a **forked causal owner** or second **`cursor`**. + +6. **`BASE ∥ WAL` WALL-CLOCK OVERLAP (`P6`, binding `G3`)** — **Extent/base bulk** (**§5**) **SHOULD overlap** WAL drain where product permits (**reduces wall-clock pin exposure**). **Does not** add a second WAL serializer — bulk channel is **not** an LSN tape (**§II 4**). + +7. **R1 — flip atomicity (`§6.4`)** — “Caught up → tail” decisions **under** the **same serializer lock** **with double-check** vs **`head`** after observing **`cursor == head`**. + +8. **R2 — saturation (`§6.6`)** — **`head − cursor` diverges**: **fail-close** **and/or** admission throttle; **never** unbounded silent wait. + +9. **`send` STRATEGY (`§6.3`)** — **`debt` non-empty ⇒ oldest (`cursor`) first**; **`debt` empty ⇒ emit tail (`incoming`)**; **timer / loop ⇒ `send(∅, debt)` drain** (**no append required**). + +**Dual-mode scope (`§13`)**: Checklist items **(3)(4)(9)** (debt-before-tail priority, periodic idle drain tied to **`send(∅, debt)`**) are **`MUST`** in **`ModeBacklog`**; **`ModeRealtime`** obeys **`§13 — E‑WALSHIPPER‑DUAL‑MODE`**. + +**Still open**: flip hysteresis (**§IV `T1`**), **`targetLSN`** vs **`head` at barrier** (**`T2`**), bearer policy while backlog exists (**`T3`**). Closing those rows **updates** checklist edge semantics but **does not** relax **(1)**–**(4)** **in **`ModeBacklog`**, **`§13`**, or **(9)** as read **`ModeBacklog`**. + +### 7. Replica routing table (**P5 operationalized**) + +| Ingress | Shape | Action | +|---------|-------|--------| +| Steady `MsgShipEntry` | WAL tuple | Substrate apply; **LWW** across LSN | +| Recover `frameWALEntry` | WAL tuple | Substrate apply **+** **`bitmap` WAL-claim** | +| Recover `frameBaseBlock` | Bytes | **`bitmap` gate** → **`ApplyBaseBlock` / synthetic frontier** | + +**`checkMonotonic`**: protocol defense on **recover WAL** stream (**gap / backward / dup / +1**). + +--- + +## III. Coordination, wiring, failure surface + +### 8. Coordinator & pin + +**`PeerShipCoordinator`**: single routing truth **per volume** (`StartSession`/`EndSession`, phase). +**Pin / recycle**: **`v3-recovery-pin-floor-wire.md`**, **`MinPinAcrossActiveSessions`**. + +### 9. Wiring & transports + +**Legacy + `port+1` dual-lane** — **`v3-recovery-wiring-plan.md`**. +Bearer differs; **§7 semantic is P3**. + +### 10. Failure taxonomy pointer + +**`Failure*` matrix / engine bypass** — **`core/recovery/failure.go`** + transport classification. **Option B** (retry budget) orthogonal. + +--- + +## IV. 未定细节 · 协议→执行 风险登记 (**explicit TBD**) + +**Policy**: every row is **allowed to stay open** until a **leaf spec or code contract** closes it. **Closing** = promote into **§II–§III** or a linked doc + **INV/test pin**. + +| ID | Open question | Why it can bite at code time | +|----|----------------|------------------------------| +| **T1** | **“Caught up” / flip policy** (**refines §6.4**): strict `cursor == head` vs hysteresis **`head − cursor < ε`** vs **time-based idle**; **product** semantics when **flip races append** — serializer already required | Premature declaring **normal** ⇒ **gap** if **outside §6.4**; late flip ⇒ long session / pin | +| **T2** | **`targetLSN` vs moving `head`**: if **`head > target`** at session start, does backlog drain stop at **target** for **barrier** only, or always emit to **head**? **Related (§6.2)**: **catch-up closes `cursor→head` gap** (operational debt), not “match **`head`’s value at an instant**”; **barrier / pin** still need an **engine-chosen cut** agreed with receiver. | **Barrier false negatives** or **double application** if sender/receiver disagree | +| **T3** | **Bearer choice while `cursor < head`**: **steady bearer silent vs dual-lane only** (**§II–§III** closes); **either way** obey **§6 P2 single serializer** | **P2** if two pumps both emit unconsumed prefix | +| **T4** | **Duplicate** `lsn == applied` on recover path: **hard error** vs **idempotent no-op** | Retransmit policy vs **checkMonotonic** | +| **T5** | **Sparse base** vs dense full-LBA — contract with **bitmap** density | Performance / correctness on huge volumes | +| **T6** | **Multi-replica RF>2**: independent **`cursor[P]` per peer**, not “modes” per replica | Out of **current** binary topology scope — RF>2 convergence **needs** explicit coordinator story + tests | +| **T7** | **Saturation policy detail** (**§6.6**): threshold timing, hysteresis vs flapping sessions, UX/ops signaling | Fail-close/throttle correctness under bursty workloads | +| **T8** | **Replica/session liveness** (**R3**): deadline, eviction, **no-progress** watchdog; **`PinUnderRetention` storms** (**R4**) — **pin_floor** vs **`S`** vs ack cadence; **distinct timers** for stall vs **Primary-restart / fencing** (**R5** — see **`v3-recovery-execution-institution.md`**) | Wrong coupling ⇒ **silent hang** or **spurious recycle** | + +**Residual engineering catalog (**does not overturn **§I**; closes in leaf specs**)**: **dual-lane ingress accounting → one RebuildSession** (**R6**); **substrate Scan contract vs pin** (**R7**); **misconnect / cross-volume tests** (**R8**). **Architect direction**: protocol-first pinning before large code divergence. + +### IV.1 Cross-cutting **progress envelope** (not recover-only) + +Steady **replicate** and **recover** share the same physical bottlenecks: **Wal `head` growth**, **ship `cursor`**, **extent/base bulk**, **receiver apply**, **`SetPinFloor` / `MinPin` / retention `S`**. **§6.6** and **§IV T7–T8** apply whenever **sustained ingress exceeds sustainable drain** — including **non-session** operation. **P5 `bitmap`** gives **WAL-vs-base** correctness at the replica; it does **not** guarantee **throughput parity** or **pin velocity**. + +**Design obligation**: product SHOULD define a single **progress envelope** (metering, throttle, watchdog, degraded / fail semantics) for **steady + recover**, then specialize budgets per phase — see **`v3-recovery-execution-institution.md` § Risk envelope**. + +**Operational analogy** (informative): distributed block systems expose **admission**, **recovery vs foreground limits**, **degraded replicas** rather than indefinite internal queue growth; our failure modes should remain **explicit** at the engine/ops boundary. + +--- + +## V. Binding gaps · QA checks · history · index + +### 11. Binding gaps (until closed) + +| ID | Statement | +|----|-----------| +| **G0** | **`WalShipper` unified loop** (**§6**, checklist **§6.8**) in **production** — **cursor + serializer**, not orphaned sender pumps | +| **G1** | **P2** proven by tests (**CHK-WALSHIPPER-SINGLE-CURSOR** or superseding suite) | +| **G2** | **P7** — dual-lane obeys engine **`fromLSN`** | +| **G3** | **P6** — base ∥ WAL overlap (**perf**) | + +### 12. QA checks + +| ID | Predicate | +|----|-----------| +| **CHK-WALSHIPPER-SINGLE-CURSOR** | **No duplicate logical `lsn` to same peer** on two live outbound paths (**§6.3**) | +| **CHK-WALSHIPPER-TIMER-DRAIN** | **`ModeBacklog` only** ( **`§13` — E‑WALSHIPPER‑DUAL‑MODE** ): no sustained **`cursor ≪ head`** while **Primary append‑idle** without timer‑driven **`emit-from-cursor`** attempts (**§6.8 item 4**); aligns **`v3-recovery-wal-shipper-spec.md` §9** `Test-Timer-Drains-Idle` | +| **CHK-BARRIER-BEFORE-CLOSE** | **P1** | +| **CHK-RECOVER-REWIND-ONCE** | First post-open recover-WAL expects **`fromLSN+1`** | + +### 13. Architect-approved exceptions (**binding**) + +Controlled relaxations of **§6** prose **without** weakening **§I**. Each row **requires** engineer + architect sign‑off **before merge** unless **already Approved** herein. + +#### E‑WALSHIPPER‑DUAL‑MODE — ModeRealtime vs ModeBacklog + +| Field | Statement | +|-------|-----------| +| **ID** | **E‑WALSHIPPER‑DUAL‑MODE** | +| **Status** | **Approved** (**2026-04-30**). | +| **Motivation (**T4a / fresh receiver**)**: Substrate‑scan WAL ship during **steady Realtime** **`NotifyAppend`** can replay **dead‑window** bytes and violate **no‑replay** session invariants asserted by tests/product. | +| **`ModeBacklog`** (recover / catch‑up Wal path on WalShipper) | **`§6.3**(A)** `send(incoming, debt)`, **`§6.3**(B)** `send(∅, debt)` / periodic **`ShipOpportunity`**, and **§6.8**(3)(4)(9) **`MUST`** hold **in full**. Timer **`MUST`** attempt **`emit-from-cursor`** when **`cursor < head`** even if Primary append‑idle (**§6.8 item 4**). **Priority (= oldest unsent)** **`MUST`** use substrate **`ScanLBAs`** (or equivalent single‑tape authoritative read). | +| **`ModeRealtime`** (steady per‑append **`NotifyAppend`**) | **Literal §6.3(B)** “idle timer **`MUST`”** **`MUST NOT`** apply — ship is **`NotifyAppend`**‑driven. **Hot path **`MUST NOT`** substrate‑replay** for ship byte selection (**caller `data`** is canonical on optimized tail emit). **`send(incoming, debt)` debt‑priority **`MUST NOT`** be re‑interpreted as “Realtime must drain before tail”**: gap **`cursor < head`** **`MUST`** be corrected only via **`ModeBacklog`** + **`DrainBacklog`** / coordinator **re‑anchor**, not Realtime substrate scan. **`CHK‑WALSHIPPER‑TIMER‑DRAIN`** applies to **Backlog** only (see **§12** predicate). | +| **Safety switch** | Implementations **SHOULD** expose **`StrictRealtimeOrdering`**: **log‑warn** on **`lsn ≠ cursor + 1`** (dense WAL) **by default**; **strict / error path** optional **production** opt‑in. **Hard opt‑in** (`strict=true`): engine **SHOULD** require prior proof of **ordering discipline** (**rebuild‑on‑gap**, fresh peer policy). Operational detail **`v3-recovery-wal-shipper-mini-plan.md` §11.6**. | +| **Not relaxed** | **P2** single serializer (**§6.8(1)**); **§6.4** (**R1**); **§6.8(5)** framing≠second tape; **§6.6** (**R2**) observability hooks; **§I P1–P7**. | + +**Doc bridge**: **`v3-recovery-wal-shipper-mini-plan.md`** **§11.2 dual‑mode contract** ⇄ this exception. + +--- + +### 14. Revision + +| Date | Change | +|------|--------| +| 2026-04-27 | v1 stub | +| 2026-04-29 | v2 WalShipper / receiver | +| 2026-04-29 | **v3.1**: **How to read** table ↔ **§I P1–P7** numbering; **§IV** TBD cross-refs | +| 2026-04-29 | **v3.2**: **Wal shipper as one tape + `cursor`/pin/timer** (**§6**); **flip = Wal+ship atomicity (**R1**)**; worst-case **1:1** scheduling table; saturation **§6.6 (**R2**)**; **§IV** T7–T8 + R-catalog; QA **CHK-WALSHIPPER-SINGLE-CURSOR** | +| 2026-04-29 | **v3.3**: **§IV.1 progress envelope** (steady+recover); **T8** disambiguates R3/R4 vs **R5 restart**; **execution institution** risk ledger **R1–R6**, spec-before-delete‑bad‑code ordering, substrate admission note (**WALStore** vs **smartwal**) | +| 2026-04-29 | **v3.4**: **§6** → **`v3-recovery-wal-shipper-spec.md`** (**implementable** WalShipper; unified wal stream docs must align) | +| 2026-04-29 | **v3.6**: **`v3-recovery-wal-shipper-mini-plan.md`** (**spec→code bridge**); **doc map** | +| **2026-04-30** | **v3.7**: **`§6.8` implementer checklist** (single serializer, priority, timer drain, framing≠tape, **`G3` pointer**); **§V** **`CHK-WALSHIPPER-TIMER-DRAIN`** | +| **2026-04-30** | **v3.8**: **§6.2–6.3** — **policy vs operational debt**, **moving-target / gap→0**, **`send(incoming, debt)` / `send(∅, debt)`** strategy abstraction + **steady debt-free sends**; **§6.8** item **(9)**; **`T2`** cross-ref | +| **2026-04-30** | **v3.9**: **§13** (**E‑WALSHIPPER‑DUAL‑MODE**); **`§14`/`§15` renumber**; **§6.3** carve‑out pointer; **`CHK‑WALSHIPPER‑TIMER‑DRAIN`** scoped per **§13** | + +### 15. Document map + +| Doc | Role | +|-----|------| +| **This** | **Foundation + index**; **`§6.8`** implementer checklist; **`§13`** architect-approved exceptions (**E‑WALSHIPPER‑DUAL‑MODE**) | +| `v3-recovery-pin-floor-wire.md` | Pin bytes | +| `v3-recovery-wal-shipper-mini-plan.md` | **Implementation bridge** (**seaweed_block** phased PR ↔ spec §INV) | +| `v3-recovery-wal-shipper-spec.md` | **WalShipper algorithm** (priority, **R1**, **INV-***) | +| `v3-recovery-unified-wal-stream-*.md` | **Align to WalShipper spec** or mark superseded | +| `v3-recovery-wiring-plan.md` | Ports / flags | +| `v3-recovery-inv-test-map.md` | INV ↔ tests | +| `v3-recovery-execution-institution.md` | Lifecycle + **risk envelope** / spec→delete-code ordering | +| `v3-recovery-live-line-backlog-spec.md` | Retired stub → **here** | +| `v3-storage-logical-pin-gate.md` | **LogicalStorage substrates + `RecycleFloorGate` / pin** | diff --git a/sw-block/design/v3-recovery-wal-shipper-mini-plan.md b/sw-block/design/v3-recovery-wal-shipper-mini-plan.md index 35d21640b..79de3c512 100644 --- a/sw-block/design/v3-recovery-wal-shipper-mini-plan.md +++ b/sw-block/design/v3-recovery-wal-shipper-mini-plan.md @@ -135,6 +135,7 @@ Wal-shipper-spec **§9** names coarse tests — here map **package + style**: | 2026-04-29 | v0.3 — **§10 P2d decision request** appended; P2c split into **slice A / B-1 / B-2** (all merged into `g7-redo/wal-shipper-impl`) | | 2026-04-30 | v0.4 — **P2d ratified + implemented** (`cb8ff1c`): per-connection format dispatch (steady=`MsgShipEntry`, dual-lane=`frameWALEntry`); resident WalShipper with EmitProfile + EmitKind; RecoverySink wired into `startRebuildDualLane` | | 2026-04-30 | v0.5 — **§11 C1–C3 hardening sequence** appended after architect review of `cb8ff1c` exposed three correctness gaps vs §6.8 checklist (writer race, missing RecordShipped, no timer drain, sequential BASE/WAL) | +| 2026-04-30 | v0.6 — **Dual-mode **`ModeBacklog`/ModeRealtime`** + consensus **`§13` E‑WALSHIPPER‑DUAL‑MODE** (**T4a carve-out**) — rewrote **§11.2**, **§11.3** `#3/#4/#9`, §11.4 anchors, §11.5 (**StrictRealtimeOrdering**, TOCTOU); **§11.6 migrations** (**replicaID**, rebuild-on-gap); CHK/timer scope = **Backlog** | --- @@ -199,7 +200,7 @@ Choice — **per-connection format dispatch**, not single-format unification. Du |---|---------|-----| | **G-WRITE-RACE** | (1) SINGLE-SERIALIZER (mechanical) | `Sender.writeFrame` (`writerMu`) and `EmitFunc.conn.Write` (no mutex) race on same dual-lane conn; can interleave header/payload of two frames | | **G-RECORDSHIPPED** | accounting | WalShipper-routed emits don't call `coord.RecordShipped`; `shipCursor` stays at `fromLSN` for whole session | -| **G-TIMER** | (4) TIMER + new (9) PRIORITY (consensus v3.8) | No periodic `emit-from-cursor` loop. Primary-idle starvation possible. NotifyAppend in Realtime emits new tail directly even when `cursor < head` (debt) — violates `send(incoming, debt)` | +| **G-TIMER** | (4) TIMER + **(9) PRIORITY** (consensus v3.8) | **Originally**: no periodic `emit-from-cursor`; Realtime **`NotifyAppend`** could skip debt. **Resolved in C2 + clarified by dual-mode (**`§13` E‑WALSHIPPER‑DUAL‑MODE**)** — **idle timer MUST** **`ModeBacklog` only**; Realtime ships **without** substrate drain per **T4a** (**§11.2a**). | | **G-PARALLEL** | (6) BASE ∥ WAL overlap | `Sender.Run` runs `streamBase` fully before `sink.DrainBacklog`. Strictly serial. P6 / `G3` requires wall-clock overlap | ### §11.2 Commit sequence @@ -213,15 +214,21 @@ Choice — **per-connection format dispatch**, not single-format unification. Du - `transport.RecoverySink` implements both. Steady (Ship) path keeps own mutex (no contention). - Tests: `-race` on concurrent Sender.writeFrame + EmitFunc; assert `coord.PinFloor` advances during dual-lane session. -**C2 — `send(∅, debt)` timer + `send(incoming, debt)` priority** (consensus v3.8 §6.3 / §6.8 #9): +**C2 — `send(∅, debt)` timer + **`ModeBacklog`** priority** (**consensus** **§6.3 / §6.8 #3,#4,#9**, scoped **`v3-recovery-algorithm-consensus.md` §13 — E‑WALSHIPPER‑DUAL‑MODE**): -- `WalShipper` gains internal goroutine driven by `IdleSleep` cadence. On tick under `shipMu`: if `cursor < head`, run one `ScanLBAs(cursor, ...)` cycle. -- `NotifyAppend` Realtime path uses **debt condition `cursor < head`** (NOT `lsn > cursor + 1` — fails the dense single-LSN-of-debt edge case where `lsn = cursor + 1 = head`). - - `cursor < head` → debt exists → nudge drain, return nil (don't emit `lsn` directly) - - `cursor == head` (no debt) + `lsn > cursor` → direct emit of caller's `data` (optimized tail path) -- `NotifyAppend(lba, lsn, data)` signature unchanged. `data` is **canonical only on the no-debt path**; debt path drain reads substrate (one source of truth). Documented in PR. -- Lifecycle: timer runs for shipper's lifetime (one per (volume, replicaID)); emit gated by session-installed context (nil-conn → silent drop). Timer doesn't autonomously decide a peer. -- Tests: `Test-Timer-Drains-Idle` (CHK-WALSHIPPER-TIMER-DRAIN); `Test-Priority-OldFirst` (§6.8 #3 / §6.3 — NOT CHK-WALSHIPPER-SINGLE-CURSOR); `Test-NoGap-DenseLSN-Edge` (cursor=10, head=11, lsn=11 — debt path required). +- **`ModeBacklog`**: `WalShipper` internal `timerLoop` (`IdleSleep`) + **`nudgeCh`**. **`drainOpportunity`**: under **`shipMu`**, **`mode`, `cursor`, `head`** read atomically (**TOCTOU fix** `377bcb0`); if **`cursor < head`**, one **`ScanLBAs(cursor, …)`** cycle (`emit-from-cursor`), then **`cursor++`**; **`assertCaughtUpAndEnableTailShipLocked`** unchanged. +- **`ModeRealtime`** (`NotifyAppend`): **does not** run substrate-scan ship on **`NotifyAppend` hot path** (**T4a / no‑replay`). **`lsn <= cursor`** ⇒ idempotent **`nil`**; else **`emit(EmitKindLive, …, data)`**, **`cursor = lsn`** (optimized steady tail — caller bytes canonical). **`cursor < head`** in Realtime ⇒ **BUG / caller contract violation** ⇒ engine **`MUST`** drive **`ModeBacklog`** **`DrainBacklog`** **or rebuild** (**§11.6**); **never** silently “repair” via Realtime scan. +- **`StrictRealtimeOrdering`**: **`wal_shipper`** config — **default log-warn** on dense **`lsn ≠ cursor + 1`**; **strict / error-return** optional **production** opt‑in (**pre-condition §11.6**). +- Signature **`NotifyAppend(lba, lsn, data)`**: **`data`** used **Realtime** optimized path **only**; **`ModeBacklog`** drain **`MUST`** read substrate (**one tape** authoritative bytes). +- **`DisableTimerDrain`**: production **`MUST false`** (test-only; see §11.5). +- Tests: **`TestC2_TimerDrainsIdle`** (**CHK-WALSHIPPER-TIMER‑DRAIN**, **`ModeBacklog`** only per **§13**); **`TestC2_PriorityOldFirst`** (§6.8 #3 / #9, Backlog); **`TestC2_NoGapDenseLSNEdge`** (**regression**) — dense **single-unsent‑LSN** debt ships only via Backlog/timer (bans naive **`lsn > cursor + 1`** debt predicate). + +**§11.2a — Dual-mode normative contract (consensus ⇄ mini-plan)** + +| Mode | Scheduling | +|------|------------| +| **`ModeBacklog`** | **`send(incoming, debt)`**, **`send(∅, debt)`**: idle timer **`MUST`** (**§13**). Oldest‑first (**§6.8 #3**) via **`ScanLBAs`**. Scope of **`CHK-WALSHIPPER-TIMER‑DRAIN`**. | +| **`ModeRealtime`** | **Append‑driven only** — **literal §6.3(B)** idle timer **`MUST NOT`** apply. **Hot path `MUST NOT` substrate‑replay** for shipped WAL bytes (**§13**). | **C3 — BASE ∥ WAL parallel in `Sender.Run`** (§6.8 #6 / P6 / G3): @@ -240,13 +247,13 @@ Choice — **per-connection format dispatch**, not single-format unification. Du |---|---|---|---| | 1 | SINGLE-SERIALIZER | `walShipperEntry.writeMu` + `WalShipper.shipMu`; lock hierarchy `shipMu → writeMu` documented | C1 `294d4bf` | | 2 | ONE TAPE | substrate canonical, single cursor, single shipMu | (P0/P1, retained) | -| 3 | PRIORITY | drain emits LSN-ascending in cursor scan order during Backlog | C2 `53c292f` | -| 4 | TIMER | `WalShipper.timerLoop` goroutine via `IdleSleep` cadence + nudgeCh | C2 `53c292f` | +| 3 | PRIORITY | **`ModeBacklog`** only (**§13**): drain emits LSN-ascending from substrate scan | C2 `53c292f` | +| 4 | TIMER | **`ModeBacklog`**: `timerLoop` + **`drainOpportunity`** — idle drain **`MUST NOT`** be suppressed while **`cursor cursor + 1` — banned regression anchor) | C2 `53c292f` | +| 9 | `send(·,·)` | **`ModeBacklog`**: **`cursor cursor+1`** predicate). **`ModeRealtime`**: **`NotifyAppend`** tail (**T4a**); **`StrictRealtimeOrdering`** + TOCTOU (`377bcb0`) | C2 `53c292f` | ### §11.4 Test anchors (CHK + regression) @@ -257,7 +264,8 @@ Choice — **per-connection format dispatch**, not single-format unification. Du | `TestC1_PostEmitHook_AdvancesShipCursor` | accounting (RecordShipped) | (same file) | | `TestC2_TimerDrainsIdle` | #4 (CHK-WALSHIPPER-TIMER-DRAIN) | `core/transport/c2_timer_drain_test.go` | | `TestC2_PriorityOldFirst` | #3 / #9 | (same file) | -| `TestC2_NoGapDenseLSNEdge` | #9 — **regression anchor**: bans `lsn > cursor + 1` debt check | (same file) | +| `TestC2_NoGapDenseLSNEdge` | #9 / **Backlog** — **regression anchor**: bans **`lsn > cursor + 1`** as sole debt detector; asserts single‑LSN debt via Backlog/timer | (same file) | +| **`StrictRealtimeOrdering`** (ordering guard) | **`§13`** safety switch (`wal_shipper` — default warn / strict opt‑in; see tests in `wal_shipper`/`NotifyAppend`) | **`377bcb0`** | | `TestC3_BaseWalParallel_FramesInterleave` | #6 (CHK-BASE-WAL-OVERLAP via interleaved frame indices) | `core/recovery/c3_base_wal_parallel_test.go` | ### §11.5 Invariants documented in code (architect-review-derived) @@ -265,9 +273,16 @@ Choice — **per-connection format dispatch**, not single-format unification. Du 1. **`DisableTimerDrain` test-only invariant** (`wal_shipper.go` config doc): production MUST leave false; tests using true MUST NOT linger in `cursor < head` without manual drain source. Today's only users: R2 saturation tests, scope-bounded. 2. **Lock hierarchy** (`walShipperEntry.writeMu` doc): `shipMu` always outermost; `writeMu` acquired under it. Reversal risks deadlock. 3. **`RecordShipped` coverage audit** (PR description): four production emit paths verified — bridging streamBacklog/flushAndSeal call inline; WalShipper-routed Backlog drain + Realtime post-R1 fire via C1 postEmit hook. Steady-state `Ship()` without session correctly skips (no session = no shipCursor to advance). +4. **`StrictRealtimeOrdering`** (**`wal_shipper`**) — default **warn** on **`lsn ≠ cursor + 1`** (dense WAL); strict opt‑in (**§13**). **Production `strict=true`**: **pre‑condition §11.6** — engine **must** rebuild / re-anchor on gap first. +5. **`drainOpportunity` atomic read** (**`377bcb0`)**: **`mode` + `cursor` + `head`** under **`shipMu`** **before** cross-checking backlog work — **fixes TOCTOU** vs **`head`/cursor** snapshot. -### §11.6 Open items (carried for hardware run / next phase) +### §11.6 Open items · production migrations (hardware + ops) + +| Item | Action | +|------|--------| +| **§IV T2 (`targetLSN` vs moving `head` / barrier)** | Architect-open at consensus **§IV**; hardware run validates **receiver `achievedLSN` ⇄ coordinator `targetLSN`** at C3 barrier. | +| **§IV T7 (R2 sustained pressure)** | Single-shot **`OnSaturation`** is structural only; throttle / escalation remains **engine** work. | +| **Engine rebuild-on-gap** (before `StrictRealtimeOrdering=true`) | If **Ship→WalShipper** sees Realtime **`NotifyAppend`** violations (out‑of‑order / **`lsn≠cursor+1`** under strict policy, or **`cursor cursor + 1` as sole debt test** — use **`cursor < head`** (**Backlog**); (**b**) **Realtime never substrate-scan for ship** (**T4a**, **§13**). | -- **§IV T2 (barrier vs frozen targetLSN)**: still architect-open. C3 barrier writes after BASE ∥ WAL goroutines complete; receiver computes achievedLSN; coord.CanEmitSessionComplete checks `achieved >= target`. With WalShipper-resident gap→0 semantics, alignment between coord's session-target and receiver's frontier is a hardware-run concern. -- **§IV T7 (R2 saturation under sustained pressure)**: hook is wired (single-shot), production-grade throttle/escalation is engine concern. Hardware run with sustained `append > ship` rate is the validation venue. -- **PR template line** (suggested by reviewer): for future PRs touching `WalShipper.NotifyAppend`, add a checklist item: "**禁止 `lsn > cursor + 1` 当 debt 判定 — 用 `cursor < head`**" (the banned-check regression anchor). diff --git a/sw-block/design/v3-recovery-wal-shipper-spec.md b/sw-block/design/v3-recovery-wal-shipper-spec.md new file mode 100644 index 000000000..ee470d07b --- /dev/null +++ b/sw-block/design/v3-recovery-wal-shipper-spec.md @@ -0,0 +1,197 @@ +# WalShipper — algorithm spec (**implementable**) + +**Status**: Normative leaf spec (**Primary-side**, per-peer Wal replication scheduling). Extends **`v3-recovery-algorithm-consensus.md` §6** without contradicting **§I P1–P7**. + +**Audience**: Engineers implementing or reviewing `WalShipper` / catch-up sender / dual-lane Wal path. + +**Goal**: Write **this** algorithm once — reviewers match code to headings **§2–§7**; no parallel “mystery shipper” in `core/recovery/sender` unless architect supersedes this doc. + +**Consensus checklist (**`seaweedfs`** / **`seaweed_block`** handoff**)**: **`v3-recovery-algorithm-consensus.md` §6.8** — **nine MUST** bullets (**item 9** = **`send` strategy**); **implementations SHOULD map each item to tests** (below §9). **Scope**: §6.8 **(3)(4)(9)** **as literal idle / debt-priority** apply **`ModeBacklog`**; **`ModeRealtime`** obey **`§13` — E‑WALSHIPPER‑DUAL‑MODE** (`**v3-recovery-wal-shipper-mini-plan.md` §11.2a**). + +--- + +## 1. Scope + +| In scope | Out of scope (link elsewhere) | +|----------|--------------------------------| +| **One ordered Wal tape** on Primary: `head`, per-peer **`cursor[P]`**, **emit** order to **P** | **Extent/base bulk** — `v3-recovery-algorithm-consensus.md` §5; bytes `v3-recovery-pin-floor-wire.md` | +| **Single serializer** for “next LSN sent to P” (**P2**) | Receiver **`checkMonotonic`**, **`bitmap`** — consensus §7; replica apply | +| **R1** transition (caught-up vs append race) | **Engine** retry / escalate after fail-close — executor / engine | +| **R2** saturation observability hooks | **`MinPin`/recycle** full story — coordinator + substrate | + +--- + +## 2. State (**per replica `P`**) + +Stored in **one owning struct** (“`WalShipper`” here); names are logical. + +| Symbol | Meaning | +|--------|---------| +| **`head`** | Next LSN the Primary Wal will assign on append (**exclusive upper bound** of assigned LSNs: valid entries satisfy `pin < lsn ≤ head` per product pinning, see **§3**). **Single source** — bumps only inside **Wal append** path owned by same package or **serialized** through **§6.1**. | +| **`cursor`** | Next LSN **to send** to peer **P** (inclusive lower bound for unsent prefix). **`cursor ∈ [fromLSN, head]`** once session valid. (**Shorthand**: `cursor` means `cursor[P]`.) | +| **`fromLSN` / `pin`** | Session lower bound from engine/coordinator (**P7**). **`cursor` initialized ≤ `fromLSN` semantics** as product defines (usually `cursor := fromLSN` before first emit). | +| **`shipMu`** | **One mutex** (or FIFO channel) guarding **§4–§6** — **`head`/`cursor` decision + emit claim** (**R1**, **§6.1 consensus**). | +| **`lastEmit`** | Optional diagnostic: last LSN **handed to transport** for P (for metrics). | + +**Vocabulary (consensus **§6.2–6.3**)** — **normative mapping**: + +| Term | Meaning | +|------|---------| +| **`debt` non-empty** | **`cursor < head`** — **oldest unsent** is always **`cursor`** on the **single** Primary Wal order. | +| **`debt` empty** | **`cursor == head`** — **no unsent prefix**; next emit is **tail** / **R1**-safe new assign. | +| **`incoming`** | **New tail** Wal LSN made visible to the serializer **this `ShipOpportunity`** (typically from **append**), **not** allowed to **skip ahead** of **`cursor..head−1`**. | +| **`send(incoming, debt)`** | **Strategy choice**: if **`debt` empty** → **emit `incoming`** (steady **debt-free** path); if **`debt` non-empty** → **emit `cursor` first**, **`incoming` waits** on the same tape. | +| **`send(∅, debt)`** | **Timer / auxiliary loop** opportunity: **no append paired to this tick** (**`∅`**) but **still drain** **`cursor < head`** — consensus **§6.3(B)** / checklist **§6.8 item 4**. | + +**Policy note (informative)**: **Upper layers** open recover **policy debt** via **`fromLSN`** / extent-anchored admission (**P7**). **Closing** that hole is **gap → 0** on **`cursor→head`**, not “freeze **`head`’s coordinate**” — **consensus §6.2**, **`T2`**. + +**Forbidden**: second concurrent “recover pump” and “steady pump” **both** calling `Emit` for the same **`P`** without **`shipMu`** — **violates INV-SINGLE**. + +--- + +## 3. Invariants (**must hold between observable steps**) + +**INV-SINGLE** — At most **one** execution context at a time runs **`PlanNextEmit` → `EmitToTransport(P, lsn)`** for a given **`P`**. + +**INV-MONOTONIC-CURSOR** — **`cursor` never decreases** unless session is torn down / re-seeded (**P7** engine reset). + +**INV-SUBSET** — **Every emitted LSN** satisfies **`pin < lsn ≤ head`** (or product’s closed/open interval choice, but **consistent** with receiver `expected`). + +**INV-NO-GAP-R1** — No LSN **`k`** with **`cursor ≤ k ≤ head`** is **skipped** across the **catch-up → tail** seam: **§5**. + +**INV-NO-DOUBLE-LIVE** — No duplicate emit of same logical **`lsn`** on two live outbound paths for **P** — matches **CHK-WALSHIPPER-SINGLE-CURSOR**. + +--- + +## 4. Core priority rule (**default scheduling**) + +On each **`ShipOpportunity`** (defined as: **after Wal append**, or **timer tick**, or **transport ready** callback — all must **serialize** through **`shipMu`**): + +### 4.1 Strategy summary (**must match consensus §6.3**) + +- **`send(incoming, debt)`**: + - **`debt` empty** (`cursor == head`): **emit `incoming`** (new tail) when this opportunity includes one — **debt-free steady send**. + - **`debt` non-empty** (`cursor < head`): **emit `cursor`** (oldest unsent). **`incoming` does not overtake** — it **extends** the tape behind the unsent prefix. +- **`send(∅, debt)`**: **Timer-only** (or loop) opportunity — even if **no `incoming`** this tick, if **`cursor < head`**, **PlanNextEmit** follows **(1)** below until budget / blocked. + +### 4.2 PlanNextEmit steps + +1. If **`cursor < head`**: **PlanNextEmit** returns **`cursor`** (the **oldest unsent** LSN). After successful handoff to transport (or committed send slot), **`cursor ← cursor + 1`** (or next present LSN if sparse — product must define; **binary Wal is dense LSN**). +2. Else **`cursor == head`**: **no debt**. If this opportunity was triggered by **a new append** that just set **`head = h`**: **PlanNextEmit** returns **`h`** (the new tail) **once** (steady **append → ship one**). +3. **Timer-only** opportunity with **`cursor < head`**: same as (1) — **`send(∅, debt)`**; **drain without requiring new appends** (consensus §6.3(B)). + +**Wrong model (non‑compliant)**: alternating “K live + M backlog” as **two mailboxes**. **Right model**: single tape, **`cursor` priority** before tail within one serializer. + +--- + +## 5. **R1** — Caught-up / transition (**mutex + double-check**) + +**Problem**: Between “observed **`cursor == head`**” and “declare tail-push only”, **`head` grows** (`head = h+1`); tail-push path never sees **`h+1`** → **gap**. + +**Procedure `AssertCaughtUpAndEnableTailShip`** (runs **inside `shipMu`**): + +1. Read **`h ← head`**. +2. If **`cursor < h`**, return **`Backlogged`** (caller runs §4 loop). +3. **Double-check**: read **`h2 ← head`** again under same lock. If **`h2 ≠ h`** or **`cursor < h2`**, go to step 1 or §4 as appropriate. +4. Only now may implementation set internal flag **`tailPushEligible := true`** (if it uses one) — **optional**; many implementations **need no flag** if **§4** always applies: when **`cursor == head`** and append occurs, next **`ShipOpportunity`** ships **exactly** the new tail. + +**Append path** (**Wal append**) **MUST**: + +- Bump **`head`** **only** while holding **`shipMu`**, **or** +- Bump **`head`** then **immediately** call **`ShipOpportunity`** under **`shipMu`**, **or** +- Use a **single threaded** Wal append + shipper goroutine communicating by ordered messages. + +Any **ordering** works iff **INV-NO-GAP-R1** is provable. + +--- + +## 6. Saturation (**R2**) — bounded progress (**not algorithm beauty**) + +**Observation** (**outside hot path mutex if needed**): + +- **`lag := head − cursor`** (dense LSN **or product-defined synthetic distance**). + +**Threshold policy** (**product**, consensus §6.6): if **`lag`**, **time‑above‑T**, or **pin pressure** satisfies fail condition: + +- **`SignalSaturation`** → executor ends session **`FailureTimeout`/product kind** **and/or** triggers **Wal admission throttle** on Primary — **never** silently infinite spin. + +Timer **priority drain** (**§4**) remains required for **liveness** but **does not** disprove **R2** math. + +--- + +## 7. Pseudocode (**reference shape**) + +```text +// One goroutine OR all entrypoints take shipMu. + +function OnWalAppend(assigns_new_lsn): + lock(shipMu) + head++ + + EmitAccordingTo§4UntilBudgetOrBlocked() // drains cursor..head-1 before new tail rules + + // Typically: while cursor < head: emit(cursor); cursor++ + // then if new tail == head prev+1 logic: emit that tail once + + unlock(shipMu) + +function OnShipTimer(): + lock(shipMu) + EmitAccordingTo§4UntilBudgetOrBlocked() + unlock(shipMu) + +function EmitAccordingTo§4UntilBudgetOrBlocked(): + while transport_has_credit() and cursor < head: + emit(cursor) + cursor++ + while transport_has_credit() and cursor == head and pending_tail_ship_from_append: + emit(head); clear pending_tail_flag // depends on refactor; simplest: append path calls Emit once after advancing head while cursor tracked +``` + +Implementations **flatten** this; the **semantic** contract is **§3–§5**, not the sketch’s loop details. + +--- + +## 8. Interface boundaries (**who owns what**) + +| Concern | Owner | +|---------|-------| +| **`fromLSN` / pin / SessionID** | Engine + **`PeerShipCoordinator`** | +| **`head` mutation** | Primary Wal append path (**substrate**/BlockVol — **serialized** vs shipper §5) | +| **`cursor` mutation** | **WalShipper** only (per **P**) | +| **`EmitToTransport`** | WalShipper calls into transport (SWRP / `frameWALEntry` / dual-lane) — **one caller** per **P** | + +--- + +## 9. Tests (**minimum**) + +| ID | Intent | +|----|--------| +| **CHK-WALSHIPPER-SINGLE-CURSOR** | consensus **§V** / **§6.8 item 1** — no duplicate live paths | +| **CHK-WALSHIPPER-TIMER-DRAIN** | consensus **§V / §6.8 item 4** | +| **Test-R1-Flip-NoGap** | After synthetic `Scan`/drain at `H0`, concurrent append `H0+1` before unlock — **receiver** must see **`H0+1`** (integration) | +| **Test-Priority-OldFirst** | With `cursor < head` and new append, **first emit** is **`cursor`**, not new tail — **`send(incoming, debt)`** with **`debt` non-empty** | +| **Test-Timer-Drains-Idle** | No append; timer moves **`cursor`** toward **`head`** — **`send(∅, debt)`** | + +--- + +## 10. Revision + +| Date | Change | +|------|--------| +| 2026-04-29 | Initial **WalShipper** implementable spec (priority queue on one tape, **R1** double-check, **R2** hook, boundaries). | +| **2026-04-30** | Pointer to consensus **§6.8** implementer checklist; **§V** **`CHK-WALSHIPPER-TIMER-DRAIN`** cross-ref in **§9**. | +| **2026-04-30** | **§2** vocabulary **`incoming` / `debt` / `send(·,·)`**; **§4.1–4.2** strategy split; tests cross-ref **§6.3**; checklist **nine** items (**§6.8(9)**). | +| **2026-04-30** | **§1** intro + **consensus §6.3(B)** cross-ref — **`ModeBacklog` vs `ModeRealtime`** (**§13** E‑WALSHIPPER‑DUAL‑MODE); mini-plan **§11.2a** bridge | + +--- + +## 11. Document map + +| Doc | Role | +|-----|------| +| **`v3-recovery-algorithm-consensus.md` §6 + §6.8** | **Axioms + stability + implementer MUST checklist** | +| **This** | **Algorithm steps + invariants for one implementation** | +| **`v3-recovery-wal-shipper-mini-plan.md`** | **Phased PR / file anchoring (**`seaweed_block`**) **↔** §INV tests** | +| `v3-recovery-wiring-plan.md` | Bearer / port | +| `v3-recovery-unified-wal-stream-*.md` | Historical kickoff — **align or supersede** with this spec |