Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 31 additions & 4 deletions RELEASE_NOTES_2026.09.1.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ This is a **file-count** bound, not a byte bound — compacted output size track

Out-of-range values fall back to the default with a startup warning rather than failing. **`1` is not usable** and is treated as out of range: compaction's adaptive retry rejects any batch below two files, so a batch size of one would fail every batch of every partition.

## Groundwork: edge-to-cloud sync (not yet usable)
## New: edge-to-cloud sync (manual)

Arc runs at the edge today — a standalone binary with local storage in a vehicle, a factory cell, or a forward deployment — with no first-class way to ship that data to a central Arc. Backup is a full DR snapshot rather than incremental sync, and Parquet import re-ingests rows through the write path, which breaks the end-to-end checksum and double-counts on retry.

Expand All @@ -88,9 +88,36 @@ Reconcile is answered from a hub-side index of received files rather than by rea

Resume is supported on local storage. On S3 and Azure a dropped transfer restarts from zero, because block objects cannot be appended to; this is a throughput cost on intermittent links, not a correctness problem.

The rest of the work is still internal: the spoke-side ledger, the `SyncTransport` interface, and the HMAC scheme. Manual export/import will be OSS when it ships; the automatic scheduled agent will be an Enterprise feature.
### The spoke side

Manual export/import will be OSS when it ships; the automatic scheduled agent will be an Enterprise feature.
An edge Arc now syncs to a hub on demand. Enable it on the spoke:

```toml
[edge_sync.spoke]
enabled = true # default false
hub_url = "https://hub.example.com" # required when enabled
spoke_id = "rocket-01" # this spoke's ID, as registered on the hub
hub_id = "ground-station" # the REMOTE hub's edge_sync.hub_id
max_attempts = 5 # attempts before a file is marked failed
max_concurrent = 2 # simultaneous transfers
batch_size = 0 # files per reconcile round-trip; 0 = hub default
```

The secret is **environment-only**: `ARC_EDGE_SYNC_SPOKE_SECRET`, the value the hub returned once at registration. A secret in the config file is **refused at startup** rather than ignored — one that is ignored still leaks, and leaving it in place makes the committed copy look load-bearing. `hub_id` is validated at load for the same reason it is easy to get wrong: it is bound into every request MAC, so a mismatch fails *every* request with a `400` that looks like a hub problem.

Three admin endpoints drive it:

| Endpoint | Purpose |
|---|---|
| `POST /api/v1/spoke-sync/run` | Run one sync pass and return what it did. |
| `GET /api/v1/spoke-sync/status` | Pending/synced/failed counts and sync lag. |
| `GET /api/v1/spoke-sync/ledger` | Per-file state, attempts, and last error. |

A pass recovers transfers interrupted by a crash, discovers new files, reconciles the backlog, and 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, not just the first batch. Conflicts are reported in full rather than counted and are not retried: the same path holding different content means a spoke-ID collision or corruption, and re-sending would either be refused or destroy evidence.

Files are hashed once at discovery, and the ledger survives restarts, so a spoke re-run after a crash neither re-hashes nor re-sends what already landed. Nothing is deleted from the spoke — sync is a copy, and local retention stays in the operator's hands.

This is the **manual** form, and it is OSS. The automatic scheduled agent will be an Enterprise feature.

## Security hardening

