Skip to content

feat(sync): spoke-side ledger for edge-to-cloud replication (#569 PR 1/9) - #570

Merged
xe-nvdk merged 1 commit into
mainfrom
feat/edge-sync-ledger
Aug 6, 2026
Merged

xe-nvdk merged 1 commit into
mainfrom
feat/edge-sync-ledger

Conversation

@xe-nvdk

@xe-nvdk xe-nvdk commented Aug 6, 2026

Copy link
Copy Markdown
Member

First of nine PRs for edge sync — see #569 for the phase-1 breakdown and docs/progress/2026-06-04-edge-sync-architecture-converged.md for the design.

Summary

Adds internal/sync with the spoke-side ledger: the durable record of which Parquet files have reached which hub, and how far a partial transfer got.

  • sync_ledger keyed UNIQUE(hub_id, path) from day one — multi-hub becomes config later, not a schema migration on a disconnected box.
  • State transitions enforced in SQL, not by trusting callers.
  • MarkFailed's retry cap decided inside a single UPDATE — no read-then-write race.
  • bytes_sent as the resume checkpoint, preserved across failure, clamped to size_bytes.
  • PruneSynced batched at 1000 rows, PK-ordered, ctx.Err() between batches, no vacuum.
  • UTC normalization on write and read.

No production callers. Nothing constructs the ledger, so its schema is never created. Wiring lands with /sync/run (PR 8).

Design decision recorded on the issue

Granularity is whole-file, not row-group — decision comment. raft.FileEntry carries SHA256/SizeBytes per file with no row-group data, which is what makes §5.1's reconcile answerable from the manifest with no I/O on parquet bytes. Row-group sync would force the spoke to open every pending file's footer, need a second integrity chain with no authority to verify against, and turn §6.1's three-case receive into a range-merge problem. Revisit later is cheap: a row-group range is an additive column.

Configuration matrix

Configuration Reaches new code? Preconditions established?
OSS standalone, any config No — grep NewLedger outside the package returns nothing n/a; no live derefs
Cluster mode, any config No — no wiring in cmd/arc/main.go n/a
Tests only Yes — ledger_test.go is the sole caller test DB via sql.Open, t.Cleanup closes

Every row is "no", which inverts the usual risk profile: no config reaches this code, so no nil-deref or subsystem-interaction bug is possible. Risk sits entirely in SQL correctness and the state machine, which is where review was pointed.

Config keys introduced: none. Deref check: NewLedger rejects a nil *sql.DB; no other pointer is dereferenced internally.

Review

Two adversarial passes. The first found 3 blockers + 6 high, all fixed; I verified each empirically with probes before fixing rather than taking them on faith:

  • Lost-update race in MarkFailed — read-then-write let a concurrent worker push attempts past the cap between statements; the losing writer resurrected the entry to pending, retrying forever. Now one atomic UPDATE with the CASE inside.
  • A factually wrong time-domain comment — probed the real stored format: "2026-08-06 22:40:43.006808+00:00", space-separated with an offset, neither RFC3339 nor datetime() output. The comment pointed a future maintainer the wrong way.
  • discovered_at not UTC-normalized — probed: the same instant stored two ways gave SQL equality 0.
  • Every illegal transition permitted — MarkInFlight on a synced row succeeded with synced_at still set, double-counting in Stats.
  • PendingBytes went negative (−4900 with an over-large offset), and since it is a SUM, one bad offset silently cancelled other files' real pending bytes.
  • Full scan + temp B-tree on the prune cutoff and status query — fixed with a second index, both re-verified with EXPLAIN QUERY PLAN.

Two reviewer claims I checked and declined: the prune's state filter is redundant-by-construction rather than untested (removing either filter alone still protects unsynced rows — deliberate belt-and-braces on a DELETE), and database as a column name doesn't actually need quoting (probed unquoted in WHERE/SELECT/GROUP BY/ORDER BY; four existing Arc tables already do this).

The second pass found no bugs. It suggested an explicit pending → synced reconcile-path test, which is added.

I also found one of my own tests was vacuous — the TrackBatch rollback test never triggered a rollback because empty strings satisfy NOT NULL. Replaced with a trigger-based abort and verified it fails pre-fix (reported 1 inserted but rollback discarded every row).

Test plan

  • go build ./cmd/... ./internal/...
  • go test ./internal/sync/ — 23 tests
  • go test -race ./internal/sync/
  • go vet / gofmt -l clean
  • Query plans verified with EXPLAIN QUERY PLAN (no full scans, no temp B-trees)
  • New regression tests verified to fail pre-fix
  • Binary run — N/A. No startup path reaches this code; there is nothing to exercise. The usual run-the-binary gate applies to PR 8, which first gives the ledger a caller.

Notes

  • Docs: no docs.basekick.net change — nothing here is configurable or observable. Docs land with PR 8's endpoints.
  • Package name: internal/sync shadows stdlib sync; importers needing sync.Mutex must alias. Renaming to edgesync is cheap now and progressively more annoying as PRs 2–9 stack. Happy to do it before PR 2.
  • Deferred to PR 8: RecoverInFlight has no age guard and Pending has no claim/lease, so two workers could hand out the same row. Both are real but belong with the agent, where concurrency actually exists — building a lease with no consumer now would be speculative.

🤖 Generated with Claude Code

First of nine PRs for edge sync (#569). Adds internal/sync with the
spoke-side ledger: the durable record of which Parquet files have
reached which hub, and how far a partial transfer got.

The unit of sync is the file, not the row — Arc already produces
immutable content-addressed Parquet with a SHA256 in the manifest, so
shipping files gives end-to-end integrity for free and makes idempotency
trivial. sync_ledger is keyed (hub_id, path) from day one so multi-hub
is later config rather than a schema migration on a disconnected box.

State transitions are enforced in SQL, not by trusting callers. Every
mutation carries `AND state = ?`, and ErrInvalidTransition is distinct
from ErrNotFound because "wrong state" and "untracked path" are
different bugs. Without the guards, MarkInFlight on a synced entry
produced a row that simultaneously claimed the hub had the file and
that it was being sent.

MarkFailed decides the retry cap inside a single UPDATE rather than
reading attempts and deciding in Go. A read-then-write pair loses
updates: a concurrent worker can push attempts past the cap between the
two statements, and the losing writer then resurrects the entry to
pending, so it retries forever. §8.2 makes concurrent workers the
default, not a hypothetical.

bytes_sent is the resume checkpoint and survives failure deliberately —
a dropped link does not invalidate the bytes the hub already accepted,
and restarting a large file from zero is exactly what a contested link
cannot afford. It is clamped to size_bytes in SQL because PendingBytes
is a SUM across the backlog, so one bad offset would go negative and
cancel out other files' real pending bytes.

Timestamps are UTC-normalized on write and read: go-sqlite3 serializes
time.Time with its offset intact, so the same instant in two zones
yields two different strings and breaks both SQL equality and ORDER BY.

PruneSynced batches at 1000 rows ordered by primary key, checks
ctx.Err() between batches, and runs no vacuum. A second index
(state, synced_at, hub_id) keeps its cutoff query and the status
endpoint's newest-synced lookup off full scans and temp B-trees —
both verified with EXPLAIN QUERY PLAN.

No production callers yet: nothing constructs the ledger, so its schema
is never created. Wiring lands with the /sync/run endpoint (PR 8).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@xe-nvdk
xe-nvdk merged commit cdc0b46 into main Aug 6, 2026
4 checks passed
@xe-nvdk
xe-nvdk deleted the feat/edge-sync-ledger branch August 6, 2026 22:59
xe-nvdk added a commit that referenced this pull request Aug 7, 2026
Completes the edge-sync story: an edge Arc now pushes its Parquet files
to a hub on demand. The hub receive side, ledger, transport interface,
HMAC scheme, reconcile index, and spoke registry landed in #570-#576;
this is the client that drives them.

A pass recovers transfers interrupted by a crash, discovers new files,
reconciles the backlog in one round-trip, then streams what the hub
lacks — newest first, so a contact window that closes mid-backlog has
already delivered the freshest telemetry. It pages until the backlog
drains, so one pass on a spoke returning from a long outage moves
everything rather than the first batch.

Operator surface (admin-only, on /api/v1/spoke-sync):
  POST /run      run one pass, return what it did
  GET  /status   pending/synced/failed counts and sync lag
  GET  /ledger   per-file state, attempts, and last error

Config lives in [edge_sync.spoke]; the secret is environment-only via
ARC_EDGE_SYNC_SPOKE_SECRET. This is the manual form and is OSS; the
scheduled agent is Enterprise and lands later.

Bugs found and fixed during review, each with a regression test verified
to fail against the pre-fix code:

- Secret guard used v.IsSet, which consults the environment under
  AutomaticEnv — so Arc refused to start in the one configuration the
  guard exists to require. Now v.InConfig, which reads only the file.
- Run fetched a single page: it sent batch_size files, reported success,
  and silently stranded the rest.
- A negative batch_size reached make()'s capacity argument and panicked
  the spoke on its first pass. Clamped in NewAgent, rejected at load.
- Fiber's Group().Use() matches by string prefix, not path segment, so
  the hub's group at /api/v1/sync also matched /api/v1/sync-spoke/*, and
  its body limit ran on the operator routes. Moved to /api/v1/spoke-sync,
  with a test pinning the property rather than the name.
- The ledger view showed only pending entries, hiding the exhausted files
  it exists to diagnose. Added Ledger.Unfinished.
- Reconcile-path conflicts stayed pending and rode along in every future
  reconcile payload forever. Added Ledger.MarkConflicted, since MarkFailed
  requires in_flight and silently matched no rows.

Also: max_concurrent is bounded (64), cancellation is checked before the
select rather than as a case (both being ready made it a coin flip),
hub_url is parsed rather than prefix-matched, and every new key has a
v.SetDefault.

Test plan:
- [x] 37 new tests across agent, spoke API, e2e, and config
- [x] e2e drives the real agent against real hub handlers over real HTTP
- [x] go test -race ./internal/edgesync/ — clean
- [x] two live processes: 20 files, batch_size=7, all byte-identical on
      the hub, idempotent re-run, no staging residue
- [x] non-default max_attempts/max_concurrent/batch_size exercised
- [x] auth boundary: forged MAC 401, unknown spoke 401, wrong secret
      writes nothing
- [x] negative batch_size refused at startup with a clear message
- [x] go vet, gofmt clean
xe-nvdk added a commit that referenced this pull request Aug 7, 2026
Completes the edge-sync story: an edge Arc now pushes its Parquet files
to a hub on demand. The hub receive side, ledger, transport interface,
HMAC scheme, reconcile index, and spoke registry landed in #570-#576;
this is the client that drives them.

A pass recovers transfers interrupted by a crash, discovers new files,
reconciles the backlog in one round-trip, then streams what the hub
lacks — newest first, so a contact window that closes mid-backlog has
already delivered the freshest telemetry. It pages until the backlog
drains, so one pass on a spoke returning from a long outage moves
everything rather than the first batch.

Operator surface (admin-only, on /api/v1/spoke-sync):
  POST /run      run one pass, return what it did
  GET  /status   pending/synced/failed counts and sync lag
  GET  /ledger   per-file state, attempts, and last error

Config lives in [edge_sync.spoke]; the secret is environment-only via
ARC_EDGE_SYNC_SPOKE_SECRET. This is the manual form and is OSS; the
scheduled agent is Enterprise and lands later.

Bugs found and fixed during review, each with a regression test verified
to fail against the pre-fix code:

- Secret guard used v.IsSet, which consults the environment under
  AutomaticEnv — so Arc refused to start in the one configuration the
  guard exists to require. Now v.InConfig, which reads only the file.
- Run fetched a single page: it sent batch_size files, reported success,
  and silently stranded the rest.
- A negative batch_size reached make()'s capacity argument and panicked
  the spoke on its first pass. Clamped in NewAgent, rejected at load.
- Fiber's Group().Use() matches by string prefix, not path segment, so
  the hub's group at /api/v1/sync also matched /api/v1/sync-spoke/*, and
  its body limit ran on the operator routes. Moved to /api/v1/spoke-sync,
  with a test pinning the property rather than the name.
- The ledger view showed only pending entries, hiding the exhausted files
  it exists to diagnose. Added Ledger.Unfinished.
- Reconcile-path conflicts stayed pending and rode along in every future
  reconcile payload forever. Added Ledger.MarkConflicted, since MarkFailed
  requires in_flight and silently matched no rows.

Also: max_concurrent is bounded (64), cancellation is checked before the
select rather than as a case (both being ready made it a coin flip),
hub_url is parsed rather than prefix-matched, and every new key has a
v.SetDefault.

Test plan:
- [x] 37 new tests across agent, spoke API, e2e, and config
- [x] e2e drives the real agent against real hub handlers over real HTTP
- [x] go test -race ./internal/edgesync/ — clean
- [x] two live processes: 20 files, batch_size=7, all byte-identical on
      the hub, idempotent re-run, no staging residue
- [x] non-default max_attempts/max_concurrent/batch_size exercised
- [x] auth boundary: forged MAC 401, unknown spoke 401, wrong secret
      writes nothing
- [x] negative batch_size refused at startup with a clear message
- [x] go vet, gofmt clean
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant