diff --git a/src/cli.ts b/src/cli.ts index cd38571..8f14fed 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -1,8 +1,10 @@ import { program } from "commander"; import { startDaemon } from "./start-daemon"; +import { consumeRecoveryStream } from "./utils/recovery-stream"; import { handleDaemonCommand, callDaemon, + callDaemonRaw, callAuth, ensureDaemonRunning, isDaemonRunning, @@ -1341,6 +1343,55 @@ program } }); +// Recover - re-derive proofs from a mint's restore endpoint +program + .command("recover") + .description("Recover wallet proofs from a mint's restore endpoint") + .requiredOption( + "-m, --mint-url ", + "Mint URL to recover proofs from", + ) + .action(async (options: { mintUrl: string }) => { + await ensureDaemonRunning(); + try { + const response = await callDaemonRaw("/wallet/recover", { + method: "POST", + body: { url: options.mintUrl }, + }); + + // Validation errors (e.g. missing field) come back as JSON before any + // streaming begins. + if (!response.ok) { + const contentType = response.headers.get("content-type") || ""; + if (contentType.includes("application/json")) { + const data = (await response + .json() + .catch(() => ({}))) as { error?: string }; + console.error(data.error || `HTTP ${response.status}`); + } else { + console.error(`HTTP ${response.status}`); + } + process.exit(1); + } + + const message = await consumeRecoveryStream(response.body, (progress) => { + console.log(progress); + }); + console.log(`\u2713 ${message}`); + } catch (error) { + const message = (error as Error).message; + if ( + message?.includes("fetch failed") || + message?.includes("Connection refused") + ) { + console.error("Daemon is not running and failed to auto-start"); + process.exit(1); + } + console.error(message); + process.exit(1); + } + }); + // Monitor - interactive TUI program .command("monitor") diff --git a/src/daemon/http/index.ts b/src/daemon/http/index.ts index 0e6401b..e525fe9 100644 --- a/src/daemon/http/index.ts +++ b/src/daemon/http/index.ts @@ -20,6 +20,7 @@ import { import { decodeCashuTokenAmount } from "../wallet"; import { getClientsFromStore } from "../../utils/clients"; import { getUsageSummary } from "./usage-summary"; +import { encodeRecoveryEvent } from "../../utils/recovery-stream"; type ClientMode = "xcashu" | "lazyrefund" | "apikeys"; @@ -167,6 +168,31 @@ async function respond( } } +export async function streamMintRecovery( + res: ServerResponse, + walletClient: CocodClient, + mintUrl: string, +): Promise { + res.writeHead(200, { + "Content-Type": "application/x-ndjson; charset=utf-8", + "Cache-Control": "no-store", + }); + try { + const message = await walletClient.recoverMint(mintUrl, (line) => { + res.write(encodeRecoveryEvent({ type: "progress", message: line })); + }); + res.end(encodeRecoveryEvent({ type: "result", ok: true, message })); + } catch (error) { + res.end( + encodeRecoveryEvent({ + type: "result", + ok: false, + error: toErrorMessage(error), + }), + ); + } +} + function requireStringField( body: Record, field: string, @@ -466,6 +492,19 @@ export function createDaemonRequestHandler(deps: { return; } + // Recover proofs from a mint's restore endpoint. Progress and the required + // terminal result are sent as NDJSON so clients can detect truncated runs. + if (req.method === "POST" && url.pathname === "/wallet/recover") { + try { + const body = await readJsonBody(req); + const mintUrl = getRequiredStringField(body, "url"); + await streamMintRecovery(res, deps.walletClient, mintUrl); + } catch (error) { + respondWithError(res, error); + } + return; + } + if (req.method === "GET" && url.pathname === "/wallet/history") { await respond(res, async () => { const offsetParam = url.searchParams.get("offset"); diff --git a/src/daemon/http/recovery.test.ts b/src/daemon/http/recovery.test.ts new file mode 100644 index 0000000..4ed8530 --- /dev/null +++ b/src/daemon/http/recovery.test.ts @@ -0,0 +1,75 @@ +import { describe, expect, it } from "bun:test"; +import type { ServerResponse } from "http"; +import type { CocodClient } from "../wallet/cocod-client"; +import { streamMintRecovery } from "./index"; + +function responseRecorder(): { + response: ServerResponse; + chunks: string[]; + headers: Record; +} { + const chunks: string[] = []; + const headers: Record = {}; + const response = { + writeHead(_status: number, values: Record) { + Object.assign(headers, values); + return this; + }, + write(chunk: string) { + chunks.push(chunk); + return true; + }, + end(chunk?: string) { + if (chunk) chunks.push(chunk); + return this; + }, + } as unknown as ServerResponse; + return { response, chunks, headers }; +} + +function clientWithRecover( + recoverMint: CocodClient["recoverMint"], +): CocodClient { + return { recoverMint } as CocodClient; +} + +describe("streamMintRecovery", () => { + it("streams progress and a terminal success event", async () => { + const recorder = responseRecorder(); + const client = clientWithRecover(async (_url, progress) => { + progress?.("Scanning keysets"); + return "Recovery complete"; + }); + + await streamMintRecovery( + recorder.response, + client, + "https://mint.example.com", + ); + + expect(recorder.headers["Content-Type"]).toContain("application/x-ndjson"); + expect(recorder.chunks.map((line) => JSON.parse(line))).toEqual([ + { type: "progress", message: "Scanning keysets" }, + { type: "result", ok: true, message: "Recovery complete" }, + ]); + }); + + it("always terminates a started stream with an error result", async () => { + const recorder = responseRecorder(); + const client = clientWithRecover(async (_url, progress) => { + progress?.("Scanning keysets"); + throw new Error("Mint restore failed"); + }); + + await streamMintRecovery( + recorder.response, + client, + "https://mint.example.com", + ); + + expect(recorder.chunks.map((line) => JSON.parse(line))).toEqual([ + { type: "progress", message: "Scanning keysets" }, + { type: "result", ok: false, error: "Mint restore failed" }, + ]); + }); +}); diff --git a/src/daemon/wallet/coco-client.npc.test.ts b/src/daemon/wallet/coco-client.npc.test.ts index d348473..d3fa5c2 100644 --- a/src/daemon/wallet/coco-client.npc.test.ts +++ b/src/daemon/wallet/coco-client.npc.test.ts @@ -162,6 +162,22 @@ describe("in-process coco client NPC integration", () => { }); }); +describe("legacy cocod mint recovery", () => { + it("reports mint recovery as unsupported", async () => { + const client = createCocodClient({ + socketPath: "/tmp/unused-cocod-recovery-test.sock", + }); + + try { + await client.recoverMint("https://mint.example.com"); + throw new Error("expected recovery to fail"); + } catch (error) { + expect(error).toBeInstanceOf(CocodHttpError); + expect((error as InstanceType).status).toBe(501); + } + }); +}); + describe("legacy cocod client NPC passthroughs", () => { type FetchImpl = ( input: string | URL | Request, diff --git a/src/daemon/wallet/coco-client.ts b/src/daemon/wallet/coco-client.ts index 2b351a3..e0d580d 100644 --- a/src/daemon/wallet/coco-client.ts +++ b/src/daemon/wallet/coco-client.ts @@ -31,6 +31,7 @@ import type { NpcUsernameResult, } from "./cocod-client"; import { logger } from "../../utils/logger"; +import { RecoveryGate } from "./recovery-gate"; const NPC_DEFAULT_BASE_URL = "https://npubx.cash"; @@ -517,6 +518,53 @@ export interface CreateCocoClientOptions { npcBaseUrl?: string; } +type CocoManager = Awaited>; + +export async function recoverMintProofs( + coco: CocoManager, + mintUrl: string, + onProgress?: (message: string) => void, +): Promise { + const normalized = normalizeMintUrl(mintUrl); + onProgress?.(`Fetching mint info for ${normalized}…`); + const { keysets } = await coco.mint.addMint(normalized, { trusted: true }); + onProgress?.( + `Recovery started for ${keysets.length} keyset(s) from ${normalized}; this may take several minutes…`, + ); + + // Background services can mutate proofs and counters independently of + // public client calls, so stop them while restore overwrites counters. + let subscriptionsPaused = false; + let mintWatcherStopped = false; + let proofWatcherStopped = false; + let mintProcessorStopped = false; + try { + await coco.pauseSubscriptions(); + subscriptionsPaused = true; + await coco.disableMintOperationWatcher(); + mintWatcherStopped = true; + await coco.disableProofStateWatcher(); + proofWatcherStopped = true; + await coco.disableMintOperationProcessor(); + mintProcessorStopped = true; + await coco.wallet.restore(normalized); + } finally { + const restarts: Promise[] = []; + if (mintWatcherStopped) restarts.push(coco.enableMintOperationWatcher()); + if (proofWatcherStopped) restarts.push(coco.enableProofStateWatcher()); + if (mintProcessorStopped) restarts.push(coco.enableMintOperationProcessor()); + if (subscriptionsPaused) restarts.push(coco.resumeSubscriptions()); + const results = await Promise.allSettled(restarts); + const failedRestart = results.find( + (result): result is PromiseRejectedResult => result.status === "rejected", + ); + if (failedRestart) throw failedRestart.reason; + } + + onProgress?.(`Recovery complete for ${normalized}`); + return `Mint ${normalized} recovery completed successfully`; +} + export async function createCocoClient( options: CreateCocoClientOptions = {}, ): Promise { @@ -632,6 +680,7 @@ export async function createCocoClient( }; let disposed = false; + const recoveryGate = new RecoveryGate(); return { async ping(): Promise { try { @@ -668,51 +717,59 @@ export async function createCocoClient( }, async receiveCashu(token: string): Promise { - await coco.wallet.receive(token); - return "Token received successfully"; + return recoveryGate.runMutation(async () => { + await coco.wallet.receive(token); + return "Token received successfully"; + }); }, async receiveBolt11(amount: number, mintUrl?: string): Promise { - const targetMint = mintUrl || walletConfig.defaultMintUrl; - if (!targetMint) { - throw new Error("No trusted mint available for Lightning invoice"); - } - const op = await coco.ops.mint.prepare({ - mintUrl: targetMint, - amount, - method: "bolt11", + return recoveryGate.runMutation(async () => { + const targetMint = mintUrl || walletConfig.defaultMintUrl; + if (!targetMint) { + throw new Error("No trusted mint available for Lightning invoice"); + } + const op = await coco.ops.mint.prepare({ + mintUrl: targetMint, + amount, + method: "bolt11", + }); + if (!("request" in op)) { + throw new Error("mint prepare did not return a payment request"); + } + return op.request as string; }); - if (!("request" in op)) { - throw new Error("mint prepare did not return a payment request"); - } - return op.request as string; }, async sendCashu(amount: number, mintUrl?: string): Promise { - const targetMint = mintUrl || walletConfig.defaultMintUrl; - if (!targetMint) { - throw new Error("No trusted mint available for sending"); - } - const prepared = await coco.ops.send.prepare({ - mintUrl: targetMint, - amount, + return recoveryGate.runMutation(async () => { + const targetMint = mintUrl || walletConfig.defaultMintUrl; + if (!targetMint) { + throw new Error("No trusted mint available for sending"); + } + const prepared = await coco.ops.send.prepare({ + mintUrl: targetMint, + amount, + }); + const { token } = await coco.ops.send.execute(prepared.id); + return getEncodedToken(token); }); - const { token } = await coco.ops.send.execute(prepared.id); - return getEncodedToken(token); }, async sendBolt11(invoice: string, mintUrl?: string): Promise { - const targetMint = mintUrl || walletConfig.defaultMintUrl; - if (!targetMint) { - throw new Error("No trusted mint available for Lightning payment"); - } - const prepared = await coco.ops.melt.prepare({ - mintUrl: targetMint, - method: "bolt11", - methodData: { invoice }, + return recoveryGate.runMutation(async () => { + const targetMint = mintUrl || walletConfig.defaultMintUrl; + if (!targetMint) { + throw new Error("No trusted mint available for Lightning payment"); + } + const prepared = await coco.ops.melt.prepare({ + mintUrl: targetMint, + method: "bolt11", + methodData: { invoice }, + }); + await coco.ops.melt.execute(prepared.id); + return "Payment sent successfully"; }); - await coco.ops.melt.execute(prepared.id); - return "Payment sent successfully"; }, async listMints(): Promise { @@ -746,6 +803,15 @@ export async function createCocoClient( return `Default mint set to ${mintUrl}`; }, + async recoverMint( + mintUrl: string, + onProgress?: (message: string) => void, + ): Promise { + return recoveryGate.runRecovery(() => + recoverMintProofs(coco, mintUrl, onProgress), + ); + }, + async dispose(): Promise { if (disposed) return; disposed = true; @@ -782,19 +848,21 @@ export async function createCocoClient( username: string, confirm?: boolean, ): Promise { - const result = await npcApi().setUsername(username, confirm === true); - if (result.success) { - return { success: true }; - } - return { - success: false, - paymentRequest: - result.pr as NpcUsernameResult["paymentRequest"], - }; + return recoveryGate.runMutation(async () => { + const result = await npcApi().setUsername(username, confirm === true); + if (result.success) { + return { success: true }; + } + return { + success: false, + paymentRequest: + result.pr as NpcUsernameResult["paymentRequest"], + }; + }); }, async syncNpc(): Promise { - await npcApi().sync(); + await recoveryGate.runMutation(() => npcApi().sync()); }, }; } diff --git a/src/daemon/wallet/cocod-client.ts b/src/daemon/wallet/cocod-client.ts index 6725812..c26a319 100644 --- a/src/daemon/wallet/cocod-client.ts +++ b/src/daemon/wallet/cocod-client.ts @@ -81,6 +81,15 @@ export interface CocodClient { getMintInfo(url: string): Promise; getDefaultMint(): Promise; setDefaultMint(url: string): Promise; + /** + * Recover deterministic proofs using the mint's NUT-09 restore endpoint. + * Recovered proofs are checked using NUT-07 before being persisted. + * `onProgress` receives human-readable status messages. + */ + recoverMint( + mintUrl: string, + onProgress?: (message: string) => void, + ): Promise; /** Release resources held by in-process wallet implementations. */ dispose?(): Promise; getHistory(offset?: number, limit?: number): Promise; @@ -398,6 +407,15 @@ export function createCocodClient( async setDefaultMint(url: string): Promise { return post("/mints/default", { url }); }, + async recoverMint( + _mintUrl: string, + _onProgress?: (message: string) => void, + ): Promise { + throw new CocodHttpError( + 501, + "Mint recovery is not supported by the legacy cocod client.", + ); + }, async getHistory(_offset?: number, _limit?: number): Promise { return []; }, diff --git a/src/daemon/wallet/recover-mint.test.ts b/src/daemon/wallet/recover-mint.test.ts new file mode 100644 index 0000000..44e4c5c --- /dev/null +++ b/src/daemon/wallet/recover-mint.test.ts @@ -0,0 +1,73 @@ +import { describe, expect, it } from "bun:test"; +import type { initializeCoco } from "@cashu/coco-core"; +import { recoverMintProofs } from "./coco-client"; + +type CocoManager = Awaited>; + +function manager(options: { restoreError?: Error } = {}): { + coco: CocoManager; + calls: string[]; +} { + const calls: string[] = []; + const record = (name: string) => async () => { + calls.push(name); + }; + const coco = { + mint: { + addMint: async (url: string, addOptions: { trusted: boolean }) => { + calls.push(`add:${url}:${addOptions.trusted}`); + return { keysets: [{}, {}] }; + }, + }, + wallet: { + restore: async (url: string) => { + calls.push(`restore:${url}`); + if (options.restoreError) throw options.restoreError; + }, + }, + pauseSubscriptions: record("pause-subscriptions"), + disableMintOperationWatcher: record("stop-mint-watcher"), + disableProofStateWatcher: record("stop-proof-watcher"), + disableMintOperationProcessor: record("stop-mint-processor"), + enableMintOperationWatcher: record("start-mint-watcher"), + enableProofStateWatcher: record("start-proof-watcher"), + enableMintOperationProcessor: async () => { + calls.push("start-mint-processor"); + return true; + }, + resumeSubscriptions: record("resume-subscriptions"), + } as unknown as CocoManager; + return { coco, calls }; +} + +describe("recoverMintProofs", () => { + it("normalizes the mint URL, restores it, and emits coarse progress", async () => { + const { coco, calls } = manager(); + const progress: string[] = []; + + const result = await recoverMintProofs( + coco, + "https://mint.example.com/", + (message) => progress.push(message), + ); + + expect(calls).toContain("add:https://mint.example.com:true"); + expect(calls).toContain("restore:https://mint.example.com"); + expect(progress).toHaveLength(3); + expect(progress[1]).toContain("2 keyset(s)"); + expect(result).toContain("https://mint.example.com"); + }); + + it("restarts background services when restore fails", async () => { + const { coco, calls } = manager({ restoreError: new Error("restore failed") }); + + await expect( + recoverMintProofs(coco, "https://mint.example.com"), + ).rejects.toThrow("restore failed"); + + expect(calls).toContain("start-mint-watcher"); + expect(calls).toContain("start-proof-watcher"); + expect(calls).toContain("start-mint-processor"); + expect(calls).toContain("resume-subscriptions"); + }); +}); diff --git a/src/daemon/wallet/recovery-gate.test.ts b/src/daemon/wallet/recovery-gate.test.ts new file mode 100644 index 0000000..4438fdf --- /dev/null +++ b/src/daemon/wallet/recovery-gate.test.ts @@ -0,0 +1,62 @@ +import { describe, expect, it } from "bun:test"; +import { RecoveryGate } from "./recovery-gate"; + +function deferred(): { + promise: Promise; + resolve: () => void; +} { + let resolve!: () => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + +describe("RecoveryGate", () => { + it("rejects mutations while recovery is running", async () => { + const gate = new RecoveryGate(); + const hold = deferred(); + const recovery = gate.runRecovery(() => hold.promise); + + await expect(gate.runMutation(async () => undefined)).rejects.toThrow( + "recovery is in progress", + ); + hold.resolve(); + await recovery; + }); + + it("rejects concurrent recovery attempts", async () => { + const gate = new RecoveryGate(); + const hold = deferred(); + const recovery = gate.runRecovery(() => hold.promise); + + await expect(gate.runRecovery(async () => undefined)).rejects.toThrow( + "already in progress", + ); + hold.resolve(); + await recovery; + }); + + it("does not start recovery during an active mutation", async () => { + const gate = new RecoveryGate(); + const hold = deferred(); + const mutation = gate.runMutation(() => hold.promise); + + await expect(gate.runRecovery(async () => undefined)).rejects.toThrow( + "another wallet operation", + ); + hold.resolve(); + await mutation; + }); + + it("releases the gate after a failed recovery", async () => { + const gate = new RecoveryGate(); + await expect( + gate.runRecovery(async () => { + throw new Error("restore failed"); + }), + ).rejects.toThrow("restore failed"); + + await expect(gate.runMutation(async () => "ok")).resolves.toBe("ok"); + }); +}); diff --git a/src/daemon/wallet/recovery-gate.ts b/src/daemon/wallet/recovery-gate.ts new file mode 100644 index 0000000..624ec16 --- /dev/null +++ b/src/daemon/wallet/recovery-gate.ts @@ -0,0 +1,49 @@ +export class WalletRecoveryBusyError extends Error { + constructor(message: string) { + super(message); + this.name = "WalletRecoveryBusyError"; + } +} + +/** + * Coordinates public wallet mutations with destructive counter restoration. + * Admission is synchronous, so no mutation can enter between the recovery + * availability check and recovery becoming exclusive. + */ +export class RecoveryGate { + private activeMutations = 0; + private recovering = false; + + async runMutation(operation: () => Promise): Promise { + if (this.recovering) { + throw new WalletRecoveryBusyError( + "Wallet recovery is in progress; try again when it completes.", + ); + } + + this.activeMutations += 1; + try { + return await operation(); + } finally { + this.activeMutations -= 1; + } + } + + async runRecovery(operation: () => Promise): Promise { + if (this.recovering) { + throw new WalletRecoveryBusyError("Wallet recovery is already in progress."); + } + if (this.activeMutations > 0) { + throw new WalletRecoveryBusyError( + "Cannot start wallet recovery while another wallet operation is in progress.", + ); + } + + this.recovering = true; + try { + return await operation(); + } finally { + this.recovering = false; + } + } +} diff --git a/src/utils/daemon-client.ts b/src/utils/daemon-client.ts index ff27906..fb2bc0d 100644 --- a/src/utils/daemon-client.ts +++ b/src/utils/daemon-client.ts @@ -95,6 +95,46 @@ export async function callDaemon( return _callUrl(baseUrl, path, options, config); } +/** + * Performs an authenticated daemon request and returns the raw fetch + * `Response` so callers can stream the body (e.g. for live progress output). + * Handles NIP-98 auth for remote daemons exactly like `callDaemon`. + */ +export async function callDaemonRaw( + path: string, + options: { method?: HttpMethod; body?: object } = {}, +): Promise { + const config = await loadConfig(); + const baseUrl = getDaemonBaseUrl(config); + const { method = "POST", body } = options; + const url = `${baseUrl}${path}`; + + const bodyString = body ? JSON.stringify(body) : undefined; + const bodyBytes = bodyString + ? new TextEncoder().encode(bodyString) + : undefined; + + let authorization: string | undefined; + if ((config.daemonUrl || config.authUrl) && config.nsec) { + const secretKey = parseSecretKey(config.nsec); + authorization = await createNIP98Authorization( + secretKey, + url, + method, + bodyBytes, + ); + } + + return fetch(url, { + method, + headers: { + ...(authorization ? { Authorization: authorization } : {}), + ...(bodyString ? { "Content-Type": "application/json" } : {}), + }, + body: bodyString, + }); +} + /** Like callDaemon but sends requests to the auth proxy URL instead. * Falls back to the daemon URL if no authUrl is configured. */ export async function callAuth( diff --git a/src/utils/recovery-stream.test.ts b/src/utils/recovery-stream.test.ts new file mode 100644 index 0000000..93e4bec --- /dev/null +++ b/src/utils/recovery-stream.test.ts @@ -0,0 +1,63 @@ +import { describe, expect, it } from "bun:test"; +import { + consumeRecoveryStream, + encodeRecoveryEvent, +} from "./recovery-stream"; + +function stream(chunks: string[]): ReadableStream { + const encoder = new TextEncoder(); + return new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(encoder.encode(chunk)); + controller.close(); + }, + }); +} + +describe("recovery NDJSON stream", () => { + it("reports progress and returns the terminal success message", async () => { + const progress: string[] = []; + const body = stream([ + encodeRecoveryEvent({ type: "progress", message: "Scanning mint" }), + encodeRecoveryEvent({ type: "result", ok: true, message: "Recovered" }), + ]); + + await expect( + consumeRecoveryStream(body, (message) => progress.push(message)), + ).resolves.toBe("Recovered"); + expect(progress).toEqual(["Scanning mint"]); + }); + + it("handles events split across transport chunks", async () => { + const encoded = encodeRecoveryEvent({ + type: "result", + ok: true, + message: "Recovered", + }); + await expect( + consumeRecoveryStream(stream([encoded.slice(0, 8), encoded.slice(8)])), + ).resolves.toBe("Recovered"); + }); + + it("propagates a terminal recovery failure", async () => { + const body = stream([ + encodeRecoveryEvent({ type: "result", ok: false, error: "Mint failed" }), + ]); + await expect(consumeRecoveryStream(body)).rejects.toThrow("Mint failed"); + }); + + it("rejects a truncated response without a terminal result", async () => { + const body = stream([ + encodeRecoveryEvent({ type: "progress", message: "Scanning mint" }), + ]); + await expect(consumeRecoveryStream(body)).rejects.toThrow( + "ended before a result", + ); + }); + + it("rejects an absent response body", async () => { + await expect(consumeRecoveryStream(null)).rejects.toThrow( + "no recovery response body", + ); + }); +}); diff --git a/src/utils/recovery-stream.ts b/src/utils/recovery-stream.ts new file mode 100644 index 0000000..2387229 --- /dev/null +++ b/src/utils/recovery-stream.ts @@ -0,0 +1,93 @@ +export type RecoveryStreamEvent = + | { type: "progress"; message: string } + | { type: "result"; ok: true; message: string } + | { type: "result"; ok: false; error: string }; + +export function encodeRecoveryEvent(event: RecoveryStreamEvent): string { + return `${JSON.stringify(event)}\n`; +} + +function parseRecoveryEvent(line: string): RecoveryStreamEvent { + let value: unknown; + try { + value = JSON.parse(line); + } catch { + throw new Error("Daemon returned malformed recovery progress data."); + } + + if (!value || typeof value !== "object") { + throw new Error("Daemon returned an invalid recovery event."); + } + + const event = value as Record; + if (event.type === "progress" && typeof event.message === "string") { + return { type: "progress", message: event.message }; + } + if ( + event.type === "result" && + event.ok === true && + typeof event.message === "string" + ) { + return { type: "result", ok: true, message: event.message }; + } + if ( + event.type === "result" && + event.ok === false && + typeof event.error === "string" + ) { + return { type: "result", ok: false, error: event.error }; + } + + throw new Error("Daemon returned an invalid recovery event."); +} + +/** + * Consume a recovery NDJSON response. A terminal result event is mandatory; + * EOF before that event is treated as a failed/truncated recovery. + */ +export async function consumeRecoveryStream( + body: ReadableStream | null, + onProgress: (message: string) => void = () => {}, +): Promise { + if (!body) { + throw new Error("Daemon returned no recovery response body."); + } + + const reader = body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + let result: Extract | undefined; + + const consumeLine = (line: string): void => { + if (!line.trim()) return; + if (result) { + throw new Error("Daemon returned data after the recovery result."); + } + const event = parseRecoveryEvent(line); + if (event.type === "progress") { + onProgress(event.message); + } else { + result = event; + } + }; + + while (true) { + const { done, value } = await reader.read(); + if (done) break; + buffer += decoder.decode(value, { stream: true }); + const lines = buffer.split("\n"); + buffer = lines.pop() ?? ""; + for (const line of lines) consumeLine(line); + } + + buffer += decoder.decode(); + consumeLine(buffer); + + if (!result) { + throw new Error("Recovery response ended before a result was received."); + } + if (!result.ok) { + throw new Error(result.error); + } + return result.message; +}