Expand Down Expand Up @@ -373,7 +400,7 @@ The forwarding-header and loop-guard changes are cluster-only; the partition-pru
2. **No API or on-disk format changes.** Reads, queries, and storage layout are untouched. Queries with a start date earlier than `1970-01-01` are pruned from the epoch forward; since Arc stores no pre-epoch data, results are unchanged.
3. **Clustered operators:** if any external tooling deliberately sets `X-Arc-Forwarded-By`, `X-Real-IP`, or `X-Forwarded-*` headers on requests to Arc and expects them to survive an inter-node forward, note that these are now stripped on the forwarding hop and re-established by Arc. Client IP has never been derived from these headers, so log/audit attribution is unchanged.
4. **Active licenses keep working.** No re-activation required.
5. **Edge sync is off by default and nothing changes unless you enable it.** With `edge_sync.enabled=false` (the default) no routes are mounted and `/api/v1/sync/file` returns 404. The `sync_ledger` and `sync_history` tables are still not created — the spoke side is not yet wired into startup.
5. **Edge sync is off by default and nothing changes unless you enable it.** With `edge_sync.enabled=false` (the default) the hub mounts no routes and `/api/v1/sync/file` returns 404. With `edge_sync.spoke.enabled=false` (the default) no spoke routes are mounted, no `sync_ledger` or `sync_history` tables are created, and nothing is read from local storage. Both sides are independent: a hub need not be a spoke, and a spoke need not be a hub.

## Dependencies

Expand Down
64 changes: 64 additions & 0 deletions cmd/arc/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -2107,6 +2107,70 @@ func main() {
}
}

// Register the edge-sync SPOKE side (#569) — this instance pushing its
// files to a hub. Independent of the hub block above: an Arc can be a hub,
// a spoke, both, or neither.
//
// The secret comes from ARC_EDGE_SYNC_SPOKE_SECRET only; config load
// already refused a config-file secret and verified every required field,
// so by here the configuration is known good.
if cfg.EdgeSync.Spoke.Enabled {
spokeLogger := logger.Get("edgesync-spoke")

// The ledger shares the auth database, like the hub index. Reuse the
// handle opened for the hub block when both roles are enabled on one
// instance, rather than opening a second one on the same file.
spokeDB, spokeOwnsDB, err := sharedSQLiteHandle(authManager, cfg.Auth.DBPath)
if err != nil {
log.Fatal().Err(err).Str("path", cfg.Auth.DBPath).Msg("Failed to open the edge sync spoke database; refusing to start")
}
if spokeOwnsDB {
shutdownCoordinator.RegisterHook("edgesync-spoke-db", func(context.Context) error {
return spokeDB.Close()
}, shutdown.PriorityDatabase)
}

spokeLedger, err := edgesync.NewLedger(spokeDB, spokeLogger)
if err != nil {
log.Fatal().Err(err).Msg("Failed to create the edge sync ledger; refusing to start")
}

spokeTransport, err := edgesync.NewHTTPTransport(edgesync.HTTPTransportConfig{
BaseURL: cfg.EdgeSync.Spoke.HubURL,
SpokeID: cfg.EdgeSync.Spoke.SpokeID,
Secret: cfg.EdgeSync.Spoke.Secret,
})
if err != nil {
log.Fatal().Err(err).Msg("Failed to create the edge sync transport; refusing to start")
}

syncAgent, err := edgesync.NewAgent(edgesync.AgentConfig{
Ledger: spokeLedger,
Transport: spokeTransport,
Backend: storageBackend,
HubID: cfg.EdgeSync.Spoke.HubID,
SpokeID: cfg.EdgeSync.Spoke.SpokeID,
MaxAttempts: cfg.EdgeSync.Spoke.MaxAttempts,
MaxConcurrent: cfg.EdgeSync.Spoke.MaxConcurrent,
BatchSize: cfg.EdgeSync.Spoke.BatchSize,
Logger: spokeLogger,
})
if err != nil {
log.Fatal().Err(err).Msg("Failed to create the edge sync agent; refusing to start")
}

spokeHandler, err := api.NewEdgeSyncSpokeHandler(syncAgent, authManager, spokeLogger)
if err != nil {
log.Fatal().Err(err).Msg("Failed to create the edge sync spoke handler; refusing to start")
}
spokeHandler.RegisterRoutes(server.GetApp())

spokeLogger.Info().
Str("hub_url", cfg.EdgeSync.Spoke.HubURL).
Str("spoke_id", cfg.EdgeSync.Spoke.SpokeID).
Msg("Edge sync spoke enabled; trigger a pass with POST /api/v1/spoke-sync/run")
}

