fix: buffer SSH readiness confirmation
This commit is contained in:
@@ -256,8 +256,13 @@ export function createConcreteDiagnosticAdapters(
|
|||||||
dependencies.sshForwardConfirmed
|
dependencies.sshForwardConfirmed
|
||||||
? dependencies.sshForwardConfirmed(tunnel, signal)
|
? dependencies.sshForwardConfirmed(tunnel, signal)
|
||||||
: new Promise<void>((resolve, reject) => {
|
: new Promise<void>((resolve, reject) => {
|
||||||
|
let stderrBuffer = "";
|
||||||
|
const confirmation = new RegExp(`Local forwarding listening on 127\\.0\\.0\\.1 port ${port}\\.?`);
|
||||||
const confirm = (data: Buffer | string) => {
|
const confirm = (data: Buffer | string) => {
|
||||||
if (new RegExp(`Local forwarding listening on 127\\.0\\.0\\.1 port ${port}\\.?`).test(data.toString())) {
|
stderrBuffer = `${stderrBuffer}${data.toString()}`.slice(-4096);
|
||||||
|
const lines = stderrBuffer.split(/\r?\n/);
|
||||||
|
stderrBuffer = lines.pop() ?? "";
|
||||||
|
if (lines.some((line) => confirmation.test(line))) {
|
||||||
child.stderr?.off?.("data", confirm);
|
child.stderr?.off?.("data", confirm);
|
||||||
resolve();
|
resolve();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -406,6 +406,27 @@ test("permits the probe only after this SSH child confirms its forwarded port",
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("accepts an SSH forward confirmation split across stderr chunks", async () => {
|
||||||
|
const directory = await mkdtemp(join(tmpdir(), "thothii-diagnostic-"));
|
||||||
|
const privateKeyFile = join(directory, "ssh-key");
|
||||||
|
await writeFile(privateKeyFile, "test-key\n", { mode: 0o600 });
|
||||||
|
const stderr = new EventEmitter();
|
||||||
|
const child = Object.assign(new EventEmitter(), { kill: vi.fn(() => true), stderr });
|
||||||
|
child.kill.mockImplementation(() => { child.emit("exit", 0); return true; });
|
||||||
|
const probe = vi.fn(async () => undefined);
|
||||||
|
try {
|
||||||
|
const adapter = createConcreteDiagnosticAdapters({ sshSpawn: vi.fn(() => child), reserveLoopbackPort: async () => 45432 } as any);
|
||||||
|
const running = adapter.withSshTunnel({ sshHost: "bastion.example.test", sshPort: 22, sshUser: "tunnel", privateKeyFile, knownHostsFile: "/run/secrets/known-hosts", targetHost: "dwh.internal", targetPort: 5432, localHost: "127.0.0.1", localPort: 0, timeoutMs: 500, signal: new AbortController().signal }, probe);
|
||||||
|
for (let attempt = 0; attempt < 20 && stderr.listenerCount("data") === 0; attempt += 1) await new Promise((resolve) => setTimeout(resolve, 1));
|
||||||
|
stderr.emit("data", "debug1: Local forwarding listening on 127.0.0.1 ");
|
||||||
|
stderr.emit("data", "port 45432.\n");
|
||||||
|
await running;
|
||||||
|
expect(probe).toHaveBeenCalledOnce();
|
||||||
|
} finally {
|
||||||
|
await rm(directory, { recursive: true, force: true });
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
test("requires a matching embedding model vector and removes its unique write probe", async () => {
|
test("requires a matching embedding model vector and removes its unique write probe", async () => {
|
||||||
const adapters = successfulAdapters();
|
const adapters = successfulAdapters();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user