Repository navigation
feat(sync): spoke-side ledger for edge-to-cloud replication (#569 PR 1/9) - #570
Merged
Merged
Conversation
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>
9 tasks done
This was referenced Aug 6, 2026
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
10 tasks done
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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
First of nine PRs for edge sync — see #569 for the phase-1 breakdown and
docs/progress/2026-06-04-edge-sync-architecture-converged.mdfor the design.Summary
Adds
internal/syncwith the spoke-side ledger: the durable record of which Parquet files have reached which hub, and how far a partial transfer got.sync_ledgerkeyedUNIQUE(hub_id, path)from day one — multi-hub becomes config later, not a schema migration on a disconnected box.MarkFailed's retry cap decided inside a single UPDATE — no read-then-write race.bytes_sentas the resume checkpoint, preserved across failure, clamped tosize_bytes.PruneSyncedbatched at 1000 rows, PK-ordered,ctx.Err()between batches, no vacuum.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.FileEntrycarriesSHA256/SizeBytesper 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
grep NewLedgeroutside the package returns nothingcmd/arc/main.goledger_test.gois the sole callersql.Open,t.CleanupclosesEvery 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:
NewLedgerrejects 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:
MarkFailed— read-then-write let a concurrent worker pushattemptspast the cap between statements; the losing writer resurrected the entry topending, retrying forever. Now one atomicUPDATEwith theCASEinside."2026-08-06 22:40:43.006808+00:00", space-separated with an offset, neither RFC3339 nordatetime()output. The comment pointed a future maintainer the wrong way.discovered_atnot UTC-normalized — probed: the same instant stored two ways gave SQL equality0.MarkInFlighton a synced row succeeded withsynced_atstill set, double-counting inStats.PendingByteswent negative (−4900 with an over-large offset), and since it is aSUM, one bad offset silently cancelled other files' real pending bytes.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
databaseas a column name doesn't actually need quoting (probed unquoted inWHERE/SELECT/GROUP BY/ORDER BY; four existing Arc tables already do this).The second pass found no bugs. It suggested an explicit
pending → syncedreconcile-path test, which is added.I also found one of my own tests was vacuous — the
TrackBatchrollback test never triggered a rollback because empty strings satisfyNOT 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 testsgo test -race ./internal/sync/go vet/gofmt -lcleanEXPLAIN QUERY PLAN(no full scans, no temp B-trees)Notes
docs.basekick.netchange — nothing here is configurable or observable. Docs land with PR 8's endpoints.internal/syncshadows stdlibsync; importers needingsync.Mutexmust alias. Renaming toedgesyncis cheap now and progressively more annoying as PRs 2–9 stack. Happy to do it before PR 2.RecoverInFlighthas no age guard andPendinghas 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