Split session daemon from web server

This commit is contained in:
Federico Jaramillo Martinez
2026-05-07 13:12:48 +02:00
parent fa55723246
commit 2ffff85313
9 changed files with 272 additions and 93 deletions
+10
View File
@@ -0,0 +1,10 @@
import { homedir } from "node:os";
import { join } from "node:path";
export function sessiondSocketPath(): string {
return process.env.PI_WEB_SESSIOND_SOCKET ?? join(homedir(), ".pi-web", "sessiond.sock");
}
export function sessiondHttpUrl(): string | undefined {
return process.env.PI_WEB_SESSIOND_URL;
}
@@ -0,0 +1,65 @@
import http from "node:http";
import { WebSocket } from "ws";
import { sessiondHttpUrl, sessiondSocketPath } from "./config.js";
export class SessionDaemonClient {
private readonly baseUrl = sessiondHttpUrl();
private readonly socketPath = sessiondSocketPath();
async request(method: string, path: string, body?: unknown): Promise<{ statusCode: number; headers: Record<string, string>; body: string }> {
const payload = body === undefined ? undefined : JSON.stringify(body);
if (this.baseUrl) return this.requestUrl(method, path, payload);
return this.requestSocket(method, path, payload);
}
connectWebSocket(path: string): WebSocket {
if (this.baseUrl) {
const url = new URL(path, this.baseUrl);
url.protocol = url.protocol === "https:" ? "wss:" : "ws:";
return new WebSocket(url);
}
return new WebSocket(`ws+unix:${this.socketPath}:${path}`);
}
private async requestUrl(method: string, path: string, payload?: string) {
const response = await fetch(new URL(path, this.baseUrl), {
method,
headers: payload ? { "content-type": "application/json" } : undefined,
body: payload,
});
return {
statusCode: response.status,
headers: Object.fromEntries(response.headers.entries()),
body: await response.text(),
};
}
private requestSocket(method: string, path: string, payload?: string): Promise<{ statusCode: number; headers: Record<string, string>; body: string }> {
return new Promise((resolve, reject) => {
const request = http.request(
{
socketPath: this.socketPath,
path,
method,
headers: payload
? { "content-type": "application/json", "content-length": Buffer.byteLength(payload) }
: undefined,
},
(response) => {
const chunks: Buffer[] = [];
response.on("data", (chunk) => chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)));
response.on("end", () => {
resolve({
statusCode: response.statusCode ?? 500,
headers: Object.fromEntries(Object.entries(response.headers).map(([key, value]) => [key, Array.isArray(value) ? value.join(", ") : value ?? ""])),
body: Buffer.concat(chunks).toString("utf8"),
});
});
},
);
request.on("error", reject);
if (payload) request.write(payload);
request.end();
});
}
}
+56
View File
@@ -0,0 +1,56 @@
import type { FastifyInstance, FastifyReply } from "fastify";
import { WebSocket, type RawData } from "ws";
import { SessionDaemonClient } from "./sessionDaemonClient.js";
export async function registerSessionProxyRoutes(app: FastifyInstance, daemon = new SessionDaemonClient()): Promise<void> {
const proxy = async (request: { method: string; url: string; body?: unknown }, reply: FastifyReply) => {
try {
const upstream = await daemon.request(request.method, stripApiPrefix(request.url), request.body);
reply.code(upstream.statusCode);
if (upstream.headers["content-type"]) reply.header("content-type", upstream.headers["content-type"]);
return upstream.body ? JSON.parse(upstream.body) : undefined;
} catch (error) {
requestFailed(reply, error);
}
};
app.get<{ Querystring: { cwd?: string } }>("/api/sessions", (request, reply) => proxy(request, reply));
app.post<{ Body: { cwd: string } }>("/api/sessions", (request, reply) => proxy(request, reply));
app.get<{ Params: { sessionId: string } }>("/api/sessions/:sessionId/messages", (request, reply) => proxy(request, reply));
app.get<{ Params: { sessionId: string } }>("/api/sessions/:sessionId/status", (request, reply) => proxy(request, reply));
app.get<{ Params: { sessionId: string } }>("/api/sessions/:sessionId/commands", (request, reply) => proxy(request, reply));
app.post<{ Params: { sessionId: string }; Body: { text: string } }>("/api/sessions/:sessionId/prompt", (request, reply) => proxy(request, reply));
app.post<{ Params: { sessionId: string }; Body: { text: string } }>("/api/sessions/:sessionId/commands/run", (request, reply) => proxy(request, reply));
app.post<{ Params: { sessionId: string }; Body: { requestId: string; value: string } }>("/api/sessions/:sessionId/commands/respond", (request, reply) => proxy(request, reply));
app.post<{ Params: { sessionId: string } }>("/api/sessions/:sessionId/abort", (request, reply) => proxy(request, reply));
app.post<{ Params: { sessionId: string } }>("/api/sessions/:sessionId/stop", (request, reply) => proxy(request, reply));
app.get<{ Params: { sessionId: string } }>("/api/sessions/:sessionId/events", { websocket: true }, (socket, request) => {
bridgeSockets(socket, daemon.connectWebSocket(`/sessions/${request.params.sessionId}/events`));
});
app.get("/api/sessions/events", { websocket: true }, (socket) => {
bridgeSockets(socket, daemon.connectWebSocket("/sessions/events"));
});
}
function stripApiPrefix(url: string): string {
return url.startsWith("/api") ? url.slice(4) || "/" : url;
}
function requestFailed(reply: FastifyReply, error: unknown) {
reply.code(502).send({ error: `Session daemon unavailable: ${error instanceof Error ? error.message : String(error)}` });
}
function bridgeSockets(client: WebSocket, upstream: WebSocket): void {
client.on("message", (data) => sendIfOpen(upstream, data));
upstream.on("message", (data) => sendIfOpen(client, data));
client.on("close", () => upstream.close());
upstream.on("close", () => client.close());
upstream.on("error", () => client.close());
client.on("error", () => upstream.close());
}
function sendIfOpen(socket: WebSocket, data: RawData): void {
if (socket.readyState === WebSocket.OPEN) socket.send(data);
}