feat(api): port Fase 5 REST API — Bun.serve router + JWT middleware + auth/session/conversation/chat handlers + composition root

This commit is contained in:
asepharyana
2026-09-02 22:50:38 +07:00
parent 1b130e69b3
commit 49e79322b4
14 changed files with 1007 additions and 0 deletions
+18
View File
@@ -0,0 +1,18 @@
{
"name": "@zesdex/api",
"version": "1.21.2",
"private": true,
"type": "module",
"scripts": {
"api": "bun src/main.ts",
"build": "true"
},
"dependencies": {
"@zesdex/domain": "workspace:*",
"@zesdex/application": "workspace:*",
"@zesdex/infrastructure": "workspace:*"
},
"devDependencies": {
"@types/bun": "^1.2.0"
}
}
+120
View File
@@ -0,0 +1,120 @@
/**
* Data Transfer Objects for the REST API.
* Wire format for request/response bodies — independent of domain entities.
* Mirrors `apps/interfaces/api/src/dto/*.rs`.
*/
/* -------------------------------------------------------------------------- */
/* Auth DTOs */
/* -------------------------------------------------------------------------- */
export interface LoginRequest {
username: string;
password: string;
}
export interface RegisterRequest {
username: string;
password: string;
/** Optional display name. */
display_name?: string;
}
export interface RefreshRequest {
refresh_token: string;
}
export interface AuthResponse {
access_token: string;
refresh_token: string;
token_type: string;
expires_in: number;
}
/* -------------------------------------------------------------------------- */
/* Session DTOs */
/* -------------------------------------------------------------------------- */
export interface CreateSessionRequest {
title: string;
}
export interface SessionResponse {
id: string;
created_at: number;
updated_at: number;
title: string;
model: string;
message_count: number;
archived: boolean;
summary?: string;
}
export interface SessionListResponse {
sessions: SessionResponse[];
total: number;
}
/* -------------------------------------------------------------------------- */
/* Conversation DTOs */
/* -------------------------------------------------------------------------- */
export interface AddMessageRequest {
role: string;
content: string;
}
export interface ToolFunctionResponse {
name: string;
arguments: string;
}
export interface ToolCallResponse {
id: string;
function: ToolFunctionResponse;
}
export interface MessageResponse {
role: string;
content?: string;
tool_calls?: ToolCallResponse[];
tool_call_id?: string;
}
export interface ConversationResponse {
session_id: string;
messages: MessageResponse[];
message_count: number;
model: string;
system_prompt: string;
max_tokens?: number;
temperature?: number;
}
export interface ChatCompletionRequest {
session_id: string;
message: string;
model?: string;
max_tokens?: number;
temperature?: number;
}
export interface ChatCompletionResponse {
reply: string;
prompt_tokens: number;
completion_tokens: number;
}
/* -------------------------------------------------------------------------- */
/* Error DTO */
/* -------------------------------------------------------------------------- */
export interface ErrorResponse {
message: string;
code: number;
}
/** Build a standardised error response body. */
export function errorBody(message: string, code: number): ErrorResponse {
return { message, code };
}
+71
View File
@@ -0,0 +1,71 @@
/**
* Typed API error with automatic HTTP response conversion.
* Mirrors `apps/interfaces/api/src/error.rs`.
*
* Each kind maps to a status code:
* - BadRequest → 400
* - Unauthorized → 401
* - NotFound → 404
* - Conflict → 409
* - TooManyRequests → 429
* - Internal → 500
* - ChatProxy → 502
*/
export type ApiErrorKind =
| "BadRequest"
| "Unauthorized"
| "NotFound"
| "Conflict"
| "TooManyRequests"
| "Internal"
| "ChatProxy";
export const HTTP_STATUS: Record<ApiErrorKind, number> = {
BadRequest: 400,
Unauthorized: 401,
NotFound: 404,
Conflict: 409,
TooManyRequests: 429,
Internal: 500,
ChatProxy: 502,
};
export class ApiError extends Error {
constructor(public kind: ApiErrorKind, message: string) {
super(message);
this.name = "ApiError";
}
get status(): number {
return HTTP_STATUS[this.kind];
}
static badRequest(msg: string): ApiError {
return new ApiError("BadRequest", msg);
}
static unauthorized(msg: string): ApiError {
return new ApiError("Unauthorized", msg);
}
static notFound(msg: string): ApiError {
return new ApiError("NotFound", msg);
}
static conflict(msg: string): ApiError {
return new ApiError("Conflict", msg);
}
static tooManyRequests(msg: string): ApiError {
return new ApiError("TooManyRequests", msg);
}
static internal(msg: string): ApiError {
return new ApiError("Internal", msg);
}
static chatProxy(msg: string): ApiError {
return new ApiError("ChatProxy", msg);
}
/** Whether an internal error's details should be hidden from the client. */
get exposeDetails(): boolean {
// Internal errors hide details; all others surface the message.
return this.kind !== "Internal";
}
}
+133
View File
@@ -0,0 +1,133 @@
/**
* Authentication handlers — login, register, and token refresh.
* Mirrors `apps/interfaces/api/src/handlers/auth.rs`.
*
* Endpoints:
* - POST /auth/login — authenticate, return JWT pair
* - POST /auth/register — create account, return JWT pair
* - POST /auth/refresh — exchange refresh token for a new pair
*
* All three are rate-limited (20 attempts / 10 min per client).
*/
import type { ApiState, RateLimiter } from "../state.ts";
import { ApiError } from "../error.ts";
import { loadUsers, saveUsers } from "../state.ts";
import type { AuthResponse, LoginRequest, RegisterRequest, RefreshRequest } from "../dto.ts";
/** Login/register brute-force protection: 20 attempts per 10-minute window. */
const AUTH_RATE_LIMIT_MAX = 20;
const AUTH_RATE_LIMIT_WINDOW_SECS = 600;
/** Extract a coarse client identity (X-Forwarded-For first hop or "unknown"). */
function clientId(headers: Headers): string {
const xff = headers.get("x-forwarded-for");
if (xff) {
const first = xff.split(",")[0]?.trim();
if (first) return first;
}
return "unknown";
}
/** Enforce the auth rate limit; throws TooManyRequests when exceeded. */
function enforceRateLimit(limiter: RateLimiter, headers: Headers): void {
if (!limiter.check(clientId(headers), AUTH_RATE_LIMIT_MAX, AUTH_RATE_LIMIT_WINDOW_SECS)) {
throw ApiError.tooManyRequests("Too many requests, try again later");
}
}
function requireFields(req: LoginRequest & RegisterRequest, reg: boolean): void {
if (req.username === "") {
throw ApiError.badRequest("Username is required");
}
if (req.password === "") {
throw ApiError.badRequest("Password is required");
}
if (reg && req.password.length < 6) {
throw ApiError.badRequest("Password must be at least 6 characters");
}
}
/** Issue an AuthResponse for the given subject. */
function authResponse(state: ApiState, sub: string): AuthResponse {
const [access, refresh] = state.token_service.generateTokens(sub);
return {
access_token: access,
refresh_token: refresh,
token_type: "Bearer",
expires_in: 3600, // access token expiry in seconds
};
}
/** POST /auth/login — authenticate and issue a JWT pair. */
export async function loginHandler(
state: ApiState,
headers: Headers,
body: LoginRequest,
): Promise<AuthResponse> {
enforceRateLimit(state.auth_rate_limiter, headers);
requireFields(body, false);
const users = loadUsers(state.store_base_dir);
if (!users) {
throw ApiError.unauthorized("Invalid username or password");
}
const storedHash = users[body.username];
if (!storedHash) {
throw ApiError.unauthorized("Invalid username or password");
}
const valid = await state.password_service.verify(body.password, storedHash);
if (!valid) {
throw ApiError.unauthorized("Invalid username or password");
}
return authResponse(state, body.username);
}
/** POST /auth/register — create a new user account and issue a JWT pair. */
export async function registerHandler(
state: ApiState,
headers: Headers,
body: RegisterRequest,
): Promise<AuthResponse> {
enforceRateLimit(state.auth_rate_limiter, headers);
requireFields(body, true);
// Hash first (async), then serialize the read-modify-write of users.json.
const hash = await state.password_service.hash(body.password);
const unlock = state.users_lock.lock();
try {
const users = loadUsers(state.store_base_dir) ?? {};
if (body.username in users) {
throw ApiError.conflict("Username already exists");
}
users[body.username] = hash;
saveUsers(state.store_base_dir, users);
} finally {
unlock();
}
return authResponse(state, body.username);
}
/** POST /auth/refresh — exchange a refresh token for a fresh pair. */
export async function refreshHandler(
state: ApiState,
headers: Headers,
body: RefreshRequest,
): Promise<AuthResponse> {
enforceRateLimit(state.auth_rate_limiter, headers);
if (body.refresh_token === "") {
throw ApiError.badRequest("Refresh token is required");
}
let sub: string;
try {
sub = state.token_service.verifyRefreshToken(body.refresh_token);
} catch {
throw ApiError.unauthorized("Invalid or expired refresh token");
}
return authResponse(state, sub);
}
+66
View File
@@ -0,0 +1,66 @@
/**
* LLM chat completion proxy handler.
* Mirrors `apps/interfaces/api/src/handlers/chat.rs`.
*
* Endpoints:
* - POST /chat/completions — proxy a completion to the LLM, persisting history.
*/
import type { ApiState } from "../state.ts";
import { ApiError } from "../error.ts";
import { newConversation, Roles, type ChatMessage } from "@zesdex/domain";
import type { ChatCompletionResponse } from "../dto.ts";
/** POST /chat/completions — proxy to LLM provider and persist the conversation. */
export async function chatCompletionsHandler(
state: ApiState,
req: {
session_id: string;
message: string;
model?: string;
max_tokens?: number;
temperature?: number;
},
): Promise<ChatCompletionResponse> {
if (req.session_id === "") throw ApiError.badRequest("session_id is required");
if (req.message === "") throw ApiError.badRequest("message is required");
// Load or create the conversation for this session.
let conversation;
try {
conversation = await state.conversation_service.loadConversation(req.session_id);
} catch {
conversation = newConversation("", req.session_id);
conversation.model = req.model ?? state.llm_client.model;
if (req.max_tokens !== undefined) conversation.max_tokens = req.max_tokens;
if (req.temperature !== undefined) conversation.temperature = req.temperature;
}
if (req.model !== undefined) conversation.model = req.model;
const userMsg: ChatMessage = { role: Roles.User, content: req.message };
conversation.messages.push(userMsg);
// Call the LLM (non-streaming), using the conversation's resolved model.
let response: ChatMessage;
let usage: [number, number] | null = null;
try {
const result = await state.llm_client.chat(
conversation.messages,
undefined,
conversation.max_tokens,
conversation.temperature,
);
response = result.message;
usage = result.usage;
} catch (e) {
throw ApiError.chatProxy(`LLM request failed: ${(e as Error).message}`);
}
const [promptTokens, completionTokens] = usage ?? [0, 0];
const replyText = response.content ?? "";
const assistantMsg: ChatMessage = { role: Roles.Assistant, content: replyText };
await state.conversation_service.addMessage(conversation, assistantMsg);
return { reply: replyText, prompt_tokens: promptTokens, completion_tokens: completionTokens };
}
@@ -0,0 +1,110 @@
/**
* Conversation message-history handlers.
* Mirrors `apps/interfaces/api/src/handlers/conversations.rs`.
*
* Endpoints:
* - GET /sessions/:id/conversations — fetch conversation
* - POST /sessions/:id/conversations — append a message
* - DELETE /sessions/:id/conversations/:cid — delete a message by index
*/
import type { ApiState } from "../state.ts";
import { ApiError } from "../error.ts";
import { Roles, type ChatMessage } from "@zesdex/domain";
import type {
AddMessageRequest,
ConversationResponse,
MessageResponse,
ToolCallResponse,
} from "../dto.ts";
import type { Conversation } from "@zesdex/domain";
/** Map a single ChatMessage to its wire DTO. */
export function toMessageResponse(m: ChatMessage): MessageResponse {
let tool_calls: ToolCallResponse[] | undefined;
if (m.tool_calls && m.tool_calls.length > 0) {
tool_calls = m.tool_calls.map((tc) => ({
id: tc.id,
function: {
name: tc.function.name,
arguments: JSON.stringify(tc.function.arguments),
},
}));
}
const out: MessageResponse = { role: m.role };
if (m.content !== null && m.content !== undefined) out.content = m.content;
if (tool_calls) out.tool_calls = tool_calls;
if (m.tool_call_id) out.tool_call_id = m.tool_call_id;
return out;
}
/** Map a domain Conversation to its wire DTO. */
export function toConversationResponse(c: Conversation): ConversationResponse {
return {
session_id: c.session_id,
messages: c.messages.map(toMessageResponse),
message_count: c.messages.length,
model: c.model,
system_prompt: c.system_prompt,
...(c.max_tokens !== undefined ? { max_tokens: c.max_tokens } : {}),
...(c.temperature !== undefined ? { temperature: c.temperature } : {}),
};
}
/** GET /sessions/:id/conversations — fetch the conversation for a session. */
export async function getConversationHandler(
state: ApiState,
id: string,
): Promise<ConversationResponse> {
if (id === "") throw ApiError.badRequest("Session ID is required");
const conversation = await state.conversation_service.loadConversation(id);
return toConversationResponse(conversation);
}
/** POST /sessions/:id/conversations — append a message to a conversation. */
export async function addMessageHandler(
state: ApiState,
id: string,
req: AddMessageRequest,
): Promise<{ status: number; body: ConversationResponse }> {
if (id === "") throw ApiError.badRequest("Session ID is required");
if (req.content === "") throw ApiError.badRequest("Message content is required");
const role = req.role.toLowerCase();
if (role !== "user" && role !== "assistant") {
throw ApiError.badRequest(`Invalid role: ${req.role}`);
}
const msg: ChatMessage = {
role: role === "user" ? Roles.User : Roles.Assistant,
content: req.content,
};
const conversation = await state.conversation_service.loadConversation(id);
await state.conversation_service.addMessage(conversation, msg);
return { status: 200, body: toConversationResponse(conversation) };
}
/** DELETE /sessions/:id/conversations/:cid — delete a message by index. */
export async function deleteMessageHandler(
state: ApiState,
id: string,
cid: string,
): Promise<{ status: number }> {
if (id === "") throw ApiError.badRequest("Session ID is required");
const index = Number.parseInt(cid, 10);
if (Number.isNaN(index)) {
throw ApiError.badRequest(`Invalid message index: ${cid}`);
}
const conversation = await state.conversation_service.loadConversation(id);
if (index >= conversation.messages.length) {
throw ApiError.notFound(
`Message index ${index} out of bounds (max: ${conversation.messages.length - 1})`,
);
}
conversation.messages.splice(index, 1);
await state.conversation_service.saveConversation(conversation);
return { status: 204 };
}
@@ -0,0 +1,73 @@
/**
* Session management handlers.
* Mirrors `apps/interfaces/api/src/handlers/sessions.rs`.
*
* Endpoints:
* - GET /sessions — list all sessions
* - POST /sessions — create a new session
* - DELETE /sessions/:id — archive/close a session
*/
import type { ApiState } from "../state.ts";
import { ApiError } from "../error.ts";
import type { CreateSessionRequest, SessionListResponse, SessionResponse } from "../dto.ts";
import { newSessionId } from "@zesdex/domain";
/** Map a domain Session to its wire DTO. */
export function toSessionResponse(s: {
id: string;
created_at: number;
updated_at: number;
title: string;
model: string;
message_count: number;
archived: boolean;
summary?: string;
}): SessionResponse {
return {
id: s.id,
created_at: s.created_at,
updated_at: s.updated_at,
title: s.title,
model: s.model,
message_count: s.message_count,
archived: s.archived,
...(s.summary !== undefined ? { summary: s.summary } : {}),
};
}
/** GET /sessions — list all (non-archived) sessions. */
export async function listSessionsHandler(
state: ApiState,
): Promise<SessionListResponse> {
const sessions = await state.session_service.listAll();
const sessionResponses = sessions.map(toSessionResponse);
return { sessions: sessionResponses, total: sessionResponses.length };
}
/** POST /sessions — create a new session. */
export async function createSessionHandler(
state: ApiState,
req: CreateSessionRequest,
): Promise<{ status: number; body: SessionResponse }> {
if (req.title.trim() === "") {
throw ApiError.badRequest("Session title is required");
}
const session = await state.session_service.createSession(req.title);
return { status: 201, body: toSessionResponse(session) };
}
/** DELETE /sessions/:id — archive/close a session (returns 204). */
export async function deleteSessionHandler(
state: ApiState,
id: string,
): Promise<{ status: number }> {
if (id === "") {
throw ApiError.badRequest("Session ID is required");
}
const parsed = newSessionId(id);
if (!parsed.ok) {
throw ApiError.badRequest(`Invalid session ID: ${parsed.error}`);
}
await state.session_service.archiveSession(parsed.value);
return { status: 204 };
}
+10
View File
@@ -0,0 +1,10 @@
/**
* Zesdex REST API interface package.
* Mirrors `apps/interfaces/api/`.
*/
export * from "./dto.ts";
export { ApiError, HTTP_STATUS } from "./error.ts";
export { RateLimiter, newApiState } from "./state.ts";
export type { ApiState } from "./state.ts";
export { authenticateRequest } from "./middleware/auth.ts";
export { startApiServer } from "./server.ts";
+21
View File
@@ -0,0 +1,21 @@
#!/usr/bin/env bun
/**
* Zesdex API server — standalone entry point.
* Mirrors the `--api` mode of the Rust gateway.
*/
import { newApiState, startApiServer } from "./index.ts";
function envOr(key: string, fallback: string): string {
const v = process.env[key];
return v && v !== "" ? v : fallback;
}
const port = Number.parseInt(envOr("ZESDEX_API_PORT", "8080"), 10);
const baseDir = envOr("ZESDEX_STORE", process.env.HOME ? `${process.env.HOME}/.local/share/zesdex` : ".");
const jwtSecret = envOr("ZESDEX_JWT_SECRET", "dev-secret-change-me");
const apiKey = envOr("OPENAI_API_KEY", "");
const model = envOr("ZESDEX_MODEL", "claude-opus-5");
const baseUrl = process.env.OPENAI_API_BASE ?? undefined;
const state = newApiState(baseDir, jwtSecret, apiKey, model, baseUrl);
startApiServer(state, port);
@@ -0,0 +1,51 @@
/**
* JWT authentication middleware for the API.
* Mirrors `apps/interfaces/api/src/middleware/auth.rs`.
*
* Validates the `Authorization: Bearer <token>` header and injects the
* validated subject into `req.user`. Refresh tokens are rejected on
* protected routes — only access tokens are accepted.
*/
import { verifyToken } from "@zesdex/infrastructure";
import type { ApiState } from "../state.ts";
/** Claims decoded from a valid JWT, attached to the request. */
export interface JwtClaims {
sub: string;
}
/** Result of authentication: either the subject or an error status/detail. */
export type AuthResult = { ok: true; sub: string } | { ok: false; status: number; detail: string };
/**
* Authenticate a request against the API state's JWT secret.
* Reads the `Authorization: Bearer <token>` header, verifies the signature
* and token type (must be `access`), and returns the subject on success.
*/
export function authenticateRequest(state: ApiState, headers: Headers): AuthResult {
const authHeader = headers.get("Authorization");
if (!authHeader) {
return { ok: false, status: 401, detail: "Missing or invalid Authorization header" };
}
const prefix = "Bearer ";
if (!authHeader.startsWith(prefix)) {
return { ok: false, status: 401, detail: "Missing or invalid Authorization header" };
}
const token = authHeader.slice(prefix.length).trim();
if (!token) {
return { ok: false, status: 401, detail: "Missing or invalid Authorization header" };
}
try {
const claims = verifyToken(state.jwt_secret, token);
if (claims.typ !== "access") {
return {
ok: false,
status: 401,
detail: "refresh tokens are not accepted on protected routes",
};
}
return { ok: true, sub: claims.sub };
} catch (e) {
return { ok: false, status: 401, detail: `Invalid token: ${(e as Error).message}` };
}
}
+183
View File
@@ -0,0 +1,183 @@
/**
* Zesdex REST API — HTTP server and router (Bun.serve).
* Mirrors `apps/interfaces/api/src/lib.rs`.
*
* Routes:
* - Public: /api/v1/auth/{login,register,refresh}, /api/v1/health
* - Protected (JWT): /api/v1/sessions..., /api/v1/chat/completions
*
* Uses Bun's native HTTP server — no external framework dependency.
*/
import type { ApiState } from "./state.ts";
import { ApiError } from "./error.ts";
import { errorBody } from "./dto.ts";
import { authenticateRequest } from "./middleware/auth.ts";
import { loginHandler, registerHandler, refreshHandler } from "./handlers/auth.ts";
import {
listSessionsHandler,
createSessionHandler,
deleteSessionHandler,
} from "./handlers/sessions.ts";
import {
getConversationHandler,
addMessageHandler,
deleteMessageHandler,
} from "./handlers/conversations.ts";
import { chatCompletionsHandler } from "./handlers/chat.ts";
/** CORS headers applied to every response (permissive for local/dev). */
const CORS_HEADERS: Record<string, string> = {
"Access-Control-Allow-Origin": "*",
"Access-Control-Allow-Methods": "GET, POST, DELETE, OPTIONS",
"Access-Control-Allow-Headers": "Content-Type, Authorization",
};
/** A handler's response: JSON body plus optional status code. */
interface Result {
status: number;
body: unknown;
}
/** Per-request context passed to handlers. */
interface Ctx {
state: ApiState;
params: string[];
body: unknown;
headers: Headers;
}
/** Route mapping: regex matches URL path after `/api/v1`. */
interface Route {
method: string;
pattern: RegExp;
protected: boolean;
handle: (ctx: Ctx) => Promise<Result>;
}
function result(body: unknown, status = 200): Result {
return { status, body };
}
function json(res: Result): Response {
return new Response(JSON.stringify(res.body), {
status: res.status,
headers: { "Content-Type": "application/json", ...CORS_HEADERS },
});
}
function noContent(): Response {
return new Response(null, { status: 204, headers: CORS_HEADERS });
}
/** Build the route table for the API. Matches against paths without the /api/v1 prefix. */
function buildRoutes(): Route[] {
return [
// Auth — public (rate-limited inside handlers)
{ method: "POST", pattern: /^\/auth\/login$/, protected: false, handle: (c) => loginHandler(c.state, c.headers, c.body as never).then((r) => result(r)) },
{ method: "POST", pattern: /^\/auth\/register$/, protected: false, handle: (c) => registerHandler(c.state, c.headers, c.body as never).then((r) => result(r)) },
{ method: "POST", pattern: /^\/auth\/refresh$/, protected: false, handle: (c) => refreshHandler(c.state, c.headers, c.body as never).then((r) => result(r)) },
// Health — public
{ method: "GET", pattern: /^\/health$/, protected: false, handle: () => Promise.resolve(result({ status: "ok" })) },
// Sessions — protected
{ method: "GET", pattern: /^\/sessions$/, protected: true, handle: (c) => listSessionsHandler(c.state).then((s) => result(s)) },
{ method: "POST", pattern: /^\/sessions$/, protected: true, handle: async (c) => {
const r = await createSessionHandler(c.state, c.body as { title: string });
return result(r.body, r.status);
} },
{ method: "DELETE", pattern: /^\/sessions\/([^/]+)$/, protected: true, handle: async (c) => {
await deleteSessionHandler(c.state, c.params[0]!);
return { status: 204, body: null };
} },
// Conversations (sub-resource of sessions) — protected
{ method: "GET", pattern: /^\/sessions\/([^/]+)\/conversations$/, protected: true, handle: async (c) => {
const conv = await getConversationHandler(c.state, c.params[0]!);
return result(conv);
} },
{ method: "POST", pattern: /^\/sessions\/([^/]+)\/conversations$/, protected: true, handle: async (c) => {
const r = await addMessageHandler(c.state, c.params[0]!, c.body as { role: string; content: string });
return result(r.body, r.status);
} },
{ method: "DELETE", pattern: /^\/sessions\/([^/]+)\/conversations\/([^/]+)$/, protected: true, handle: async (c) => {
await deleteMessageHandler(c.state, c.params[0]!, c.params[1]!);
return { status: 204, body: null };
} },
// Chat — protected
{ method: "POST", pattern: /^\/chat\/completions$/, protected: true, handle: (c) => chatCompletionsHandler(c.state, c.body as never).then((s) => result(s)) },
];
}
/** Handle a request against the router; maps ApiError → HTTP response. */
async function handleRequest(state: ApiState, routes: Route[], req: Request): Promise<Response> {
const url = new URL(req.url);
// Strip the /api/v1 version prefix so route patterns match cleanly.
const fullPath = url.pathname;
const path = fullPath.startsWith("/api/v1") ? fullPath.slice("/api/v1".length) : fullPath;
const method = req.method;
// CORS preflight
if (method === "OPTIONS") {
return new Response(null, { status: 204, headers: CORS_HEADERS });
}
const route = routes.find((r) => r.method === method && r.pattern.test(path));
if (!route) {
return json(result(errorBody("Not found", 404), 404));
}
// JWT auth for protected routes (subject is validated but not currently exposed).
if (route.protected) {
const auth = authenticateRequest(state, req.headers);
if (!auth.ok) {
return json(result({ error: "Unauthorized", detail: auth.detail }, auth.status));
}
}
// Parse JSON body if present; read once and reuse.
let body: unknown = undefined;
const raw = await req.text();
if (raw) {
try {
body = JSON.parse(raw);
} catch {
return json(result(errorBody("Bad request: malformed JSON body", 400), 400));
}
}
const match = path.match(route.pattern);
const ctx: Ctx = {
state,
params: match ? match.slice(1) : [],
body,
headers: req.headers,
};
try {
const res = await route.handle(ctx);
if (res.status === 204) return noContent();
return json(res);
} catch (e) {
if (e instanceof ApiError) {
const msg = e.exposeDetails ? e.message : "An internal error occurred";
return json(result(errorBody(msg, e.status), e.status));
}
console.error("Unhandled error:", e);
return json(result(errorBody("An internal error occurred", 500)));
}
}
/** Start the API server on the given port. Returns a handle. */
export function startApiServer(state: ApiState, port: number): { stop: () => void; port: number } {
const routes = buildRoutes();
const server = Bun.serve({
port,
async fetch(req) {
return handleRequest(state, routes, req);
},
});
console.log(`zesdex-api listening on http://0.0.0.0:${server.port}`);
return { stop: () => server.stop(true), port: server.port ?? port };
}
+134
View File
@@ -0,0 +1,134 @@
/**
* Shared application state for the REST API server.
* Mirrors `apps/interfaces/api/src/state.rs`.
*
* `ApiState` holds all service implementations wired to concrete
* infrastructure adapters. Constructed once at startup (composition root)
* and shared across all requests.
*/
import * as fs from "node:fs";
import * as path from "node:path";
import {
SessionServiceImpl,
ConversationServiceImpl,
SettingsServiceImpl,
MemoryServiceImpl,
} from "@zesdex/application";
import {
JsonSettingsRepository,
JsonAppConfigRepository,
JsonConversationRepository,
MarkdownMemoryRepository,
FileSystemSessionRepository,
FileSystemSessionLockRepository,
Argon2PasswordService,
Hs256TokenService,
LlmClient,
} from "@zesdex/infrastructure";
/**
* Sliding-window rate limiter for auth endpoints.
* Mirrors `infrastructure::middleware::rate_limit::RateLimiter`.
*/
export class RateLimiter {
private hits = new Map<string, number[]>();
/** Check whether `key` is within `max` requests per `windowSecs`. Returns true if allowed. */
check(key: string, max: number, windowSecs: number): boolean {
const now = Date.now();
const cutoff = now - windowSecs * 1000;
const list = (this.hits.get(key) ?? []).filter((t) => t > cutoff);
if (list.length >= max) {
this.hits.set(key, list);
return false;
}
list.push(now);
this.hits.set(key, list);
return true;
}
}
/** Concrete session repository wired from infrastructure. */
const sessionRepo = new FileSystemSessionRepository();
const sessionLockRepo = new FileSystemSessionLockRepository();
const conversationRepo = new JsonConversationRepository();
const settingsRepo = new JsonSettingsRepository();
const appConfigRepo = new JsonAppConfigRepository();
const memoryRepo = new MarkdownMemoryRepository();
/** Rooted, long-lived shared API state. */
export interface ApiState {
store_base_dir: string;
jwt_secret: string;
session_service: SessionServiceImpl;
conversation_service: ConversationServiceImpl;
settings_service: SettingsServiceImpl;
memory_service: MemoryServiceImpl;
password_service: Argon2PasswordService;
token_service: Hs256TokenService;
auth_rate_limiter: RateLimiter;
/** Serializes read-modify-write of `users.json` (TOCTOU race guard). */
users_lock: { lock: () => () => void };
llm_client: LlmClient;
}
/**
* A no-op lock that returns an unlock function.
* Bun is single-threaded per process, so the JS event loop already
* serialises synchronous read-modify-write of users.json — the lock is
* retained for structural parity with the Rust mutex guard.
*/
function noopLock(): { lock: () => () => void } {
return { lock: () => () => {} };
}
/** Construct a new API state with all services wired to their defaults. */
export function newApiState(
baseDir: string,
jwtSecret: string,
llmApiKey: string,
llmModel: string,
llmBaseUrl?: string,
): ApiState {
const sessionsDir = path.join(baseDir, "sessions");
const memoryDir = path.join(baseDir, "memories");
const session_service = new SessionServiceImpl(sessionRepo, sessionLockRepo, baseDir);
const conversation_service = new ConversationServiceImpl(conversationRepo, sessionsDir);
const settings_service = new SettingsServiceImpl(settingsRepo, appConfigRepo, baseDir);
const memory_service = new MemoryServiceImpl(memoryRepo, memoryDir);
const token_service = new Hs256TokenService(jwtSecret);
const llm = new LlmClient(llmApiKey, llmModel, llmBaseUrl);
return {
store_base_dir: baseDir,
jwt_secret: jwtSecret,
session_service,
conversation_service,
settings_service,
memory_service,
password_service: new Argon2PasswordService(),
token_service,
auth_rate_limiter: new RateLimiter(),
users_lock: noopLock(),
llm_client: llm,
};
}
/** Load the `users.json` map (username → password hash), or null if absent. */
export function loadUsers(baseDir: string): Record<string, string> | null {
const p = path.join(baseDir, "users.json");
try {
return JSON.parse(fs.readFileSync(p, "utf8")) as Record<string, string>;
} catch {
return null;
}
}
/** Persist the users map to `users.json` (pretty-printed). */
export function saveUsers(baseDir: string, users: Record<string, string>): void {
const p = path.join(baseDir, "users.json");
fs.mkdirSync(path.dirname(p), { recursive: true });
fs.writeFileSync(p, JSON.stringify(users, null, 2));
}
@@ -6,6 +6,9 @@ export * from "./utils.ts";
export * from "./llm/index.ts";
export * from "./persistence/index.ts";
export * from "./tools/mod.ts";
export * from "./auth/index.ts";
export * from "./ipc/index.ts";
export * from "./bgbash/index.ts";
export { InfrastructureToolExecutor } from "./tools/executor.ts";
export { toolDefs, allTools, toolIsRisky, toolIsParallelSafe } from "./tools/registry.ts";
export { BestPracticeEngine } from "./best_practice/engine.ts";