// Register TLE handler (streaming TLE ingestion)
tleHandler := api.NewTLEHandler(arrowBuffer, logger.Get("tle"))
if authManager != nil && rbacManager != nil {
Expand Down
214 changes: 214 additions & 0 deletions internal/api/edgesync_e2e_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,214 @@
package api

import (
"context"
"crypto/sha256"
"database/sql"
"encoding/hex"
"fmt"
"io"
"net/http"
"net/http/httptest"
"os"
"testing"

"github.com/basekick-labs/arc/internal/cluster/security"
"github.com/basekick-labs/arc/internal/edgesync"
"github.com/basekick-labs/arc/internal/storage"
"github.com/gofiber/fiber/v2"
"github.com/rs/zerolog"
)

// TestEdgeSync_SpokeToHubEndToEnd drives the real spoke agent against the real
// hub handlers over real HTTP.
//
// Every other test in this sequence exercises one side against an in-process
// stand-in. This is the only one that proves the two halves interoperate — the
// header names, the HMAC field order, the status-code mapping, and the resume
// offsets all have to agree, and each was written from the design doc rather
// than from the other side's code.
func TestEdgeSync_SpokeToHubEndToEnd(t *testing.T) {
ctx := context.Background()

// --- Hub ---
hubDir, err := os.MkdirTemp("", "e2e-hub-*")
if err != nil {
t.Fatalf("hub dir: %v", err)
}
t.Cleanup(func() { os.RemoveAll(hubDir) })

hubBackend, err := storage.NewLocalBackend(hubDir, zerolog.Nop())
if err != nil {
t.Fatalf("hub backend: %v", err)
}
t.Cleanup(func() { hubBackend.Close() })

hubIndex := newTestAPIHubIndex(t)
receiver, err := edgesync.NewReceiver(edgesync.ReceiverConfig{
Backend: hubBackend, Index: hubIndex, Logger: zerolog.Nop(),
})
if err != nil {
t.Fatalf("receiver: %v", err)
}
reconciler, err := edgesync.NewReconciler(edgesync.ReconcilerConfig{
Index: hubIndex, Backend: hubBackend, MaxEntries: 100,
})
if err != nil {
t.Fatalf("reconciler: %v", err)
}

const spokeID, hubID, secret = "rocket-01", "ground-station", "e2e-shared-secret"
handler, err := NewEdgeSyncHandler(EdgeSyncHandlerConfig{
Receiver: receiver,
Reconciler: reconciler,
SpokeSecrets: StaticSpokeSecrets(map[string]string{spokeID: secret}),
Replay: security.NewNonceCache(security.HMACTimestampTolerance),
HubID: hubID,
MaxFileBytes: 8 << 20,
Logger: zerolog.Nop(),
})
if err != nil {
t.Fatalf("handler: %v", err)
}

app := fiber.New(fiber.Config{DisableStartupMessage: true, BodyLimit: 32 << 20})
handler.RegisterRoutes(app)

// A real HTTP listener, so the transport exercises actual network I/O
// rather than Fiber's in-process test harness.
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
resp, err := app.Test(r, testRequestTimeoutMS)
if err != nil {
w.WriteHeader(http.StatusInternalServerError)
return
}
defer resp.Body.Close()
for k, vs := range resp.Header {
for _, v := range vs {
w.Header().Add(k, v)
}
}
w.WriteHeader(resp.StatusCode)
_, _ = copyBody(w, resp.Body)
}))
t.Cleanup(srv.Close)

// --- Spoke ---
spokeDir, err := os.MkdirTemp("", "e2e-spoke-*")
if err != nil {
t.Fatalf("spoke dir: %v", err)
}
t.Cleanup(func() { os.RemoveAll(spokeDir) })

spokeBackend, err := storage.NewLocalBackend(spokeDir, zerolog.Nop())
if err != nil {
t.Fatalf("spoke backend: %v", err)
}
t.Cleanup(func() { spokeBackend.Close() })

