fix(localdb): snapshot a wiped peer whose deltas were already acked and pruned

A converged pair prunes its oplog on ack, so "no oplog rows at all" is the
steady state of a healthy pair — not an anomaly. ComputeSnapshotRequiredAsync
measured the peer's gap against the oldest surviving oplog row and treated an
empty oplog as "no gap possible", which made that steady state the one state
from which a wiped peer could never be healed: it reconnected with an empty
database and a zero watermark, was told it needed no snapshot, and stayed
empty indefinitely.

With an empty oplog the oldest seq we could still stream is one past
last_acked_seq — ack-pruning is the only thing that empties the oplog without
also flagging needs_snapshot, and it deletes exactly seq <= last_acked_seq.
Measuring the gap against that heals the wiped peer and leaves every other
case identical: a converged peer's watermark equals last_acked_seq (no
snapshot), and a pair that has never written anything has both at zero.

Found by the OtOpcUa LocalDb Phase 1 live gate (check 4), which recorded it as
a documented limitation. Covered by a test that reproduces the live-gate
scenario through the real session stack: RED before this fix (the wiped side
never converges, 15 s timeout), green after. 148/148 pass.

Version 0.1.2.

Claude-Session: https://claude.ai/code/session_01GASWkNEi68FSCtvr6rLoEW
This commit is contained in:
Joseph Doherty
2026-07-21 01:17:32 -04:00
parent c42348eff5
commit cad3bcb6bf
4 changed files with 66 additions and 6 deletions
+1 -1
View File
@@ -5,7 +5,7 @@
<Nullable>enable</Nullable>
<ImplicitUsings>enable</ImplicitUsings>
<LangVersion>latest</LangVersion>
<Version>0.1.1</Version>
<Version>0.1.2</Version>
<ManagePackageVersionsCentrally>true</ManagePackageVersionsCentrally>
<TreatWarningsAsErrors>true</TreatWarningsAsErrors>
<!-- NU1510 (package-pruning advice): the Replication and Tests projects carry a
+9 -2
View File
@@ -10,7 +10,7 @@ Change capture uses SQLite triggers writing a single HLC-stamped oplog inside th
consumer's own write transaction (crash-consistent, non-bypassable). Convergence is
last-writer-wins over a hybrid logical clock. All timestamps are UTC.
## Packages (all `0.1.0`)
## Packages (all `0.1.2`)
| Package | Contents |
|---|---|
@@ -18,7 +18,14 @@ last-writer-wins over a hybrid logical clock. All timestamps are UTC.
| `ZB.MOM.WW.LocalDb.Contracts` | `localdb_sync.v1` proto + generated gRPC types (packable; a future peer could be non-.NET). |
| `ZB.MOM.WW.LocalDb.Replication` | gRPC sync engine: passive `LocalDbSyncService`, initiating `SyncBackgroundService`, `MaintenanceBackgroundService`, LWW apply, watermark acks + pruning, snapshot resync, `ISyncStatus`, `LocalDbMetrics`, `AddZbLocalDbReplication()` / `MapZbLocalDbSync()`. |
> Status: **built and published to the Gitea feed at `0.1.0` (2026-07-18); no app has adopted it yet.**
> Status: **published to the Gitea feed; adopted by ScadaBridge (Phases 1+2) and OtOpcUa (Phase 1).**
>
> - `0.1.2` (2026-07-21) — fixes snapshot resync for a **wiped peer**: a converged pair prunes its
> oplog to empty on ack, and the handshake read an empty oplog as "no gap possible", so a peer
> that came back with an empty database and a zero watermark was never snapshotted. The gap is
> now measured against `last_acked_seq` when the oplog is empty.
> - `0.1.1` (2026-07-20) — creates the parent directory of `LocalDb:Path` (a missing directory was
> a hard boot failure, `SQLite Error 14`).
## Quickstart
@@ -418,9 +418,20 @@ internal sealed class SyncSession
var minRows = await _db.QueryAsync(
"SELECT COALESCE(MIN(seq), 0) FROM __localdb_oplog", static r => r.GetInt64(0), null, ct);
var minSeq = minRows[0];
// minSeq 0 => oplog empty. A gap exists when the peer's applied watermark falls below the
// oldest seq we can still stream (pruning discarded everything at or below its horizon).
return minSeq > 0 && peerHandshake.LastAppliedRemoteSeq < minSeq - 1;
// The oldest seq we could still stream as a delta. With rows in the oplog that is simply the
// oldest one; with an EMPTY oplog it is the next seq we will ever write — one past
// last_acked_seq, because ack-pruning is the only thing that empties the oplog without
// flagging needs_snapshot, and it deletes exactly seq <= last_acked_seq.
//
// Treating an empty oplog as "no gap possible" is what let a wiped peer starve: a converged
// pair prunes to empty on ack, so the healthy node held no delta at all, and a peer that
// came back with a zero watermark and an empty database was told it needed nothing. The
// steady state of a healthy pair was precisely the state that could not heal one.
var streamableFrom = minSeq > 0 ? minSeq : peerState.LastAckedSeq + 1;
// A gap exists when the peer's applied watermark falls below what we can still stream.
return peerHandshake.LastAppliedRemoteSeq < streamableFrom - 1;
}
private static async Task<SyncMessage> ReceiveOneAsync(SyncDuplex duplex, CancellationToken ct)
@@ -146,6 +146,12 @@ public sealed class SnapshotResyncTests : IDisposable
return r[0];
}
private static async Task<long> OplogCount(SqliteLocalDb db)
{
var r = await db.QueryAsync("SELECT COUNT(*) FROM __localdb_oplog", x => x.GetInt64(0));
return r[0];
}
private static SnapshotRow SnapRow(long id, string sku, long hlc, string nodeId = "peer") =>
new()
{
@@ -156,6 +162,42 @@ public sealed class SnapshotResyncTests : IDisposable
IsTombstone = false,
};
[Fact]
public async Task WipedPeer_AfterEverythingWasAckedAndPruned_StillGetsSnapshot()
{
// The gap the OtOpcUa Phase 1 live gate found (check 4). A converged pair prunes its oplog
// on ack, so the steady state is "no oplog at all". If the peer is then wiped — a rebuilt
// node, a lost volume — it comes back with an empty database and a zero watermark, and the
// healthy side has no delta left to send it. Snapshot detection compared the peer's
// watermark against the oldest oplog row, which does not exist, so it concluded no snapshot
// was needed and the wiped node stayed empty indefinitely: the one case where replication
// most obviously should heal a node was the case it could not.
var a = await NewSideAsync();
for (var i = 1; i <= 3; i++)
await a.Db.ExecuteAsync("INSERT INTO orders (id, sku, qty) VALUES (@id, 'A', @id)", new { id = i });
await a.Store.RecordPeerAckAsync(await a.Store.GetMaxSeqAsync());
Assert.Equal(0, await OplogCount(a.Db));
// Not the already-covered cap-exceeded path: nothing here flagged a resync.
Assert.Equal(0, await NeedsSnapshot(a.Db));
// A brand-new database is exactly what a wiped node is: no rows, watermark zero.
var b = await NewSideAsync();
Assert.Empty(await ReadOrders(b.Db));
var (da, db) = DuplexPair();
using var cts = new CancellationTokenSource(RunTimeout);
var runA = a.Session.RunAsync(da, cts.Token);
var runB = b.Session.RunAsync(db, cts.Token);
await WaitForAsync(async () => (await ReadOrders(b.Db)).Count == 3, RunTimeout);
Assert.Equal(await ReadOrders(a.Db), await ReadOrders(b.Db));
cts.Cancel();
await SwallowAsync(runA);
await SwallowAsync(runB);
}
[Fact]
public async Task PrunedPeer_GetsSnapshot_ThenConverges()
{