From b9c384f87cbce23524c861567b26cac0c7ee0e97 Mon Sep 17 00:00:00 2001 From: Adesh Nalpet Adimurthy <390.adesh@gmail.com> Date: Sat, 15 Aug 2026 15:32:01 -0400 Subject: [PATCH] Raft Metadata Object-Store Archive and Restore From --- .github/workflows/e2e.yml | 43 +++ dashboard/src/api.ts | 28 ++ dashboard/src/app.tsx | 42 ++- docs/pages/deployment.md | 27 +- harness/README.md | 2 + harness/docker/entrypoint.sh | 6 + harness/local/docker-compose.ds.yml | 2 + .../metadata/raft/MetadataNode.java | 59 ++++ .../metadata/raft/MetadataStateMachine.java | 36 +++ .../raft/ObjectStorageSnapshotArchive.java | 254 ++++++++++++++++++ .../metadata/raft/SnapshotArchive.java | 27 ++ .../raft/MetadataStateMachineRestoreTest.java | 48 ++++ .../ObjectStorageSnapshotArchiveTest.java | 113 ++++++++ .../io/streamstack/server/AdminServer.java | 37 +++ .../streamstack/server/StreamStackNode.java | 28 +- .../server/model/config/ServerConfig.java | 37 +++ .../streamstack/server/AdminServerTest.java | 36 +++ ...MetadataArchiveRestoreIntegrationTest.java | 180 +++++++++++++ 18 files changed, 1002 insertions(+), 3 deletions(-) create mode 100644 metadata/src/main/java/io/streamstack/metadata/raft/ObjectStorageSnapshotArchive.java create mode 100644 metadata/src/main/java/io/streamstack/metadata/raft/SnapshotArchive.java create mode 100644 metadata/src/test/java/io/streamstack/metadata/raft/MetadataStateMachineRestoreTest.java create mode 100644 metadata/src/test/java/io/streamstack/metadata/raft/ObjectStorageSnapshotArchiveTest.java create mode 100644 server/src/test/java/io/streamstack/server/MetadataArchiveRestoreIntegrationTest.java diff --git a/.github/workflows/e2e.yml b/.github/workflows/e2e.yml index 907292e..262d523 100644 --- a/.github/workflows/e2e.yml +++ b/.github/workflows/e2e.yml @@ -54,6 +54,28 @@ jobs: java -jar cli/target/streamstack.jar ds bench --endpoint http://127.0.0.1:4437 \ -b 1024 -n 64 -w 8 -d 8 + - name: Archive metadata snapshot + run: | + curl -sf -X PUT -H 'Content-Type: application/json' -d '' http://127.0.0.1:4437/e2e/restore + curl -sf -X POST -H 'Content-Type: application/json' -d '{"payload":"restore-me"}' http://127.0.0.1:4437/e2e/restore + curl -sf -X POST http://127.0.0.1:9090/admin/snapshot | jq -e '.appliedIndex > 0' + timeout 60 bash -c "until curl -sf http://127.0.0.1:9090/admin/snapshots | jq -e '.snapshots | length > 0' > /dev/null; do sleep 2; done" + + - name: Restore from storage after data dir loss + run: | + docker compose --env-file harness/local/.env \ + -f harness/local/docker-compose.minio.yml \ + -f harness/local/docker-compose.ds.yml \ + rm -sf node1 + docker volume rm local_node1-data + RESTORE_FROM_STORAGE=true docker compose --env-file harness/local/.env \ + -f harness/local/docker-compose.minio.yml \ + -f harness/local/docker-compose.ds.yml \ + up -d node1 + timeout 180 bash -c 'until curl -sf http://127.0.0.1:9090/ready; do sleep 2; done' + curl -sf http://127.0.0.1:9090/admin/streams/e2e/restore | jq -e '.ownerLocal == true and .streamId != null' + curl -sf -m 30 http://127.0.0.1:4437/e2e/restore | grep -q '"payload":"restore-me"' + - name: Collect logs if: always() run: | @@ -169,6 +191,27 @@ jobs: curl -sf -X PUT -H 'Content-Type: application/json' -d '' "http://127.0.0.1:${http_port}/e2e/failover" curl -sf -X POST -H 'Content-Type: application/json' -d '{"after":"failover"}' "http://127.0.0.1:${http_port}/e2e/failover" curl -sf -m 30 "http://127.0.0.1:${http_port}/e2e/failover" | grep -q '"after":"failover"' + echo "$leader" > killed-leader.txt + + - name: Replace node with empty data dir + run: | + node=$(cat killed-leader.txt) + admin_port=$((9090 + node)) + echo "replacing node ${node} with a fresh data volume" + docker compose --env-file harness/local/.env \ + -f harness/local/docker-compose.minio.yml \ + -f harness/local/docker-compose.cluster.ds.yml \ + rm -sf "node${node}" + docker volume rm "local_node${node}-data" + docker compose --env-file harness/local/.env \ + -f harness/local/docker-compose.minio.yml \ + -f harness/local/docker-compose.cluster.ds.yml \ + up -d "node${node}" + timeout 180 bash -c "until curl -sf http://127.0.0.1:${admin_port}/ready; do sleep 2; done" + curl -sf "http://127.0.0.1:${admin_port}/admin/cluster" | jq -e '.registered == true and .raft.appliedIndex > 0' + curl -sf "http://127.0.0.1:${admin_port}/admin/nodes" | jq -e '.nodes | length == 3' + http_port=$((4436 + node)) + curl -sfL -m 30 "http://127.0.0.1:${http_port}/e2e/failover" | grep -q '"after":"failover"' - name: Collect logs if: always() diff --git a/dashboard/src/api.ts b/dashboard/src/api.ts index 8e6386c..8c9f5d2 100644 --- a/dashboard/src/api.ts +++ b/dashboard/src/api.ts @@ -31,6 +31,20 @@ export interface Readiness { registered: boolean } +export interface ArchivedSnapshot { + key: string + appliedIndex: number + timestampMs: number + size: number +} + +export interface SnapshotArchiveInfo { + archiveSuccessCount: number + archiveFailureCount: number + lastArchivedIndex: number + snapshots: ArchivedSnapshot[] +} + async function get(path: string): Promise { const res = await fetch(path, { headers: { Accept: 'application/json' } }) const body = (await res.json()) as T @@ -45,3 +59,17 @@ async function get(path: string): Promise { export const fetchCluster = () => get('/admin/cluster') export const fetchNodes = () => get<{ nodes: NodeInfo[] }>('/admin/nodes') export const fetchReady = () => get('/ready') + +export async function fetchSnapshots(): Promise { + const res = await fetch('/admin/snapshots', { headers: { Accept: 'application/json' } }) + + if (res.status === 404) { + return null + } + + if (!res.ok) { + throw new Error(`GET /admin/snapshots failed: ${res.status}`) + } + + return (await res.json()) as SnapshotArchiveInfo +} diff --git a/dashboard/src/app.tsx b/dashboard/src/app.tsx index 475ae93..d3c0b58 100644 --- a/dashboard/src/app.tsx +++ b/dashboard/src/app.tsx @@ -3,9 +3,11 @@ import { fetchCluster, fetchNodes, fetchReady, + fetchSnapshots, type ClusterInfo, type NodeInfo, type Readiness, + type SnapshotArchiveInfo, } from './api' const POLL_INTERVAL_MS = 2000 @@ -44,6 +46,7 @@ export function App() { const [cluster, setCluster] = useState(null) const [nodes, setNodes] = useState([]) const [ready, setReady] = useState(null) + const [snapshots, setSnapshots] = useState(null) const [error, setError] = useState(null) const [updatedAt, setUpdatedAt] = useState(null) @@ -52,7 +55,12 @@ export function App() { async function poll() { try { - const [c, n, r] = await Promise.all([fetchCluster(), fetchNodes(), fetchReady()]) + const [c, n, r, s] = await Promise.all([ + fetchCluster(), + fetchNodes(), + fetchReady(), + fetchSnapshots(), + ]) if (!alive) { return @@ -61,6 +69,7 @@ export function App() { setCluster(c) setNodes(n.nodes) setReady(r) + setSnapshots(s) setError(null) setUpdatedAt(new Date()) } catch (e) { @@ -185,6 +194,37 @@ export function App() { )} +
+

Metadata snapshots

+ {!snapshots ? ( +
Snapshot archive disabled
+ ) : snapshots.snapshots.length === 0 ? ( +
No archived snapshots yet
+ ) : ( + + + + + + + + + + {[...snapshots.snapshots].reverse().map((s) => ( + + + + + + ))} + +
Applied indexArchived atSize
{s.appliedIndex}{new Date(s.timestampMs).toLocaleString()}{s.size}
+ )} + {snapshots && snapshots.archiveFailureCount > 0 ? ( +
Archive failures: {snapshots.archiveFailureCount}
+ ) : null} +
+