sw-block/design: §13 E-WALSHIPPER-DUAL-MODE + mini-plan §11.2a + spec sync

Architect-approved exception (consensus v3.9, 2026-04-30): WalShipper
dual-mode contract carves Realtime out of §6.3(B) timer + §6.8(3)(9)
send(incoming, debt) normative scope. Backlog mode is fully normative;
Realtime is per-append (T4a no-replay), with StrictRealtimeOrdering
as the production safety switch.

Files:

  v3-recovery-algorithm-consensus.md
    - NEW §13: Architect-approved exceptions, with E-WALSHIPPER-
      DUAL-MODE entry. Old §13 Revision → §14; old §14 Document
      map → §15. Cross-refs in §I, §6.3(A)(B), §6.8 checklist,
      §V §12 (CHK-WALSHIPPER-TIMER-DRAIN explicitly limited to
      ModeBacklog), §15 (index gains §13).
    - §14 revision row v3.9.

  v3-recovery-wal-shipper-mini-plan.md
    - §11.1 G-TIMER row updated to "C2 + dual-mode/§13".
    - §11.2 replaced (was: pre-merge "Realtime drains on cursor<head"
      bullets) with C2 dual-mode contract + §11.2a normative table:
        * Backlog: §6.3 / §6.8(3)(4)(9) fully apply (timer, scan,
          oldest-first, send(·, debt)).
        * Realtime: NotifyAppend per-LSN-in-order; lsn==cursor+1
          invariant; no substrate replay (§13 carve-out).
        * Production safety switch: StrictRealtimeOrdering.
    - §11.3 compliance receipt rebuilt as 1–9 ordered table with
      commit SHAs (P2d cb8ff1c, C1 294d4bf, C2 53c292f, C3 40f2935,
      review-fix 377bcb0).
    - §11.4 test anchor table gains StrictRealtimeOrdering opt-in.
    - §11.5 invariants gain (4) §13 dual-mode + (5) replicaID drift.
    - §11.6 reshaped as production-migration table:
        * Engine drives rebuild-on-gap before flipping
          StrictRealtimeOrdering=true.
        * Fix replicaID drift in startRebuildDualLane (sink uses
          engine replicaID; bridge.coord uses dl.ReplicaID).
        * Hardware-run carries: §IV T2 (barrier vs targetLSN),
          §IV T7 (R2 saturation under sustained pressure).
        * PR template line: ban `lsn > cursor + 1` debt detection
          (banned regression anchor; use `cursor < head` if it
          ever needs to come back, but only inside Backlog).
    - §8 revision: v0.6.

  v3-recovery-wal-shipper-spec.md
    - Goal/checklist front-matter cross-refs §13 + mini-plan §11.2a.
    - §10 revision row aligning with consensus v3.9.

Companion implementation on seaweed_block g7-redo/wal-shipper-impl:
  - 377bcb0 review-fix: T4a Realtime sequence guard + drainOpportunity TOCTOU
  - cb17338 docs: drainOpportunity comment matches §13 dual-mode contract
This commit is contained in:
pingqiu committed 2026-04-30 09:48:10 -07:00
1 parent aeeb71f0e7
commit 91a313e1e1
3 files changed
+535 -17

No files matched your search

@@ -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<head` predicate** clears — **protects pin/retention** at latency cost. |
**None** of these create a **second Wal**; they reshape **who runs when inside the same entity**.
#### 6.6 Saturation (**R2**) — **normative fail-close**
If **sustained append rate > 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** |
@@ -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<head`** (**§13**) | C2 `53c292f` |
| 5 | FRAMING ≠ SECOND TAPE | `EmitProfile` (Steady / DualLane) selects encoder per conn | P2d `cb8ff1c` |
| 6 | BASE ∥ WAL | `Sender.Run` two-goroutine errgroup + shared `writeMu` | C3 `40f2935` |
| 7 | R1 | `assertCaughtUpAndEnableTailShipLocked` under shipMu + double-check | (P0, retained) |
| 8 | R2 | `OnSaturation` single-shot per session | (P0, retained) |
| 9 | `send(·,·)` | NotifyAppend uses **`cursor < head`** debt condition (NOT `lsn > cursor + 1` — banned regression anchor) | C2 `53c292f` |
| 9 | `send(·,·)` | **`ModeBacklog`**: **`cursor<head`** + priority + **`send(∅, debt)`** (no banned **`lsn > 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<head`** stuck in wrong mode): engine **MUST** drive **rebuild** or session reopen to **re‑anchor **`cursor`** (§13)** before tightening production ordering. |
| **startRebuildDualLane `replicaID` consistency** | **Caller obligation**: **RecoverySink** / WalShipper key (**engine **`replicaID`**)** **MUST** match **`bridge.coord`** / dual‑lane **`dl.ReplicaID`** used for **`RecordShipped`** and cursor polls. Drift surfaced during C2 debugging — **bridge‑side cleanup** (same ID end‑to‑end). |
| **PR template guard** | For PRs touching **`WalShipper.NotifyAppend`**: (**a**) **Ban `lsn > 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).
@@ -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 |