From 7f4bc790f1cfb1cd058a1488dfdc66684fce9d98 Mon Sep 17 00:00:00 2001 From: vivek1504 Date: Fri, 17 Jul 2026 16:54:12 +0530 Subject: [PATCH 1/4] added types for ai agent execution --- src/runtime/protocol.ts | 14 ++++++++++ src/runtime/transport.ts | 17 ++++++++++++ src/types/types.ts | 56 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 87 insertions(+) diff --git a/src/runtime/protocol.ts b/src/runtime/protocol.ts index 91635ce..c822329 100644 --- a/src/runtime/protocol.ts +++ b/src/runtime/protocol.ts @@ -17,6 +17,7 @@ export function buildPayload(req: any, subPath: string): string { export function readVsockResponse( socket: Socket, timeout: number, + onStreamChunk?: (chunk: any) => void ): Promise<{ type: string; data: any; error?: string }> { return new Promise((resolve, reject) => { let buffer = ""; @@ -41,6 +42,14 @@ export function readVsockResponse( onData = (chunk: Buffer) => { buffer += chunk.toString(); + if (buffer.length > 10 * 1024 * 1024) { + clearTimeout(timer); + cleanup(); + socket.destroy(); + reject(new Error("Response too large")); + return; + } + let index; while ((index = buffer.indexOf("\n")) >= 0) { @@ -53,6 +62,11 @@ export function readVsockResponse( try { const msg = JSON.parse(line); + if (msg.type === "stream") { + onStreamChunk?.(msg); + continue; + } + if (msg.type === "response" || msg.type === "error") { clearTimeout(timer); cleanup(); diff --git a/src/runtime/transport.ts b/src/runtime/transport.ts index f853189..e4a3c47 100644 --- a/src/runtime/transport.ts +++ b/src/runtime/transport.ts @@ -90,3 +90,20 @@ export async function sendRequest(subPath: string, req: any, res: any, vm: Vm) { res.status(statusCode).json(msg.data ?? { error: msg.error }); } } + +export async function sendMessage( + vm: Vm, + message: Record, + onStreamChunk?: (chunk: any) => void, + timeout: number = 60000, +): Promise { + const socket = await getVmSocket(vm); + socket.write(JSON.stringify(message) + "\n"); + + transportLogger.debug( + { vmId: vm.id, messageType: message.type, messageId: message.id }, + "message sent to VM" + ); + + return readVsockResponse(socket, timeout, onStreamChunk); +} diff --git a/src/types/types.ts b/src/types/types.ts index a5ce129..cdb539d 100644 --- a/src/types/types.ts +++ b/src/types/types.ts @@ -34,3 +34,59 @@ export interface RuntimeFunction { vms: Vm[]; readyVms: Set; } + +export interface ExecuteMessage { + type: "execute"; + id: string; + command: string; + args: string[]; + cwd?: string; + env?: Record; + timeout?: number; +} + +export interface WriteFileMessage { + type: "write_file"; + id: string; + path: string; + content: string; + mode?: number; +} + +export interface ReadFileMessage { + type: "read_file"; + id: string; + path: string; +} + +export interface ListFilesMessage { + type: "list_files"; + id: string; + path?: string; + recursive?: boolean; +} + +export interface CancelMessage { + type: "cancel"; + id: string; +} + +export interface StreamMessage { + type: "stream"; + id: string; + stream: "stdout" | "stderr"; + data: string; +} + +export interface ResponseMessage { + type: "response"; + id: string; + data: any; +} + +export interface ErrorMessage { + type: "error"; + id: string; + error: string; + code?: number; +} From 33112027040553a2c9fe7a13107a442609007e63 Mon Sep 17 00:00:00 2001 From: vivek1504 Date: Fri, 17 Jul 2026 17:33:53 +0530 Subject: [PATCH 2/4] added session life cycle management --- src/exec/session.ts | 73 +++++++++++++++++++++++++++++++++++++++++++++ src/utils/logger.ts | 1 + 2 files changed, 74 insertions(+) create mode 100644 src/exec/session.ts diff --git a/src/exec/session.ts b/src/exec/session.ts new file mode 100644 index 0000000..c0fac0f --- /dev/null +++ b/src/exec/session.ts @@ -0,0 +1,73 @@ +import { runtimeStore } from "../runtime/store.js"; +import { cleanupVm } from "../runtime/cleanup.js"; +import { sessionLogger } from "../utils/logger.js"; + +export interface Session { + sessionId: string; + createdAt: number; + lastActivityAt: number; + state: "creating" | "active" | "destroying"; +} + +const sessions = new Map(); + +export function getSession(sessionId: string): Session | undefined { + return sessions.get(sessionId); +} + +export function createSession(sessionId: string): Session { + const session: Session = { + sessionId, + createdAt: Date.now(), + lastActivityAt: Date.now(), + state: "creating", + }; + sessions.set(sessionId, session); + return session; +} + +export function touchSession(sessionId: string): void { + const session = sessions.get(sessionId); + if (session) session.lastActivityAt = Date.now(); +} + +export async function destroySession(sessionId: string): Promise { + const session = sessions.get(sessionId); + if (!session) return false; + + session.state = "destroying"; + + const fn = runtimeStore.functions.get(sessionId); + if (fn) { + for (const vm of [...fn.vms]) { + await cleanupVm(fn, vm); + } + runtimeStore.functions.delete(sessionId); + } + + sessions.delete(sessionId); + sessionLogger.info({ sessionId }, "session destroyed"); + return true; +} + +export function getActiveSessions(): number { + return sessions.size; +} + +export function getAllSessions(): Session[] { + return [...sessions.values()]; +} + +export function startSessionReaper(ttlMs: number = 30 * 60 * 1000): void { + setInterval(() => { + const now = Date.now(); + for (const [id, session] of sessions) { + if (now - session.lastActivityAt > ttlMs && session.state === "active") { + sessionLogger.info({ sessionId: id, idleMs: now - session.lastActivityAt }, "reaping idle session"); + destroySession(id).catch(err => { + sessionLogger.error({ sessionId: id, err }, "session reap failed"); + }); + } + } + }, 60_000); +} diff --git a/src/utils/logger.ts b/src/utils/logger.ts index 7c7d20b..cdc7d0c 100644 --- a/src/utils/logger.ts +++ b/src/utils/logger.ts @@ -33,6 +33,7 @@ export const vmManagerLogger = logger.child({ module: "vm-manager" }); export const transportLogger = logger.child({ module: "transport" }); export const protocolLogger = logger.child({ module: "protocol" }); export const cleanupLogger = logger.child({ module: "cleanup" }); +export const sessionLogger = logger.child({module: "session"}) export const httpLoggerOptions: Options = { logger: logger.child({ module: "http" }), From c6704c15d18a7887eb9e2f3c7f0a3fd32bad55f6 Mon Sep 17 00:00:00 2001 From: vivek1504 Date: Fri, 17 Jul 2026 18:12:42 +0530 Subject: [PATCH 3/4] added gateway.ts --- src/exec/gateway.ts | 66 +++++++++++++++++++++++++++++++++++++++ src/runtime/vm-manager.ts | 3 +- src/utils/logger.ts | 3 +- 3 files changed, 70 insertions(+), 2 deletions(-) create mode 100644 src/exec/gateway.ts diff --git a/src/exec/gateway.ts b/src/exec/gateway.ts new file mode 100644 index 0000000..dc0c849 --- /dev/null +++ b/src/exec/gateway.ts @@ -0,0 +1,66 @@ +import { runtimeStore } from "../runtime/store.js"; +import { readVsockResponse } from "../runtime/protocol.js"; +import { getVmSocket } from "../runtime/transport.js"; +import { getSession, createSession, touchSession } from "./session.js"; +import { gatewayLogger } from "../utils/logger.js"; +import crypto from "crypto"; +import type { Vm } from "../types/types.js"; + +import { createVm } from "../runtime/vm-manager.js"; +import { Deque } from "../runtime/deque.js"; + +export async function ensureSession(sessionId: string): Promise { + let session = getSession(sessionId); + if (session?.state === "active") { + const fn = runtimeStore.functions.get(sessionId); + if (fn && fn.vms.length > 0) return fn.vms[0]!; + } + + session = session || createSession(sessionId); + + let fn = runtimeStore.functions.get(sessionId); + if (!fn) { + fn = { + functionId: sessionId, + queue: new Deque(), + vms: [], + readyVms: new Set(), + weight: 1, + inflightCount: 0, + deficit: 0, + pendingCreations: 0, + }; + runtimeStore.functions.set(sessionId, fn); + } + + if (fn.vms.length === 0) { + await createVm(sessionId, fn, "__exec__"); + } + + session.state = "active"; + return fn.vms[0]!; +} + +export async function sendSessionMessage( + sessionId: string, + message: Record, + onStream?: (chunk: any) => void, + timeout: number = 60000, +): Promise { + const vm = await ensureSession(sessionId); + touchSession(sessionId); + + const id = message.id || crypto.randomUUID(); + const fullMessage = { ...message, id }; + + const socket = await getVmSocket(vm); + socket.write(JSON.stringify(fullMessage) + "\n"); + + gatewayLogger.debug( + { sessionId, messageType: message.type, messageId: id }, + "message sent to VM" + ); + + const result = await readVsockResponse(socket, timeout, onStream); + return { ...result, messageId: id }; +} diff --git a/src/runtime/vm-manager.ts b/src/runtime/vm-manager.ts index 172b3e7..5dca8dc 100644 --- a/src/runtime/vm-manager.ts +++ b/src/runtime/vm-manager.ts @@ -11,6 +11,7 @@ import type { RuntimeFunction, Vm } from "../types/types.js"; export async function createVm( functionId: string, fn: RuntimeFunction, + snapshotId?: string, ): Promise { const instanceId = crypto.randomBytes(4).toString("hex"); const apiSock = `/tmp/firecracker-${functionId}-${instanceId}.socket`; @@ -38,7 +39,7 @@ export async function createVm( await waitForFirecrackerApiSocket(apiSock); const client = createFcClient(apiSock); - await restoreVm(client, functionId, vsock); + await restoreVm(client, snapshotId || functionId, vsock); const vm: Vm = { id: instanceId, diff --git a/src/utils/logger.ts b/src/utils/logger.ts index cdc7d0c..7cf2fa6 100644 --- a/src/utils/logger.ts +++ b/src/utils/logger.ts @@ -33,7 +33,8 @@ export const vmManagerLogger = logger.child({ module: "vm-manager" }); export const transportLogger = logger.child({ module: "transport" }); export const protocolLogger = logger.child({ module: "protocol" }); export const cleanupLogger = logger.child({ module: "cleanup" }); -export const sessionLogger = logger.child({module: "session"}) +export const sessionLogger = logger.child({ module: "session" }) +export const gatewayLogger = logger.child({module: "gateway"}) export const httpLoggerOptions: Options = { logger: logger.child({ module: "http" }), From 469c3e0f049cf55ce26527603057b0cd4bafd783 Mon Sep 17 00:00:00 2001 From: vivek1504 Date: Fri, 17 Jul 2026 20:12:43 +0530 Subject: [PATCH 4/4] added exec routes --- src/routes/exec.ts | 90 +++++++++++++++++++++++++++++++++++++++++++++ src/utils/logger.ts | 3 +- 2 files changed, 92 insertions(+), 1 deletion(-) create mode 100644 src/routes/exec.ts diff --git a/src/routes/exec.ts b/src/routes/exec.ts new file mode 100644 index 0000000..7eab70e --- /dev/null +++ b/src/routes/exec.ts @@ -0,0 +1,90 @@ +import { Router } from "express"; +import { sendSessionMessage, ensureSession } from "../exec/gateway.js"; +import { destroySession, getSession, getAllSessions } from "../exec/session.js"; +import { execLogger } from "../utils/logger.js"; + +export const execRouter = Router(); + +execRouter.post("/:sessionId/execute", async (req, res) => { + const { sessionId } = req.params; + const { command, args, cwd, env, timeout } = req.body; + + if (!command) return res.status(400).json({ error: "command is required" }); + + const output: { stream: string; data: string; ts: number }[] = []; + + try { + const result = await sendSessionMessage( + sessionId, + { type: "execute", command, args, cwd, env, timeout }, + (chunk) => { + output.push({ stream: chunk.stream, data: chunk.data, ts: Date.now() }); + }, + ); + + res.json({ + exitCode: result.data?.exitCode, + signal: result.data?.signal, + duration: result.data?.duration, + output, + }); + } catch (err: any) { + res.status(500).json({ error: err.message }); + } +}); + +execRouter.post("/:sessionId/write", async (req, res) => { + const { sessionId } = req.params; + const { path, content, mode } = req.body; + + if (!path || !content) return res.status(400).json({ error: "path and content required" }); + + try { + const result = await sendSessionMessage(sessionId, { + type: "write_file", path, content, mode, + }); + res.json(result.data); + } catch (err: any) { + res.status(500).json({ error: err.message }); + } +}); + +execRouter.get("/:sessionId/read", async (req, res) => { + const { sessionId } = req.params; + const { path } = req.query; + + if (!path) return res.status(400).json({ error: "path query param required" }); + + try { + const result = await sendSessionMessage(sessionId, { + type: "read_file", path, + }); + res.json(result.data); + } catch (err: any) { + res.status(500).json({ error: err.message }); + } +}); + +execRouter.get("/:sessionId/files", async (req, res) => { + const { sessionId } = req.params; + const { path, recursive } = req.query; + + try { + const result = await sendSessionMessage(sessionId, { + type: "list_files", path, recursive: recursive === "true", + }); + res.json(result.data); + } catch (err: any) { + res.status(500).json({ error: err.message }); + } +}); + +execRouter.delete("/:sessionId", async (req, res) => { + const { sessionId } = req.params; + const destroyed = await destroySession(sessionId); + res.json({ destroyed }); +}); + +execRouter.get("/", (_req, res) => { + res.json({ sessions: getAllSessions() }); +}); diff --git a/src/utils/logger.ts b/src/utils/logger.ts index 7cf2fa6..606b4a9 100644 --- a/src/utils/logger.ts +++ b/src/utils/logger.ts @@ -34,7 +34,8 @@ export const transportLogger = logger.child({ module: "transport" }); export const protocolLogger = logger.child({ module: "protocol" }); export const cleanupLogger = logger.child({ module: "cleanup" }); export const sessionLogger = logger.child({ module: "session" }) -export const gatewayLogger = logger.child({module: "gateway"}) +export const gatewayLogger = logger.child({ module: "gateway" }) +export const execLogger= logger.child({module: "exec" }) export const httpLoggerOptions: Options = { logger: logger.child({ module: "http" }),