feat(backend): RpcClient (spawn/send/request/event) over JSONL
This commit is contained in:
@@ -0,0 +1,38 @@
|
||||
import type { ChildProcessWithoutNullStreams } from "node:child_process";
|
||||
import { attachJsonlReader } from "./line-splitter.js";
|
||||
|
||||
type Listener = (evt: any) => void;
|
||||
|
||||
export class RpcClient {
|
||||
private seq = 0;
|
||||
private pending = new Map<string, (resp: any) => void>();
|
||||
private listeners = new Set<Listener>();
|
||||
|
||||
constructor(private child: ChildProcessWithoutNullStreams) {
|
||||
attachJsonlReader(child.stdout, (line) => {
|
||||
if (!line) return;
|
||||
let msg: any;
|
||||
try { msg = JSON.parse(line); } catch { return; }
|
||||
if (msg.type === "response" && msg.id && this.pending.has(msg.id)) {
|
||||
this.pending.get(msg.id)!(msg);
|
||||
this.pending.delete(msg.id);
|
||||
return;
|
||||
}
|
||||
for (const l of this.listeners) l(msg);
|
||||
});
|
||||
}
|
||||
|
||||
nextId(): string { return `c${++this.seq}`; }
|
||||
|
||||
send(cmd: object): void { this.child.stdin.write(JSON.stringify(cmd) + "\n"); }
|
||||
|
||||
request(cmd: object & { type: string }): Promise<any> {
|
||||
const id = this.nextId();
|
||||
return new Promise((resolve) => {
|
||||
this.pending.set(id, resolve);
|
||||
this.send({ ...cmd, id });
|
||||
});
|
||||
}
|
||||
|
||||
on(_event: "event", cb: Listener): void { this.listeners.add(cb); }
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
import { test, expect } from "vitest";
|
||||
import { spawn } from "node:child_process";
|
||||
import path from "node:path";
|
||||
import { RpcClient } from "../src/rpc/rpc-client.js";
|
||||
|
||||
const FAKE = path.resolve("../harness/tests/fake_pi/fake_pi_rpc.mjs");
|
||||
const SCRIPT = path.resolve("../harness/tests/fake_pi/scripts/f1_disambiguation.json");
|
||||
|
||||
test("request(get_available_models) correla la response", async () => {
|
||||
const child = spawn("node", [FAKE, SCRIPT]);
|
||||
const rpc = new RpcClient(child as any);
|
||||
const res = await rpc.request({ type: "get_available_models" });
|
||||
expect(res.data.models[0].provider).toBe("zai");
|
||||
child.stdin.end();
|
||||
});
|
||||
|
||||
test("prompt emette un evento extension_ui_request", async () => {
|
||||
const child = spawn("node", [FAKE, SCRIPT]);
|
||||
const rpc = new RpcClient(child as any);
|
||||
const got = new Promise<any>((resolve) => rpc.on("event", (e) => e.type === "extension_ui_request" && resolve(e)));
|
||||
rpc.send({ type: "prompt", message: "/nuova-domanda \"x\"" });
|
||||
const evt = await got;
|
||||
expect(evt.method).toBe("input"); // shape nativa di Pi
|
||||
expect(JSON.parse(evt.title).widget).toBe("select"); // descriptor nel title
|
||||
child.stdin.end();
|
||||
});
|
||||
Reference in New Issue
Block a user