Skip to content

feat(edgesync): hub receive endpoint with verify-before-commit (#569 PR 4/9) - #574

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

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

Conversation

@xe-nvdk

@xe-nvdk xe-nvdk commented Aug 7, 2026 •

Copy link
Copy Markdown
Member

Fourth of nine PRs for edge sync — see #569. Follows #570 (ledger), #571 (rename), #572 (transport), #573 (HMAC).

First PR in this sequence with a real HTTP endpoint, so it carries the auth non-negotiables and the binary-run gate.

Summary

POST /api/v1/sync/file accepts an authenticated file push from a spoke. Bytes stream into a staging area while Arc hashes them, the digest is compared against what the spoke declared, and only on a match is the file promoted to its final namespaced location. A mismatch is discarded, so corrupt bytes never appear where a reader would find them.

New config: edge_sync.enabled (default false), edge_sync.hub_id, edge_sync.max_file_bytes (512MiB default).

A design constraint §5.2 did not anticipate

The design specifies "atomic-rename into the final namespaced path". storage.Backend has no rename, and resume needs AppendingBackend, which only LocalBackend implements.

Resolution (recorded on the issue): promote by streaming copy rather than rename — the same pipe pattern tiering/migrator.go already uses. Verify-before-commit is preserved on every backend; the cost is an extra read+write on object storage. Resume degrades rather than failing: where the backend cannot append, a dropped transfer restarts from zero, matching what the peer-fetch puller already does for ErrResumeNotSupported.

An earlier version of this analysis concluded the hub had to be local-only, mirroring Iceberg. That was wrong — it confused "no rename" with "no verify-before-commit" and would have refused a working configuration.

Two blockers found by adversarial review

Both reproduced before fixing; both have regression tests that fail against the pre-fix code.

A failed promote wedged the path permanently. LocalBackend.StatFile falls back to {path}.part when the final file is absent; ReadTo does not. An interrupted promote left a .part that the identity check read as "the file exists", and every retry then failed hashing it — forever, surviving restart, with no cleanup path anywhere. Fixed by using Exists (not .part-aware) and clearing stale .part at the destination before promoting.

A pre-auth memory DoS. Fiber buffers the body before routing — so before auth — and server.max_payload_size defaults to 1GB. Anyone able to reach the port could pin 1GB per connection with no token and no spoke secret. edge_sync.max_file_bytes is now checked on Content-Length and on the declared size, and config load refuses a cap above the server limit because the larger value could never take effect.

Other findings fixed

  • The hub returned a result its own validator rejects — a staged prefix longer than a re-declared size gave BytesAccepted > SizeBytes, which PutResult.Validate refuses, so a spoke would hard-error on a well-formed reply.
  • A Raft election told the spoke its request was malformed — every non-resume error mapped to 400. Hub-side failures are now 503 with no internal error text.
  • AlreadyPresent never re-attempted registration — a file whose manifest write failed stayed on disk but invisible to every reader permanently, because the spoke marks it synced and never sends it again.
  • .. was only rejected as a whole segment, but LocalBackend replaces every .. substring with _, so a..b.parquet and a_b.parquet collided on disk and produced an unresolvable 409 naming a digest the spoke never uploaded.
  • PartitionTime was unset in the cluster wiring (found before either reviewer reported it) — tiering/router.go filters time-ranged queries by it, so synced files would have vanished from those queries.
  • Duplicate component JSON key in every log line.
  • SweepStaging added — staging had no reclamation path at all, so an abandoned transfer leaked disk forever.

Configuration matrix

Configuration Reaches new code? Preconditions established?
OSS standalone, enabled=false (default) No — block skipped, no routes n/a. Verified live: endpoint 404
OSS standalone, enabled=true + hub_id Yes clusterCoordinator nil → registerFile nil, guarded in Receive. Verified live
hub_id empty No — config load fails Verified live: refuses to start
hub_id with / or NUL No — config load fails Verified live
max_file_bytes > server.max_payload_size No — config load fails Verified live
Cluster + enabled=true Yes — registerFile proposes to Raft RegisterFileInManifest handles leader and follower (forwards). Not exercised live
auth.enabled=true Yes — RequireAdmin mounted authManager nil-guarded
S3/Azure hub Yes — SupportsResume() false Restarts from zero. Not exercised live

Two rows not exercised live; both reviewers were pointed at them specifically.

Test plan

  • go build ./cmd/... ./internal/..., go vet, gofmt -l clean
  • go test + -race clean across edgesync, api, config
  • Binary runs at non-default values: enabled + hub_id, max_file_bytes=256MiB, plus three configs that must refuse to start — all verified
  • Real end-to-end over HTTP: authenticated upload committed and landed under rocket-01/metrics/...; replay refused with reason: replay; unauthenticated and forged requests rejected with nothing written; oversized upload 413
  • Every new regression test verified to fail pre-fix

Review

Security-focused pass, deep pass, and a third review. Both of the first two independently found the wedge.

Notes

  • Docs: new docs/advanced/edge-sync.md, with a limitations section stating plainly that no spoke can be registered yet.
  • Deliberately shipping enable-able but not usable: spoke registration lands in PR 7, so an enabled hub rejects every request and warns at startup — visible and fails closed.
  • SweepStaging is unwired: scheduling it needs the agent's scheduler (PR 8).

Fourth of nine PRs for edge sync (#569). Adds the hub side of a file
transfer: POST /api/v1/sync/file accepts an authenticated push from a
spoke, verifies it, and only then makes it visible.

The ordering is the point. Bytes stream into a staging area under
.sync-staging/ while a SHA-256 runs over them; the digest is compared
against what the spoke declared; and only on a match is the file
promoted to its final namespaced location. A mismatch is discarded, so
corrupt bytes never appear where a reader or a later reconcile would
find them.

§5.2 specifies "atomic-rename into the final path", which storage.Backend
cannot express — there is no rename, and resume needs AppendingBackend,
which only LocalBackend implements. Promotion is therefore a streaming
copy (the same pipe pattern tiering/migrator.go uses), which preserves
verify-before-commit on every backend at the cost of an extra read+write
on object storage. Resume degrades rather than failing: where the
backend cannot append, a dropped transfer restarts from zero, matching
what the peer-fetch puller already does for ErrResumeNotSupported. The
reasoning is recorded on #569 so it is not re-derived.

Two independent auth layers, both mandatory: Arc's API token middleware
on the route group, and a per-spoke HMAC binding spoke, hub, path, and
content digest. The handler constructor refuses to start without any
security dependency — a nil secret lookup or replay guard would silently
downgrade authentication, and a handler that half-works is worse than
one that will not start.

Adversarial review found two blockers, both fixed with regression tests
that fail against the pre-fix code:

LocalBackend.StatFile falls back to the "{path}.part" staging file when
the final file is absent; ReadTo does not. An interrupted promote left a
.part that the identity check read as "the file exists", and every retry
then failed hashing it — permanently, surviving restart, with no cleanup
path. The check now uses Exists, and promote clears a stale .part at the
destination first.

The Fiber app buffers request bodies before routing, so before auth runs,
and server.max_payload_size defaults to 1GB — anyone able to reach the
port could pin 1GB per connection with no token and no spoke secret.
edge_sync.max_file_bytes (512MiB default) is checked on Content-Length
and on the declared size, and config load refuses a cap above the server
limit because the larger value could never take effect.

Also fixed: a staged prefix longer than a re-declared size produced
BytesAccepted > SizeBytes, which PutResult.Validate rejects, so a spoke
would hard-error on the hub's own reply. A manifest write failing during
a Raft election mapped to 400, telling the spoke its request was
malformed when it should retry — hub-side failures are now 503. And
AlreadyPresent now re-attempts registration: without it, a file whose
manifest write failed stayed on disk but invisible to every reader
forever, because the spoke marks it synced and never sends it again.

Path validation rejects ".." anywhere rather than only as a whole
segment: LocalBackend replaces every ".." substring with "_", so
"a..b.parquet" and "a_b.parquet" collided on disk and produced an
unresolvable 409 naming a digest the spoke never uploaded.

SweepStaging reclaims abandoned partials but is not yet scheduled —
wiring it needs the agent's scheduler (PR 8). Spoke registration lands
in PR 7, so an enabled hub currently rejects every request and warns at
startup: visible and fails closed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@xe-nvdk
xe-nvdk merged commit 3747eeb into main Aug 7, 2026
4 checks passed
@xe-nvdk
xe-nvdk deleted the feat/edge-sync-receive branch August 7, 2026 16:51
xe-nvdk added a commit that referenced this pull request Aug 7, 2026
Combines PRs 6 and 7 of nine for edge sync (#569). PR 6's spoke
namespacing and 409-never-overwrite shipped in #574 and #575 — the
receive and reconcile endpoints could not be correct without them — so
what remained was the per-spoke secret binding, which lives here.

This makes the hub usable. Until now it rejected every request for want
of a way to register a spoke.

§8.1 specifies a `secret_hash` column. That cannot work: HMAC
verification recomputes the MAC from the secret, so the hub needs the
plaintext at request time — unlike an API token, which is only ever
checked against a value the caller presents. Secrets are therefore
encrypted at rest with AES-256-GCM under ARC_ENCRYPTION_KEY, the same
mechanism MQTT uses for broker passwords, and the column is named
secret_encrypted so the storage model cannot be misread from the schema.

The hub generates the secret and returns it once, from registration or
rotation. It is never readable again — Spoke carries no secret field, so
metadata paths cannot leak it by construction. Generating rather than
accepting one removes the failure mode where an operator or an
automation picks something weak or reuses it across a fleet.

Adversarial review found two blockers, both in claims I had made and not
properly checked.

The comment asserting the admin routes inherit the sync group's
middleware was wrong. I had tested it with the sync group registered
first; main.go registers admin first, and under that order the
middleware does not run. Worse, the behavior flips on reordering two
lines, at which point a registration whose JSON exceeds the sync body
limit would get a 413 from the wrong group before auth ran. Admin routes
now sit at /api/v1/sync-spokes, so the isolation is structural rather
than a property of registration order.

The cipher guard did not guard. mqtt.NewPasswordEncryptor returns a
non-nil pass-through for an empty key, which satisfies the interface, so
a caller that forgot to check the key would have stored every secret in
plaintext while construction reported success. NewRegistry now
round-trips a canary and refuses a cipher that does not encrypt.

Three mutations survived because the tests used a single spoke: a
RecordActivity that updates every row, a List that drops its ORDER BY,
and a scan that always reports enabled. All three are invisible with one
registration. Multi-spoke tests now catch each, and List has a stable
tiebreaker so bulk provisioning cannot produce arbitrary order.

Also: spoke names are bounded and reject control characters, since they
are echoed into logs and every list response; a decrypt failure now logs
at Warn rather than Debug, because it is a hub misconfiguration
affecting every spoke rather than the scan noise an unknown spoke ID is;
and the startup path verifies the configured key still decrypts what is
stored, refusing to start on a changed key rather than letting the hub
look healthy while nothing can authenticate.

With auth.enabled=false the admin routes are unauthenticated, matching
every other Arc admin endpoint in that mode. Review's recommendation was
to keep the pattern rather than introduce a lone inconsistency, and to
make the posture visible: startup now warns explicitly that anyone
reaching the port can mint write credentials.
xe-nvdk added a commit that referenced this pull request Sep 11, 2026
#737) (#742)

validateSpokeID rejected separators, NUL, the exact values "." and "..",
and any leading dot, but allowed ".." anywhere inside an identifier. The
local backend then resolves paths through sanitizePath, which replaces
every ".." with "_", so "rocket..01" and "rocket_01" named one directory.
Since the spoke ID is the first path segment of everything a spoke
writes, a spoke authenticated as one could write into the other's
namespace. The HMAC proves which spoke is asking; the namespace mapping
decides where it may write, and that mapping was not injective.

validateSyncPath already rejects ".." anywhere, with a comment recording
exactly this many-to-one reasoning, added by #574 for the source path.
The identity that prefixes it was left open. Same rule, same place, and
it covers all seven callers at once.

Reported privately by rexpository. Recorded as hardening rather than an
advisory because every path to it runs through an administrator creating
the colliding ID: registration is admin-only, secrets are generated
server-side and shown once, and there is no self-registration. It is
local-backend only, since S3 and Azure preserve key spelling, and edge
sync is opt-in.

The validator also gains the bounds it was missing. A 128-byte cap, no
control characters, no leading or trailing whitespace: the same rules the
registry already applies to the spoke NAME, which unlike the ID never
becomes a path component. An over-long ID registered successfully and
then failed every write with ENAMETOOLONG, which is the same "hand out a
credential that only surfaces later as an unexplained refusal" shape this
change exists to close.

Tightening a validator does nothing for rows already stored, so
Registry.InvalidStoredIDs re-checks them and the hub names each offender
at startup. Without it the only signal is an edge box failing at the far
end with no matching hub log line.

Four tests. The rejection table gains the traversal spellings plus the
new bounds. Registry.Register is covered, since a receive-side rejection
alone still lets an operator issue a credential that resolves into
somebody else's namespace. InvalidStoredIDs is covered against a legacy
row inserted behind Register's back, which is the state an upgraded hub
is actually in.

The fourth is the invariant that was violated: every accepted spoke ID
must resolve to a namespace no other accepted ID can reach, asserted
through the real local backend rather than by reasoning about the
validator, since the collision lived in the gap between them. Its first
draft was vacuous and passed with the fix reverted, because reading each
ID back straight after its own write goes through the same sanitizer as
the write and returns the bytes it just stored. Writing every candidate
before reading any of them makes the collision visible as an earlier
spoke's payload having been replaced. It now fails without the fix on two
pairs, including a..b against a_b, which no hand-written case listed.

Two follow-ups filed rather than folded in. #740: the same collision
survives through case folding and Unicode normalization on macOS and
Windows volumes, which the injective test deliberately excludes because
including it would pass on Linux CI and fail on a macOS dev machine;
rejecting mixed case would break identifiers that are legitimate today,
so it needs a deprecation path rather than a validator line. #741:
sanitizePath rewrites rather than rejects, which is the root of this
collision class, but it backs every LocalBackend operation in Arc and
needs its own audit.

Release notes and the README bug-reporter list credit the reporter. The
migration guidance says what actually has to happen: for an ID like
rocket..01 the already-synced data sits under rocket_01/, the hub's file
index is keyed on spoke ID, so a re-pointed spoke finds no index rows and
re-uploads its whole backlog. Upgrade note 2 no longer claims no
configuration change is required, because for such a spoke there is one.

Fixes #737
xe-nvdk added a commit that referenced this pull request Sep 12, 2026
…741) (#745)

sanitizePath repaired its input: every ".." became "_" and NUL bytes were
stripped, with a comment claiming this prevented directory traversal. It
did not. Containment was enforced separately by validatePath, resolving
the path and checking it still sat under the root.

What the rewrite did do was make the mapping from key to file
many-to-one. "a..b" and "a_b" were one location, as were "..foo" and
"_foo". That is the root of two bugs already fixed in this release, #574
on the edge-sync source path and #737 on the spoke ID, each closed by
teaching one caller not to send ".." and each leaving the next caller to
rediscover it. Rejecting makes the mapping injective for every caller of
this backend. Not for every backend: S3 and Azure still validate nothing,
which is #743.

Two audits informed it. Instrumenting the backend and running the full
suite found no path that the rewrite altered, and none that
filepath.Clean altered. Both are evidence, not proof, and review found
two things neither could see.

The first wedges ingest, and would have shipped. The writer builds
"{database}/{measurement}/{year}/..." by interpolation, so an empty
measurement yields "db//2026/...". filepath.Join used to collapse that,
landing the rows a level up under the database itself; this change
refuses the key instead. A refused flush keeps its data in the WAL and
retries forever, so one such record would poison that buffer key
permanently. Two producers could reach it: extractMeasurements SKIPPED
empty names, so the isValidMeasurementName loop right after it never saw
them, and MQTT took "m": "" as present because it satisfies the type
assertion, so the "mqtt" default never fired. Both now reject. MQTT
subscriptions also validate their database and every topic-mapping
target at config load, which was the one ingest surface with no handler
in front of it.

The second is stored operator config. Retention builds a prefix from
policy.Database and policy.Measurement, and neither was validated:
create checked only non-empty, update checked nothing, and the scheduler
replays the row forever. Both validate now, and rows already stored are
named at startup. The old behaviour there was worse than a rejection: an
empty database folded to the storage root, so retention enumerated every
database as a measurement and deleted aged files instance-wide.

Those API checks use a storage-segment rule, not isValidDatabaseName.
The create-time rule is much stricter than the storage contract, and the
storage root legitimately holds directories that never went through it:
an edge-sync hub writes each spoke namespace there, and validateSpokeID
permits a leading digit, interior dots and 128 bytes. Gating the
per-name routes on the create rule would have returned 400 for
directories GET /api/v1/databases had just listed, and would have made
the _internal 403 unreachable behind a 400.

A latent hole the first draft would have shipped: the validator scans
"/" segments while the join used filepath.Separator, so on a platform
whose OS also splits on backslash a segment built from them would pass
the scan whole and then escape. Arc builds only linux and darwin, so it
was never live. Backslash is rejected now, as edgesync already did.

Rejections carry ErrInvalidPath, because a permanent error delivered
into loops written for transient ones is how a compaction manifest gets
retried forever. No caller branches on it yet; the consumers are #744.

validateSyncPath checked its leading dot on the whole string rather than
per segment, so "db/./cpu/x.parquet" was accepted. path.Join cleaned it
before storage, but Measurement is derived from the RAW path, so "."
reached the Raft manifest and reconciliation then built "db/./" from it.

GetFullPath is deleted: no callers, and it swallowed the validation
error, which is the behaviour this change exists to remove. The eight
surviving copies of "validate and sanitize the path to prevent path
traversal" are corrected, since that comment is the lie the change is
about, and the two edgesync comments justifying their rules by the fold
now say why the rules remain without it.

Performance, since validatePath runs on every local read and write. Once
the path is known clean and relative the three path/filepath calls after
it were redundant, so the check is a single pass with no allocation:

  BenchmarkValidatePath      605 ns/op -> 110 ns/op   1 alloc, unchanged
  BenchmarkValidatePathDeep  757 ns/op -> 150 ns/op   1 alloc, unchanged

That is roughly 0.06% of a 4 MiB Parquet write and about half of a bare
existence check, so it is real in the stat-heavy reconciliation and
tiering loops and noise on ingest. An earlier draft cited 1565 ns/op as
the baseline; that measurement was noisy and 605 is correct. The change
is carried on correctness, not on the speedup.

Containment is guaranteed by construction rather than recomputed, so the
proof is a property test instead of a filepath.Rel call per request. Its
generator is weighted to produce mostly ACCEPTED keys and asserts it
reached 10,000 of them across 2,000 distinct values, because a generator
that mostly rejects skips almost every iteration and asserts nothing.
Every accepted key is also checked against filepath.Join.

Fixes #741
xe-nvdk added a commit that referenced this pull request Sep 12, 2026
…) (#748)

#741 made LocalBackend refuse malformed keys and claimed the mapping from
key to object was now injective "for every caller at once". That held for
one implementation out of three. S3 prepended its prefix with no
validation and Azure passed the key straight to NewBlockBlobClient, so
one key behaved three ways and the next reporter would refile #737
against whichever backend they ran.

The rule is now a property of the Backend interface, and what it
guarantees is the thing whose absence produced #574, #737 and #741: two
different keys can never name one object.

It also fixes a collision #741 introduced. "coll" and "coll/" resolved to
one local file: two keys, one object, first write silently lost.
checkStoragePath accepted a trailing separator because ~30 List callers
pass db+"/" and 13 pass "", and validatePath then stripped it. That
exemption is right for a prefix and wrong for a key, so ValidateKey and
ValidateListPrefix are now separate.

The contract test changed shape too. #741 asserted matching accept and
reject sets, which cannot see this: every backend accepts "coll/", and
what differed was the mapping, not the verdict. The property is
injectivity, so the test is injectivity.

Everything about the object stores was measured against a live MinIO and
Azurite. The first draft was wrong because I reasoned instead:

  - A leading separator is refused. MinIO STRIPS it, so "/a/x" and "a/x"
    were one object there, while S3 and Azure keep them distinct. I first
    wrote this up as a collision live on S3 and Azure; it is MinIO's
    alone, and on S3 the failure is an orphan rather than an overwrite.
  - A backslash is refused because Azure treats it as a separator: "a\b"
    and "a/b" are ONE blob. That is the live collision, and it is not the
    one I first found.
  - "." and ".." segments and empty interior segments are refused
    locally, naming the key, instead of as a remote 400 describing an S3
    API error. Azure stored them literally, so there it is a real
    behaviour change.
  - "a..b" and "..foo" are accepted everywhere and stored as asked.

A listing never returns a key the backend would refuse. Object stores
carry directory-marker objects ending in a separator, written by consoles
and sync tools, and Arc feeds List output straight into Read, Exists and
Delete. Returning one is worse than the original bug: restore skips the
file and still reports success (backup/restore.go:124), compaction's
"already compacted" escape becomes a hard failure (job.go:605), and
manifest recovery retries a permanent error forever (manifest.go:309).
Filtering them makes List output always usable, which is what those
callers already assumed. Found by writing a marker through the raw S3 API
and watching a read of it fail.

SanitizeS3Prefix becomes ValidateS3Prefix. Its damage was on the SUCCESS
path, so returning an error instead of "" would have touched none of it:
"/" stayed "/", "a//b" became "a//b/" and "." became "./", each of which
400s every write, while "a/..b" was replaced with "", which is the bucket
root rather than a safe default. A prefix that cannot form usable keys
now stops the backend from being built, which is a startup-breaking
upgrade and is called out in the upgrade notes.

prefixedKey returns an error so a new S3 method cannot skip validation.
Not a total chokepoint, and the comment says so: the read path builds
s3:// URIs in storage/util.go for DuckDB read_parquet, and iceberg-go
writes metadata through its own FileIO. Filed as #746 along with the five
unreconciled path validators.

raft.ValidateManifestPath is deliberately NOT delegated to ValidateKey,
and now says why. It runs inside FSM Apply, which includes log replay, so
tightening it would make a node reject an entry an older binary accepted
and two versions would build different state from one log. The gap closes
at the consumer instead: #747 tracks the four loops that should
quarantine on ErrInvalidPath rather than retry it.

DeleteBatch stays best-effort. Two callers return its error with no
per-file fallback, so failing a whole batch on one bad key would make a
backup permanently undeletable and leave compaction inputs beside their
output.

Two Iceberg callers would have started failing: exporter.go passes "./"
when a metadata location sits directly under the warehouse root and "/"
when warehouseRelKey returns empty. Both listed empty before and are
skipped now; the second would otherwise have failed DROP DATABASE.

ValidateKey gains the bounds it never had, 1024 bytes overall and 255 per
segment, and validateSyncPath gains the same minus the ".part" staging
suffix, so a spoke cannot send a path that passes both validators and
then fails only once that suffix is appended.

CI now runs a MinIO and the tagged contract tests against it. None of the
behaviour above is reproducible in a unit test, because it is a property
of the servers. docker run rather than a service container, since
quay.io/minio/minio needs `server /data` and services cannot supply
arguments.

Fixes #743
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