{"version":3,"file":"server.d.ts","sourceRoot":"","sources":["../src/server.ts"],"names":[],"mappings":"AAAA,OAAO,EAAgB,KAAK,MAAM,IAAI,UAAU,EAA6C,MAAM,WAAW,CAAC;AAY/G,OAAO,EAGN,KAAK,OAAO,EAGZ,KAAK,OAAO,EACZ,KAAK,KAAK,EAEV,KAAK,mBAAmB,EACxB,MAAM,uBAAuB,CAAC;AAE/B,OAAO,KAAK,EAAE,YAAY,EAAE,MAAM,aAAa,CAAC;AAOhD,OAAO,EAWN,KAAK,YAAY,EACjB,KAAK,oBAAoB,EAGzB,MAAM,oBAAoB,CAAC;AAE5B,OAAO,EAAE,UAAU,EAAE,KAAK,YAAY,EAAE,MAAM,aAAa,CAAC;AAc5D,UAAU,iBAAiB;IAC1B,SAAS,EAAE,MAAM,CAAC;IAClB,KAAK,CAAC,EAAE,MAAM,CAAC;IACf,KAAK,EAAE,KAAK,CAAC,GAAG,CAAC,CAAC;IAClB,OAAO,CAAC,EAAE,mBAAmB,CAAC;IAC9B,aAAa,CAAC,EAAE,oBAAoB,CAAC;IACrC,iBAAiB,CAAC,EAAE,OAAO,EAAE,CAAC;IAC9B,cAAc,CAAC,EAAE,OAAO,EAAE,CAAC;CAC3B;AA+4BD,wBAAgB,kBAAkB,CACjC,OAAO,EAAE,YAAY,EACrB,IAAI,EAAE,IAAI,CAAC,iBAAiB,EAAE,gBAAgB,GAAG,mBAAmB,CAAC,GACnE,OAAO,CAOT;AAoMD,wBAAgB,cAAc,CAAC,cAAc,CAAC,EAAE,OAAO,CAAC,YAAY,CAAC,GAAG,UAAU,CAoIjF;AAqED,wBAAgB,WAAW,CAAC,cAAc,CAAC,EAAE,OAAO,CAAC,YAAY,CAAC,GAAG,UAAU,CAU9E","sourcesContent":["import { createServer, type Server as HttpServer, type IncomingMessage, type ServerResponse } from \"node:http\";\nimport { createRequire } from \"node:module\";\nimport {\n\ttype CompactionPreparationOptions,\n\ttype CompactionSettings,\n\tcompactLegacy,\n\tDEFAULT_COMPACTION_SETTINGS,\n\ttype LegacyCompactResult,\n\ttype ProxyAssistantMessageEvent,\n\tprepareLegacyCompaction,\n\ttype SessionTreeEntry,\n} from \"@earendil-works/pi-agent-core\";\nimport {\n\ttype AssistantMessage,\n\ttype AssistantMessageEvent,\n\ttype Context,\n\tcreateModels,\n\tcreateProvider,\n\ttype Message,\n\ttype Model,\n\ttype ProviderStreams,\n\ttype SimpleStreamOptions,\n} from \"@earendil-works/pi-ai\";\nimport { streamSimple } from \"@earendil-works/pi-ai/compat\";\nimport type { ServerConfig } from \"./config.ts\";\nimport { loadConfig } from \"./config.ts\";\nimport { PiServerError, PiServerErrorCode, type PiServerErrorResponse } from \"./error-codes.ts\";\nimport { encodeErrorEvent, encodeProxyEvent } from \"./event-encoding.ts\";\nimport { ReceiveUploadError, receiveUpload } from \"./receive-upload.ts\";\nimport { CHUNK_ENDPOINT, type RequestChunkBody, receiveRequestChunk } from \"./request-chunks.ts\";\nimport { deletePersistedSession, loadPersistedSessions, savePersistedSession } from \"./session-persistence.ts\";\nimport {\n\tappendCompactionEntry,\n\tappendMessages,\n\tappendSessionEntries,\n\tdeleteSession as deleteSessionFromStore,\n\tgetOrCreateSession,\n\tgetSession,\n\tgetSessionBranch,\n\tlistSessions,\n\treplaceMessages,\n\treplaceSessionTree,\n\ttype SessionState,\n\ttype SessionStaticContext,\n\tsetStaticContext,\n\tswitchSessionLeaf,\n} from \"./session-store.ts\";\n\nexport { loadConfig, type ServerConfig } from \"./config.ts\";\n\ninterface PackageMetadata {\n\tversion: string;\n}\n\nconst packageMetadata = createRequire(import.meta.url)(\"../package.json\") as PackageMetadata;\nconst PI_SERVER_VERSION = packageMetadata.version;\n\ninterface SessionInitBody {\n\tsessionId: string;\n\tstaticContext?: SessionStaticContext;\n}\n\ninterface StreamRequestBody {\n\tsessionId: string;\n\trunId?: string;\n\tmodel: Model<any>;\n\toptions?: SimpleStreamOptions;\n\tstaticContext?: SessionStaticContext;\n\tephemeralMessages?: Message[];\n\tcontextOverlay?: Message[];\n}\n\ninterface SessionSyncBody {\n\tsessionId: string;\n\tmessages: Message[];\n\tstaticContext?: SessionStaticContext;\n}\n\ninterface SessionAppendBody {\n\tsessionId: string;\n\tmessages: Message[];\n\tstaticContext?: SessionStaticContext;\n}\n\ninterface SessionTreeSyncBody {\n\tsessionId: string;\n\tentries: SessionTreeEntry[];\n\tleafId: string | null;\n\tstaticContext?: SessionStaticContext;\n}\n\ninterface SessionTreeSwitchBody {\n\tsessionId: string;\n\tleafId: string | null;\n}\n\ninterface SessionCompactBody {\n\tsessionId: string;\n\trunId?: string;\n\tmodel: Model<any>;\n\toptions?: SimpleStreamOptions;\n\tsettings?: CompactionSettings;\n\tpreparation?: CompactionPreparationOptions;\n\tcustomInstructions?: string;\n\tbaseTreeHash?: string;\n\tfullResponse?: boolean;\n\tstreamResponse?: boolean;\n}\n\nfunction createRequestModels(model: Model<any>, options: SimpleStreamOptions) {\n\tconst models = createModels();\n\tconst requestStream: ProviderStreams[\"streamSimple\"] = (requestModel, context, streamOptions) =>\n\t\tstreamSimple(requestModel, context, {\n\t\t\t...streamOptions,\n\t\t\t...(options.timeoutMs !== undefined ? { timeoutMs: options.timeoutMs } : {}),\n\t\t});\n\tmodels.setProvider(\n\t\tcreateProvider({\n\t\t\tid: model.provider,\n\t\t\tname: model.provider,\n\t\t\tmodels: [model],\n\t\t\tauth: {\n\t\t\t\tapiKey: {\n\t\t\t\t\tname: \"pi-server request auth\",\n\t\t\t\t\tresolve: async () => ({ auth: { apiKey: options.apiKey, headers: options.headers } }),\n\t\t\t\t},\n\t\t\t},\n\t\t\tapi: {\n\t\t\t\tstream: requestStream,\n\t\t\t\tstreamSimple: requestStream,\n\t\t\t},\n\t\t}),\n\t);\n\treturn models;\n}\n\ninterface StreamRunRecord {\n\tsessionId: string;\n\trunId: string;\n\tkind: \"stream\" | \"compact\";\n\tstatus: \"running\" | \"completed\" | \"failed\" | \"aborted\";\n\tmessage?: AssistantMessage;\n\tcompactResult?: SessionCompactSuccessBody;\n\terrorMessage?: string;\n\tcreatedAt: number;\n\tupdatedAt: number;\n\tcancelRequested: boolean;\n\tcontroller: AbortController;\n\tsettled: Promise<void>;\n\tresolveSettled: (() => void) | undefined;\n}\n\nconst STREAM_HEARTBEAT = \": keep-alive\\n\\n\";\nconst JSON_HEARTBEAT = \" \\n\";\nconst STREAM_HEARTBEAT_INTERVAL_MS = 25_000;\nconst STREAM_RUN_TTL_MS = 60 * 60 * 1000; // 1 hour\nconst streamRuns = new Map<string, StreamRunRecord>();\n\nfunction readBody(req: IncomingMessage): Promise<string> {\n\treturn new Promise((resolve, reject) => {\n\t\tconst chunks: Buffer[] = [];\n\t\treq.on(\"data\", (chunk: Buffer) => chunks.push(chunk));\n\t\treq.on(\"end\", () => resolve(Buffer.concat(chunks).toString(\"utf-8\")));\n\t\treq.on(\"error\", reject);\n\t});\n}\n\nfunction sendJson(res: ServerResponse, status: number, body: unknown): void {\n\tconst data = JSON.stringify(body);\n\tres.writeHead(status, { \"Content-Type\": \"application/json\", \"Content-Length\": Buffer.byteLength(data) });\n\tres.end(data);\n}\n\nfunction sendError(\n\tres: ServerResponse,\n\tstatus: number,\n\tmessage: string,\n\tcode?: PiServerErrorCode,\n\tdetails?: Record<string, unknown>,\n): void {\n\tconst body: PiServerErrorResponse = { error: message };\n\tif (code) body.code = code;\n\tif (details) body.details = details;\n\tsendJson(res, status, body);\n}\n\nfunction logRequestError(req: IncomingMessage, error: unknown): void {\n\tconst message = error instanceof Error ? error.stack || error.message : String(error);\n\tconsole.error(`${req.method ?? \"UNKNOWN\"} ${req.url ?? \"/\"} failed: ${message}`);\n}\n\nfunction authenticate(config: ServerConfig, req: IncomingMessage): boolean {\n\tif (!config.authToken) return true;\n\tconst header = req.headers.authorization;\n\tif (!header) return false;\n\tconst token = header.startsWith(\"Bearer \") ? header.slice(7) : header;\n\treturn token === config.authToken;\n}\n\nfunction sessionResponseBody(session: SessionState) {\n\treturn {\n\t\tsessionId: session.sessionId,\n\t\tstaticContextHash: session.staticContextHash,\n\t\ttreeHash: session.treeHash,\n\t\tmessageCount: session.messages.length,\n\t\tentryCount: session.entries.length,\n\t\tleafId: session.leafId,\n\t\trevision: session.revision,\n\t};\n}\n\nfunction persistSession(config: ServerConfig, session: SessionState): void {\n\tsavePersistedSession(config.sessionStoreDir, session);\n}\n\nfunction runKey(sessionId: string, runId: string): string {\n\treturn `${sessionId}\\0${runId}`;\n}\n\nfunction getStreamRun(sessionId: string, runId: string): StreamRunRecord | undefined {\n\treturn streamRuns.get(runKey(sessionId, runId));\n}\n\nfunction createRunRecord(sessionId: string, runId: string, kind: StreamRunRecord[\"kind\"]): StreamRunRecord {\n\tlet resolveSettled: (() => void) | undefined;\n\tconst settled = new Promise<void>((resolve) => {\n\t\tresolveSettled = resolve;\n\t});\n\treturn {\n\t\tsessionId,\n\t\trunId,\n\t\tkind,\n\t\tstatus: \"running\",\n\t\tcreatedAt: Date.now(),\n\t\tupdatedAt: Date.now(),\n\t\tcancelRequested: false,\n\t\tcontroller: new AbortController(),\n\t\tsettled,\n\t\tresolveSettled,\n\t};\n}\n\nfunction settleRun(run: StreamRunRecord | undefined): void {\n\tif (!run) return;\n\tconst resolve = run.resolveSettled;\n\trun.resolveSettled = undefined;\n\tresolve?.();\n}\n\nfunction startStreamRun(sessionId: string, runId: string, kind: StreamRunRecord[\"kind\"] = \"stream\"): StreamRunRecord {\n\tconst existing = getStreamRun(sessionId, runId);\n\tif (existing?.status === \"completed\" || existing?.status === \"aborted\") return existing;\n\tconst run = existing ?? createRunRecord(sessionId, runId, kind);\n\tif (existing) {\n\t\trun.message = undefined;\n\t\trun.errorMessage = undefined;\n\t\trun.cancelRequested = false;\n\t\trun.controller = new AbortController();\n\t\tlet resolveSettled: (() => void) | undefined;\n\t\trun.settled = new Promise<void>((resolve) => {\n\t\t\tresolveSettled = resolve;\n\t\t});\n\t\trun.resolveSettled = resolveSettled;\n\t}\n\trun.kind = kind;\n\trun.status = \"running\";\n\trun.updatedAt = Date.now();\n\tstreamRuns.set(runKey(sessionId, runId), run);\n\treturn run;\n}\n\nfunction streamRunResponseBody(run: StreamRunRecord) {\n\treturn {\n\t\tsessionId: run.sessionId,\n\t\trunId: run.runId,\n\t\tstatus: run.status,\n\t\tcreatedAt: run.createdAt,\n\t\tupdatedAt: run.updatedAt,\n\t\t...(run.message ? { message: run.message } : {}),\n\t\t...(run.compactResult ? { compactResult: run.compactResult } : {}),\n\t\t...(run.errorMessage ? { errorMessage: run.errorMessage } : {}),\n\t};\n}\n\nfunction replayStreamRun(run: StreamRunRecord): ProxyAssistantMessageEvent[] {\n\tif (!run.message) return [];\n\tif (\n\t\trun.message.stopReason !== \"stop\" &&\n\t\trun.message.stopReason !== \"length\" &&\n\t\trun.message.stopReason !== \"toolUse\" &&\n\t\trun.message.stopReason !== \"deferred\"\n\t) {\n\t\treturn [];\n\t}\n\tconst events: ProxyAssistantMessageEvent[] = [{ type: \"start\" }];\n\tfor (const [contentIndex, content] of run.message.content.entries()) {\n\t\tswitch (content.type) {\n\t\t\tcase \"text\":\n\t\t\t\tevents.push(\n\t\t\t\t\t{ type: \"text_start\", contentIndex },\n\t\t\t\t\t{ type: \"text_delta\", contentIndex, delta: content.text },\n\t\t\t\t\t{ type: \"text_end\", contentIndex, contentSignature: content.textSignature },\n\t\t\t\t);\n\t\t\t\tbreak;\n\t\t\tcase \"thinking\":\n\t\t\t\tevents.push(\n\t\t\t\t\t{ type: \"thinking_start\", contentIndex },\n\t\t\t\t\t{ type: \"thinking_delta\", contentIndex, delta: content.thinking },\n\t\t\t\t\t{ type: \"thinking_end\", contentIndex, contentSignature: content.thinkingSignature },\n\t\t\t\t);\n\t\t\t\tbreak;\n\t\t\tcase \"toolCall\":\n\t\t\t\tevents.push(\n\t\t\t\t\t{ type: \"toolcall_start\", contentIndex, id: content.id, toolName: content.name },\n\t\t\t\t\t{ type: \"toolcall_delta\", contentIndex, delta: JSON.stringify(content.arguments) },\n\t\t\t\t\t{ type: \"toolcall_end\", contentIndex, toolCall: content },\n\t\t\t\t);\n\t\t\t\tbreak;\n\t\t\tdefault:\n\t\t\t\tbreak;\n\t\t}\n\t}\n\tevents.push({\n\t\ttype: \"done\",\n\t\treason: run.message.stopReason,\n\t\tusage: run.message.usage,\n\t\tdeferred: run.message.deferred,\n\t\tmessage: run.message,\n\t});\n\treturn events;\n}\n\nfunction completeStreamRun(run: StreamRunRecord | undefined, message: AssistantMessage): void {\n\tif (!run) return;\n\tif (run.cancelRequested) {\n\t\trun.status = \"aborted\";\n\t\trun.message = undefined;\n\t\trun.errorMessage = \"Stream run aborted\";\n\t\trun.updatedAt = Date.now();\n\t\treturn;\n\t}\n\trun.status = \"completed\";\n\trun.message = message;\n\trun.errorMessage = undefined;\n\trun.updatedAt = Date.now();\n}\n\nfunction failStreamRun(run: StreamRunRecord | undefined, errorMessage: string): void {\n\tif (!run) return;\n\tif (run.cancelRequested) {\n\t\trun.status = \"aborted\";\n\t\trun.errorMessage = \"Stream run aborted\";\n\t\trun.updatedAt = Date.now();\n\t\treturn;\n\t}\n\trun.status = \"failed\";\n\trun.errorMessage = errorMessage;\n\trun.updatedAt = Date.now();\n}\n\nfunction touchStreamRun(run: StreamRunRecord | undefined): void {\n\tif (run) run.updatedAt = Date.now();\n}\n\nfunction writeStreamEvent(res: ServerResponse, event: ProxyAssistantMessageEvent): void {\n\tif (!res.writableEnded && !res.destroyed) {\n\t\tres.write(encodeProxyEvent(event));\n\t}\n}\n\nfunction writeStreamError(res: ServerResponse, message: string): void {\n\tif (!res.writableEnded && !res.destroyed) {\n\t\tres.write(encodeErrorEvent(message));\n\t}\n}\n\nfunction endStreamResponse(res: ServerResponse): void {\n\tif (!res.writableEnded && !res.destroyed) {\n\t\tres.end();\n\t}\n}\n\nfunction cleanupExpiredStreamRuns(nowMs: number): void {\n\tfor (const [key, run] of streamRuns) {\n\t\tif (nowMs - run.updatedAt > STREAM_RUN_TTL_MS) {\n\t\t\tstreamRuns.delete(key);\n\t\t}\n\t}\n}\n\nfunction deleteStreamRunsForSession(sessionId: string): void {\n\tfor (const [key, run] of streamRuns) {\n\t\tif (run.sessionId === sessionId) streamRuns.delete(key);\n\t}\n}\n\nconst streamRunCleanupTimer = setInterval(() => cleanupExpiredStreamRuns(Date.now()), STREAM_RUN_TTL_MS / 4);\nstreamRunCleanupTimer.unref();\n\nfunction sessionHistoryFullResponseBody(session: SessionState, baseMessageCount: number) {\n\treturn {\n\t\tsessionId: session.sessionId,\n\t\tstaticContext: session.staticContext,\n\t\tstaticContextHash: session.staticContextHash,\n\t\ttreeHash: session.treeHash,\n\t\tmessageCount: session.messages.length,\n\t\tentryCount: session.entries.length,\n\t\tleafId: session.leafId,\n\t\trevision: session.revision,\n\t\tentries: session.entries,\n\t\tbaseMessageCount,\n\t\tmessages: session.messages.slice(baseMessageCount),\n\t};\n}\n\nfunction sessionTreePatchResponseBody(\n\tsession: SessionState,\n\tbaseMessageCount: number,\n\tentriesFrom: number,\n\tbaseRevision: number | undefined,\n) {\n\treturn {\n\t\tsessionId: session.sessionId,\n\t\tstaticContext: session.staticContext,\n\t\tstaticContextHash: session.staticContextHash,\n\t\ttreeHash: session.treeHash,\n\t\tmessageCount: session.messages.length,\n\t\tentryCount: session.entries.length,\n\t\tleafId: session.leafId,\n\t\trevision: session.revision,\n\t\tbaseMessageCount,\n\t\tmessages: session.messages.slice(baseMessageCount),\n\t\ttreePatch: {\n\t\t\tentriesFrom,\n\t\t\tbaseRevision,\n\t\t\tentries: session.entries.slice(entriesFrom),\n\t\t\tleafId: session.leafId,\n\t\t\trevision: session.revision,\n\t\t},\n\t};\n}\n\ntype PreparedCompaction = Parameters<typeof compactLegacy>[0];\n\ninterface PreparedSessionCompact {\n\tsession: SessionState;\n\tsessionId: string;\n\trevision: number;\n\ttreeHash: string;\n\tleafId: string | null;\n\tentryCount: number;\n\tpreparation: PreparedCompaction;\n\toptions: SimpleStreamOptions;\n}\n\ninterface SessionCompactSuccessBody {\n\tsuccess: true;\n\tcompaction: Omit<LegacyCompactResult, \"retainedTail\">;\n\tcompactionEntry: SessionTreeEntry;\n\tsessionId: string;\n\tstaticContextHash: string;\n\ttreeHash: string;\n\tmessageCount: number;\n\tentryCount: number;\n\tleafId: string | null;\n\trevision: number;\n\tstaticContext?: SessionStaticContext;\n\ttreePatch?: {\n\t\tbaseTreeHash: string;\n\t\tentriesFrom: number;\n\t\tentries: SessionTreeEntry[];\n\t\tleafId: string | null;\n\t\trevision: number;\n\t};\n\tentries?: SessionTreeEntry[];\n\tmessages?: Message[];\n}\n\ninterface SessionCompactHttpResponse {\n\tstatus: number;\n\tbody: PiServerErrorResponse | SessionCompactSuccessBody;\n}\n\nfunction handleSessionInit(config: ServerConfig, body: SessionInitBody, res: ServerResponse): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (body.staticContext) {\n\t\tsetStaticContext(body.sessionId, body.staticContext);\n\t} else {\n\t\tgetOrCreateSession(body.sessionId);\n\t}\n\tconst session = getSession(body.sessionId)!;\n\tpersistSession(config, session);\n\tsendJson(res, 200, sessionResponseBody(session));\n}\n\nfunction handleSessionUpdate(\n\tconfig: ServerConfig,\n\tbody: SessionInitBody & { staticContext: SessionStaticContext },\n\tres: ServerResponse,\n): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (!body.staticContext) {\n\t\tsendError(res, 400, \"staticContext is required for update\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tsetStaticContext(body.sessionId, body.staticContext);\n\tconst session = getSession(body.sessionId)!;\n\tpersistSession(config, session);\n\tsendJson(res, 200, sessionResponseBody(session));\n}\n\nfunction handleSessionSync(config: ServerConfig, body: SessionSyncBody, res: ServerResponse): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (!Array.isArray(body.messages)) {\n\t\tsendError(res, 400, \"messages is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (body.staticContext) {\n\t\tsetStaticContext(body.sessionId, body.staticContext);\n\t}\n\tconst session = replaceMessages(body.sessionId, body.messages);\n\tpersistSession(config, session);\n\tsendJson(res, 200, sessionResponseBody(session));\n}\n\nfunction handleSessionAppend(config: ServerConfig, body: SessionAppendBody, res: ServerResponse): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (!Array.isArray(body.messages)) {\n\t\tsendError(res, 400, \"messages is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (body.staticContext) {\n\t\tsetStaticContext(body.sessionId, body.staticContext);\n\t}\n\tconst session = appendMessages(body.sessionId, body.messages);\n\tpersistSession(config, session);\n\tsendJson(res, 200, sessionResponseBody(session));\n}\n\nfunction handleSessionTreeSync(config: ServerConfig, body: SessionTreeSyncBody, res: ServerResponse): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (!Array.isArray(body.entries)) {\n\t\tsendError(res, 400, \"entries is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (body.staticContext) {\n\t\tsetStaticContext(body.sessionId, body.staticContext);\n\t}\n\tconst session = replaceSessionTree(body.sessionId, body.entries, body.leafId ?? null);\n\tpersistSession(config, session);\n\tsendJson(res, 200, sessionResponseBody(session));\n}\n\nfunction handleSessionTreeAppend(config: ServerConfig, body: SessionTreeSyncBody, res: ServerResponse): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (!Array.isArray(body.entries)) {\n\t\tsendError(res, 400, \"entries is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\ttry {\n\t\tif (body.staticContext) {\n\t\t\tsetStaticContext(body.sessionId, body.staticContext);\n\t\t}\n\t\tconst session = appendSessionEntries(body.sessionId, body.entries, body.leafId ?? null);\n\t\tpersistSession(config, session);\n\t\tsendJson(res, 200, sessionResponseBody(session));\n\t} catch (error) {\n\t\tif (error instanceof PiServerError) {\n\t\t\tsendError(res, 400, error.message, error.code, error.details);\n\t\t} else {\n\t\t\tconst message = error instanceof Error ? error.message : String(error);\n\t\t\tsendError(res, 500, message, PiServerErrorCode.INTERNAL_ERROR);\n\t\t}\n\t}\n}\n\nfunction handleSessionTreeSwitch(config: ServerConfig, body: SessionTreeSwitchBody, res: ServerResponse): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\ttry {\n\t\tconst session = switchSessionLeaf(body.sessionId, body.leafId ?? null);\n\t\tpersistSession(config, session);\n\t\tsendJson(res, 200, sessionResponseBody(session));\n\t} catch (error) {\n\t\tif (error instanceof PiServerError) {\n\t\t\tsendError(res, 400, error.message, error.code, error.details);\n\t\t} else {\n\t\t\tconst message = error instanceof Error ? error.message : String(error);\n\t\t\tsendError(res, 500, message, PiServerErrorCode.INTERNAL_ERROR);\n\t\t}\n\t}\n}\n\nfunction prepareSessionCompact(body: SessionCompactBody): PreparedSessionCompact | SessionCompactHttpResponse {\n\tif (!body.sessionId) {\n\t\treturn { status: 400, body: { error: \"sessionId is required\" } };\n\t}\n\tif (!body.model) {\n\t\treturn { status: 400, body: { error: \"model is required\" } };\n\t}\n\n\tconst session = getSession(body.sessionId);\n\tif (!session) {\n\t\treturn { status: 404, body: { error: \"session not found\", code: PiServerErrorCode.SESSION_NOT_FOUND } };\n\t}\n\n\tconst entries = getSessionBranch(session);\n\tconst leaf = entries.at(-1);\n\tif (leaf?.type === \"compaction\") {\n\t\tif (!leaf.firstKeptEntryId) {\n\t\t\treturn { status: 400, body: { error: \"Stored compaction is missing firstKeptEntryId\" } };\n\t\t}\n\t\treturn {\n\t\t\tstatus: 200,\n\t\t\tbody: {\n\t\t\t\tsuccess: true,\n\t\t\t\tcompaction: {\n\t\t\t\t\tsummary: leaf.summary,\n\t\t\t\t\tfirstKeptEntryId: leaf.firstKeptEntryId,\n\t\t\t\t\ttokensBefore: leaf.tokensBefore,\n\t\t\t\t\tdetails: leaf.details,\n\t\t\t\t},\n\t\t\t\tcompactionEntry: leaf,\n\t\t\t\t...sessionResponseBody(session),\n\t\t\t\tentries: session.entries,\n\t\t\t\tmessages: session.messages,\n\t\t\t},\n\t\t};\n\t}\n\tconst preparationResult = prepareLegacyCompaction(\n\t\tentries,\n\t\tbody.settings ?? DEFAULT_COMPACTION_SETTINGS,\n\t\tbody.preparation,\n\t);\n\tif (!preparationResult.ok) {\n\t\treturn { status: 400, body: { error: preparationResult.error.message } };\n\t}\n\tif (!preparationResult.value) {\n\t\treturn { status: 400, body: { error: \"Nothing to compact\" } };\n\t}\n\n\tconst options = body.options ?? {};\n\treturn {\n\t\tsession,\n\t\tsessionId: session.sessionId,\n\t\trevision: session.revision,\n\t\ttreeHash: session.treeHash,\n\t\tleafId: session.leafId,\n\t\tentryCount: session.entries.length,\n\t\tpreparation: preparationResult.value,\n\t\toptions,\n\t};\n}\n\nfunction sessionCompactConflict(\n\tbody: SessionCompactBody,\n\tprepared: PreparedSessionCompact,\n\tcurrentSession: SessionState | undefined,\n): SessionCompactHttpResponse {\n\treturn {\n\t\tstatus: 409,\n\t\tbody: {\n\t\t\terror: \"Session changed while compaction was running\",\n\t\t\tcode: PiServerErrorCode.SESSION_STATE_CONFLICT,\n\t\t\tdetails: {\n\t\t\t\tsessionId: body.sessionId,\n\t\t\t\texpectedRevision: prepared.revision,\n\t\t\t\texpectedTreeHash: prepared.treeHash,\n\t\t\t\texpectedLeafId: prepared.leafId,\n\t\t\t\tactualRevision: currentSession?.revision ?? null,\n\t\t\t\tactualTreeHash: currentSession?.treeHash ?? null,\n\t\t\t\tactualLeafId: currentSession?.leafId ?? null,\n\t\t\t},\n\t\t},\n\t};\n}\n\nasync function completeSessionCompact(\n\tconfig: ServerConfig,\n\tbody: SessionCompactBody,\n\tprepared: PreparedSessionCompact,\n\trun?: StreamRunRecord,\n): Promise<SessionCompactHttpResponse> {\n\tconst result = await compactLegacy(\n\t\tprepared.preparation,\n\t\tcreateRequestModels(body.model, prepared.options),\n\t\tbody.model,\n\t\tbody.customInstructions,\n\t\trun?.controller.signal,\n\t\tprepared.options.reasoning,\n\t);\n\tif (!result.ok) {\n\t\treturn { status: 500, body: { error: result.error.message, code: PiServerErrorCode.INTERNAL_ERROR } };\n\t}\n\tif (run?.cancelRequested || run?.controller.signal.aborted) {\n\t\treturn { status: 409, body: { error: \"Compaction aborted\", code: PiServerErrorCode.INVALID_REQUEST } };\n\t}\n\n\tconst currentSession = getSession(prepared.sessionId);\n\tif (\n\t\t!currentSession ||\n\t\tcurrentSession !== prepared.session ||\n\t\tcurrentSession.revision !== prepared.revision ||\n\t\tcurrentSession.treeHash !== prepared.treeHash ||\n\t\tcurrentSession.leafId !== prepared.leafId\n\t) {\n\t\treturn sessionCompactConflict(body, prepared, currentSession);\n\t}\n\n\tconst baseTreeHash = prepared.treeHash;\n\tconst baseEntryCount = prepared.entryCount;\n\tconst compaction = result.value;\n\tconst { session: updatedSession, entry: compactionEntry } = appendCompactionEntry(body.sessionId, compaction);\n\tpersistSession(config, updatedSession);\n\tif (!body.fullResponse && body.baseTreeHash === baseTreeHash) {\n\t\treturn {\n\t\t\tstatus: 200,\n\t\t\tbody: {\n\t\t\t\tsuccess: true,\n\t\t\t\tcompaction: compaction satisfies LegacyCompactResult,\n\t\t\t\tcompactionEntry,\n\t\t\t\t...sessionResponseBody(updatedSession),\n\t\t\t\tstaticContext: updatedSession.staticContext,\n\t\t\t\ttreePatch: {\n\t\t\t\t\tbaseTreeHash,\n\t\t\t\t\tentriesFrom: baseEntryCount,\n\t\t\t\t\tentries: [compactionEntry],\n\t\t\t\t\tleafId: updatedSession.leafId,\n\t\t\t\t\trevision: updatedSession.revision,\n\t\t\t\t},\n\t\t\t},\n\t\t};\n\t}\n\treturn {\n\t\tstatus: 200,\n\t\tbody: {\n\t\t\tsuccess: true,\n\t\t\tcompaction: compaction satisfies LegacyCompactResult,\n\t\t\tcompactionEntry,\n\t\t\t...sessionResponseBody(updatedSession),\n\t\t\tstaticContext: updatedSession.staticContext,\n\t\t\tentries: updatedSession.entries,\n\t\t\tmessages: updatedSession.messages,\n\t\t},\n\t};\n}\n\nfunction writeServerSentEvent(res: ServerResponse, event: string, body: unknown): void {\n\tres.write(`event: ${event}\\n`);\n\tres.write(`data: ${JSON.stringify(body)}\\n\\n`);\n}\n\nasync function handleSessionCompactStream(\n\tconfig: ServerConfig,\n\tbody: SessionCompactBody,\n\tprepared: PreparedSessionCompact,\n\tres: ServerResponse,\n\trun?: StreamRunRecord,\n): Promise<void> {\n\tres.writeHead(200, {\n\t\t\"Content-Type\": \"text/event-stream\",\n\t\t\"Cache-Control\": \"no-cache\",\n\t\tConnection: \"keep-alive\",\n\t});\n\tres.flushHeaders();\n\tres.write(STREAM_HEARTBEAT);\n\n\tconst heartbeat = setInterval(() => {\n\t\ttouchStreamRun(run);\n\t\tif (!res.writableEnded) {\n\t\t\tres.write(STREAM_HEARTBEAT);\n\t\t}\n\t}, STREAM_HEARTBEAT_INTERVAL_MS);\n\theartbeat.unref();\n\n\ttry {\n\t\tconst result = await completeSessionCompact(config, body, prepared, run);\n\t\tfinishCompactRun(run, result);\n\t\twriteServerSentEvent(res, result.status >= 400 ? \"error\" : \"result\", result.body);\n\t} catch (err) {\n\t\tconst message = err instanceof Error ? err.message : String(err);\n\t\tif (run?.status === \"running\") {\n\t\t\trun.status = run.cancelRequested ? \"aborted\" : \"failed\";\n\t\t\trun.errorMessage = message;\n\t\t\trun.updatedAt = Date.now();\n\t\t}\n\t\twriteServerSentEvent(res, \"error\", { error: message });\n\t} finally {\n\t\tclearInterval(heartbeat);\n\t\tsettleRun(run);\n\t\tres.end();\n\t}\n}\n\nasync function handleSessionCompactJsonStream(\n\tconfig: ServerConfig,\n\tbody: SessionCompactBody,\n\tprepared: PreparedSessionCompact,\n\tres: ServerResponse,\n\trun?: StreamRunRecord,\n): Promise<void> {\n\tres.writeHead(200, {\n\t\t\"Content-Type\": \"application/json\",\n\t\t\"Cache-Control\": \"no-cache\",\n\t\tConnection: \"keep-alive\",\n\t});\n\tres.flushHeaders();\n\tres.write(JSON_HEARTBEAT);\n\n\tconst heartbeat = setInterval(() => {\n\t\ttouchStreamRun(run);\n\t\tif (!res.writableEnded) {\n\t\t\tres.write(JSON_HEARTBEAT);\n\t\t}\n\t}, STREAM_HEARTBEAT_INTERVAL_MS);\n\theartbeat.unref();\n\n\ttry {\n\t\tconst result = await completeSessionCompact(config, body, prepared, run);\n\t\tfinishCompactRun(run, result);\n\t\tres.write(JSON.stringify(result.body));\n\t} catch (err) {\n\t\tconst message = err instanceof Error ? err.message : String(err);\n\t\tif (run?.status === \"running\") {\n\t\t\trun.status = run.cancelRequested ? \"aborted\" : \"failed\";\n\t\t\trun.errorMessage = message;\n\t\t\trun.updatedAt = Date.now();\n\t\t}\n\t\tres.write(JSON.stringify({ error: message }));\n\t} finally {\n\t\tclearInterval(heartbeat);\n\t\tsettleRun(run);\n\t\tres.end();\n\t}\n}\n\nfunction finishCompactRun(run: StreamRunRecord | undefined, result: SessionCompactHttpResponse): void {\n\tif (!run) return;\n\tif (run.cancelRequested || run.controller.signal.aborted) {\n\t\trun.status = \"aborted\";\n\t\trun.errorMessage = \"Compaction aborted\";\n\t} else {\n\t\trun.status = result.status >= 400 ? \"failed\" : \"completed\";\n\t\tif (result.status >= 400 && \"error\" in result.body) {\n\t\t\trun.errorMessage = result.body.error;\n\t\t} else if (\"success\" in result.body) {\n\t\t\trun.compactResult = { ...result.body, entries: result.body.entries?.slice() };\n\t\t}\n\t}\n\trun.updatedAt = Date.now();\n}\n\nasync function handleSessionCompact(\n\tconfig: ServerConfig,\n\tbody: SessionCompactBody,\n\tres: ServerResponse,\n): Promise<void> {\n\tconst existingRun = body.runId ? getStreamRun(body.sessionId, body.runId) : undefined;\n\tif (existingRun?.status === \"aborted\") {\n\t\tsendJson(res, 409, { error: \"Compaction run was aborted\", code: PiServerErrorCode.INVALID_REQUEST });\n\t\treturn;\n\t}\n\tif (existingRun?.status === \"completed\" && existingRun.compactResult) {\n\t\tsendJson(res, 200, existingRun.compactResult);\n\t\treturn;\n\t}\n\tconst prepared = prepareSessionCompact(body);\n\tif (\"status\" in prepared) {\n\t\tsendJson(res, prepared.status, prepared.body);\n\t\treturn;\n\t}\n\n\tlet run: StreamRunRecord | undefined;\n\tif (body.runId) {\n\t\tconst existingRun = getStreamRun(body.sessionId, body.runId);\n\t\tif (existingRun?.status === \"aborted\") {\n\t\t\tsendJson(res, 409, { error: \"Compaction run was aborted\", code: PiServerErrorCode.INVALID_REQUEST });\n\t\t\treturn;\n\t\t}\n\t\tif (existingRun?.status === \"running\") {\n\t\t\tsendError(res, 409, \"A run with this runId is already in progress\", PiServerErrorCode.RUN_IN_PROGRESS);\n\t\t\treturn;\n\t\t}\n\t\trun = startStreamRun(body.sessionId, body.runId, \"compact\");\n\t}\n\n\tif (body.streamResponse) {\n\t\tawait handleSessionCompactStream(config, body, prepared, res, run);\n\t\treturn;\n\t}\n\n\tawait handleSessionCompactJsonStream(config, body, prepared, res, run);\n}\n\nfunction handleSessionHistory(\n\tsessionId: string,\n\tfrom: number | undefined,\n\tentriesFrom: number | undefined,\n\trevision: number | undefined,\n\tbaseTreeHash: string | undefined,\n\tres: ServerResponse,\n): void {\n\tconst session = getSession(sessionId);\n\tif (!session) {\n\t\tsendError(res, 404, \"session not found\", PiServerErrorCode.SESSION_NOT_FOUND);\n\t\treturn;\n\t}\n\tconst baseMessageCount = from ?? 0;\n\tif (\n\t\tentriesFrom !== undefined &&\n\t\tentriesFrom <= session.entries.length &&\n\t\t(revision === undefined || revision <= session.revision) &&\n\t\t(baseTreeHash === undefined || baseTreeHash === session.prefixHashes[entriesFrom])\n\t) {\n\t\tsendJson(res, 200, sessionTreePatchResponseBody(session, baseMessageCount, entriesFrom, revision));\n\t\treturn;\n\t}\n\tsendJson(res, 200, sessionHistoryFullResponseBody(session, baseMessageCount));\n}\n\nfunction handleSessionRun(sessionId: string, runId: string, res: ServerResponse): void {\n\tconst run = getStreamRun(sessionId, runId);\n\tif (!run) {\n\t\tsendError(res, 404, \"run not found\", PiServerErrorCode.RUN_NOT_FOUND);\n\t\treturn;\n\t}\n\tsendJson(res, 200, streamRunResponseBody(run));\n}\n\nasync function handleSessionRunAbort(sessionId: string, runId: string, res: ServerResponse): Promise<void> {\n\tif (!sessionId || !runId) {\n\t\tsendError(res, 400, \"sessionId and runId are required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\n\tlet run = getStreamRun(sessionId, runId);\n\tif (!run) {\n\t\trun = createRunRecord(sessionId, runId, \"stream\");\n\t\trun.status = \"aborted\";\n\t\trun.cancelRequested = true;\n\t\trun.errorMessage = \"Stream run aborted before it started\";\n\t\trun.updatedAt = Date.now();\n\t\tsettleRun(run);\n\t\tstreamRuns.set(runKey(sessionId, runId), run);\n\t\tsendJson(res, 200, streamRunResponseBody(run));\n\t\treturn;\n\t}\n\n\tif (run.status !== \"running\") {\n\t\tawait run.settled;\n\t\tsendJson(res, 200, streamRunResponseBody(run));\n\t\treturn;\n\t}\n\n\trun.cancelRequested = true;\n\trun.updatedAt = Date.now();\n\trun.controller.abort();\n\tawait run.settled;\n\tsendJson(res, 200, streamRunResponseBody(run));\n}\n\nexport function buildStreamContext(\n\tsession: SessionState,\n\tbody: Pick<StreamRequestBody, \"contextOverlay\" | \"ephemeralMessages\">,\n): Context {\n\tconst messages = body.contextOverlay ?? [...session.messages, ...(body.ephemeralMessages ?? [])];\n\treturn {\n\t\tsystemPrompt: session.staticContext?.systemPrompt,\n\t\tmessages,\n\t\ttools: session.staticContext?.tools,\n\t};\n}\n\nfunction handleStream(config: ServerConfig, body: StreamRequestBody, res: ServerResponse): void {\n\tif (!body.sessionId) {\n\t\tsendError(res, 400, \"sessionId is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\tif (!body.model) {\n\t\tsendError(res, 400, \"model is required\", PiServerErrorCode.REQUIRED_FIELD_MISSING);\n\t\treturn;\n\t}\n\n\tcleanupExpiredStreamRuns(Date.now());\n\n\tconst session = getOrCreateSession(body.sessionId);\n\n\tif (body.staticContext) {\n\t\tsetStaticContext(body.sessionId, body.staticContext);\n\t\tpersistSession(config, session);\n\t}\n\n\tif (!session.staticContext && !body.staticContext) {\n\t\tsendError(\n\t\t\tres,\n\t\t\t400,\n\t\t\t\"Session has no static context. Initialize with /api/session/init first.\",\n\t\t\tPiServerErrorCode.SESSION_NO_STATIC_CONTEXT,\n\t\t);\n\t\treturn;\n\t}\n\n\tif (body.ephemeralMessages !== undefined && !Array.isArray(body.ephemeralMessages)) {\n\t\tsendError(res, 400, \"ephemeralMessages must be an array\", PiServerErrorCode.INVALID_REQUEST);\n\t\treturn;\n\t}\n\tif (body.contextOverlay !== undefined && !Array.isArray(body.contextOverlay)) {\n\t\tsendError(res, 400, \"contextOverlay must be an array\", PiServerErrorCode.INVALID_REQUEST);\n\t\treturn;\n\t}\n\n\tconst context = buildStreamContext(session, body);\n\tconst existingRun = body.runId ? getStreamRun(body.sessionId, body.runId) : undefined;\n\n\tconst resolvedModel = body.model;\n\n\tif (existingRun?.status === \"completed\") {\n\t\tres.writeHead(200, {\n\t\t\t\"Content-Type\": \"text/event-stream\",\n\t\t\t\"Cache-Control\": \"no-cache\",\n\t\t\tConnection: \"keep-alive\",\n\t\t});\n\t\tres.flushHeaders();\n\t\tres.write(STREAM_HEARTBEAT);\n\t\tfor (const event of replayStreamRun(existingRun)) {\n\t\t\twriteStreamEvent(res, event);\n\t\t}\n\t\tendStreamResponse(res);\n\t\treturn;\n\t}\n\tif (existingRun?.status === \"aborted\") {\n\t\tsendError(res, 409, \"The stream run was aborted\", PiServerErrorCode.INVALID_REQUEST);\n\t\treturn;\n\t}\n\tif (existingRun?.status === \"running\") {\n\t\tsendError(res, 409, \"A stream with this runId is already in progress\", PiServerErrorCode.RUN_IN_PROGRESS);\n\t\treturn;\n\t}\n\n\tconst run = body.runId ? startStreamRun(body.sessionId, body.runId) : undefined;\n\tconst streamOptions: SimpleStreamOptions = {\n\t\t...(body.options ?? {}),\n\t\t...(run ? { signal: run.controller.signal } : {}),\n\t};\n\n\tres.writeHead(200, {\n\t\t\"Content-Type\": \"text/event-stream\",\n\t\t\"Cache-Control\": \"no-cache\",\n\t\tConnection: \"keep-alive\",\n\t});\n\tres.flushHeaders();\n\tres.write(STREAM_HEARTBEAT);\n\n\tconst heartbeat = setInterval(() => {\n\t\ttouchStreamRun(run);\n\t\tif (!res.writableEnded && !res.destroyed) {\n\t\t\tres.write(STREAM_HEARTBEAT);\n\t\t}\n\t}, STREAM_HEARTBEAT_INTERVAL_MS);\n\theartbeat.unref();\n\n\tlet stream: AsyncIterable<AssistantMessageEvent>;\n\ttry {\n\t\tstream = streamSimple(resolvedModel, context, streamOptions);\n\t} catch (err) {\n\t\tclearInterval(heartbeat);\n\t\tconst message = err instanceof Error ? err.message : String(err);\n\t\tfailStreamRun(run, message);\n\t\tsettleRun(run);\n\t\twriteStreamError(res, message);\n\t\tendStreamResponse(res);\n\t\treturn;\n\t}\n\n\tvoid (async () => {\n\t\ttry {\n\t\t\tfor await (const event of stream) {\n\t\t\t\tconst proxyEvent = toProxyEvent(event);\n\t\t\t\tif (proxyEvent) {\n\t\t\t\t\ttouchStreamRun(run);\n\t\t\t\t\twriteStreamEvent(res, proxyEvent);\n\t\t\t\t}\n\t\t\t\tif (event.type === \"done\") {\n\t\t\t\t\tcompleteStreamRun(run, event.message);\n\t\t\t\t} else if (event.type === \"error\") {\n\t\t\t\t\tfailStreamRun(run, event.error.errorMessage ?? event.reason);\n\t\t\t\t}\n\t\t\t}\n\t\t} catch (err) {\n\t\t\tconst message = err instanceof Error ? err.message : String(err);\n\t\t\tfailStreamRun(run, message);\n\t\t\twriteStreamError(res, message);\n\t\t} finally {\n\t\t\tclearInterval(heartbeat);\n\t\t\tif (run?.status === \"running\") {\n\t\t\t\tfailStreamRun(run, run.cancelRequested ? \"Stream run aborted\" : \"Stream ended before a terminal event\");\n\t\t\t}\n\t\t\tsettleRun(run);\n\t\t\tendStreamResponse(res);\n\t\t}\n\t})();\n}\n\nasync function handlePostRequest(\n\tconfig: ServerConfig,\n\tpathname: string,\n\tbody: unknown,\n\tres: ServerResponse,\n): Promise<boolean> {\n\tif (pathname === \"/api/receive\") {\n\t\ttry {\n\t\t\tsendJson(res, 200, receiveUpload(config.uploadDir, body));\n\t\t} catch (error) {\n\t\t\tif (!(error instanceof ReceiveUploadError)) throw error;\n\t\t\tsendJson(res, error.status, { error: error.message });\n\t\t}\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/init\") {\n\t\thandleSessionInit(config, body as SessionInitBody, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/update\") {\n\t\thandleSessionUpdate(config, body as SessionInitBody & { staticContext: SessionStaticContext }, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/sync\") {\n\t\thandleSessionSync(config, body as SessionSyncBody, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/append\") {\n\t\thandleSessionAppend(config, body as SessionAppendBody, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/tree/sync\") {\n\t\thandleSessionTreeSync(config, body as SessionTreeSyncBody, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/tree/append\") {\n\t\thandleSessionTreeAppend(config, body as SessionTreeSyncBody, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/tree/switch\") {\n\t\thandleSessionTreeSwitch(config, body as SessionTreeSwitchBody, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/session/compact\") {\n\t\tawait handleSessionCompact(config, body as SessionCompactBody, res);\n\t\treturn true;\n\t}\n\n\tif (pathname === \"/api/stream\") {\n\t\thandleStream(config, body as StreamRequestBody, res);\n\t\treturn true;\n\t}\n\n\treturn false;\n}\n\nexport function createPiServer(configOverride?: Partial<ServerConfig>): HttpServer {\n\tconst config = loadConfig(configOverride);\n\tloadPersistedSessions(config.sessionStoreDir);\n\n\tconst server = createServer(async (req, res) => {\n\t\tconst url = new URL(req.url ?? \"/\", \"http://localhost\");\n\n\t\tif (req.method === \"GET\" && url.pathname === \"/health\") {\n\t\t\tsendJson(res, 200, { status: \"ok\" });\n\t\t\treturn;\n\t\t}\n\n\t\tif (!authenticate(config, req)) {\n\t\t\tconst body =\n\t\t\t\treq.method === \"GET\" && url.pathname === \"/\"\n\t\t\t\t\t? { error: \"Unauthorized\", version: PI_SERVER_VERSION }\n\t\t\t\t\t: { error: \"Unauthorized\" };\n\t\t\tsendJson(res, 401, body);\n\t\t\treturn;\n\t\t}\n\n\t\tif (req.method === \"GET\" && url.pathname === \"/api/sessions\") {\n\t\t\tsendJson(res, 200, { sessions: listSessions() });\n\t\t\treturn;\n\t\t}\n\n\t\tconst runMatch = /^\\/api\\/session\\/([^/]+)\\/runs\\/([^/]+)$/.exec(url.pathname);\n\t\tif (req.method === \"GET\" && runMatch) {\n\t\t\thandleSessionRun(decodeURIComponent(runMatch[1]), decodeURIComponent(runMatch[2]), res);\n\t\t\treturn;\n\t\t}\n\n\t\tconst runAbortMatch = /^\\/api\\/session\\/([^/]+)\\/runs\\/([^/]+)\\/abort$/.exec(url.pathname);\n\t\tif (req.method === \"POST\" && runAbortMatch) {\n\t\t\tawait readBody(req);\n\t\t\tawait handleSessionRunAbort(decodeURIComponent(runAbortMatch[1]), decodeURIComponent(runAbortMatch[2]), res);\n\t\t\treturn;\n\t\t}\n\n\t\tif (req.method === \"GET\" && url.pathname.startsWith(\"/api/session/\") && url.pathname.endsWith(\"/history\")) {\n\t\t\tconst encodedSessionId = url.pathname.slice(\"/api/session/\".length, -\"/history\".length);\n\t\t\tconst fromParam = url.searchParams.get(\"from\");\n\t\t\tconst from = fromParam === null ? undefined : Number(fromParam);\n\t\t\tif (from !== undefined && (!Number.isInteger(from) || from < 0)) {\n\t\t\t\tsendError(res, 400, \"from must be a non-negative integer\", PiServerErrorCode.INVALID_REQUEST);\n\t\t\t\treturn;\n\t\t\t}\n\t\t\tconst entriesFromParam = url.searchParams.get(\"entriesFrom\");\n\t\t\tconst entriesFrom = entriesFromParam === null ? undefined : Number(entriesFromParam);\n\t\t\tif (entriesFrom !== undefined && (!Number.isInteger(entriesFrom) || entriesFrom < 0)) {\n\t\t\t\tsendError(res, 400, \"entriesFrom must be a non-negative integer\", PiServerErrorCode.INVALID_REQUEST);\n\t\t\t\treturn;\n\t\t\t}\n\t\t\tconst revisionParam = url.searchParams.get(\"revision\");\n\t\t\tconst revision = revisionParam === null ? undefined : Number(revisionParam);\n\t\t\tif (revision !== undefined && (!Number.isInteger(revision) || revision < 0)) {\n\t\t\t\tsendError(res, 400, \"revision must be a non-negative integer\", PiServerErrorCode.INVALID_REQUEST);\n\t\t\t\treturn;\n\t\t\t}\n\t\t\thandleSessionHistory(\n\t\t\t\tdecodeURIComponent(encodedSessionId),\n\t\t\t\tfrom,\n\t\t\t\tentriesFrom,\n\t\t\t\trevision,\n\t\t\t\turl.searchParams.get(\"baseTreeHash\") ?? undefined,\n\t\t\t\tres,\n\t\t\t);\n\t\t\treturn;\n\t\t}\n\n\t\tif (req.method === \"POST\" && url.pathname === CHUNK_ENDPOINT) {\n\t\t\ttry {\n\t\t\t\tconst body = JSON.parse(await readBody(req)) as RequestChunkBody;\n\t\t\t\tconst chunkResult = receiveRequestChunk(body);\n\t\t\t\tif (!chunkResult.complete) {\n\t\t\t\t\tsendJson(res, 200, chunkResult.ack);\n\t\t\t\t\treturn;\n\t\t\t\t}\n\t\t\t\tconst handled = await handlePostRequest(\n\t\t\t\t\tconfig,\n\t\t\t\t\tchunkResult.target,\n\t\t\t\t\tJSON.parse(chunkResult.bodyJson) as unknown,\n\t\t\t\t\tres,\n\t\t\t\t);\n\t\t\t\tif (!handled && !res.headersSent) {\n\t\t\t\t\tsendError(res, 404, \"Not found\", PiServerErrorCode.INVALID_REQUEST);\n\t\t\t\t}\n\t\t\t} catch (err) {\n\t\t\t\tlogRequestError(req, err);\n\t\t\t\tif (!res.headersSent) {\n\t\t\t\t\tsendError(res, 400, err instanceof Error ? err.message : String(err), PiServerErrorCode.INVALID_REQUEST);\n\t\t\t\t} else {\n\t\t\t\t\tres.write(encodeErrorEvent(err instanceof Error ? err.message : String(err)));\n\t\t\t\t\tres.end();\n\t\t\t\t}\n\t\t\t}\n\t\t\treturn;\n\t\t}\n\n\t\tif (req.method === \"POST\") {\n\t\t\ttry {\n\t\t\t\tconst body = JSON.parse(await readBody(req)) as unknown;\n\t\t\t\tif (await handlePostRequest(config, url.pathname, body, res)) return;\n\t\t\t\tif (!res.headersSent) {\n\t\t\t\t\tsendError(res, 404, \"Not found\", PiServerErrorCode.INVALID_REQUEST);\n\t\t\t\t\treturn;\n\t\t\t\t}\n\t\t\t} catch (err) {\n\t\t\t\tlogRequestError(req, err);\n\t\t\t\tif (!res.headersSent) {\n\t\t\t\t\tsendError(res, 500, err instanceof Error ? err.message : String(err), PiServerErrorCode.INTERNAL_ERROR);\n\t\t\t\t} else {\n\t\t\t\t\tres.write(encodeErrorEvent(err instanceof Error ? err.message : String(err)));\n\t\t\t\t\tres.end();\n\t\t\t\t}\n\t\t\t\treturn;\n\t\t\t}\n\t\t}\n\n\t\tif (req.method === \"DELETE\" && url.pathname.startsWith(\"/api/session/\")) {\n\t\t\tconst sessionId = decodeURIComponent(url.pathname.slice(\"/api/session/\".length));\n\t\t\tdeleteSessionFromStore(sessionId);\n\t\t\tdeleteStreamRunsForSession(sessionId);\n\t\t\tdeletePersistedSession(config.sessionStoreDir, sessionId);\n\t\t\tsendJson(res, 200, { deleted: sessionId });\n\t\t\treturn;\n\t\t}\n\n\t\tsendError(res, 404, \"Not found\", PiServerErrorCode.INVALID_REQUEST);\n\t});\n\n\treturn server;\n}\n\nfunction toProxyEvent(event: AssistantMessageEvent): ProxyAssistantMessageEvent | undefined {\n\tswitch (event.type) {\n\t\tcase \"start\":\n\t\t\treturn { type: \"start\" };\n\t\tcase \"text_start\":\n\t\t\treturn { type: \"text_start\", contentIndex: event.contentIndex };\n\t\tcase \"text_delta\":\n\t\t\treturn { type: \"text_delta\", contentIndex: event.contentIndex, delta: event.delta };\n\t\tcase \"text_end\":\n\t\t\treturn {\n\t\t\t\ttype: \"text_end\",\n\t\t\t\tcontentIndex: event.contentIndex,\n\t\t\t\tcontentSignature:\n\t\t\t\t\tevent.partial.content[event.contentIndex]?.type === \"text\"\n\t\t\t\t\t\t? (event.partial.content[event.contentIndex] as { textSignature?: string }).textSignature\n\t\t\t\t\t\t: undefined,\n\t\t\t};\n\t\tcase \"thinking_start\":\n\t\t\treturn { type: \"thinking_start\", contentIndex: event.contentIndex };\n\t\tcase \"thinking_delta\":\n\t\t\treturn { type: \"thinking_delta\", contentIndex: event.contentIndex, delta: event.delta };\n\t\tcase \"thinking_end\":\n\t\t\treturn {\n\t\t\t\ttype: \"thinking_end\",\n\t\t\t\tcontentIndex: event.contentIndex,\n\t\t\t\tcontentSignature:\n\t\t\t\t\tevent.partial.content[event.contentIndex]?.type === \"thinking\"\n\t\t\t\t\t\t? (event.partial.content[event.contentIndex] as { thinkingSignature?: string }).thinkingSignature\n\t\t\t\t\t\t: undefined,\n\t\t\t};\n\t\tcase \"toolcall_start\":\n\t\t\treturn {\n\t\t\t\ttype: \"toolcall_start\",\n\t\t\t\tcontentIndex: event.contentIndex,\n\t\t\t\tid:\n\t\t\t\t\tevent.partial.content[event.contentIndex]?.type === \"toolCall\"\n\t\t\t\t\t\t? (event.partial.content[event.contentIndex] as { id: string }).id\n\t\t\t\t\t\t: \"\",\n\t\t\t\ttoolName:\n\t\t\t\t\tevent.partial.content[event.contentIndex]?.type === \"toolCall\"\n\t\t\t\t\t\t? (event.partial.content[event.contentIndex] as { name: string }).name\n\t\t\t\t\t\t: \"\",\n\t\t\t};\n\t\tcase \"toolcall_delta\":\n\t\t\treturn { type: \"toolcall_delta\", contentIndex: event.contentIndex, delta: event.delta };\n\t\tcase \"toolcall_end\":\n\t\t\treturn { type: \"toolcall_end\", contentIndex: event.contentIndex, toolCall: event.toolCall };\n\t\tcase \"done\":\n\t\t\treturn {\n\t\t\t\ttype: \"done\",\n\t\t\t\treason: event.reason,\n\t\t\t\tusage: event.message.usage,\n\t\t\t\tdeferred: event.message.deferred,\n\t\t\t\tmessage: event.message,\n\t\t\t};\n\t\tcase \"error\":\n\t\t\treturn {\n\t\t\t\ttype: \"error\",\n\t\t\t\treason: event.reason,\n\t\t\t\terrorMessage: event.error.errorMessage,\n\t\t\t\tusage: event.error.usage,\n\t\t\t};\n\t\tdefault:\n\t\t\treturn undefined;\n\t}\n}\n\nexport function startServer(configOverride?: Partial<ServerConfig>): HttpServer {\n\tconst config = loadConfig(configOverride);\n\tconst server = createPiServer(configOverride);\n\tserver.listen(config.port, config.host, () => {\n\t\tconst address = server.address();\n\t\tconst port = typeof address === \"object\" && address !== null ? address.port : config.port;\n\t\tconst host = config.host.includes(\":\") && !config.host.startsWith(\"[\") ? `[${config.host}]` : config.host;\n\t\tconsole.log(`pi-server listening on ${host}:${port}`);\n\t});\n\treturn server;\n}\n"]}