diff --git a/docs/plans/2026-07-17-localdb-implementation.md b/docs/plans/2026-07-17-localdb-implementation.md new file mode 100644 index 0000000..e515734 --- /dev/null +++ b/docs/plans/2026-07-17-localdb-implementation.md @@ -0,0 +1,511 @@ +# ZB.MOM.WW.LocalDb Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers-extended-cc:executing-plans to implement this plan task-by-task. + +**Goal:** Build `ZB.MOM.WW.LocalDb` — an embedded SQLite local cache library with optional bidirectional async 2-node replication over gRPC (triggers+oplog CDC, HLC last-writer-wins, all-UTC). + +**Architecture:** Three packages in a plain scadaproj directory (NOT a nested git repo — commits go to scadaproj `main`): core (`ZB.MOM.WW.LocalDb`: DB host, consumer SQL API, table registration, capture triggers, HLC), replication (`ZB.MOM.WW.LocalDb.Replication`: gRPC sync engine, LWW apply, watermarks, snapshot resync), contracts (`ZB.MOM.WW.LocalDb.Contracts`: `localdb_sync.v1` proto). Design of record: `docs/plans/2026-07-17-localdb-design.md` — read it first; it is the spec. + +**Tech Stack:** .NET 10, Microsoft.Data.Sqlite 10.0.7 (WAL mode), Grpc (Google.Protobuf 3.34.1 / Grpc.Core.Api + Grpc.Tools 2.76.0 in Contracts, Grpc.AspNetCore + Grpc.Net.Client in Replication), xunit 2.9.3. Conventions mirror `ZB.MOM.WW.Secrets/` exactly (Directory.Build.props / Directory.Packages.props / .slnx / src+tests / `IsPackable` metadata / `InternalsVisibleTo` tests). + +**Key design invariants (do not violate):** +1. All timestamps UTC. HLC = 64-bit `(utcUnixMs << 16) | counter`. No local time anywhere. +2. Capture is transactional: triggers write oplog + row-version in the SAME consumer transaction. Never capture outside the transaction. +3. Replication apply NEVER re-fires capture (the `__localdb_applying` guard) and NEVER re-forwards applied entries (2-node topology: no oplog write on apply). +4. LWW compare key is `(hlc, node_id)` lexicographic — node_id tie-break makes both nodes converge deterministically. +5. Local writes never block on the peer. Everything replication-side is async + batched. +6. Fail closed: schema-digest mismatch at handshake refuses sync; PK-less registration throws at startup; a connection without the `zb_hlc_next` UDF fails writes to registered tables (this is intentional — capture cannot be silently skipped). + +**Documented limitations (state in README, don't "fix"):** updating PK columns of a replicated row is unsupported (capture treats it as an update under the NEW pk; the OLD pk row is not tombstoned); exactly two nodes; LWW only. + +--- + +### Task 1: Scaffold the library + +**Classification:** small +**Estimated implement time:** ~5 min +**Parallelizable with:** none + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/Directory.Build.props` +- Create: `ZB.MOM.WW.LocalDb/Directory.Packages.props` +- Create: `ZB.MOM.WW.LocalDb/ZB.MOM.WW.LocalDb.slnx` +- Create: `ZB.MOM.WW.LocalDb/README.md` (stub: name + one-paragraph purpose + link to design doc) +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/ZB.MOM.WW.LocalDb.csproj` +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Contracts/ZB.MOM.WW.LocalDb.Contracts.csproj` +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/ZB.MOM.WW.LocalDb.Replication.csproj` +- Create: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/ZB.MOM.WW.LocalDb.Tests.csproj` + +**Step 1: Copy conventions from Secrets.** `Directory.Build.props` = same as `ZB.MOM.WW.Secrets/Directory.Build.props` but `0.1.0` and `true` added to the PropertyGroup. `Directory.Packages.props`: + +```xml + + + true + + + + + + + + + + + + + + + + + + + + + + + + +``` + +Core csproj: mirror Secrets' packable metadata (`PackageId ZB.MOM.WW.LocalDb`, Description "Embedded SQLite local cache with optional bidirectional async 2-node gRPC replication for the ZB.MOM.WW SCADA family.", Project/Repository URL `https://gitea.dohertylan.com/dohertj2/zb-mom-ww-localdb`), refs: Microsoft.Data.Sqlite, Options(+Config), DI.Abstractions, Configuration.Abstractions, Hosting.Abstractions, Logging.Abstractions; `InternalsVisibleTo ZB.MOM.WW.LocalDb.Tests`. Contracts csproj: copy HistorianGateway Contracts pattern — ``, Google.Protobuf + Grpc.Core.Api + Grpc.Tools(PrivateAssets=all), packable, AnyCPU note not needed (no repo x64 contract here). Add a placeholder `Protos/localdb_sync.proto` (header only, filled in Task 7). Replication csproj: packable, ProjectRefs to core + Contracts, PackageRefs Grpc.AspNetCore, Grpc.Net.Client, Hosting.Abstractions, Logging.Abstractions; `InternalsVisibleTo ZB.MOM.WW.LocalDb.Tests`. Tests csproj: mirror Secrets tests csproj + `Microsoft.AspNetCore.TestHost`, ProjectRefs to all three src projects; add ``. + +`.slnx`: same shape as Secrets' (`/src/` folder with 3 projects, `/tests/` with 1). + +**Step 2: Build.** Run from `ZB.MOM.WW.LocalDb/`: `dotnet build ZB.MOM.WW.LocalDb.slnx` → expect 0 warnings 0 errors. `dotnet test` → 0 tests, green. + +**Step 3: Commit.** `git add ZB.MOM.WW.LocalDb && git commit -m "feat(localdb): scaffold ZB.MOM.WW.LocalDb (core/contracts/replication + tests)"` + +--- + +### Task 2: HybridLogicalClock + +**Classification:** small +**Estimated implement time:** ~4 min +**Parallelizable with:** Task 7 + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Hlc/HybridLogicalClock.cs` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/HybridLogicalClockTests.cs` + +**Step 1: Write failing tests** — `Next_IsStrictlyMonotonic` (1000 calls, each > previous), `Next_UsesUtcPhysicalMs` (decode top 48 bits ≈ injected now), `Next_SameMs_IncrementsCounter`, `Observe_RemoteAhead_AdvancesPastRemote`, `Observe_RemoteBehind_NoRegression`, `Resume_FromPersistedValue_NeverRegresses` (construct with initial=persisted high-water > wall clock → Next() still > persisted), `CounterOverflow_RollsPhysicalForward` (counter forced to 0xFFFF → next physical+1, counter 0). + +**Step 2: Run** `dotnet test --filter "FullyQualifiedName~HybridLogicalClockTests"` → FAIL (type missing). + +**Step 3: Implement:** + +```csharp +namespace ZB.MOM.WW.LocalDb.Hlc; + +/// 64-bit HLC: (UTC unix-ms << 16) | 16-bit logical counter. Thread-safe. +public sealed class HybridLogicalClock +{ + private readonly Func _utcNowMs; + private readonly Lock _lock = new(); + private long _last; + + public HybridLogicalClock(long initialValue = 0, Func? utcNowMs = null) + { + _utcNowMs = utcNowMs ?? (() => DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()); + _last = initialValue; + } + + public long Next() + { + lock (_lock) + { + var candidate = _utcNowMs() << 16; + _last = candidate > _last ? candidate : _last + 1; // +1 rolls counter; counter overflow naturally carries into physical + return _last; + } + } + + /// Merge a peer's HLC so our next stamp orders after everything we've seen. + public void Observe(long remote) { lock (_lock) { if (remote > _last) _last = remote; } } + + public long Current { get { lock (_lock) return _last; } } + + public static long PhysicalMs(long hlc) => hlc >> 16; + public static int Counter(long hlc) => (int)(hlc & 0xFFFF); +} +``` + +(Note `_last + 1` handles both same-ms increment AND counter overflow — the carry propagates into the physical bits, which is exactly "roll physical forward". Verify the overflow test agrees.) + +**Step 4: Run tests** → PASS. **Step 5: Commit** `feat(localdb): hybrid logical clock (UTC-ms<<16 | counter)`. + +--- + +### Task 3: Lib-owned schema + LocalDbSchema migrator + +**Classification:** standard +**Estimated implement time:** ~4 min +**Parallelizable with:** Task 2, Task 7 + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/LocalDbSchema.cs` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/LocalDbSchemaTests.cs` + +**Step 1: Failing tests** — open a temp-file `SqliteConnection`, call `LocalDbSchema.EnsureCreated(conn)`: asserts all five objects exist (`__localdb_meta`, `__localdb_oplog`, `__localdb_peer_state`, `__localdb_row_version`, `__localdb_applying`); `EnsureCreated_IsIdempotent` (call twice); `Meta_NodeId_MintedOnceAndStable` (reopen → same GUID); `Applying_SeededWithZeroRow`. + +**Step 2:** Run → FAIL. + +**Step 3: Implement.** `internal static class LocalDbSchema` with `EnsureCreated(SqliteConnection)` executing (one transaction): + +```sql +CREATE TABLE IF NOT EXISTS __localdb_meta ( + id INTEGER PRIMARY KEY CHECK (id = 1), + node_id TEXT NOT NULL, + hlc_high_water INTEGER NOT NULL DEFAULT 0, + schema_version INTEGER NOT NULL +); +CREATE TABLE IF NOT EXISTS __localdb_oplog ( + seq INTEGER PRIMARY KEY AUTOINCREMENT, + table_name TEXT NOT NULL, + pk_json TEXT NOT NULL, + row_json TEXT NULL, + hlc INTEGER NOT NULL, + node_id TEXT NOT NULL, + is_tombstone INTEGER NOT NULL DEFAULT 0 +); +CREATE INDEX IF NOT EXISTS __localdb_oplog_hlc ON __localdb_oplog(table_name, pk_json, hlc); +CREATE TABLE IF NOT EXISTS __localdb_peer_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + peer_node_id TEXT NULL, + last_acked_seq INTEGER NOT NULL DEFAULT 0, -- highest OUR seq the peer has applied (drives pruning) + last_applied_remote_seq INTEGER NOT NULL DEFAULT 0, -- highest PEER seq we have applied (our resume point) + last_seen_hlc INTEGER NOT NULL DEFAULT 0, + needs_snapshot INTEGER NOT NULL DEFAULT 0, + last_sync_utc TEXT NULL +); +CREATE TABLE IF NOT EXISTS __localdb_row_version ( + table_name TEXT NOT NULL, + pk_json TEXT NOT NULL, + hlc INTEGER NOT NULL, + node_id TEXT NOT NULL, + is_tombstone INTEGER NOT NULL DEFAULT 0, + tombstone_utc TEXT NULL, + PRIMARY KEY (table_name, pk_json) +) WITHOUT ROWID; +CREATE TABLE IF NOT EXISTS __localdb_applying ( + id INTEGER PRIMARY KEY CHECK (id = 1), + applying INTEGER NOT NULL DEFAULT 0 +); +CREATE TABLE IF NOT EXISTS __localdb_dead_letter ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + received_utc TEXT NOT NULL, + table_name TEXT NOT NULL, + pk_json TEXT NOT NULL, + row_json TEXT NULL, + hlc INTEGER NOT NULL, + node_id TEXT NOT NULL, + error TEXT NOT NULL +); +INSERT INTO __localdb_meta (id, node_id, schema_version) + SELECT 1, , 1 WHERE NOT EXISTS (SELECT 1 FROM __localdb_meta); +INSERT INTO __localdb_applying (id, applying) SELECT 1, 0 WHERE NOT EXISTS (SELECT 1 FROM __localdb_applying); +INSERT INTO __localdb_peer_state (id) SELECT 1 WHERE NOT EXISTS (SELECT 1 FROM __localdb_peer_state); +``` + +Also expose `internal static (string NodeId, long HlcHighWater) ReadMeta(conn)` and `internal static void FlushHlcHighWater(conn, long value)`. + +**Step 4:** tests PASS. **Step 5: Commit** `feat(localdb): lib-owned schema (meta/oplog/peer_state/row_version/applying/dead_letter)`. + +--- + +### Task 4: LocalDbConnectionFactory + ILocalDb consumer API + +**Classification:** standard +**Estimated implement time:** ~5 min +**Parallelizable with:** none (needs Tasks 2+3) + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/ILocalDb.cs` +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/LocalDbOptions.cs` +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/SqliteLocalDb.cs` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/SqliteLocalDbTests.cs` + +**Step 1: Failing tests** — `Open_CreatesDbFile_WalMode` (query `PRAGMA journal_mode` == "wal"), `ExecuteAsync_ParameterizedInsert_Works`, `QueryAsync_MapsRows`, `Transaction_CommitAndRollback`, `HlcUdf_RegisteredOnConnection` (`SELECT zb_hlc_next()` returns increasing values), `Hlc_ResumesFromDurableMax_OnReopen` (write rows→reopen→`Current` ≥ max hlc seen). + +**Step 2:** FAIL. + +**Step 3: Implement.** +- `LocalDbOptions`: `Path` (required), `BusyTimeoutMs` (default 5000), `Synchronous` (default "NORMAL"). +- `ILocalDb`: `Task ExecuteAsync(string sql, object? params, CancellationToken)`; `Task> QueryAsync(string sql, Func map, object? params, CancellationToken)`; `Task BeginTransactionAsync(CancellationToken)` (tx exposes same Execute/Query + Commit/Rollback); `SqliteConnection CreateConnection()` (escape hatch — documented: still UDF-registered, safe for capture). +- `SqliteLocalDb : ILocalDb, IDisposable`: on construction — open master connection, `PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL; PRAGMA busy_timeout=...; PRAGMA foreign_keys=ON;`, `LocalDbSchema.EnsureCreated`, recover HLC: `initial = MAX(meta.hlc_high_water, COALESCE((SELECT MAX(hlc) FROM __localdb_oplog),0), COALESCE((SELECT MAX(hlc) FROM __localdb_row_version),0))`, create the shared `HybridLogicalClock`. Every connection the factory hands out gets `conn.CreateFunction("zb_hlc_next", () => _clock.Next());` (deterministic: false). `Dispose` flushes `hlc_high_water = _clock.Current`. Parameter binding: anonymous-object properties → `@name` parameters (small reflection helper, cache `PropertyInfo[]` per type). + +**Step 4:** PASS. **Step 5: Commit** `feat(localdb): SqliteLocalDb host (WAL, UDF-registered connections, HLC recovery, consumer SQL API)`. + +--- + +### Task 5: Table registration + capture triggers + +**Classification:** high-risk (data contract — trigger SQL is the CDC heart) +**Estimated implement time:** ~5 min +**Parallelizable with:** none + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/ReplicatedTable.cs` (name, PK cols, all cols+declared types, `Digest` = lowercase hex SHA-256 of `"{table}|pk:{p1,p2}|cols:{c1:TYPE,c2:TYPE,...}"` — column order = PRAGMA order) +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/TriggerSqlGenerator.cs` +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Registration/LocalDbRegistrationException.cs` +- Modify: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Internal/SqliteLocalDb.cs` (add `RegisterReplicated(string tableName)` + `IReadOnlyDictionary ReplicatedTables`) +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/RegistrationTests.cs`, `.../CaptureTriggerTests.cs` + +**Step 1: Failing tests.** Registration: `Register_TableWithPk_InstallsThreeTriggers` (sqlite_master count), `Register_NoPk_Throws LocalDbRegistrationException`, `Register_MissingTable_Throws`, `Register_IsIdempotent_RegeneratesTriggers`, `Digest_ChangesWhenSchemaChanges`. Capture: create `orders(id INTEGER PRIMARY KEY, sku TEXT, qty INTEGER)`, register, then: `Insert_WritesOplogAndRowVersion_SameHlc` (both rows exist, hlc equal, row_json = full row via json comparison, pk_json = `{"id":1}`), `Update_WritesNewFullRow`, `Delete_WritesTombstone_NullRowJson_RowVersionTombstoned`, `MultiRowStatement_OneOplogEntryPerRow_DistinctHlcs`, `ApplyingFlag_SuppressesCapture` (set applying=1 in txn, write, applying=0 → no oplog rows), `UnregisteredTable_NoCapture`. + +**Step 2:** FAIL. + +**Step 3: Implement.** `RegisterReplicated`: PRAGMA `table_info('T')` → cols, types, pk ordinals (pk>0). No PK → throw. Quote all identifiers with `"`. Generated SQL (template for table T, pk cols P, all cols C — `DROP TRIGGER IF EXISTS` each first): + +```sql +CREATE TRIGGER "__localdb_T_ai" AFTER INSERT ON "T" +WHEN (SELECT applying FROM __localdb_applying WHERE id = 1) = 0 +BEGIN + INSERT INTO __localdb_oplog (table_name, pk_json, row_json, hlc, node_id, is_tombstone) + VALUES ('T', + json_object('p1', NEW."p1" /*, ...pk cols */), + json_object('c1', NEW."c1", 'c2', NEW."c2" /*, ...all cols */), + zb_hlc_next(), + (SELECT node_id FROM __localdb_meta WHERE id = 1), + 0); + INSERT OR REPLACE INTO __localdb_row_version (table_name, pk_json, hlc, node_id, is_tombstone, tombstone_utc) + SELECT table_name, pk_json, hlc, node_id, 0, NULL FROM __localdb_oplog WHERE seq = last_insert_rowid(); +END; +``` + +`_au` (AFTER UPDATE): identical body with NEW values. `_ad` (AFTER DELETE): oplog row with OLD pk json, `row_json` NULL, `is_tombstone` 1; row_version upsert via the same `last_insert_rowid()` SELECT but with `1, strftime('%Y-%m-%dT%H:%M:%fZ','now')` for tombstone fields. (The `last_insert_rowid()` trick binds row_version to the SAME hlc as the oplog entry — one `zb_hlc_next()` call per row. SQLite restores last_insert_rowid after the trigger completes, so consumers are unaffected.) + +`ReplicatedTables` is the in-memory registry the replication package reads (digests for handshake, column lists for apply/snapshot). + +**Step 4:** PASS. **Step 5: Commit** `feat(localdb): opt-in table registration + capture triggers (oplog + row_version, applying guard)`. + +--- + +### Task 6: Options validation + AddZbLocalDb DI + +**Classification:** small +**Estimated implement time:** ~3 min +**Parallelizable with:** Task 7, Task 8 + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/DependencyInjection/LocalDbServiceCollectionExtensions.cs` +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/LocalDbOptionsValidator.cs` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/DependencyInjectionTests.cs` + +**Steps (TDD):** tests: `AddZbLocalDb_ResolvesILocalDb_Singleton`, `AddZbLocalDb_WithRegistrations_TablesRegisteredOnStart`, `Validator_EmptyPath_Fails`, `Validator_NonPositiveBusyTimeout_Fails`. Implement `AddZbLocalDb(IConfiguration config, Action? register = null)`: binds `LocalDb` section, `IValidateOptions` validator (plain `IValidateOptions` — do NOT take a package dep on `ZB.MOM.WW.Configuration`; keep this lib dependency-light like Audit), registers `SqliteLocalDb` as singleton `ILocalDb`; the builder collects table names, applied at first resolution. Commit `feat(localdb): options validation + AddZbLocalDb DI`. + +--- + +### Task 7: Contracts — localdb_sync.v1 proto + +**Classification:** standard (wire contract) +**Estimated implement time:** ~3 min +**Parallelizable with:** Task 2, Task 3, Task 6 + +**Files:** +- Modify: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Contracts/Protos/localdb_sync.proto` + +**Step 1: Write the proto** (no TDD; compile is the test): + +```proto +syntax = "proto3"; +package localdb_sync.v1; +option csharp_namespace = "ZB.MOM.WW.LocalDb.Contracts.V1"; + +// One long-lived bidirectional stream carries BOTH directions of replication. +service LocalDbSync { + rpc Sync (stream SyncMessage) returns (stream SyncMessage); +} + +message SyncMessage { + oneof msg { + Handshake handshake = 1; // first message each direction + HandshakeAck handshake_ack = 2; + DeltaBatch delta_batch = 3; + DeltaAck delta_ack = 4; + SnapshotBegin snapshot_begin = 5; // sender switches to snapshot mode + SnapshotBatch snapshot_batch = 6; + SnapshotComplete snapshot_complete = 7; + } +} + +message Handshake { + string node_id = 1; + uint32 lib_schema_version = 2; + repeated TableDigest tables = 3; // fail-closed on any mismatch + int64 last_applied_remote_seq = 4; // highest of YOUR oplog seqs I have applied +} +message TableDigest { string table_name = 1; string digest = 2; } +message HandshakeAck { string node_id = 1; bool snapshot_required = 2; } + +message OplogEntry { + int64 seq = 1; + string table_name = 2; + string pk_json = 3; + optional string row_json = 4; // absent => tombstone + int64 hlc = 5; // UTC-based HLC (physical ms << 16 | counter) + string node_id = 6; + bool is_tombstone = 7; +} +message DeltaBatch { repeated OplogEntry entries = 1; } +message DeltaAck { int64 applied_thru_seq = 1; } + +message SnapshotBegin { int64 as_of_seq = 1; } +message SnapshotBatch { string table_name = 1; repeated SnapshotRow rows = 2; } +message SnapshotRow { + string pk_json = 1; + optional string row_json = 2; + int64 hlc = 3; + string node_id = 4; + bool is_tombstone = 5; +} +message SnapshotComplete { int64 as_of_seq = 1; } +``` + +**Step 2:** `dotnet build` → generated types compile, 0 warnings. **Step 3: Commit** `feat(localdb): localdb_sync.v1 wire contract`. + +--- + +### Task 8: OplogStore — read / ack / prune / peer state + +**Classification:** standard +**Estimated implement time:** ~5 min +**Parallelizable with:** Task 6 + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/Internal/OplogStore.cs` +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/ReplicationOptions.cs` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/OplogStoreTests.cs` + +**Step 1: Failing tests** — `ReadBatchAbove_ReturnsInSeqOrder_RespectsMaxBatch`, `RecordPeerAck_PrunesBelowWatermark`, `Prune_KeepsTombstoneRowVersions_WithinRetention` / `..._DropsExpired` (row_version tombstones older than `TombstoneRetention` removed; live-row versions NEVER pruned), `OplogDepth_ReportsBacklog`, `CapExceeded_SetsNeedsSnapshotAndPrunes` (backlog rows > `MaxOplogRows` OR oldest entry age > `MaxOplogAge` → prune to cap + `peer_state.needs_snapshot=1`), `PeerState_RoundTrips`. + +**Step 2:** FAIL. **Step 3: Implement** `internal sealed class OplogStore(ILocalDb db, ReplicationOptions opts)` — plain SQL over the Task-3 tables. `ReplicationOptions`: `PeerAddress`, `ApiKey` (nullable), `FlushInterval` (default 250 ms), `MaxBatchSize` (default 500), `MaxOplogRows` (default 1_000_000), `MaxOplogAge` (default 7 d), `TombstoneRetention` (default 7 d), `ReconnectBackoffMax` (default 60 s), `MaxHlcDriftAhead` (default 5 min, warn-only bool `FailClosedOnDrift`). All time values UTC. **Step 4:** PASS. **Step 5: Commit** `feat(localdb): oplog store (batch read, ack watermark, pruning caps, tombstone retention)`. + +--- + +### Task 9: LwwApplier + +**Classification:** high-risk (convergence correctness) +**Estimated implement time:** ~5 min +**Parallelizable with:** none + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/Internal/LwwApplier.cs` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/LwwApplierTests.cs` + +**Step 1: Failing tests** — `Apply_NewerRemote_Wins` (row replaced, row_version updated to remote (hlc,node)), `Apply_OlderRemote_Discarded_StillAcked`, `Apply_EqualHlc_NodeIdTieBreak_Deterministic`, `Apply_Tombstone_DeletesRow`, `Apply_DoesNotFireCaptureTriggers` (no oplog rows appear), `Apply_Batch_Atomic_WatermarkSameTxn` (fault injection mid-batch → nothing applied, watermark unmoved), `Apply_PoisonEntry_DeadLettered_BatchContinues` (entry for a column that doesn't exist → dead_letter row + counter, rest of batch applies), `Apply_IsIdempotent` (re-apply same batch → no-op), `Apply_ObservesRemoteHlc` (local clock advanced past remote hlc). + +**Step 2:** FAIL. **Step 3: Implement.** `ApplyBatchAsync(IReadOnlyList entries, CancellationToken)` — single transaction: + +1. `UPDATE __localdb_applying SET applying = 1 WHERE id = 1;` +2. Per entry: read `(hlc, node_id)` from `__localdb_row_version` for (table, pk_json). Remote wins iff `remote.hlc > local.hlc || (remote.hlc == local.hlc && string.CompareOrdinal(remote.node_id, local.node_id) > 0)` (no local version row = remote wins). Loser → skip (count `conflictsDiscarded`). Winner → tombstone ? parameterized `DELETE FROM "T" WHERE "p1" = json_extract(@pk,'$.p1') AND ...` : `INSERT OR REPLACE INTO "T" ("c1",...) VALUES (json_extract(@row,'$.c1'), ...)` using the `ReplicatedTables` column registry; then upsert `__localdb_row_version` with remote (hlc,node_id,is_tombstone). Per-entry `SqliteException` → dead-letter insert + continue. +3. `UPDATE __localdb_peer_state SET last_applied_remote_seq = @maxSeq, last_seen_hlc = @maxHlc, last_sync_utc = @utcNow WHERE id = 1;` +4. `UPDATE __localdb_applying SET applying = 0 WHERE id = 1;` commit. +5. After commit: `clock.Observe(maxHlc)`. Return `(appliedThruSeq, applied, discarded, deadLettered)`. + +**Step 4:** PASS. **Step 5: Commit** `feat(localdb): LWW applier (hlc+node_id tie-break, applying guard, dead-letter, atomic watermark)`. + +--- + +### Task 10: SyncSession — the shared stream state machine + +**Classification:** high-risk (protocol) +**Estimated implement time:** ~5 min +**Parallelizable with:** none + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/Internal/SyncSession.cs` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/SyncSessionTests.cs` + +Both sides run the SAME session logic over an abstracted duplex pipe (`Func send`, `ChannelReader inbox`) — the gRPC service and client are thin adapters (Tasks 11–12). That makes the protocol fully unit-testable with two in-memory sessions wired inbox↔send. + +**Step 1: Failing tests** (two `SyncSession`s over in-memory channels, each with its own temp DB): `Handshake_Exchanged_BothProceed`, `Handshake_DigestMismatch_FailsClosed_TypedError` (`LocalDbSchemaMismatchException`, stream ends, nothing applied), `Handshake_SchemaVersionMismatch_FailsClosed`, `Deltas_FlowBothDirections_AcksAdvanceWatermarks`, `Ack_DrivesPruning`, `HlcDriftAhead_Warns` (fake clock peer 10 min ahead → logged warning; with `FailClosedOnDrift` → session ends), `PeerNeedsSnapshot_TriggersSnapshotPath` (assert `SnapshotBegin` sent when `last_applied_remote_seq` < pruned horizon or `needs_snapshot=1`). + +**Step 2:** FAIL. **Step 3: Implement** `SyncSession.RunAsync(ct)`: send Handshake (node_id, version=1, digests from `ReplicatedTables`, `last_applied_remote_seq` from peer_state); await peer Handshake → validate → HandshakeAck (with `snapshot_required` decision). Then two concurrent loops until ct/stream end: **pump** (every `FlushInterval` or on oplog signal: `OplogStore.ReadBatchAbove(peerAckedSeq)` → send `DeltaBatch`s) and **receive** (DeltaBatch → `LwwApplier` → send `DeltaAck`; DeltaAck → `OplogStore.RecordPeerAck` → prune; snapshot messages → Task 12 handlers). All exceptions surface as typed session-end reasons; no silent catch. + +**Step 4:** PASS. **Step 5: Commit** `feat(localdb): SyncSession protocol state machine (handshake fail-closed, bidirectional delta pump)`. + +--- + +### Task 11: gRPC service + client + BackgroundService adapters + +**Classification:** standard +**Estimated implement time:** ~5 min +**Parallelizable with:** none + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/LocalDbSyncService.cs` (passive side: `LocalDbSync.LocalDbSyncBase.Sync` → wrap request/response streams into a `SyncSession` pipe; reject a second concurrent stream with `ALREADY_EXISTS`) +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/LocalDbSyncClient.cs` + `SyncBackgroundService.cs` (initiator: `GrpcChannel.ForAddress(PeerAddress)`, optional `Authorization: Bearer ` metadata per family convention, open duplex call → `SyncSession`; on fault: capped exponential backoff `1s→ReconnectBackoffMax`, jittered, forever) +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/DependencyInjection/ReplicationServiceCollectionExtensions.cs` — `AddZbLocalDbReplication(IConfiguration)` (binds `LocalDb:Replication`, validator: initiator requires `PeerAddress` absolute http(s) URI, positive intervals/caps) + `MapZbLocalDbSync(IEndpointRouteBuilder)` +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/GrpcAdapterTests.cs` + +**Steps (TDD):** tests use `Microsoft.AspNetCore.TestHost` (in-process Kestrel-less server hosting `MapZbLocalDbSync`) + `GrpcChannel` over the TestServer handler: `EndToEnd_TwoDbs_InsertOnA_AppearsOnB` (the first real wire round-trip), `ClientReconnects_AfterServerRestart_WithBackoff`, `SecondConcurrentStream_Rejected`, `Validator_MissingPeerAddress_FailsOnInitiator`. Commit `feat(localdb): gRPC sync adapters, background service, replication DI`. + +--- + +### Task 12: Snapshot resync + +**Classification:** high-risk +**Estimated implement time:** ~5 min +**Parallelizable with:** none + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/Internal/SnapshotStreamer.cs` (sender: per registered table, page `__localdb_row_version LEFT JOIN "T"` in chunks of `MaxBatchSize` → `SnapshotBatch` rows incl. tombstones; wrap in `SnapshotBegin(as_of_seq=current max oplog seq)` / `SnapshotComplete`) +- Modify: `.../Internal/SyncSession.cs` (receive path: snapshot rows go through the SAME `LwwApplier` LWW compare — a snapshot NEVER truncates, it merges, because the receiver may hold newer local writes; after `SnapshotComplete`: clear `needs_snapshot`, set `last_applied_remote_seq = as_of_seq`, resume deltas) +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/SnapshotResyncTests.cs` + +**Steps (TDD):** `PrunedPeer_GetsSnapshot_ThenConverges` (A accumulates past cap while B offline → reconnect → snapshot → both identical), `Snapshot_MergesNotTruncates` (B has newer local row during snapshot → B's row survives on both), `Snapshot_CarriesTombstones` (row deleted on A while B offline+pruned → gone on B after resync), `DeltasResume_AfterSnapshot`. Commit `feat(localdb): snapshot resync (LWW-merging, tombstone-carrying, delta resume)`. + +--- + +### Task 13: Telemetry meters + sync status surface + +**Classification:** small +**Estimated implement time:** ~3 min +**Parallelizable with:** Task 14 + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb.Replication/LocalDbMetrics.cs` +- Modify: touchpoints in OplogStore / LwwApplier / SyncSession / SyncBackgroundService +- Test: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/MetricsTests.cs` + +**Steps (TDD via `MetricCollector`):** `System.Diagnostics.Metrics.Meter` named `zb.mom.ww.localdb`: observable gauges `localdb.oplog.depth`, `localdb.sync.lag.seconds` (UTC now − last_sync_utc; follows METRIC-CONVENTIONS seconds-unit rule), counters `localdb.sync.applied`, `localdb.sync.conflicts_discarded`, `localdb.sync.dead_lettered`, `localdb.sync.reconnects`, `localdb.sync.snapshots`. Plus `ISyncStatus` (connected?, last sync UTC, backlog, peer node id) registered in DI for consumer health checks. Commit `feat(localdb): telemetry meters + ISyncStatus`. + +--- + +### Task 14: Convergence integration + property tests + +**Classification:** high-risk (this is the proof the design works) +**Estimated implement time:** ~5 min +**Parallelizable with:** Task 13 + +**Files:** +- Create: `ZB.MOM.WW.LocalDb/tests/ZB.MOM.WW.LocalDb.Tests/Convergence/ConvergenceFixture.cs` (two full stacks — DB + replication — wired over TestHost; helper `AssertConvergedAsync()` = quiesce (both oplogs drained + acked), then dump each registered table + row_version ordered by pk from both DBs and assert byte-identical) +- Create: `.../Convergence/BidirectionalConvergenceTests.cs` +- Create: `.../Convergence/RandomOpsConvergenceTests.cs` + +**Tests:** `ConcurrentWrites_BothSides_Converge`, `UpdateDeleteRace_Converges_TombstoneVsUpdate_LwwDecides`, `OfflineAccumulation_Reconnect_Converges`, `RestartMidStream_NoLoss_NoDuplication` (kill background service mid-batch, restart → converges; proves idempotent replay), and the property test: seeded `Random` (seed logged), 500 mixed ops (insert/update/delete, overlapping keys) interleaved on both nodes with random disconnect/reconnect → `AssertConvergedAsync`. Run each ≥3 seeds. Commit `test(localdb): bidirectional convergence + randomized property tests`. + +--- + +### Task 15: README, pack verification, scadaproj index row + +**Classification:** small +**Estimated implement time:** ~4 min +**Parallelizable with:** none + +**Files:** +- Modify: `ZB.MOM.WW.LocalDb/README.md` (full: what it is, quickstart both modes, config reference table for `LocalDb` + `LocalDb:Replication` sections, LWW semantics, documented limitations from the plan header, ops notes: oplog caps → snapshot fallback, dead-letter table, metrics list) +- Modify: `CLAUDE.md` (scadaproj root — add LocalDb row to the Component normalization table: Status "Built (lib 0.1.0, not yet published/adopted)") +- Modify: `docs/plans/2026-07-17-localdb-design.md` (append "Implemented" note + final test count) + +**Steps:** write docs; `dotnet pack ZB.MOM.WW.LocalDb.slnx -c Release -o artifacts` → expect **3 nupkgs @ 0.1.0** each containing the README; full `dotnet test` → all green, 0 warnings. Do NOT publish to the Gitea feed (that's a separate user decision, per component-status-claims-are-optimistic policy: never claim published without doing it). Commit `docs(localdb): README + scadaproj index row; pack-verified 3 nupkgs @ 0.1.0`. + +--- + +## Execution notes + +- Work from `ZB.MOM.WW.LocalDb/` for build/test; commits go to the scadaproj repo root. +- Every task: `dotnet build ZB.MOM.WW.LocalDb.slnx` must stay 0-warning (TreatWarningsAsErrors). +- Test DBs: temp files under the test's own temp dir (`Path.GetTempPath()`/guid), deleted in Dispose — never `:memory:` for capture/replication tests (WAL + multiple connections need a real file). +- Task order is 1→15 linear except the noted parallelizable pairs. diff --git a/docs/plans/2026-07-17-localdb-implementation.md.tasks.json b/docs/plans/2026-07-17-localdb-implementation.md.tasks.json new file mode 100644 index 0000000..8a0f4c7 --- /dev/null +++ b/docs/plans/2026-07-17-localdb-implementation.md.tasks.json @@ -0,0 +1,21 @@ +{ + "planPath": "docs/plans/2026-07-17-localdb-implementation.md", + "tasks": [ + {"id": 7, "planTask": 1, "subject": "Task 1: Scaffold the library", "status": "pending"}, + {"id": 8, "planTask": 2, "subject": "Task 2: HybridLogicalClock", "status": "pending", "blockedBy": [7]}, + {"id": 9, "planTask": 3, "subject": "Task 3: Lib-owned schema + LocalDbSchema migrator", "status": "pending", "blockedBy": [7]}, + {"id": 10, "planTask": 4, "subject": "Task 4: LocalDbConnectionFactory + ILocalDb API", "status": "pending", "blockedBy": [8, 9]}, + {"id": 11, "planTask": 5, "subject": "Task 5: Table registration + capture triggers", "status": "pending", "blockedBy": [10]}, + {"id": 12, "planTask": 6, "subject": "Task 6: Options validation + AddZbLocalDb DI", "status": "pending", "blockedBy": [10]}, + {"id": 13, "planTask": 7, "subject": "Task 7: Contracts — localdb_sync.v1 proto", "status": "pending", "blockedBy": [7]}, + {"id": 14, "planTask": 8, "subject": "Task 8: OplogStore — read/ack/prune/peer state", "status": "pending", "blockedBy": [10]}, + {"id": 15, "planTask": 9, "subject": "Task 9: LwwApplier", "status": "pending", "blockedBy": [11, 13]}, + {"id": 16, "planTask": 10, "subject": "Task 10: SyncSession state machine", "status": "pending", "blockedBy": [14, 15]}, + {"id": 17, "planTask": 11, "subject": "Task 11: gRPC adapters + BackgroundService + replication DI", "status": "pending", "blockedBy": [16]}, + {"id": 18, "planTask": 12, "subject": "Task 12: Snapshot resync", "status": "pending", "blockedBy": [17]}, + {"id": 19, "planTask": 13, "subject": "Task 13: Telemetry meters + ISyncStatus", "status": "pending", "blockedBy": [18]}, + {"id": 20, "planTask": 14, "subject": "Task 14: Convergence integration + property tests", "status": "pending", "blockedBy": [18]}, + {"id": 21, "planTask": 15, "subject": "Task 15: README, pack verify, scadaproj index row", "status": "pending", "blockedBy": [19, 20]} + ], + "lastUpdated": "2026-07-17" +}