ledgerDB, err := sql.Open("sqlite3", spokeDir+"/ledger.db")
if err != nil {
t.Fatalf("ledger db: %v", err)
}
t.Cleanup(func() { ledgerDB.Close() })
ledger, err := edgesync.NewLedger(ledgerDB, zerolog.Nop())
if err != nil {
t.Fatalf("ledger: %v", err)
}
transport, err := edgesync.NewHTTPTransport(edgesync.HTTPTransportConfig{
BaseURL: srv.URL, SpokeID: spokeID, Secret: secret,
})
if err != nil {
t.Fatalf("transport: %v", err)
}
agent, err := edgesync.NewAgent(edgesync.AgentConfig{
Ledger: ledger, Transport: transport, Backend: spokeBackend,
HubID: hubID, SpokeID: spokeID, Logger: zerolog.Nop(),
})
if err != nil {
t.Fatalf("agent: %v", err)
}

// --- The spoke produces three files ---
contents := map[string][]byte{}
for i := 0; i < 3; i++ {
p := fmt.Sprintf("metrics/cpu/2026/08/07/1%d/cpu_%d.parquet", i, i)
c := []byte(fmt.Sprintf("parquet payload number %d", i))
contents[p] = c
if err := spokeBackend.Write(ctx, p, c); err != nil {
t.Fatalf("write %s: %v", p, err)
}
}

res, err := agent.Run(ctx)
if err != nil {
t.Fatalf("first sync: %v", err)
}
if res.Discovered != 3 {
t.Errorf("discovered = %d, want 3", res.Discovered)
}
if res.Sent != 3 {
t.Fatalf("sent = %d, want 3 (failed=%d partial=%d)", res.Sent, res.Failed, res.Partial)
}

// The hub must hold each file under the spoke's namespace, byte-identical.
for p, want := range contents {
got, err := hubBackend.Read(ctx, edgesync.NamespacedPath(spokeID, p))
if err != nil {
t.Errorf("hub is missing %s: %v", p, err)
continue
}
if string(got) != string(want) {
t.Errorf("%s: hub content differs from the spoke's", p)
}
// The digest the hub indexed must be the file's real one, or reconcile
// would report it missing on the next pass.
wantSHA := sha256.Sum256(want)
held, err := hubIndex.Lookup(ctx, spokeID, []string{p})
if err != nil {
t.Errorf("%s: index lookup: %v", p, err)
} else if held[p] != hex.EncodeToString(wantSHA[:]) {
t.Errorf("%s: hub indexed digest %q, want %q", p, held[p], hex.EncodeToString(wantSHA[:]))
}
}

// --- A second pass must be free ---
res2, err := agent.Run(ctx)
if err != nil {
t.Fatalf("second sync: %v", err)
}
if res2.Sent != 0 || res2.BytesSent != 0 {
t.Errorf("second pass sent %d files / %d bytes; an already-synced corpus must cost nothing",
res2.Sent, res2.BytesSent)
}

// --- A new file syncs incrementally ---
newPath := "metrics/cpu/2026/08/07/19/cpu_new.parquet"
newContent := []byte("a freshly compacted file")
if err := spokeBackend.Write(ctx, newPath, newContent); err != nil {
t.Fatalf("write new: %v", err)
}
res3, err := agent.Run(ctx)
if err != nil {
t.Fatalf("third sync: %v", err)
}
if res3.Discovered != 1 || res3.Sent != 1 {
t.Errorf("incremental pass: discovered=%d sent=%d, want 1/1", res3.Discovered, res3.Sent)
}

// --- Status reflects reality ---
st, err := agent.Status(ctx)
if err != nil {
t.Fatalf("status: %v", err)
}
if st.Pending != 0 {
t.Errorf("pending = %d after a full sync, want 0", st.Pending)
}
if st.Synced != 4 {
t.Errorf("synced = %d, want 4", st.Synced)
}
}

func copyBody(w http.ResponseWriter, r io.Reader) (int64, error) {
return io.Copy(w, r)
}
Loading