diff --git a/CLAUDE.md b/CLAUDE.md index 27b6e7f..34dde4f 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -144,6 +144,7 @@ each project's **code-verified current state**, and the **gaps** between. See | Audit (event model + writer seam) | Adopted (lib `0.1.0`; all 3 apps, merged to **local default** main/master + **pushed to origin** (gitea)) | Shared `ZB.MOM.WW.Audit` lib | [`components/audit/`](components/audit/) | [`ZB.MOM.WW.Audit/`](ZB.MOM.WW.Audit/) | | Galaxy Repository (object-hierarchy SQL browse + gRPC service) | Built (lib `0.1.0`, **published to the Gitea feed**; consumed by HistorianGateway as a feed `PackageReference`) | Shared `ZB.MOM.WW.GalaxyRepository` lib | _(design in histsdk + design doc 2026-06-23)_ | [`ZB.MOM.WW.GalaxyRepository/`](ZB.MOM.WW.GalaxyRepository/) | | Secrets (encrypted store + `${secret:}` resolution) | Built (libs `0.1.2`, **published to the Gitea feed**; **HistorianGateway adopted + live-proven** 2026-07-16; **mxaccessgw adopted G-4/G-5/G-6 + merged to `origin/main` @ `e088dfa`** 2026-07-16, box-verified fail-closed + encrypt-at-rest; **ScadaBridge adopted G-3/G-4/G-5/G-6 + merged to `origin/main` @ `128f1596`** 2026-07-16, **G-3 live-proven vs the real production MxGateway gateway** `wonder-app-vd03:5120` (secret-ref ApiKey → Connected + browsed real Galaxy); **OtOpcUa adopted G-2/G-4/G-5/G-6 + merged to `origin/master` @ `872cf7e3`** (lmxopcua) 2026-07-16 — the last app, so **all four now adopted**; Layer-B driver secrets (Galaxy API key + OpcUaClient `Password`/`UserCertificatePassword`) resolve fail-closed at driver session-open, `AddZbSecrets` registered unconditionally so driver-only nodes work; **G-2 live-proven vs the real production MxGateway** `wonder-app-vd03:5120` — a `secret:`-ref Galaxy ApiKey resolved through OtOpcUa's real `GalaxyDriverBrowser` → dummy key `MxGatewayAuthenticationException`, real key → browsed the real Galaxy root (10 nodes); temp key minted+revoked+deleted); **G-8 KEK-rotation BUILT (lib bumped `0.1.2`→`0.1.3`, 2026-07-17):** `Rewrap` DEK primitive + CAS-guarded `ApplyRewrapAsync` + `KekRotationService.RewrapAllAsync` + `secret rewrap-all` CLI + operator runbook (`ZB.MOM.WW.Secrets/docs/operations/kek-rotation.md`), full suite green + CLI smoke + adversarial crypto review PASS (TOCTOU closed via compare-and-swap); republish to feed is opt-in. **G-7 clustered replication DESIGNED + PLANNED** — fork resolved to **build the shared SQL-Server `ISecretStore`** (Akka replicator deferred phase-2), executable plan `docs/plans/2026-07-17-secrets-g7-*`. Both not yet committed/pushed (awaiting go-ahead) | Shared `ZB.MOM.WW.Secrets` lib (3 packages + CLI) | [`components/secrets/`](components/secrets/) | [`ZB.MOM.WW.Secrets/`](ZB.MOM.WW.Secrets/) | +| LocalDb (embedded cache + 2-node sync) | Built (3 pkgs `0.1.0`, **NOT yet published to the feed**; no app adoption yet) | Shared `ZB.MOM.WW.LocalDb` lib | [`docs/plans/2026-07-17-localdb-design.md`](docs/plans/2026-07-17-localdb-design.md) | [`ZB.MOM.WW.LocalDb/`](ZB.MOM.WW.LocalDb/) | The auth component is fully populated: a normalized [`spec`](components/auth/spec/SPEC.md), a proposed [`shared-contract`](components/auth/shared-contract/ZB.MOM.WW.Auth.md), three @@ -333,6 +334,11 @@ dotnet test ZB.MOM.WW.HistorianGateway.slnx --filter "Category=LiveIntegration" # runtime-only aspnet:10.0 image; login + gRPC API-key auth disabled; points at the REAL wonder datasources. dotnet publish src/ZB.MOM.WW.HistorianGateway.Server/ZB.MOM.WW.HistorianGateway.Server.csproj -c Release -o docker/publish -p:UseAppHost=false cd docker && docker compose up -d --build # dashboard/health/metrics :5220, gRPC h2c :5221 (needs host VPN for egress) + +# ZB.MOM.WW.LocalDb (~/Desktop/scadaproj/ZB.MOM.WW.LocalDb — a hosted shared lib) +dotnet build ZB.MOM.WW.LocalDb.slnx +dotnet test ZB.MOM.WW.LocalDb.slnx # 144 tests, fully offline (no live deps) +dotnet pack ZB.MOM.WW.LocalDb.slnx -c Release -o artifacts # 3 nupkgs @ 0.1.0 ``` > **Shared GLAuth (all three apps + HistorianGateway):** LDAP auth for every local dev/test stack is provided by a diff --git a/ZB.MOM.WW.LocalDb/README.md b/ZB.MOM.WW.LocalDb/README.md index 6e847ab..cf5440c 100644 --- a/ZB.MOM.WW.LocalDb/README.md +++ b/ZB.MOM.WW.LocalDb/README.md @@ -1,8 +1,211 @@ # ZB.MOM.WW.LocalDb -An embedded SQLite local cache with optional bidirectional async 2-node replication -over gRPC, for the ZB.MOM.WW SCADA / OT family. It gives an application a fast, -crash-safe local store that can optionally keep a second node in sync asynchronously, -without requiring a central database on the hot path. +An embedded **SQLite local cache** with **optional bidirectional async 2-node +replication over gRPC**, for the ZB.MOM.WW SCADA / OT family. It gives an application a +fast, crash-safe local store that can optionally keep a second node in sync +asynchronously — local writes **never** block on the peer, and a node that never enables +replication is simply a fast local DB. + +Change capture uses SQLite triggers writing a single HLC-stamped oplog inside the +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`) + +| Package | Contents | +|---|---| +| `ZB.MOM.WW.LocalDb` | Embedded DB host: `ILocalDb` SQL surface (`ExecuteAsync` / `QueryAsync` / transactions), `RegisterReplicated(table)` (PK validation + trigger install + oplog schema), `HybridLogicalClock`, `AddZbLocalDb()` DI. | +| `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, not yet published to the Gitea feed; no app has adopted it yet.** + +## Quickstart + +### Local-only (no replication) + +```csharp +builder.Services.AddZbLocalDb(builder.Configuration, db => +{ + // onReady: create your own tables, then opt them into (future) replication. + db.ExecuteAsync( + "CREATE TABLE IF NOT EXISTS widget (id INTEGER PRIMARY KEY, name TEXT, qty INTEGER)") + .GetAwaiter().GetResult(); + db.RegisterReplicated("widget"); // PK required; installs capture triggers +}); + +// Consume ILocalDb anywhere: +await db.ExecuteAsync("INSERT INTO widget (id, name, qty) VALUES (@id, @n, @q)", + new { id = 1, n = "bolt", q = 42 }); +var rows = await db.QueryAsync("SELECT name, qty FROM widget", + r => (Name: r.GetString(0), Qty: r.GetInt32(1))); +``` + +`AddZbLocalDb` is first-registration-wins: a second call is a no-op. `RegisterReplicated` +installs capture triggers even without the replication package present — you can register +tables now and turn on sync later. + +### Replicated (2 nodes) + +Every node registers the passive side (so either can accept a stream) and the initiator +(idle unless a `PeerAddress` is set). Configure exactly **one** side with the peer's +address; the other stays passive. + +```csharp +// Both nodes: +builder.Services.AddZbLocalDb(builder.Configuration, db => { /* create + RegisterReplicated */ }); +builder.Services.AddGrpc(); +builder.Services.AddZbLocalDbReplication(builder.Configuration); // passive host + initiator + maintenance + +var app = builder.Build(); +app.MapZbLocalDbSync(); // passive gRPC endpoint — host must have AddGrpc() +app.Run(); +``` + +```jsonc +// Node A — the INITIATOR: dials node B +"LocalDb": { + "Path": "nodeA.db", + "Replication": { "PeerAddress": "https://nodeB:5001", "ApiKey": "…" } +} +// Node B — PASSIVE: no PeerAddress, its SyncBackgroundService idles and only serves the endpoint +"LocalDb": { "Path": "nodeB.db", "Replication": { } } +``` + +Both directions flow on the single bidirectional stream the initiator opens; the passive +node still pushes its own deltas back. Only one side needs `PeerAddress`. + +## Configuration + +### `LocalDb` section (`LocalDbOptions`) + +| Key | Default | Meaning | +|---|---|---| +| `Path` | `""` (required) | Filesystem path to the SQLite database file. | +| `BusyTimeoutMs` | `5000` | SQLite `busy_timeout` pragma, in milliseconds. | +| `Synchronous` | `"NORMAL"` | SQLite `synchronous` pragma (`OFF` / `NORMAL` / `FULL`). | + +### `LocalDb:Replication` section (`ReplicationOptions`) + +| Key | Default | Meaning | +|---|---|---| +| `PeerAddress` | `""` | gRPC address to dial. Empty ⇒ passive (initiator idles). | +| `ApiKey` | `null` | Bearer token the initiator presents on its channel; `null` ⇒ no key header. | +| `FlushInterval` | `250 ms` | How often the change pump flushes pending oplog deltas. | +| `MaxBatchSize` | `500` | Max oplog entries per outbound batch (also the snapshot page size). | +| `MaxOplogRows` | `1_000_000` | Backlog ceiling; exceeding it flags a snapshot resync and prunes oldest rows. | +| `MaxOplogAge` | `7 days` | Age ceiling; an unpruned entry older than this flags a snapshot resync. | +| `TombstoneRetention` | `7 days` | How long deleted-row tombstones are kept before pruning. | +| `ReconnectBackoffMax` | `60 s` | Upper bound on exponential reconnect backoff (validator caps ≤ 1 day). | +| `MaintenanceInterval` | `60 s` | How often maintenance enforces caps + prunes expired tombstones. | +| `MaxHlcDriftAhead` | `5 min` | Max tolerated HLC drift where the peer's clock runs ahead of ours. | +| `FailClosedOnDrift` | `false` | When true, an over-drift entry is dead-lettered instead of applied. | + +All options are validated at startup (`ZB.MOM.WW.Configuration` validators). + +## How it works + +- **Capture** — three AFTER triggers per registered table write each change into one + shared, HLC-stamped `__localdb_oplog` plus a `__localdb_row_version` ledger, in the + consumer's own transaction. A DB-scoped `__localdb_applying` guard makes replicated + applies skip re-capture. +- **Convergence** — last-writer-wins: compare incoming HLC, then `node_id` ordinal as the + tie-break; the strict-greater winner is applied, the loser discarded (but still acked). +- **Streaming** — one long-lived bidirectional gRPC stream. The initiator dials; both + sides push oplog batches above the peer's watermark and ack the applied watermark back, + which drives pruning. Outbound uses a single writer fed by two lanes (an unbounded + control lane for handshake/acks that always preempts a bounded data lane for + delta/snapshot batches), so receiving never blocks on the ability to ack. +- **Resync** — when a peer's watermark falls below the pruned oplog horizon, it is sent a + full snapshot of the registered tables (streamed, keyset-paged) and deltas resume above + the snapshot's as-of seq. **Each node decides from its OWN state** whether it owes the + peer a snapshot. +- **Maintenance** — a background service on **both** roles enforces the oplog row/age caps + and prunes expired tombstones every `MaintenanceInterval`. + +## Semantics worth knowing + +- **PK-column updates** are captured as delete-of-old-PK + insert-of-new — a tombstone on + the old key plus a fresh insert (full delete+insert support). +- **Upserts** (`INSERT … ON CONFLICT DO UPDATE`) are fully supported (the row-version + trigger uses an explicit `ON CONFLICT` clause so it is not rewritten by the outer + statement's conflict algorithm). +- **Clock-drift guard** — warn-once per session when a peer entry's HLC physical time is + more than `MaxHlcDriftAhead` ahead. With `FailClosedOnDrift`, each over-drift entry is + dead-lettered; otherwise it is applied (the HLC merge absorbs the skew). The **session + never dies from drift** — no livelock on reconnect. +- **Writes never block on the peer**; local commit is independent of replication. +- **All timestamps are UTC** (HLC physical component is UTC Unix milliseconds). + +## Limitations (honest list) + +- **Exactly two nodes.** No mesh / N-node topologies. +- **LWW only.** No per-table pluggable conflict resolvers. +- **BLOB columns are rejected at registration** (a column declared `BLOB` fails + `RegisterReplicated`, because `json_object` cannot capture it). A *typeless* column can + still hold a BLOB value at runtime, which then surfaces as a write-time error, not a + registration error. +- **CDC semantics** — rows that existed **before** `RegisterReplicated` are not captured + and are not snapshotted; only changes after registration are replicated. Re-baselining + pre-existing rows requires re-writing them. +- **REAL ±Infinity** is carried in SQLite's JSON dialect on the wire (fine for the + SQLite↔SQLite peers this library targets). +- `AddZbLocalDb` / `AddZbLocalDbReplication` are **first-registration-wins** (idempotent). +- A **mutual simultaneous snapshot** buffers the peer's inbound snapshot in memory while + this node is still sending its own (acceptable for the small datasets this lib targets). + +## Observability + +Metrics publish under the single meter **`ZB.MOM.WW.LocalDb.Replication`** (all instruments +null-guarded — replication runs identically whether or not metrics are collected): + +| Instrument | Kind | Notes | +|---|---|---| +| `localdb.oplog.depth` | ObservableGauge | Unacked oplog backlog (polled cheaply at scrape time). | +| `localdb.sync.lag.seconds` | ObservableGauge | Seconds since the last successful sync (no measurement when never synced). | +| `localdb.sync.applied` | Counter | Rows applied from inbound **delta** batches. | +| `localdb.sync.conflicts_discarded` | Counter | Inbound rows that lost LWW. | +| `localdb.sync.dead_lettered` | Counter | Poison / over-drift rows routed to the dead-letter table. | +| `localdb.sync.reconnects` | Counter | Reconnect attempts (initiator). | +| `localdb.sync.snapshots` | Counter | Snapshots, tagged `direction=sent`/`received` (one increment per snapshot). | + +**Counter split to note:** rows merged from a **snapshot** are *not* counted under +`localdb.sync.applied`; a completed snapshot instead increments `localdb.sync.snapshots` +once (with its `direction` tag). Only delta-batch rows increment `localdb.sync.applied`. + +`ISyncStatus` (registered singleton) surfaces live state for dashboards / health: +`Connected` (true while **any** session is active — a node can run both roles at once), +`PeerNodeId`, `LastSyncUtc`, `OplogBacklog` (**nullable** — null means unknown because no +provider is wired or the poll threw, so a DB failure never reads as "0 backlog"), and +`ConnectionAttempts`. + +## Ops notes + +- When the oplog row/age caps are exceeded, maintenance prunes to the ceiling and flags + `needs_snapshot` (disk safety over delta continuity); the peer then snapshot-resyncs. +- The **`__localdb_dead_letter`** table holds poison rows (bad JSON, null-row non-tombstone + entries) and over-drift rejects (when `FailClosedOnDrift`). Inspect it to diagnose + rejects. +- **Peer auth**: `ApiKey` is sent by the initiator as a static `Authorization: Bearer …` + header on its channel. The library's passive service does **not** itself verify the key + — the passive host enforces whatever gRPC auth the host application configures (e.g. + `ZB.MOM.WW.Auth`). Configure host-side authorization to actually gate inbound streams. + +## Testing + +144 tests, all offline / no live dependencies — unit (HLC ordering + restart +persistence, trigger SQL generation, LWW decision table, watermark/pruning, digest +computation), integration over two real SQLite DBs on an in-memory (and real) gRPC pipe +(bidirectional convergence, delete/update races, offline-accumulate → reconnect, +prune → snapshot resync), plus seeded randomized property tests (random interleaved op +sequences → both DBs identical after quiesce). Family conventions apply +(`TreatWarningsAsErrors`). + +Build / test from this directory: + +```bash +dotnet build ZB.MOM.WW.LocalDb.slnx +dotnet test ZB.MOM.WW.LocalDb.slnx +``` See the design doc: [`docs/plans/2026-07-17-localdb-design.md`](../docs/plans/2026-07-17-localdb-design.md). diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/DependencyInjection/LocalDbServiceCollectionExtensions.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/DependencyInjection/LocalDbServiceCollectionExtensions.cs index d8415f5..6fe0999 100644 --- a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/DependencyInjection/LocalDbServiceCollectionExtensions.cs +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/DependencyInjection/LocalDbServiceCollectionExtensions.cs @@ -18,6 +18,7 @@ public static class LocalDbServiceCollectionExtensions /// caller. This is where consumers create their tables (via the passed ) and call /// to opt them into replication. /// + /// Subsequent calls are no-ops if an registration already exists (first registration wins). public static IServiceCollection AddZbLocalDb( this IServiceCollection services, IConfiguration configuration, @@ -44,8 +45,16 @@ public static class LocalDbServiceCollectionExtensions catch { // A throwing callback must not leak the open master connection: MS.DI does not - // dispose instances a failed factory never returned. - db.Dispose(); + // dispose instances a failed factory never returned. Guard the cleanup so a Dispose + // failure can never replace the consumer's original onReady exception. + try + { + db.Dispose(); + } + catch + { + // Swallow: the onReady exception below is the one the caller must see. + } throw; } return db; diff --git a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Hlc/HybridLogicalClock.cs b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Hlc/HybridLogicalClock.cs index 3d80201..a086f41 100644 --- a/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Hlc/HybridLogicalClock.cs +++ b/ZB.MOM.WW.LocalDb/src/ZB.MOM.WW.LocalDb/Hlc/HybridLogicalClock.cs @@ -13,6 +13,7 @@ public sealed class HybridLogicalClock _last = initialValue; } + /// Returns the next monotonically increasing HLC stamp: the current UTC physical time when it advances the clock, otherwise the last value with its logical counter incremented by one. public long Next() { lock (_lock) @@ -26,8 +27,12 @@ public sealed class HybridLogicalClock /// 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; } } + /// The most recently issued/observed HLC value, without advancing the clock. public long Current { get { lock (_lock) return _last; } } + /// Extracts the physical component (UTC Unix milliseconds) from a packed HLC value. public static long PhysicalMs(long hlc) => hlc >> 16; + + /// Extracts the 16-bit logical counter from a packed HLC value. public static int Counter(long hlc) => (int)(hlc & 0xFFFF); } diff --git a/docs/plans/2026-07-17-localdb-design.md b/docs/plans/2026-07-17-localdb-design.md index cebea38..38dc6ce 100644 --- a/docs/plans/2026-07-17-localdb-design.md +++ b/docs/plans/2026-07-17-localdb-design.md @@ -86,3 +86,16 @@ dead-letter counters. `ZB.MOM.WW.Configuration` validators on all options. - More than two nodes / mesh topologies - Per-table pluggable conflict resolvers (LWW only) - TTL/eviction (consumer concern via their own SQL) + +## Implemented — deltas from this design + +Shipped deviations from the design above (all built + tested; see `ZB.MOM.WW.LocalDb/README.md`): + +- **Engine:** Microsoft.Data.Sqlite is the default behind an engine seam; libSQL deferred (already noted above). +- **PK-column updates:** fully supported as OLD-pk tombstone + new-row insert (delete+insert), an upgrade from the original per-design limitation. +- **Clock drift:** an over-drift entry is **per-entry dead-lettered** (when `FailClosedOnDrift`) rather than faulting the session — the session never dies from drift. +- **Snapshot trigger:** each node computes its **own** `needs_snapshot` (pruned-horizon / flag); `HandshakeAck.snapshot_required` is now informational, not the decision authority. +- **`onReady` API shape:** `AddZbLocalDb(..., Action onReady)` for table creation + `RegisterReplicated`, instead of a registration-builder object. +- **Outbound path:** a single-writer, two-lane design (unbounded control lane preempts a bounded data lane) rather than an undifferentiated write path. +- **Maintenance:** a dedicated `MaintenanceBackgroundService` (both roles) enforces oplog caps + tombstone retention on an interval. +- **Snapshot sender** returns the **as-of oplog seq** it covers, so the pump skips already-covered tail deltas.