fix(cli): never replay writes after ambiguous connection failures

This commit is contained in:
redshift
2026-10-04 21:33:39 +08:00
parent 726112aabe
commit b1a346c6b6
2 changed files with 137 additions and 17 deletions
+103
View File
@@ -0,0 +1,103 @@
import { describe, expect, test } from "bun:test";
import { mkdtempSync, rmSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
// Each case loads config in an isolated process, without touching user config
// or sharing fetch mocks with other test files.
async function runCase(config: object, action: string, responses: number[]) {
const dir = mkdtempSync(join(tmpdir(), "routstrd-retry-"));
try {
writeFileSync(join(dir, "config.json"), JSON.stringify(config));
const script = `
const { callDaemon, callAuth, isDaemonRunning } = await import(${JSON.stringify(join(import.meta.dir, "daemon-client.ts"))});
const responses = ${JSON.stringify(responses)};
const requests = [];
globalThis.fetch = async (url, init) => {
requests.push({ url: String(url), method: init?.method ?? "GET" });
const status = responses.shift();
if (status === 0) throw new TypeError("Connection reset");
if (status === undefined) throw new Error("Unexpected request");
return Response.json(status >= 400 ? {error: "denied"} : {output: "ok"}, {status});
};
let result, error;
try { result = await (${action}); } catch (e) { error = e.message; }
console.log(JSON.stringify({ requests, result, error }));
`;
const proc = Bun.spawn([process.execPath, "--eval", script], {
env: { ...process.env, ROUTSTRD_DIR: dir }, stdout: "pipe", stderr: "pipe",
});
const output = await new Response(proc.stdout).text();
const stderr = await new Response(proc.stderr).text();
expect(await proc.exited, stderr).toBe(0);
return JSON.parse(output) as {
requests: { url: string; method: string }[];
result?: unknown;
error?: string;
};
} finally {
rmSync(dir, { recursive: true, force: true });
}
}
const remote = { daemonUrl: "http://localhost:18008" };
describe("daemon loopback retries", () => {
test("GET retries a connection failure", async () => {
const r = await runCase(remote, 'callDaemon("/status")', [0, 200]);
expect(r.result).toEqual({ output: "ok" });
expect(r.requests.map(r => r.url)).toEqual([
"http://localhost:18008/status", "http://127.0.0.1:18008/status",
]);
});
for (const method of ["POST", "PATCH", "DELETE"]) {
for (const [name, config, client] of [
["remote", remote, "callDaemon"],
["local", { host: "0.0.0.0", port: 18008 }, "callDaemon"],
["auth", { authUrl: "http://localhost:18008" }, "callAuth"],
] as const) {
test(`${name} ${method} is never replayed after a reset`, async () => {
const r = await runCase(config, `${client}("/wallet/send/cashu", {method: "${method}", body: {amount: 10}})`, [200, 0, 200]);
expect(r.error).toContain("outcome is unknown");
expect(r.requests.map(r => r.method)).toEqual(["GET", method]);
});
}
}
test("selects IPv4 with health probes before sending POST once", async () => {
const r = await runCase(remote, 'callDaemon("/wallet/send/cashu", {method: "POST"})', [0, 200, 200]);
expect(r.result).toEqual({ output: "ok" });
expect(r.requests).toEqual([
{ url: "http://localhost:18008/health", method: "GET" },
{ url: "http://127.0.0.1:18008/health", method: "GET" },
{ url: "http://127.0.0.1:18008/wallet/send/cashu", method: "POST" },
]);
});
test("a single-address write also reports an unknown outcome", async () => {
const r = await runCase({daemonUrl: "http://127.0.0.1:18008"}, 'callDaemon("/wallet/send/cashu", {method: "POST"})', [0]);
expect(r.error).toContain("outcome is unknown");
expect(r.requests).toHaveLength(1);
});
test("write address selection stops on HTTP errors", async () => {
const r = await runCase(remote, 'callDaemon("/wallet/send/cashu", {method: "POST"})', [403, 200]);
expect(r.error).toBe("denied");
expect(r.requests).toHaveLength(1);
});
for (const config of [remote, {host: "0.0.0.0", port: 18008}]) {
test(`health HTTP errors stop fallback (${JSON.stringify(config)})`, async () => {
const r = await runCase(config, "isDaemonRunning()", [403, 200]);
expect(r.result).toBe(false);
expect(r.requests).toHaveLength(1);
});
}
test("health connection failures still allow fallback", async () => {
const r = await runCase(remote, "isDaemonRunning()", [0, 200]);
expect(r.result).toBe(true);
expect(r.requests).toHaveLength(2);
});
});
+34 -17
View File
@@ -204,42 +204,59 @@ export async function callDaemonUrl(
} }
} }
async function callLocalDaemon( /** Resolve loopback addresses with safe reads, never by replaying a write. */
async function callDaemonCandidates(
candidates: string[],
path: string, path: string,
options: { method?: "GET" | "POST" | "PATCH" | "DELETE"; body?: object }, options: { method?: "GET" | "POST" | "PATCH" | "DELETE"; body?: object },
config: RoutstrdConfig, config: RoutstrdConfig,
): Promise<CommandResponse> { ): Promise<CommandResponse> {
// Retry only connection failures; HTTP errors prove that a server answered. const isRead = (options.method ?? "GET") === "GET";
let connectionError: DaemonConnectionError | undefined; let connectionError: DaemonConnectionError | undefined;
for (const baseUrl of localDaemonBaseUrls(config)) { for (const baseUrl of candidates) {
if (!isRead && candidates.length > 1) {
try {
// Select an address before dispatching a potentially money-moving
// request. An HTTP error is authoritative, not a reason to fall back.
await callDaemonUrl(baseUrl, "/health", { method: "GET" }, config);
} catch (error) {
if (!(error instanceof DaemonConnectionError)) throw error;
connectionError = error;
continue;
}
}
try { try {
return await callDaemonUrl(baseUrl, path, options, config); return await callDaemonUrl(baseUrl, path, options, config);
} catch (error) { } catch (error) {
if (!(error instanceof DaemonConnectionError)) throw error; if (!(error instanceof DaemonConnectionError)) throw error;
if (!isRead) {
// fetch can reject after the server has committed the operation.
throw new Error(
"Connection lost; operation outcome is unknown — check its status before retrying",
{ cause: error },
);
}
connectionError = error; connectionError = error;
} }
} }
throw connectionError ?? new Error("No daemon host candidates available"); throw connectionError ?? new Error("No daemon host candidates available");
} }
/** Call a configured remote endpoint, retrying connection failures against the async function callLocalDaemon(
* other loopback family when the host is `localhost`. */ path: string,
options: { method?: "GET" | "POST" | "PATCH" | "DELETE"; body?: object },
config: RoutstrdConfig,
): Promise<CommandResponse> {
return callDaemonCandidates(localDaemonBaseUrls(config), path, options, config);
}
async function callRemoteDaemon( async function callRemoteDaemon(
baseUrl: string, baseUrl: string,
path: string, path: string,
options: { method?: "GET" | "POST" | "PATCH" | "DELETE"; body?: object }, options: { method?: "GET" | "POST" | "PATCH" | "DELETE"; body?: object },
config: RoutstrdConfig, config: RoutstrdConfig,
): Promise<CommandResponse> { ): Promise<CommandResponse> {
let connectionError: DaemonConnectionError | undefined; return callDaemonCandidates(baseUrlCandidates(baseUrl), path, options, config);
for (const candidate of baseUrlCandidates(baseUrl)) {
try {
return await callDaemonUrl(candidate, path, options, config);
} catch (error) {
if (!(error instanceof DaemonConnectionError)) throw error;
connectionError = error;
}
}
throw connectionError ?? new Error("No daemon host candidates available");
} }
export async function callDaemon( export async function callDaemon(
@@ -282,7 +299,7 @@ export async function isDaemonRunning(): Promise<boolean> {
const response = await fetch(url, { const response = await fetch(url, {
headers: authorization ? { Authorization: authorization } : {}, headers: authorization ? { Authorization: authorization } : {},
}); });
if (response.ok) return true; return response.ok;
} catch { } catch {
// Try the next candidate host. // Try the next candidate host.
} }
@@ -298,7 +315,7 @@ export async function isDaemonRunning(): Promise<boolean> {
const response = await fetch(`${baseUrl}/health`, { const response = await fetch(`${baseUrl}/health`, {
signal: controller.signal, signal: controller.signal,
}); });
if (response.ok) return true; return response.ok;
} catch { } catch {
// Try the next candidate host. // Try the next candidate host.
} finally { } finally {