From 7e6ed34f0d96794fe3821aeea07b318a59076e98 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Thu, 3 Sep 2026 01:19:27 +0700 Subject: [PATCH] =?UTF-8?q?refactor:=20rombak=20jadi=20full=20CLI/TUI=20on?= =?UTF-8?q?ly=20=E2=80=94=20hapus=20semua=20interface=20server=20(api/ws/g?= =?UTF-8?q?rpc/web/daemon)=20+=20wire=20TUI=20ke=20agent=20beneran?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Hapus total interfaces: api (REST), ws (WebSocket), grpc, web, daemon (IPC) - Hapus infra server-only yang jadi dead code: infrastructure/ipc + seluruh stack auth (domain/auth, application/auth, persistence/iam, infrastructure/auth jwt/password/oauth_loopback) - TUI (OpenTUI) jadi interface utama: wire ke runSingleProcess() real runtime, ganti mock '(no LLM configured)' dengan turnService.runTurn asli + streaming stream_token/reasoning/tool_result/system_note/error ke transcript - Tampilkan model asli di header (mis. claude-opus-5), abort turn via Ctrl+C - CLI jadi mode headless tambahan: zesdex --headless '' one-shot turn - Bersihkan package.json (script tui/headless), tsconfig paths, docker-compose, AGENTS.md; hapus StartOAuth/login stub - verified: tsc 0 err, 70 test pass, build dist/zesdex; headless 'pong' real; TUI interactive render + real turn 'say pong in one word' -> assistant 'pong' --- AGENTS.md | 61 +++-- docker-compose.yml | 6 +- package.json | 12 +- src/application/auth/index.ts | 3 - src/application/auth/oauth_service.test.ts | 141 ---------- src/application/auth/oauth_service.ts | 121 --------- src/application/auth/session_service.ts | 44 ---- src/application/index.ts | 1 - src/domain/auth/commands.ts | 11 - src/domain/auth/error.ts | 47 ---- src/domain/auth/index.ts | 20 -- src/domain/auth/oauth.ts | 40 --- src/domain/auth/repository.ts | 39 --- src/domain/auth/service.ts | 35 --- src/domain/auth/session.ts | 56 ---- src/domain/auth/session_id.test.ts | 25 -- src/domain/auth/session_id.ts | 41 --- src/domain/auth/session_lock.ts | 120 --------- src/domain/index.ts | 1 - src/infrastructure/auth/index.ts | 8 - src/infrastructure/auth/jwt.ts | 134 ---------- src/infrastructure/auth/oauth_loopback.ts | 151 ----------- src/infrastructure/auth/password.ts | 27 -- src/infrastructure/index.ts | 2 - src/infrastructure/ipc/client.ts | 43 --- src/infrastructure/ipc/conn.ts | 90 ------- src/infrastructure/ipc/frame.ts | 101 ------- src/infrastructure/ipc/index.ts | 17 -- src/infrastructure/ipc/protocol.ts | 75 ------ src/infrastructure/ipc/server.ts | 60 ----- src/infrastructure/persistence/iam/index.ts | 4 - .../persistence/iam/oauth_repo.ts | 21 -- .../persistence/iam/session_lock_repo.ts | 78 ------ .../persistence/iam/session_repo.ts | 61 ----- src/infrastructure/persistence/index.ts | 5 +- src/interfaces/api/dto.ts | 120 --------- src/interfaces/api/error.ts | 71 ----- src/interfaces/api/handlers/auth.ts | 133 ---------- src/interfaces/api/handlers/chat.ts | 66 ----- src/interfaces/api/handlers/conversations.ts | 110 -------- src/interfaces/api/handlers/sessions.ts | 73 ----- src/interfaces/api/index.ts | 10 - src/interfaces/api/main.ts | 21 -- src/interfaces/api/middleware/auth.ts | 51 ---- src/interfaces/api/server.ts | 183 ------------- src/interfaces/api/state.ts | 134 ---------- src/interfaces/cli/args.ts | 126 ++------- src/interfaces/cli/index.ts | 144 +++------- src/interfaces/daemon/client.ts | 61 ----- src/interfaces/daemon/index.ts | 7 - src/interfaces/daemon/main.ts | 20 -- src/interfaces/daemon/server.ts | 249 ------------------ src/interfaces/grpc/index.ts | 6 - src/interfaces/grpc/main.ts | 9 - src/interfaces/grpc/server.ts | 41 --- src/interfaces/tui/action.ts | 6 +- src/interfaces/tui/command.ts | 5 - src/interfaces/tui/main.ts | 8 - src/interfaces/tui/run.ts | 63 ++++- src/interfaces/tui/state.ts | 22 +- src/interfaces/tui/ui.tsx | 94 +++++-- src/interfaces/web/index.ts | 6 - src/interfaces/web/main.ts | 10 - src/interfaces/web/server.ts | 121 --------- src/interfaces/ws/index.ts | 6 - src/interfaces/ws/main.ts | 9 - src/interfaces/ws/server.ts | 220 ---------------- tsconfig.json | 5 - 68 files changed, 250 insertions(+), 3661 deletions(-) delete mode 100644 src/application/auth/index.ts delete mode 100644 src/application/auth/oauth_service.test.ts delete mode 100644 src/application/auth/oauth_service.ts delete mode 100644 src/application/auth/session_service.ts delete mode 100644 src/domain/auth/commands.ts delete mode 100644 src/domain/auth/error.ts delete mode 100644 src/domain/auth/index.ts delete mode 100644 src/domain/auth/oauth.ts delete mode 100644 src/domain/auth/repository.ts delete mode 100644 src/domain/auth/service.ts delete mode 100644 src/domain/auth/session.ts delete mode 100644 src/domain/auth/session_id.test.ts delete mode 100644 src/domain/auth/session_id.ts delete mode 100644 src/domain/auth/session_lock.ts delete mode 100644 src/infrastructure/auth/index.ts delete mode 100644 src/infrastructure/auth/jwt.ts delete mode 100644 src/infrastructure/auth/oauth_loopback.ts delete mode 100644 src/infrastructure/auth/password.ts delete mode 100644 src/infrastructure/ipc/client.ts delete mode 100644 src/infrastructure/ipc/conn.ts delete mode 100644 src/infrastructure/ipc/frame.ts delete mode 100644 src/infrastructure/ipc/index.ts delete mode 100644 src/infrastructure/ipc/protocol.ts delete mode 100644 src/infrastructure/ipc/server.ts delete mode 100644 src/infrastructure/persistence/iam/index.ts delete mode 100644 src/infrastructure/persistence/iam/oauth_repo.ts delete mode 100644 src/infrastructure/persistence/iam/session_lock_repo.ts delete mode 100644 src/infrastructure/persistence/iam/session_repo.ts delete mode 100644 src/interfaces/api/dto.ts delete mode 100644 src/interfaces/api/error.ts delete mode 100644 src/interfaces/api/handlers/auth.ts delete mode 100644 src/interfaces/api/handlers/chat.ts delete mode 100644 src/interfaces/api/handlers/conversations.ts delete mode 100644 src/interfaces/api/handlers/sessions.ts delete mode 100644 src/interfaces/api/index.ts delete mode 100644 src/interfaces/api/main.ts delete mode 100644 src/interfaces/api/middleware/auth.ts delete mode 100644 src/interfaces/api/server.ts delete mode 100644 src/interfaces/api/state.ts delete mode 100644 src/interfaces/daemon/client.ts delete mode 100644 src/interfaces/daemon/index.ts delete mode 100644 src/interfaces/daemon/main.ts delete mode 100644 src/interfaces/daemon/server.ts delete mode 100644 src/interfaces/grpc/index.ts delete mode 100644 src/interfaces/grpc/main.ts delete mode 100644 src/interfaces/grpc/server.ts delete mode 100644 src/interfaces/tui/main.ts delete mode 100644 src/interfaces/web/index.ts delete mode 100644 src/interfaces/web/main.ts delete mode 100644 src/interfaces/web/server.ts delete mode 100644 src/interfaces/ws/index.ts delete mode 100644 src/interfaces/ws/main.ts delete mode 100644 src/interfaces/ws/server.ts diff --git a/AGENTS.md b/AGENTS.md index a410561..1dfec32 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,38 +1,55 @@ # AGENTS.md -This file provides guidance to Kilo when working with code in the zesdex repository. +This file provides guidance to AI agents (Kilo / Hermes) working in the zesdex repository. + +## Overview + +Zesdex is an **AI-native autonomous coding agent**, implemented in **TypeScript / Bun**. +It was rebuilt to be a **pure CLI/TUI tool** — there are **no network server interfaces** +(REST API / WebSocket / gRPC / web / daemon were all removed). The only two entry points: + +- `zesdex` — interactive **OpenTUI** (primary interface, wired to the real agent). +- `zesdex --headless ""` — one-shot non-interactive turn (headless mode). ## Best Practice Conventions -Zesdex follows Kana Engineering Best Practices: +1. **Clean Architecture** — strict `domain` / `application` / `infrastructure` layering + under `src/`, with thin `interfaces/` for the TUI + CLI. + - `domain` has zero I/O and zero framework dependencies. + - `application` depends only on `domain` (ports + turn service). + - `infrastructure` implements domain ports (LLM client, tool executor, file repos). + - DTOs/value objects cross layer boundaries, not entities. -1. **Clean Architecture** — Strict domain/application/infrastructure/presentation layering. - - Domain has ZERO framework dependencies. - - Application depends only on domain. - - Infrastructure implements domain traits. - - DTOs cross layer boundaries, NOT entities. +2. **Clean Code** — functions under ~40 lines, one level of abstraction per function, + descriptive names, no flag arguments, no commented-out dead code. -2. **Clean Code** — Functions under ~40 lines, one level of abstraction per function, descriptive names, no flag arguments, no commented-out code. +3. **Documentation** — every exported `function`, `class`, `interface`, and `type` + needs a `/** */` doc comment explaining what, flow, why, and return value. -3. **Documentation** — Every pub fn, struct, enum, and trait needs a doc comment (///) explaining what, flow, why, and return value. +4. **Commit Convention** — Conventional Commits in Bahasa Indonesia: + `feat(scope):`, `fix(scope):`, `chore:`, `docs:`. -4. **Commit Convention** — Conventional Commits in Bahasa Indonesia: `feat(scope):`, `fix(scope):`, `chore:`, `docs:`. +5. **Error Handling** — throw `Error` with clear messages; the turn loop wraps LLM + errors and recovers. Log with `console.*` (redirected to `~/.local/share/zesdex/zesdex.log`). -5. **Error Handling** — `anyhow::Result` and `anyhow::bail!` throughout. Log with `tracing` (never stderr). +6. **Testing** — Bun tests (`bun test`) using `describe/test/expect`. Tests are + F.I.R.S.T. Logic-pure modules (TUI state/action/command/controller) are unit-tested. -6. **Testing** — `#[cfg(test)] mod tests` blocks inline in production files. Tests are F.I.R.S.T. (Fast, Independent, Repeatable, Self-validating, Timely). +7. **No Compiler Bypasses** — no `@ts-ignore` / `any` where a real type exists. -7. **No Compiler Bypasses** — Never use `#[allow(...)]`, `#[expect(...)]`, or `#[allow(dead_code)]`. Fix the underlying code. +8. **Boy Scout Rule** — leave every module cleaner than you found it. -8. **Boy Scout Rule** — Leave every module cleaner than you found it. +## Stack -## Available Agents +- Bun 1.3.14 (packageManager pinned in `package.json`). +- TypeScript strict mode, `tsconfig.json` uses `paths` aliases (`@zesdex/*`) → + `src/*`. `bun build ... --compile` bundles into a single `dist/zesdex` binary. -- `@rust-engineer` — Rust clean architecture specialist (subagent). -- `@code-reviewer` — Code review specialist (subagent). +## Commands -## Available Commands - -- `/check` — Run cargo check, clippy, and tests. -- `/audit` — Code quality audit against clean-architecture best practices. -- `/doc` — Generate or update doc comments. +- `bun run check` — `tsc --noEmit` typecheck. +- `bun run test` — run all Bun tests. +- `bun run build` — compile `dist/zesdex` binary. +- `bun run tui` — run the interactive TUI. +- `bun run headless ""` — one-shot agent turn. +- `bun run bootstrap` — seed default settings/app_config. diff --git a/docker-compose.yml b/docker-compose.yml index 99b4e6d..d9429e5 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,5 +1,7 @@ version: "3.9" +# Zesdex is now a pure CLI/TUI tool (no REST/WS/gRPC/web/daemon servers). +# This compose file just provides an interactive TTY shell into the binary. services: zesdex: build: @@ -12,12 +14,8 @@ services: - ./workspace:/workspace:ro environment: - ZESDEX_DATA_DIR=/data - - ZESDEX_API_PORT=${ZESDEX_API_PORT:-8080} stdin_open: true tty: true - # For daemon mode, expose the IPC socket directory - # ports: - # - "127.0.0.1:${ZESDEX_PORT:-0}:${ZESDEX_PORT:-0}" volumes: zesdex-data: diff --git a/package.json b/package.json index 7f413d9..20a68ee 100644 --- a/package.json +++ b/package.json @@ -10,14 +10,10 @@ "build": "bun build src/interfaces/cli/index.ts --compile --outfile=dist/zesdex", "check": "tsc --noEmit", "test": "bun test", - "cli": "bun src/interfaces/cli/index.ts", - "bootstrap": "bun src/interfaces/cli/bootstrap.ts", - "api": "bun src/interfaces/api/main.ts", - "ws": "bun src/interfaces/ws/main.ts", - "grpc": "bun src/interfaces/grpc/main.ts", - "web": "bun src/interfaces/web/main.ts", - "daemon": "bun src/interfaces/daemon/main.ts", - "tui": "bun src/interfaces/tui/main.ts" + "tui": "bun src/interfaces/cli/index.ts", + "cli": "bun src/interfaces/cli/index.ts --headless", + "headless": "bun src/interfaces/cli/index.ts --headless", + "bootstrap": "bun src/interfaces/cli/bootstrap.ts" }, "devDependencies": { "@types/bun": "^1.2.0", diff --git a/src/application/auth/index.ts b/src/application/auth/index.ts deleted file mode 100644 index 5fc8418..0000000 --- a/src/application/auth/index.ts +++ /dev/null @@ -1,3 +0,0 @@ -/** Auth application module — OAuth PKCE + session management use-cases. */ -export * from "./oauth_service.ts"; -export * from "./session_service.ts"; \ No newline at end of file diff --git a/src/application/auth/oauth_service.test.ts b/src/application/auth/oauth_service.test.ts deleted file mode 100644 index 244f63f..0000000 --- a/src/application/auth/oauth_service.test.ts +++ /dev/null @@ -1,141 +0,0 @@ -import { describe, expect, it } from "bun:test"; -import { generatePkcePair, OAuthUseCase } from "./oauth_service.ts"; -import type { OAuthFlowStore } from "./oauth_service.ts"; -import { createHash } from "node:crypto"; -import type { OAuthToken } from "@zesdex/domain"; - -describe("generatePkcePair", () => { - it("produces a verifier ≥43 chars and a valid S256 challenge", () => { - const [verifier, challenge] = generatePkcePair(); - expect(verifier.length).toBeGreaterThanOrEqual(43); - // Challenge = base64url(sha256(verifier)). - const expected = Buffer.from(createHash("sha256").update(verifier, "utf8").digest()).toString("base64url"); - expect(challenge).toBe(expected); - }); - - it("is unique across calls", () => { - const [a] = generatePkcePair(); - const [b] = generatePkcePair(); - expect(a).not.toBe(b); - }); -}); - -class MemoryFlowStore implements OAuthFlowStore { - private verifier = ""; - private state = ""; - async saveFlowState(v: string, s: string): Promise { - this.verifier = v; - this.state = s; - } - async loadVerifier(): Promise { - return this.verifier; - } - async loadState(): Promise { - return this.state; - } - async clear(): Promise { - this.verifier = ""; - this.state = ""; - } -} - -function makeRepo() { - let token: OAuthToken | null = null; - return { - repo: { - async saveToken(_path: string, t: OAuthToken): Promise { - token = t; - }, - async loadToken(): Promise { - return token; - }, - }, - getToken: () => token, - }; -} - -describe("OAuthUseCase", () => { - it("builds an auth URL with PKCE params and persists flow state", async () => { - const { repo } = makeRepo(); - const store = new MemoryFlowStore(); - const exchanger = { - async exchangeCode(): Promise { - return { access_token: "at", refresh_token: "rt", expires_at: 9999, token_type: "Bearer" }; - }, - }; - const useCase = new OAuthUseCase(repo, store, exchanger, "/tmp/token.json"); - - const { authUrl, state } = await useCase.startFlow( - { auth_url: "https://provider.example/oauth/authorize", token_url: "https://provider.example/oauth/token", client_id: "cid", scopes: ["openid", "profile"] }, - "http://localhost:9999/callback", - ); - - const url = new URL(authUrl); - expect(url.searchParams.get("response_type")).toBe("code"); - expect(url.searchParams.get("client_id")).toBe("cid"); - expect(url.searchParams.get("scope")).toBe("openid profile"); - expect(url.searchParams.get("code_challenge_method")).toBe("S256"); - expect(url.searchParams.get("state")).toBe(state); - expect(url.searchParams.get("code_challenge")).toBeTruthy(); - // Flow state persisted. - expect(await store.loadState()).toBe(state); - expect(await store.loadVerifier()).toBeTruthy(); - }); - - it("enforces CSRF state match in completeFlow", async () => { - const { repo } = makeRepo(); - const store = new MemoryFlowStore(); - const exchanger = { - async exchangeCode(): Promise { - return { access_token: "at", refresh_token: "rt", expires_at: 9999, token_type: "Bearer" }; - }, - }; - const useCase = new OAuthUseCase(repo, store, exchanger, "/tmp/token.json"); - await useCase.startFlow( - { auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] }, - "http://localhost:9999/callback", - ); - - await expect( - useCase.completeFlow( - { auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] }, - "http://localhost:9999/callback", - "code123", - "wrong-state", - ), - ).rejects.toMatchObject({ kind: "state_mismatch" }); - }); - - it("completes flow successfully with matching state", async () => { - const { repo, getToken } = makeRepo(); - const store = new MemoryFlowStore(); - let exchangedVerifier = ""; - const exchanger = { - async exchangeCode(_tokenUrl: string, _clientId: string, _secret: string | null, _redirectUri: string, code: string, codeVerifier: string): Promise { - exchangedVerifier = codeVerifier; - return { access_token: `at-${code}`, refresh_token: "rt", expires_at: 9999, token_type: "Bearer" }; - }, - }; - const useCase = new OAuthUseCase(repo, store, exchanger, "/tmp/token.json"); - const { state } = await useCase.startFlow( - { auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] }, - "http://localhost:9999/callback", - ); - const savedVerifier = await store.loadVerifier(); - - const token = await useCase.completeFlow( - { auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] }, - "http://localhost:9999/callback", - "code123", - state, - ); - - expect(token.access_token).toBe("at-code123"); - // The code verifier passed to the exchanger is the one saved at start. - expect(exchangedVerifier).toBe(savedVerifier); - // Flow state cleared after completion. - expect(await store.loadState()).toBe(""); - // Token persisted via repo. - expect(getToken()?.access_token).toBe("at-code123"); - }); -}); \ No newline at end of file diff --git a/src/application/auth/oauth_service.ts b/src/application/auth/oauth_service.ts deleted file mode 100644 index af9bebe..0000000 --- a/src/application/auth/oauth_service.ts +++ /dev/null @@ -1,121 +0,0 @@ -/** - * OAuth 2.0 authorization-code + PKCE flow use-case. Mirrors - * `apps/application/src/auth/oauth_service.rs`. - */ -import { - type OAuthConfig, - type OAuthRepository, - type OAuthToken, - authInvalidConfig, - authStateMismatch, -} from "@zesdex/domain"; -import { createHash, randomUUID } from "node:crypto"; - -/* -------------------------------------------------------------------------- */ -/* Port traits */ -/* -------------------------------------------------------------------------- */ - -/** Persistence contract for ephemeral OAuth flow state. */ -export interface OAuthFlowStore { - saveFlowState(verifier: string, state: string): Promise; - loadVerifier(): Promise; - loadState(): Promise; - clear(): Promise; -} - -/** Abstraction for exchanging an authorization code for tokens. */ -export interface TokenExchanger { - exchangeCode( - tokenUrl: string, - clientId: string, - clientSecret: string | null, - redirectUri: string, - code: string, - codeVerifier: string, - ): Promise; -} - -/* -------------------------------------------------------------------------- */ -/* PKCE helpers */ -/* -------------------------------------------------------------------------- */ - -function base64url(input: string | Uint8Array): string { - const buf = typeof input === "string" ? new TextEncoder().encode(input) : Buffer.from(input as Uint8Array); - return Buffer.from(buf).toString("base64url"); -} - -/** Generate a PKCE code-verifier and its S256 code-challenge. */ -export function generatePkcePair(): [string, string] { - const bytes = new Uint8Array(32); - crypto.getRandomValues(bytes); - const verifier = base64url(bytes); - const challenge = base64url(createHash("sha256").update(verifier, "utf8").digest()); - return [verifier, challenge]; -} - -/** Generate a random CSRF state token (UUID-based). */ -export function generateStateToken(): string { - return randomUUID(); -} - -/* -------------------------------------------------------------------------- */ -/* OAuthUseCase */ -/* -------------------------------------------------------------------------- */ - -/** Concrete OAuth flow use-case over injected repositories. */ -export class OAuthUseCase { - constructor( - private tokenRepo: OAuthRepository, - private flowStore: OAuthFlowStore, - private tokenExchanger: TokenExchanger, - private tokenPath: string, - ) {} - - async startFlow(config: OAuthConfig, redirectUri: string): Promise<{ authUrl: string; state: string }> { - if (config.auth_url === "") { - throw authInvalidConfig("OAuth auth_url is empty"); - } - - const [verifier, challenge] = generatePkcePair(); - const state = generateStateToken(); - - await this.flowStore.saveFlowState(verifier, state); - - const url = new URL(config.auth_url); - url.searchParams.set("response_type", "code"); - url.searchParams.set("client_id", config.client_id); - url.searchParams.set("redirect_uri", redirectUri); - url.searchParams.set("scope", config.scopes.join(" ")); - url.searchParams.set("state", state); - url.searchParams.set("code_challenge_method", "S256"); - url.searchParams.set("code_challenge", challenge); - - return { authUrl: url.toString(), state }; - } - - async completeFlow(config: OAuthConfig, redirectUri: string, code: string, state: string): Promise { - // CSRF check. - const expectedState = await this.flowStore.loadState(); - if (expectedState !== state) throw authStateMismatch(); - - // Read the PKCE verifier saved in start_flow. - const verifier = await this.flowStore.loadVerifier(); - - const token = await this.tokenExchanger.exchangeCode( - config.token_url, - config.client_id, - config.client_secret ?? null, - redirectUri, - code, - verifier, - ); - - await this.tokenRepo.saveToken(this.tokenPath, token); - await this.flowStore.clear(); - return token; - } - - async getToken(): Promise { - return this.tokenRepo.loadToken(this.tokenPath); - } -} \ No newline at end of file diff --git a/src/application/auth/session_service.ts b/src/application/auth/session_service.ts deleted file mode 100644 index d40d0f1..0000000 --- a/src/application/auth/session_service.ts +++ /dev/null @@ -1,44 +0,0 @@ -/** - * Session management use-case. Mirrors `apps/application/src/auth/session_service.rs`. - */ -import { randomUUID } from "node:crypto"; -import { - type Session, - type SessionId, - type SessionLockRepository, - type SessionRepository, - newSession, - newSessionId, - authOtherError as otherErr, -} from "@zesdex/domain"; - -/** Concrete session service backed by injected repositories. */ -export class SessionServiceImpl { - constructor( - private sessionRepo: SessionRepository, - /** Lock repository is injected for future lock acquire/release (matches Rust contract). */ - // @ts-expect-error -- kept for structural parity with the Rust `SessionServiceImpl` - private lockRepo: SessionLockRepository, - private baseDir: string, - ) {} - - async createSession(title: string): Promise { - const idRes = newSessionId(randomUUID()); - if (!idRes.ok) throw otherErr(idRes.error); - const titleOwned = title === "" ? "New Session" : title; - const session = newSession(idRes.value, titleOwned); - await this.sessionRepo.saveSession(this.baseDir, session); - return session; - } - - async listAll(): Promise { - return this.sessionRepo.listSessions(this.baseDir); - } - - async archiveSession(id: SessionId): Promise { - const session = await this.sessionRepo.loadSession(this.baseDir, id); - session.archived = true; - session.updated_at = Date.now(); - await this.sessionRepo.saveSession(this.baseDir, session); - } -} \ No newline at end of file diff --git a/src/application/index.ts b/src/application/index.ts index 29b7809..3ec6e86 100644 --- a/src/application/index.ts +++ b/src/application/index.ts @@ -5,5 +5,4 @@ */ export * from "./ports/index.ts"; export * from "./agent/index.ts"; -export * from "./auth/index.ts"; export * from "./cms/index.ts"; \ No newline at end of file diff --git a/src/domain/auth/commands.ts b/src/domain/auth/commands.ts deleted file mode 100644 index 079bc28..0000000 --- a/src/domain/auth/commands.ts +++ /dev/null @@ -1,11 +0,0 @@ -/** - * Command types for IAM domain operations. Mirrors `commands.rs`. - * Carries only the data needed to construct a session entity. - */ -export interface NewSession { - title: string; -} - -export function newSessionCommand(title: string): NewSession { - return { title }; -} \ No newline at end of file diff --git a/src/domain/auth/error.ts b/src/domain/auth/error.ts deleted file mode 100644 index 648a060..0000000 --- a/src/domain/auth/error.ts +++ /dev/null @@ -1,47 +0,0 @@ -/** - * Domain error types for the IAM (auth) module. - * Mirrors `apps/domain/src/auth/error.rs`. - */ -import type { DomainError } from "../core/error.ts"; - -/** Shared repository error type for IAM persistence. */ -export type RepositoryError = DomainError; - -/** Errors from service / use-case operations in the IAM domain. */ -export type ServiceError = - | { kind: "repository"; error: DomainError } - | { kind: "invalid_config"; message: string } - | { kind: "state_mismatch" } - | { kind: "oauth_provider"; message: string } - | { kind: "other"; message: string }; - -export function repositoryErr(err: DomainError): ServiceError { - return { kind: "repository", error: err }; -} -export function invalidConfig(message: string): ServiceError { - return { kind: "invalid_config", message }; -} -export function stateMismatch(): ServiceError { - return { kind: "state_mismatch" }; -} -export function oauthProvider(message: string): ServiceError { - return { kind: "oauth_provider", message }; -} -export function authOther(message: string): ServiceError { - return { kind: "other", message }; -} - -export function serviceErrorToString(e: ServiceError): string { - switch (e.kind) { - case "repository": - return `repository error: ${e.error.message}`; - case "invalid_config": - return `invalid configuration: ${e.message}`; - case "state_mismatch": - return `OAuth state mismatch — possible CSRF attack`; - case "oauth_provider": - return `OAuth provider error: ${e.message}`; - case "other": - return e.message; - } -} \ No newline at end of file diff --git a/src/domain/auth/index.ts b/src/domain/auth/index.ts deleted file mode 100644 index 4e85570..0000000 --- a/src/domain/auth/index.ts +++ /dev/null @@ -1,20 +0,0 @@ -/** Auth domain module — sessions, locks, OAuth, and their interfaces. */ -export * from "./session.ts"; -export * from "./session_id.ts"; -export * from "./session_lock.ts"; -export * from "./oauth.ts"; -export * from "./repository.ts"; -export * from "./service.ts"; -export * from "./commands.ts"; -export { - type ServiceError as AuthServiceError, - type RepositoryError as AuthRepositoryError, -} from "./error.ts"; -export { - repositoryErr as authRepositoryErr, - invalidConfig as authInvalidConfig, - stateMismatch as authStateMismatch, - oauthProvider as authOauthProvider, - authOther as authOtherError, - serviceErrorToString as authServiceErrorToString, -} from "./error.ts"; \ No newline at end of file diff --git a/src/domain/auth/oauth.ts b/src/domain/auth/oauth.ts deleted file mode 100644 index 7b07b17..0000000 --- a/src/domain/auth/oauth.ts +++ /dev/null @@ -1,40 +0,0 @@ -/** - * OAuth entities. Mirrors `apps/domain/src/auth/oauth.rs`. - */ -export interface OAuthToken { - /** The OAuth 2.0 access token string. */ - access_token: string; - /** Optional refresh token for long-lived access. */ - refresh_token?: string; - /** Absolute expiry timestamp (epoch seconds). */ - expires_at: number; - /** Token type, e.g. `"Bearer"`. */ - token_type: string; -} - -/** Static configuration for an OAuth provider. */ -export interface OAuthConfig { - /** Authorization endpoint URL. */ - auth_url: string; - /** Token exchange endpoint URL. */ - token_url: string; - /** OAuth client identifier. */ - client_id: string; - /** Optional client secret. */ - client_secret?: string; - /** Space-separated list of requested scopes. */ - scopes: string[]; -} - -/** Default scopes for a new OAuthConfig. */ -export const DEFAULT_OAUTH_SCOPES = ["openid", "profile", "email"]; - -export function newOAuthConfig(): OAuthConfig { - return { - auth_url: "", - token_url: "", - client_id: "", - client_secret: undefined, - scopes: [...DEFAULT_OAUTH_SCOPES], - }; -} \ No newline at end of file diff --git a/src/domain/auth/repository.ts b/src/domain/auth/repository.ts deleted file mode 100644 index 0ce4b52..0000000 --- a/src/domain/auth/repository.ts +++ /dev/null @@ -1,39 +0,0 @@ -/** - * Repository trait definitions (interfaces) — pure, no impls. - * Mirrors `apps/domain/src/auth/repository.rs`. Infrastructure adapters - * implement these. - */ -import type { OAuthToken } from "./oauth.ts"; -import type { Session } from "./session.ts"; -import type { SessionId } from "./session_id.ts"; - -/** Repository for loading, saving, listing, and deleting sessions. */ -export interface SessionRepository { - /** List all loadable sessions under `/sessions/`. */ - listSessions(baseDir: string): Promise | Session[]; - /** Load a single session by id. */ - loadSession(baseDir: string, id: SessionId): Promise | Session; - /** Save a session's metadata to disk. */ - saveSession(baseDir: string, session: Session): Promise | void; - /** Delete a session directory and all its contents. */ - deleteSession(baseDir: string, id: SessionId): Promise | void; -} - -/** Repository for per-session PID-file advisory locks. */ -export interface SessionLockRepository { - /** Try to acquire the lock. `true` if acquired, `false` if a live process holds it. */ - tryLock(sessionDir: string): Promise | boolean; - /** Release the lock. */ - unlock(sessionDir: string): Promise | void; - /** Check whether a process with the given PID is alive. */ - isAlive(pid: number): boolean; -} - -/** Repository for persisting and loading OAuth tokens. */ -export interface OAuthRepository { - /** Persist an OAuth token to a JSON file. */ - saveToken(path: string, token: OAuthToken): Promise | void; - /** Load an OAuth token, returning `null` if the file does not exist. */ - loadToken(path: string): Promise | OAuthToken | null; -} - diff --git a/src/domain/auth/service.ts b/src/domain/auth/service.ts deleted file mode 100644 index 7a3a37d..0000000 --- a/src/domain/auth/service.ts +++ /dev/null @@ -1,35 +0,0 @@ -/** - * Service trait definitions — use-case boundaries for sessions and OAuth. - * Mirrors `apps/domain/src/auth/service.rs`. Implementations live in the - * application layer. - */ -import type { OAuthConfig, OAuthToken } from "./oauth.ts"; -import type { Session } from "./session.ts"; -import type { SessionId } from "./session_id.ts"; - -/** Session management use-case boundary. */ -export interface SessionService { - /** Create a new session with a generated UUID and the given title. */ - createSession(title: string): Promise | Session; - /** List all available sessions. */ - listAll(): Promise | Session[]; - /** Archive a session by id (sets `archived = true`). */ - archiveSession(id: SessionId): Promise | void; -} - -/** OAuth flow use-case boundary. */ -export interface OAuthService { - /** - * Start an OAuth authorization-code + PKCE flow. Returns `{ authUrl, state }`: - * the URL to send the user to, and the CSRF state token that must be passed - * back into `completeFlow` unchanged. - */ - startFlow(config: OAuthConfig, redirectUri: string): Promise<{ authUrl: string; state: string }> | { authUrl: string; state: string }; - /** - * Complete the flow: validate `state`, then exchange `code` for a token. - * Validates `state` against the persisted value (CSRF check). - */ - completeFlow(config: OAuthConfig, redirectUri: string, code: string, state: string): Promise | OAuthToken; - /** Retrieve the currently stored OAuth token (if any). */ - getToken(): Promise | OAuthToken | null; -} \ No newline at end of file diff --git a/src/domain/auth/session.ts b/src/domain/auth/session.ts deleted file mode 100644 index a39cd6f..0000000 --- a/src/domain/auth/session.ts +++ /dev/null @@ -1,56 +0,0 @@ -/** - * Session metadata. Mirrors `apps/domain/src/auth/session.rs`. - */ -import * as path from "node:path"; - -export interface Session { - /** Unique session identifier. */ - id: string; - /** Epoch-millis timestamp of creation. */ - created_at: number; - /** Epoch-millis timestamp of last update. */ - updated_at: number; - /** Human-readable title for the conversation. */ - title: string; - /** Model identifier string. */ - model: string; - /** Workspace root directories associated with this session. */ - workspace_roots: string[]; - /** Running count of messages in the conversation. */ - message_count: number; - /** Running count of tokens consumed. */ - token_count: number; - /** Soft-delete flag. */ - archived: boolean; - /** Optional AI-generated conversation summary. */ - summary?: string; -} - -const DEFAULT_MODEL = "anthropic/claude-opus-4-8"; - -/** Create a new session with the given id/title and current dir as root. */ -export function newSession(id: string, title: string): Session { - const now = Date.now(); - return { - id, - created_at: now, - updated_at: now, - title, - model: DEFAULT_MODEL, - workspace_roots: [process.cwd()], - message_count: 0, - token_count: 0, - archived: false, - summary: undefined, - }; -} - -/** Compute this session's directory under `/sessions/`. */ -export function sessionDir(session: Session, baseDir: string): string { - return path.join(baseDir, "sessions", session.id); -} - -/** Compute this session's `conversation.json` path. */ -export function conversationPath(session: Session, baseDir: string): string { - return path.join(sessionDir(session, baseDir), "conversation.json"); -} \ No newline at end of file diff --git a/src/domain/auth/session_id.test.ts b/src/domain/auth/session_id.test.ts deleted file mode 100644 index 9ce3a96..0000000 --- a/src/domain/auth/session_id.test.ts +++ /dev/null @@ -1,25 +0,0 @@ -import { describe, expect, it } from "bun:test"; -import { newSessionId } from "./session_id.ts"; - -describe("newSessionId", () => { - it("accepts valid UUIDs", () => { - expect(newSessionId("550e8400-e29b-41d4-a716-446655440000").ok).toBe(true); - expect(newSessionId("my-session_123").ok).toBe(true); - }); - - it("rejects path traversal", () => { - expect(newSessionId("../etc/passwd").ok).toBe(false); - expect(newSessionId("foo/../../bar").ok).toBe(false); - expect(newSessionId("foo\\..\\bar").ok).toBe(false); - }); - - it("rejects empty", () => { - expect(newSessionId("").ok).toBe(false); - }); - - it("returns the id string on success", () => { - const res = newSessionId("abc-123"); - expect(res.ok).toBe(true); - if (res.ok) expect(res.value).toBe("abc-123"); - }); -}); \ No newline at end of file diff --git a/src/domain/auth/session_id.ts b/src/domain/auth/session_id.ts deleted file mode 100644 index 0ede37e..0000000 --- a/src/domain/auth/session_id.ts +++ /dev/null @@ -1,41 +0,0 @@ -/** - * Validated session identifier. Mirrors `apps/domain/src/auth/session_id.rs`. - * - * Guarantees the inner string is non-empty and contains no path-traversal - * characters (`/`, `\`, `..`) or other unsafe delimiters. - */ - -const SESSION_ID_RE = /^[A-Za-z0-9_.-]+$/; - -export type SessionId = string & { __sessionId?: true }; - -/** - * Validate and construct a `SessionId`. - * Returns `err` message if the input contains path separators, `..`, or is empty. - */ -export function newSessionId(id: string): { ok: true; value: SessionId } | { ok: false; error: string } { - if (id === "") { - return { ok: false, error: "session id must not be empty" }; - } - if (id.includes("/") || id.includes("\\") || id.includes("..") || !SESSION_ID_RE.test(id)) { - return { ok: false, error: `session id '${id}' must not contain path separators` }; - } - return { ok: true, value: id as SessionId }; -} - -/** Assert that a session id is valid (throws if not). */ -export function assertSessionId(id: string): SessionId { - const res = newSessionId(id); - if (!res.ok) throw new Error(res.error); - return res.value; -} - -/** Return the underlying string. */ -export function sessionIdAsString(id: SessionId): string { - return id; -} - -/** Classic Rust-style UUID format (as used by the Rust app). */ -export function isUuidLike(id: string): boolean { - return /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i.test(id); -} \ No newline at end of file diff --git a/src/domain/auth/session_lock.ts b/src/domain/auth/session_lock.ts deleted file mode 100644 index 4e7cbdf..0000000 --- a/src/domain/auth/session_lock.ts +++ /dev/null @@ -1,120 +0,0 @@ -/** - * PID-file based advisory lock preventing two processes from operating on - * the same session directory concurrently. Mirrors `session_lock.rs`. - * - * Strategy: atomic `O_CREAT|O_EXCL` acquire → on existing file, liveness-check - * the owning PID (Unix: read `/proc//exe` and compare to self, plus - * `kill(pid, 0)`) → stale locks overwritten atomically (temp + rename + fsync). - */ -import * as fs from "node:fs"; -import * as path from "node:path"; - -/** A PID-file lock (`/.lock`) tied to the current process. */ -export class SessionLock { - readonly path: string; - readonly pid: number; - - constructor(sessionDir: string) { - this.path = path.join(sessionDir, ".lock"); - this.pid = process.pid; - } - - /** Attempt to acquire the lock. Returns `true` on success, `false` if held by a live process. */ - tryLock(): boolean { - // Phase 1: atomic create (O_CREAT|O_EXCL equivalent via flag 'wx'). - try { - const fd = fs.openSync(this.path, "wx"); - try { - fs.writeFileSync(fd, String(this.pid)); - fs.fsyncSync(fd); - } finally { - fs.closeSync(fd); - } - return true; - } catch (e) { - const code = (e as NodeJS.ErrnoException).code; - if (code !== "EEXIST") throw e; - // Lock file exists — check staleness. - } - - // Phase 2: read owning PID and check liveness. - const content = fs.readFileSync(this.path, "utf8").trim(); - if (content !== "") { - const pid = Number(content); - if (Number.isFinite(pid) && pid > 0) { - if (pidIsAlive(pid)) return false; - } - } - - // Phase 3: stale lock — overwrite atomically. - const tmp = this.path + ".tmp"; - const fd = fs.openSync(tmp, "w"); - try { - fs.writeFileSync(fd, String(this.pid)); - fs.fsyncSync(fd); - } finally { - fs.closeSync(fd); - } - fs.renameSync(tmp, this.path); - // Best-effort parent dir fsync. - const parent = path.dirname(this.path); - try { - const dirFd = fs.openSync(parent, "r"); - try { - fs.fsyncSync(dirFd); - } finally { - fs.closeSync(dirFd); - } - } catch { - /* best-effort */ - } - return true; - } - - /** Explicitly release the lock by removing the lock file. */ - unlock(): void { - try { - fs.unlinkSync(this.path); - } catch { - /* ignore */ - } - } -} - -/** - * Check whether a process with the given PID is alive. Conservative on - * non-Unix: returns true. - */ -export function pidIsAlive(pid: number): boolean { - if (process.platform === "win32") return true; - try { - const selfExe = fs.readlinkSync("/proc/self/exe"); - let target: string; - try { - target = fs.readlinkSync(`/proc/${pid}/exe`); - } catch { - return false; - } - if (target !== selfExe) return false; - // kill(pid, 0): throws if the process is gone / not permitted. - try { - process.kill(pid, 0); - } catch { - return false; - } - // Re-check to close the TOCTOU window. - try { - return fs.readlinkSync(`/proc/${pid}/exe`) === selfExe; - } catch { - return false; - } - } catch { - // No /proc (macOS) — fall back to kill(pid,0) only. - try { - process.kill(pid, 0); - return true; - } catch { - return false; - } - } -} \ No newline at end of file diff --git a/src/domain/index.ts b/src/domain/index.ts index 35d85de..6142f70 100644 --- a/src/domain/index.ts +++ b/src/domain/index.ts @@ -3,7 +3,6 @@ * Zero framework dependencies, zero I/O. Mirrors the Rust `zesdex-domain` crate. */ export * from "./core/index.ts"; -export * from "./auth/index.ts"; export * from "./cms/index.ts"; export * from "./agent/index.ts"; export * from "./subagent/mod.ts"; diff --git a/src/infrastructure/auth/index.ts b/src/infrastructure/auth/index.ts deleted file mode 100644 index 450823d..0000000 --- a/src/infrastructure/auth/index.ts +++ /dev/null @@ -1,8 +0,0 @@ -/** - * Auth infrastructure implementations — Argon2 passwords + HS256 JWT + OAuth loopback. - * Mirrors `apps/infrastructure/src/auth/mod.rs`. - */ -export { Argon2PasswordService } from "./password.ts"; -export { Hs256TokenService, createToken, verifyToken } from "./jwt.ts"; -export type { JwtClaims } from "./jwt.ts"; -export { LoopbackServer } from "./oauth_loopback.ts"; diff --git a/src/infrastructure/auth/jwt.ts b/src/infrastructure/auth/jwt.ts deleted file mode 100644 index 8110441..0000000 --- a/src/infrastructure/auth/jwt.ts +++ /dev/null @@ -1,134 +0,0 @@ -/** - * JWT token utilities for HMAC-SHA256 (HS256) signing and verification. - * Mirrors `apps/infrastructure/src/auth/jwt.rs`. - * - * Uses `node:crypto` HMAC — no external JWT library needed. - * Implements compact JWT encoding/decoding per RFC 7515. - */ -import { createHmac, timingSafeEqual } from "node:crypto"; -import type { TokenService } from "@zesdex/application"; - -/** Standard JWT claims. */ -export interface JwtClaims { - sub: string; - exp: number; - iat: number; - /** Token purpose: "access" or "refresh". */ - typ: "access" | "refresh"; - /** Optional session binding. */ - session_id?: string; -} - -/** JWT header (always HS256). */ -interface JwtHeader { - alg: "HS256"; - typ: "JWT"; -} - -const ACCESS_TOKEN_EXPIRY_SECS = 3600; // 1 hour -const REFRESH_TOKEN_EXPIRY_SECS = 604800; // 7 days - -function base64url(data: string | Buffer): string { - const str = - typeof data === "string" ? Buffer.from(data, "utf8").toString("base64") : data.toString("base64"); - return str.replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); -} - -function base64urlDecode(str: string): Buffer { - let s = str.replace(/-/g, "+").replace(/_/g, "/"); - while (s.length % 4) s += "="; - return Buffer.from(s, "base64"); -} - -function sign(secret: string, data: string): string { - return base64url(createHmac("sha256", secret).update(data).digest()); -} - -function encodeJson(obj: object): string { - return base64url(JSON.stringify(obj)); -} - -/** Create a JWT token with the given claims and secret. */ -export function createToken(secret: string, claims: JwtClaims): string { - const header: JwtHeader = { alg: "HS256", typ: "JWT" }; - const encHeader = encodeJson(header); - const encPayload = encodeJson(claims); - const sig = sign(secret, `${encHeader}.${encPayload}`); - return `${encHeader}.${encPayload}.${sig}`; -} - -/** Verify a JWT token and return its claims. Throws if invalid/expired. */ -export function verifyToken(secret: string, token: string): JwtClaims { - const parts = token.split("."); - if (parts.length !== 3) throw new Error("Invalid JWT: expected 3 parts"); - - const [encHeader, encPayload, sig] = parts as [string, string, string]; - - // Verify signature (constant-time) - const expectedSig = sign(secret, `${encHeader}.${encPayload}`); - const sigBuf = Buffer.from(sig, "base64"); - const expectedBuf = Buffer.from(expectedSig, "base64"); - if (sigBuf.length !== expectedBuf.length || !timingSafeEqual(sigBuf, expectedBuf)) { - throw new Error("Invalid JWT signature"); - } - - // Decode header - const header: JwtHeader = JSON.parse(base64urlDecode(encHeader).toString("utf8")); - if (header.alg !== "HS256") throw new Error(`Unsupported JWT algorithm: ${header.alg}`); - - // Decode payload - const claims: JwtClaims = JSON.parse(base64urlDecode(encPayload).toString("utf8")); - - // Validate required claims - if (!claims.sub || typeof claims.exp !== "number" || typeof claims.typ !== "string") { - throw new Error("Invalid JWT: missing required claims (sub, exp, typ)"); - } - - // Validate expiry with 60s leeway - const nowSecs = Math.floor(Date.now() / 1000); - if (claims.exp + 60 < nowSecs) { - throw new Error("JWT token expired"); - } - - return claims; -} - -/** - * Concrete TokenService implementing HS256 JWT. - * Generates access + refresh token pairs. - */ -export class Hs256TokenService implements TokenService { - constructor(private readonly secret: string) {} - - /** Generate an [access, refresh] token pair for the given subject. */ - generateTokens(sub: string): [string, string] { - const nowSecs = Math.floor(Date.now() / 1000); - const access: JwtClaims = { - sub, - exp: nowSecs + ACCESS_TOKEN_EXPIRY_SECS, - iat: nowSecs, - typ: "access", - }; - const refresh: JwtClaims = { - sub, - exp: nowSecs + REFRESH_TOKEN_EXPIRY_SECS, - iat: nowSecs, - typ: "refresh", - }; - return [createToken(this.secret, access), createToken(this.secret, refresh)]; - } - - /** Verify an access token and return the subject. Throws if invalid/expired/wrong type. */ - verifyAccessToken(token: string): string { - const claims = verifyToken(this.secret, token); - if (claims.typ !== "access") throw new Error(`Expected access token, got ${claims.typ}`); - return claims.sub; - } - - /** Verify a refresh token and return the subject. Throws if invalid/expired/wrong type. */ - verifyRefreshToken(token: string): string { - const claims = verifyToken(this.secret, token); - if (claims.typ !== "refresh") throw new Error(`Expected refresh token, got ${claims.typ}`); - return claims.sub; - } -} diff --git a/src/infrastructure/auth/oauth_loopback.ts b/src/infrastructure/auth/oauth_loopback.ts deleted file mode 100644 index 85616b7..0000000 --- a/src/infrastructure/auth/oauth_loopback.ts +++ /dev/null @@ -1,151 +0,0 @@ -/** - * Minimal loopback HTTP server for capturing OAuth authorization-code redirects. - * Mirrors `apps/infrastructure/src/auth/oauth_loopback.rs`. - * - * Binds to `127.0.0.1:`, exposes a redirect URI, and waits for a - * single callback carrying `?code=...&state=...`. Serves a static confirmation - * page and validates the state to prevent CSRF. - */ -import * as net from "node:net"; - -/** A single-use HTTP listener on 127.0.0.1 that receives the OAuth redirect. */ -export class LoopbackServer { - private server: net.Server | null = null; - private port = 0; - - /** Bind to 127.0.0.1 on an ephemeral port. */ - bind(): Promise { - return new Promise((resolve, reject) => { - const server = net.createServer(); - server.once("error", reject); - server.listen(0, "127.0.0.1", () => { - const addr = server.address(); - if (addr && typeof addr === "object") this.port = addr.port; - this.server = server; - resolve(); - }); - }); - } - - /** The redirect URI the browser should be sent to. */ - redirectUri(): string { - return `http://127.0.0.1:${this.port}/callback`; - } - - /** - * Wait for a single authorization-code callback (with timeout). - * Validates the CSRF state and returns the `code` on success. - */ - waitForCode(timeoutMs: number, expectedState: string): Promise { - return new Promise((resolve, reject) => { - const server = this.server; - if (!server) { - reject(new Error("server not bound")); - return; - } - - const timer = setTimeout(() => { - cleanup(); - reject(new Error("OAuth callback timed out")); - }, timeoutMs); - - const cleanup = () => { - clearTimeout(timer); - server.removeAllListeners("connection"); - }; - - server.once("connection", (socket) => { - socket.setTimeout(timeoutMs); - let buf = ""; - socket.on("data", (data) => { - buf += data.toString("utf8"); - // Wait for the end of the request headers - if (!buf.includes("\r\n\r\n")) return; - - const result = this.handleRequest(buf, expectedState); - socket.end(result.response); - cleanup(); - - // Give the response a chance to flush before closing. - setTimeout(() => { - if (result.error) reject(result.error); - else resolve(result.code!); - }, 20); - }); - socket.on("error", () => { - cleanup(); - reject(new Error("OAuth socket error")); - }); - }); - }); - } - - /** Close the server and release the port. */ - close(): Promise { - return new Promise((resolve) => { - const server = this.server; - this.server = null; - if (!server) { - resolve(); - return; - } - server.close(() => resolve()); - }); - } - - /** Parse the HTTP request, build the response, and return code/error. */ - private handleRequest(request: string, expectedState: string): { - response: string; - code: string | null; - error: Error | null; - } { - const code = extractQueryParam(request, "code"); - const state = extractQueryParam(request, "state"); - const stateOk = state === expectedState; - - let response: string; - if (code !== null && stateOk) { - response = - "HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\n\r\n" + - "Authorization complete. You may close this tab."; - } else if (code !== null && !stateOk) { - response = - "HTTP/1.1 400 Bad Request\r\nContent-Type: text/plain\r\n\r\n" + - "State mismatch — possible CSRF attack."; - } else { - response = - "HTTP/1.1 400 Bad Request\r\nContent-Type: text/plain\r\n\r\n" + - "Missing authorization code."; - } - - if (!stateOk) { - return { response, code: null, error: new Error("state mismatch") }; - } - if (code === null) { - return { response, code: null, error: new Error("code not found in callback") }; - } - return { response, code, error: null }; - } -} - -/** Extract a URL query parameter from an HTTP request line (percent-decoded). */ -function extractQueryParam(request: string, key: string): string | null { - const line = request.split("\r\n")[0] ?? ""; - const parts = line.split(" "); - if (parts.length < 2) return null; - const path = parts[1]!; - const query = path.split("?")[1]; - if (!query) return null; - for (const pair of query.split("&")) { - const eq = pair.indexOf("="); - const k = eq === -1 ? pair : pair.slice(0, eq); - const v = eq === -1 ? "" : pair.slice(eq + 1); - if (k === key) return urlDecode(v); - } - return null; -} - -/** Percent-decode a string (e.g. `%20` -> space). */ -function urlDecode(s: string): string { - return decodeURIComponent(s.replace(/\+/g, "%20")); -} diff --git a/src/infrastructure/auth/password.ts b/src/infrastructure/auth/password.ts deleted file mode 100644 index b680493..0000000 --- a/src/infrastructure/auth/password.ts +++ /dev/null @@ -1,27 +0,0 @@ -/** - * Argon2 password hashing and verification. - * Mirrors `apps/infrastructure/src/auth/password.rs`. - * - * Uses @node-rs/argon2 (Argon2id) for password hashing, matching the Rust argon2 crate. - */ -import { hash, verify } from "@node-rs/argon2"; -import type { PasswordService } from "@zesdex/application"; - -/** Argon2id password hashing service. */ -export class Argon2PasswordService implements PasswordService { - /** - * Hash a plaintext password using Argon2id with a random salt. - * The PHC-encoded hash string is returned. - */ - async hash(password: string): Promise { - return hash(password); - } - - /** - * Verify a plaintext password against a previously-hashed PHC string. - * Returns `true` if the password matches. - */ - async verify(password: string, hashStr: string): Promise { - return verify(hashStr, password); - } -} diff --git a/src/infrastructure/index.ts b/src/infrastructure/index.ts index 4522378..f73b097 100644 --- a/src/infrastructure/index.ts +++ b/src/infrastructure/index.ts @@ -6,8 +6,6 @@ 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"; diff --git a/src/infrastructure/ipc/client.ts b/src/infrastructure/ipc/client.ts deleted file mode 100644 index c56bba1..0000000 --- a/src/infrastructure/ipc/client.ts +++ /dev/null @@ -1,43 +0,0 @@ -/** - * IPC client — connects to the daemon's Unix socket and sends/receives - * framed JSON messages. - * Mirrors `apps/infrastructure/src/ipc/client.rs`. - */ -import * as net from "node:net"; -import { Connection } from "./conn.ts"; - -/** A thread-safe IPC client connected to a Zesdex daemon over a Unix socket. */ -export class IpcClient { - private conn: Connection; - - private constructor(conn: Connection) { - this.conn = conn; - } - - /** Connect to a Unix socket at `path`. */ - static connectUnix(path: string): Promise { - return new Promise((resolve, reject) => { - const socket = net.createConnection(path); - socket.once("error", reject); - socket.once("connect", () => { - socket.removeListener("error", reject); - resolve(new IpcClient(new Connection(socket))); - }); - }); - } - - /** Send a JSON-serializable message as one framed frame. */ - send(msg: T): void { - this.conn.send(msg); - } - - /** Receive the next message and JSON-decode it (null on close). */ - async receive(): Promise { - return this.conn.receiveJson(); - } - - /** Close the connection. */ - close(): void { - this.conn.close(); - } -} diff --git a/src/infrastructure/ipc/conn.ts b/src/infrastructure/ipc/conn.ts deleted file mode 100644 index c10a0e7..0000000 --- a/src/infrastructure/ipc/conn.ts +++ /dev/null @@ -1,90 +0,0 @@ -/** - * IPC connection — a framed, JSON-message socket over node:net. - * Mirrors `apps/infrastructure/src/ipc/conn.rs`. - * - * Wraps a `net.Socket` (or `net.Server` connection) with the length-prefixed - * framing from `frame.ts`. Provides `send` (write one framed message) and - * an async-style receive via an event emitter / queue. - */ -import * as net from "node:net"; -import { FrameReassembler, encodeFrame } from "./frame.ts"; - -/** A connection over which framed JSON messages are exchanged. */ -export class Connection { - private socket: net.Socket; - private reassembler = new FrameReassembler(); - private messageQueue: Buffer[] = []; - private waiters: Array<(buf: Buffer | null) => void> = []; - - constructor(socket: net.Socket) { - this.socket = socket; - socket.on("data", (chunk: Buffer) => { - let frames: Buffer[]; - try { - frames = this.reassembler.feed(chunk); - } catch (e) { - // oversized frame — fail the connection - const err = e as Error; - this._wakeWaiters(null); - this.socket.destroy(new Error(err.message)); - return; - } - for (const frame of frames) { - this._enqueue(frame); - } - }); - socket.on("end", () => { - this.messageQueue = []; - this._wakeWaiters(null); - }); - socket.on("error", () => { - this._wakeWaiters(null); - }); - } - - /** Serialize and send a message as one framed frame. */ - send(msg: T): void { - this.socket.write(encodeFrame(JSON.stringify(msg))); - } - - /** - * Receive the next complete message as a Buffer (parsed JSON). - * Returns `null` when the connection closes. - */ - receive(): Promise { - if (this.messageQueue.length > 0) { - return Promise.resolve(this.messageQueue.shift()!); - } - return new Promise((resolve) => { - this.waiters.push(resolve); - }); - } - - /** Receive the next message and JSON-decode it. */ - async receiveJson(): Promise { - const buf = await this.receive(); - if (buf === null) return null; - return JSON.parse(buf.toString("utf8")) as T; - } - - /** Close the underlying socket. */ - close(): void { - this.socket.end(); - } - - private _enqueue(buf: Buffer): void { - const waiter = this.waiters.shift(); - if (waiter) { - waiter(buf); - } else { - this.messageQueue.push(buf); - } - } - - private _wakeWaiters(buf: Buffer | null): void { - while (this.waiters.length > 0) { - const w = this.waiters.shift()!; - w(buf); - } - } -} diff --git a/src/infrastructure/ipc/frame.ts b/src/infrastructure/ipc/frame.ts deleted file mode 100644 index db88cf0..0000000 --- a/src/infrastructure/ipc/frame.ts +++ /dev/null @@ -1,101 +0,0 @@ -/** - * Length-prefixed framing for Unix-socket IPC. - * Mirrors `apps/infrastructure/src/ipc/frame.rs`. - * - * Every message on the wire is encoded as: - * ```text - * [ 4-byte big-endian payload length ][ payload bytes (JSON) ] - * ``` - * - * In TS this operates over `node:net` Socket (byte stream). A helper tracks - * a partial buffer so a single `data` event may contain multiple frames or a - * partial frame. - */ - -/** Maximum payload size accepted (64 MiB). */ -export const MAX_PAYLOAD = 64 * 1024 * 1024; - -/** Buffer for accumulating a frame across multiple `data` events. */ -export class FrameReassembler { - private chunks: Buffer[] = []; - private buffered = 0; - private pendingLength: number | null = null; - - /** - * Feed raw bytes and return any complete frames that were buffered. - * Each returned element is one complete JSON payload. - */ - feed(chunk: Buffer): Buffer[] { - this.chunks.push(chunk); - this.buffered += chunk.length; - return this.tryFrames(); - } - - private tryFrames(): Buffer[] { - const frames: Buffer[] = []; - for (;;) { - // Ensure we have at least the 4-byte length prefix. - if (this.buffered < 4) break; - if (this.pendingLength === null) { - this.pendingLength = this.peekLength(); - } - const len = this.pendingLength; - if (len === null) break; - if (len > MAX_PAYLOAD) { - throw new Error(`frame payload too large: ${len} bytes (max ${MAX_PAYLOAD})`); - } - if (this.buffered < 4 + len) break; // wait for full payload - // Consume the full frame. - const frame = this.take(4 + len); - const payload = frame.subarray(4, 4 + len); - frames.push(Buffer.from(payload)); - this.pendingLength = null; - } - return frames; - } - - private peekLength(): number | null { - if (this.buffered < 4) return null; - const first = this.take(4); - // Put the 4 bytes back. - this.unshift(first); - return first.readUInt32BE(0); - } - - /** Take `n` bytes from the front of the buffer. */ - private take(n: number): Buffer { - let remaining = n; - const parts: Buffer[] = []; - while (remaining > 0 && this.chunks.length > 0) { - const head = this.chunks[0]!; - if (head.length <= remaining) { - parts.push(head); - this.chunks.shift(); - remaining -= head.length; - } else { - parts.push(head.subarray(0, remaining)); - this.chunks[0] = head.subarray(remaining); - remaining = 0; - } - } - this.buffered -= n; - return Buffer.concat(parts); - } - - /** Push `n` bytes back to the front of the buffer. */ - private unshift(buf: Buffer): void { - this.chunks.unshift(buf); - this.buffered += buf.length; - } -} - -/** Encode a JSON payload into a single length-prefixed frame buffer. */ -export function encodeFrame(plainText: string | Buffer): Buffer { - const payload = typeof plainText === "string" ? Buffer.from(plainText, "utf8") : plainText; - if (payload.length > MAX_PAYLOAD) { - throw new Error(`frame payload too large: ${payload.length} bytes (max ${MAX_PAYLOAD})`); - } - const lenBuf = Buffer.alloc(4); - lenBuf.writeUInt32BE(payload.length, 0); - return Buffer.concat([lenBuf, payload]); -} diff --git a/src/infrastructure/ipc/index.ts b/src/infrastructure/ipc/index.ts deleted file mode 100644 index f2cf25a..0000000 --- a/src/infrastructure/ipc/index.ts +++ /dev/null @@ -1,17 +0,0 @@ -/** - * IPC module — frame, protocol, connection, server, client. - * Mirrors `apps/infrastructure/src/ipc/mod.rs`. - */ -export { FrameReassembler, encodeFrame, MAX_PAYLOAD } from "./frame.ts"; -export { Connection } from "./conn.ts"; -export { IpcServer } from "./server.ts"; -export { IpcClient } from "./client.ts"; -export type { - KeyAction, - KeyActionName, - ClientRequest, - MessageEntry, - ToastEntry, - StatePayload, - DaemonFrame, -} from "./protocol.ts"; diff --git a/src/infrastructure/ipc/protocol.ts b/src/infrastructure/ipc/protocol.ts deleted file mode 100644 index 5835b26..0000000 --- a/src/infrastructure/ipc/protocol.ts +++ /dev/null @@ -1,75 +0,0 @@ -/** - * Wire types for the Zesdex IPC protocol. - * Mirrors `apps/infrastructure/src/ipc/protocol.rs`. - */ - -/** A resolved key press sent from the daemon to the client. */ -export type KeyActionName = - | "Char" - | "Enter" - | "Escape" - | "Backspace" - | "Delete" - | "Tab" - | "Up" - | "Down" - | "Left" - | "Right" - | "Home" - | "End" - | "PageUp" - | "PageDown" - | "Function"; - -export interface KeyAction { - action: KeyActionName; - char?: string; - n?: number; -} - -/** A message sent from the TUI client to the daemon over the IPC socket. */ -export type ClientRequest = - | { kind: "Tick" } - | { kind: "KeyPress"; key: KeyAction; ctrl: boolean; alt: boolean; shift: boolean } - | { kind: "Submit"; text: string } - | { kind: "Paste"; text: string } - | { kind: "Resize"; cols: number; rows: number } - | { kind: "Close" } - | { kind: "ScrollUp" } - | { kind: "ScrollDown" }; - -/** A single chat message within a session. */ -export interface MessageEntry { - role: string; - content: string; - timestamp: number; -} - -/** A transient toast notification sent to the client. */ -export interface ToastEntry { - kind: string; - message: string; - created_at: number; - lifetime_ms: number; -} - -/** Full UI state snapshot pushed from the daemon to the client. */ -export interface StatePayload { - session_id: string; - messages: MessageEntry[]; - edit_count: number; - message_count: number; - overlay: string | null; - toasts: ToastEntry[]; - dirty: boolean; - input_buffer: string; - input_cursor: number; -} - -/** A frame sent from the daemon to the client. */ -export type DaemonFrame = - | { kind: "StateUpdate"; payload: StatePayload } - | { kind: "StreamToken"; token: string } - | { kind: "SystemNote"; noteKind: string; message: string } - | { kind: "ClipboardCopy"; text: string } - | { kind: "Closed" }; diff --git a/src/infrastructure/ipc/server.ts b/src/infrastructure/ipc/server.ts deleted file mode 100644 index bae2055..0000000 --- a/src/infrastructure/ipc/server.ts +++ /dev/null @@ -1,60 +0,0 @@ -/** - * IPC server — binds a Unix socket and accepts incoming client connections. - * Mirrors `apps/infrastructure/src/ipc/server.rs`. - */ -import * as net from "node:net"; -import * as fs from "node:fs"; -import { Connection } from "./conn.ts"; - -/** A Unix-socket IPC server. */ -export class IpcServer { - private server: net.Server; - private path: string; - - private constructor(path: string, server: net.Server) { - this.path = path; - this.server = server; - } - - /** Bind a Unix-socket server at `path`, removing any stale socket file. */ - static bindUnix(path: string): Promise { - return new Promise((resolve, reject) => { - if (fs.existsSync(path)) { - fs.unlinkSync(path); - } - const server = net.createServer(); - server.once("error", reject); - server.listen(path, () => { - resolve(new IpcServer(path, server)); - }); - }); - } - - /** Accept the next incoming connection as a `Connection`. */ - accept(): Promise { - return new Promise((resolve) => { - const onConnection = (socket: net.Socket) => { - serverCleanup(); - resolve(new Connection(socket)); - }; - const serverCleanup = () => { - this.server.removeListener("connection", onConnection); - }; - this.server.once("connection", onConnection); - }); - } - - /** Close the server and remove the socket file. */ - close(): Promise { - return new Promise((resolve) => { - this.server.close(() => { - try { - fs.unlinkSync(this.path); - } catch { - /* ignore */ - } - resolve(); - }); - }); - } -} diff --git a/src/infrastructure/persistence/iam/index.ts b/src/infrastructure/persistence/iam/index.ts deleted file mode 100644 index 0dd6e0b..0000000 --- a/src/infrastructure/persistence/iam/index.ts +++ /dev/null @@ -1,4 +0,0 @@ -/** IAM persistence — filesystem sessions, locks, tokens. */ -export * from "./session_repo.ts"; -export * from "./session_lock_repo.ts"; -export * from "./oauth_repo.ts"; \ No newline at end of file diff --git a/src/infrastructure/persistence/iam/oauth_repo.ts b/src/infrastructure/persistence/iam/oauth_repo.ts deleted file mode 100644 index 4d0237d..0000000 --- a/src/infrastructure/persistence/iam/oauth_repo.ts +++ /dev/null @@ -1,21 +0,0 @@ -/** - * Filesystem-backed `OAuthRepository`. Token written to JSON with mode 0o600. - */ -import * as fs from "node:fs"; -import * as path from "node:path"; -import { type OAuthRepository, type OAuthToken } from "@zesdex/domain"; -import { writeJsonAtomic } from "../../utils.ts"; - -/** Concrete filesystem OAuth token repository. */ -export class FileSystemOAuthRepository implements OAuthRepository { - async saveToken(file: string, token: OAuthToken): Promise { - fs.mkdirSync(path.dirname(file), { recursive: true }); - writeJsonAtomic(file, token, 0o600); - } - - async loadToken(file: string): Promise { - if (!fs.existsSync(file)) return null; - const data = fs.readFileSync(file, "utf8"); - return JSON.parse(data) as OAuthToken; - } -} \ No newline at end of file diff --git a/src/infrastructure/persistence/iam/session_lock_repo.ts b/src/infrastructure/persistence/iam/session_lock_repo.ts deleted file mode 100644 index b22954c..0000000 --- a/src/infrastructure/persistence/iam/session_lock_repo.ts +++ /dev/null @@ -1,78 +0,0 @@ -/** - * Filesystem-backed `SessionLockRepository` using a PID file with atomic - * `O_CREAT|O_EXCL` acquisition. Mirrors `session_lock_repo.rs`. - */ -import * as fs from "node:fs"; -import * as path from "node:path"; -import { type SessionLockRepository, pidIsAlive } from "@zesdex/domain"; - -/** Concrete filesystem session-lock repository. */ -export class FileSystemSessionLockRepository implements SessionLockRepository { - tryLock(sessionDir: string): boolean { - const file = path.join(sessionDir, ".lock"); - const pid = process.pid; - - // Phase 1: atomic create. - try { - const fd = fs.openSync(file, "wx"); - try { - fs.writeFileSync(fd, String(pid)); - fs.fsyncSync(fd); - } finally { - fs.closeSync(fd); - } - return true; - } catch (e) { - const code = (e as NodeJS.ErrnoException).code; - if (code !== "EEXIST") throw e; - } - - // Phase 2: liveness check. - const content = fs.readFileSync(file, "utf8").trim(); - const existing = Number(content); - if (Number.isFinite(existing) && existing > 0 && this.isAlive(existing)) { - return false; - } - - // Phase 3: stale lock overwrite. - const tmp = `${file}.tmp`; - let tmpFd: number; - try { - tmpFd = fs.openSync(tmp, "wx"); - } catch { - throw new Error("another process is replacing the lock"); - } - try { - fs.writeFileSync(tmpFd, String(pid)); - fs.fsyncSync(tmpFd); - } finally { - fs.closeSync(tmpFd); - } - fs.renameSync(tmp, file); - try { - const parent = path.dirname(file); - const dirFd = fs.openSync(parent, "r"); - try { - fs.fsyncSync(dirFd); - } finally { - fs.closeSync(dirFd); - } - } catch { - /* best-effort */ - } - return true; - } - - unlock(sessionDir: string): void { - const file = path.join(sessionDir, ".lock"); - try { - fs.unlinkSync(file); - } catch { - /* ignore */ - } - } - - isAlive(pid: number): boolean { - return pidIsAlive(pid); - } -} \ No newline at end of file diff --git a/src/infrastructure/persistence/iam/session_repo.ts b/src/infrastructure/persistence/iam/session_repo.ts deleted file mode 100644 index 572c97d..0000000 --- a/src/infrastructure/persistence/iam/session_repo.ts +++ /dev/null @@ -1,61 +0,0 @@ -/** - * Filesystem-backed `SessionRepository`. Each session is - * `/sessions//session.json`. - */ -import * as fs from "node:fs"; -import * as path from "node:path"; -import { - type Session, - type SessionId, - type SessionRepository, - newSessionId, - notFound, -} from "@zesdex/domain"; -import { writeJsonAtomic } from "../../utils.ts"; - -/** Concrete filesystem session repository. */ -export class FileSystemSessionRepository implements SessionRepository { - async listSessions(baseDir: string): Promise { - const dir = path.join(baseDir, "sessions"); - let entries: string[]; - try { - entries = fs.readdirSync(dir, { withFileTypes: true }) - .filter((e) => e.isDirectory()) - .map((e) => e.name); - } catch { - return []; - } - const sessions: Session[] = []; - for (const name of entries) { - const res = newSessionId(name); - if (!res.ok) continue; - try { - sessions.push(await this.loadSession(baseDir, res.value)); - } catch { - /* skip unreadable session */ - } - } - return sessions; - } - - async loadSession(baseDir: string, id: SessionId): Promise { - const file = path.join(baseDir, "sessions", id, "session.json"); - if (!fs.existsSync(file)) { - throw notFound(`session not found: ${id}`); - } - const data = fs.readFileSync(file, "utf8"); - return JSON.parse(data) as Session; - } - - async saveSession(baseDir: string, session: Session): Promise { - const dir = path.join(baseDir, "sessions", session.id); - fs.mkdirSync(dir, { recursive: true }); - const file = path.join(dir, "session.json"); - writeJsonAtomic(file, session); - } - - async deleteSession(baseDir: string, id: SessionId): Promise { - const dir = path.join(baseDir, "sessions", id); - if (fs.existsSync(dir)) fs.rmSync(dir, { recursive: true, force: true }); - } -} \ No newline at end of file diff --git a/src/infrastructure/persistence/index.ts b/src/infrastructure/persistence/index.ts index 94bf773..25223c7 100644 --- a/src/infrastructure/persistence/index.ts +++ b/src/infrastructure/persistence/index.ts @@ -1,3 +1,2 @@ -/** Persistence layer — CMS + IAM file repositories. */ -export * from "./cms/index.ts"; -export * from "./iam/index.ts"; \ No newline at end of file +/** Persistence layer — CMS file repositories. */ +export * from "./cms/index.ts"; \ No newline at end of file diff --git a/src/interfaces/api/dto.ts b/src/interfaces/api/dto.ts deleted file mode 100644 index 8825d10..0000000 --- a/src/interfaces/api/dto.ts +++ /dev/null @@ -1,120 +0,0 @@ -/** - * 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 }; -} diff --git a/src/interfaces/api/error.ts b/src/interfaces/api/error.ts deleted file mode 100644 index f3db4bd..0000000 --- a/src/interfaces/api/error.ts +++ /dev/null @@ -1,71 +0,0 @@ -/** - * 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 = { - 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"; - } -} diff --git a/src/interfaces/api/handlers/auth.ts b/src/interfaces/api/handlers/auth.ts deleted file mode 100644 index 101aaff..0000000 --- a/src/interfaces/api/handlers/auth.ts +++ /dev/null @@ -1,133 +0,0 @@ -/** - * 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 { - 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 { - 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 { - 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); -} diff --git a/src/interfaces/api/handlers/chat.ts b/src/interfaces/api/handlers/chat.ts deleted file mode 100644 index 703bacb..0000000 --- a/src/interfaces/api/handlers/chat.ts +++ /dev/null @@ -1,66 +0,0 @@ -/** - * 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 { - 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 }; -} diff --git a/src/interfaces/api/handlers/conversations.ts b/src/interfaces/api/handlers/conversations.ts deleted file mode 100644 index 2c97985..0000000 --- a/src/interfaces/api/handlers/conversations.ts +++ /dev/null @@ -1,110 +0,0 @@ -/** - * 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 { - 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 }; -} diff --git a/src/interfaces/api/handlers/sessions.ts b/src/interfaces/api/handlers/sessions.ts deleted file mode 100644 index 7d16d7c..0000000 --- a/src/interfaces/api/handlers/sessions.ts +++ /dev/null @@ -1,73 +0,0 @@ -/** - * 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 { - 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 }; -} diff --git a/src/interfaces/api/index.ts b/src/interfaces/api/index.ts deleted file mode 100644 index 2b9faa9..0000000 --- a/src/interfaces/api/index.ts +++ /dev/null @@ -1,10 +0,0 @@ -/** - * 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"; diff --git a/src/interfaces/api/main.ts b/src/interfaces/api/main.ts deleted file mode 100644 index 32f1b34..0000000 --- a/src/interfaces/api/main.ts +++ /dev/null @@ -1,21 +0,0 @@ -#!/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); diff --git a/src/interfaces/api/middleware/auth.ts b/src/interfaces/api/middleware/auth.ts deleted file mode 100644 index d34a967..0000000 --- a/src/interfaces/api/middleware/auth.ts +++ /dev/null @@ -1,51 +0,0 @@ -/** - * JWT authentication middleware for the API. - * Mirrors `apps/interfaces/api/src/middleware/auth.rs`. - * - * Validates the `Authorization: Bearer ` 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 ` 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}` }; - } -} diff --git a/src/interfaces/api/server.ts b/src/interfaces/api/server.ts deleted file mode 100644 index cbb3d10..0000000 --- a/src/interfaces/api/server.ts +++ /dev/null @@ -1,183 +0,0 @@ -/** - * 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 = { - "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; -} - -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 { - 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 }; -} diff --git a/src/interfaces/api/state.ts b/src/interfaces/api/state.ts deleted file mode 100644 index 2fbed63..0000000 --- a/src/interfaces/api/state.ts +++ /dev/null @@ -1,134 +0,0 @@ -/** - * 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(); - - /** 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 | null { - const p = path.join(baseDir, "users.json"); - try { - return JSON.parse(fs.readFileSync(p, "utf8")) as Record; - } catch { - return null; - } -} - -/** Persist the users map to `users.json` (pretty-printed). */ -export function saveUsers(baseDir: string, users: Record): void { - const p = path.join(baseDir, "users.json"); - fs.mkdirSync(path.dirname(p), { recursive: true }); - fs.writeFileSync(p, JSON.stringify(users, null, 2)); -} diff --git a/src/interfaces/cli/args.ts b/src/interfaces/cli/args.ts index a3cdd18..b29bc04 100644 --- a/src/interfaces/cli/args.ts +++ b/src/interfaces/cli/args.ts @@ -1,59 +1,32 @@ /** - * CLI argument parsing and mode dispatch. - * Mirrors `apps/gateway/src/main.rs` (clap → manual Bun.argv parsing). + * CLI argument parsing for Zesdex. * - * Modes (mutually exclusive): --daemon, --attach , --api, --ws, - * --grpc, --web; default is TUI single-process mode. + * Zesdex is now a pure CLI/TUI tool: `zesdex` boots the interactive TUI, and + * `--headless ` runs a single one-shot agent turn non-interactively. + * All server modes (REST API / WebSocket / gRPC / web / daemon) were removed. */ export interface CliOptions { - daemon: boolean; - attach: string | null; - api: boolean; - apiPort: number; - ws: boolean; - wsPort: number; - grpc: boolean; - grpcPort: number; - web: boolean; - webPort: number; + /** Non-interactive one-shot mode: run a single turn with the given prompt. */ + headless: string | null; + tos: boolean; } -const HELP = `zesdex — autonomous AI coding agent with TUI +const HELP = `zesdex — autonomous AI coding agent (TUI + headless CLI) -Usage: zesdex [OPTIONS] +Usage: + zesdex Launch the interactive TUI + zesdex --headless "" Run a single one-shot agent turn and exit Options: - --daemon Run as background daemon with IPC socket - --attach Attach TUI client to a running daemon session - --api Run REST API server - --api-port REST API port [default: 8080] - --ws Run WebSocket server - --ws-port WebSocket port [default: 8081] - --grpc Run gRPC server - --grpc-port gRPC port [default: 50051] - --web Serve web frontend - --web-port Web frontend port [default: 3000] - -h, --help Print help - -V, --version Print version + -h, --help Print help + -V, --version Print version `; /** Parse process argv (excluding node/bun) into CliOptions. */ export function parseCli(argv: string[]): CliOptions | { help: true; text: string } { - const opts: CliOptions = { - daemon: false, - attach: null, - api: false, - apiPort: 8080, - ws: false, - wsPort: 8081, - grpc: false, - grpcPort: 50051, - web: false, - webPort: 3000, - }; + const opts: CliOptions = { headless: null, tos: false }; - let sawFlag = false; for (let i = 0; i < argv.length; i++) { const arg = argv[i]!; if (arg === "--help" || arg === "-h") { @@ -62,71 +35,14 @@ export function parseCli(argv: string[]): CliOptions | { help: true; text: strin if (arg === "--version" || arg === "-V") { return { help: true, text: `zesdex 1.21.2` }; } - const assignString = (key: keyof CliOptions, next: string | undefined, name: string) => { - if (next === undefined) throw new Error(`flag '${name}' requires a value`); - (opts as unknown as Record)[key] = next; - }; - const assignPort = (key: keyof CliOptions, next: string | undefined, name: string) => { - if (next === undefined) throw new Error(`flag '${name}' requires a value`); - const n = Number(next); - if (!Number.isInteger(n) || n < 0 || n > 65535) { - throw new Error(`invalid port for '${name}': ${next}`); - } - (opts as unknown as Record)[key] = n; - }; - - switch (arg) { - case "--daemon": - opts.daemon = true; - sawFlag = true; - break; - case "--attach": - assignString("attach", argv[++i], "--attach"); - sawFlag = true; - break; - case "--api": - opts.api = true; - sawFlag = true; - break; - case "--api-port": - assignPort("apiPort", argv[++i], "--api-port"); - break; - case "--ws": - opts.ws = true; - sawFlag = true; - break; - case "--ws-port": - assignPort("wsPort", argv[++i], "--ws-port"); - break; - case "--grpc": - opts.grpc = true; - sawFlag = true; - break; - case "--grpc-port": - assignPort("grpcPort", argv[++i], "--grpc-port"); - break; - case "--web": - opts.web = true; - sawFlag = true; - break; - case "--web-port": - assignPort("webPort", argv[++i], "--web-port"); - break; - default: - throw new Error(`unexpected argument '${arg}'\n\n${HELP}`); + if (arg === "--headless") { + const next = argv[i + 1]; + if (next === undefined) throw new Error("flag '--headless' requires a value"); + opts.headless = next; + i++; + continue; } + throw new Error(`unexpected argument '${arg}'\n\n${HELP}`); } - void sawFlag; return opts; } - -/** Validate that at most one mode flag is set. */ -export function modeCount(opts: CliOptions): number { - let count = opts.daemon ? 1 : 0; - if (opts.api) count++; - if (opts.ws) count++; - if (opts.grpc) count++; - if (opts.web) count++; - if (opts.attach !== null) count++; - return count; -} diff --git a/src/interfaces/cli/index.ts b/src/interfaces/cli/index.ts index c5b095a..9e9ef35 100644 --- a/src/interfaces/cli/index.ts +++ b/src/interfaces/cli/index.ts @@ -1,17 +1,19 @@ #!/usr/bin/env bun /** * Zesdex CLI — main entry point. - * Mirrors `apps/gateway/src/main.rs`. * - * Parses argv, sets up file logging, and dispatches to the requested mode. - * Default mode runs the single-process REPL loop. + * Zesdex is a pure CLI/TUI tool: + * - `zesdex` → boots the interactive OpenTUI (primary interface) + * - `zesdex --headless "

"` → runs a single one-shot agent turn and exits + * + * All server modes (REST API / WebSocket / gRPC / web / daemon) were removed in + * the "full cli tui" rombak. The TUI owns the interactive agent runtime; this + * entry point only routes to it or runs a one-shot headless turn. */ import * as fs from "node:fs"; import * as os from "node:os"; import * as path from "node:path"; -import { newStore } from "@zesdex/domain"; -import { parseCli, modeCount } from "./args.ts"; -import { runSingleProcess } from "./compose.ts"; +import { parseCli } from "./args.ts"; /** Initialise file logging (append to /zesdex/zesdex.log). */ function initLogging(): void { @@ -22,8 +24,6 @@ function initLogging(): void { const stream = fs.createWriteStream(logPath, { flags: "a" }); const originalLog = console.log; const originalError = console.error; - const originalWarn = console.warn; - const originalInfo = console.info; const write = (level: string, args: unknown[]) => { const line = `[${new Date().toISOString()}] ${level} ${args @@ -31,8 +31,6 @@ function initLogging(): void { .join(" ")}`; stream.write(line + "\n"); if (level === "ERROR") originalError(...args); - else if (level === "WARN") originalWarn(...args); - else if (level === "INFO") originalInfo(...args); else originalLog(...args); }; @@ -42,66 +40,28 @@ function initLogging(): void { console.error = (...a: unknown[]) => write("ERROR", a); } -/** Run the REPL loop (default single-process mode). */ -async function runRepl(): Promise { +/** Run a single one-shot agent turn non-interactively, printing the outcome. */ +async function runHeadless(prompt: string): Promise { + const { runSingleProcess, buildTurnParams } = await import("./compose.ts"); const runtime = await runSingleProcess(); - console.info("zesdex ready — single-process REPL mode. Type `exit` to quit."); - // Minimal REPL over stdin. - // Use node readline. - const readline = (await import("node:readline")).default; - const rl = readline.createInterface({ input: process.stdin, output: process.stdout }); + const events: import("@zesdex/domain").TurnEvent[] = []; + const abort = new AbortController(); + const params = buildTurnParams(runtime, prompt, events, abort); - rl.on("line", async (line: string) => { - const text = line.trim(); - if (text === "exit" || text === "quit" || text === "/exit") { - rl.close(); - process.exit(0); - } - if (text === "" || text.startsWith("/")) { - if (text === "/help") { - console.log("Commands: type a message to talk to the agent; `exit` to quit."); - } - return; - } + await runtime.turnService.runTurn(params); - const events: import("@zesdex/domain").TurnEvent[] = []; - const abort = new AbortController(); - const params = buildTurnParamsFrom(runtime, text, events, abort); - - process.stdout.write("…"); - try { - await runtime.turnService.runTurn(params); - // Print the final assistant content from events. - const lastTokens: string[] = []; - for (const ev of events) { - if (ev.kind === "stream_token") lastTokens.push(ev.content); - } - const full = lastTokens.join(""); - process.stdout.write("\n"); - if (full.trim() !== "") console.log(full.trim()); - else console.log("(no response content)"); - } catch (e) { - console.error(`turn failed: ${(e as Error).message}`); - } - rl.prompt(); - }); - - rl.prompt(); + const tokens: string[] = []; + for (const ev of events) { + if (ev.kind === "stream_token") tokens.push(ev.content); + } + const full = tokens.join(""); + process.stdout.write("\n"); + if (full.trim() !== "") console.log(full.trim()); + else console.log("(no response content)"); } -// Imported lazily to avoid a circular type-only dependency. -import { buildTurnParams } from "./compose.ts"; -function buildTurnParamsFrom( - runtime: Parameters[0], - text: string, - events: import("@zesdex/domain").TurnEvent[], - abort: AbortController, -) { - return buildTurnParams(runtime, text, events, abort); -} - -/** Main entry. */ +/** Main entry. `zesdex` → TUI; `--headless ` → one-shot turn. */ async function main(): Promise { initLogging(); @@ -119,57 +79,15 @@ async function main(): Promise { process.exit(0); } - const opts = parsed; - if (modeCount(opts) > 1) { - console.error( - "Cannot specify multiple modes: --daemon, --attach, --api, --ws, --grpc, --web are mutually exclusive", - ); - process.exit(2); + if (parsed.headless !== null) { + await runHeadless(parsed.headless); + process.exit(0); } - if (opts.daemon) { - const { runDaemon } = await import("@zesdex/daemon"); - console.info("daemon mode — creating per-session IPC socket"); - const handle = await runDaemon(); - const stop = () => void handle.stop().then(() => process.exit(0)); - process.on("SIGINT", stop); - process.on("SIGTERM", stop); - await new Promise(() => {}); - } else if (opts.attach !== null) { - const { runAttach } = await import("@zesdex/daemon"); - console.info(`attach mode for daemon session '${opts.attach}'`); - const store = newStore(); - const socketPath = path.join(store.base_dir, "run", `${opts.attach}.sock`); - await runAttach(socketPath); - } else if (opts.api) { - const { newApiState, startApiServer } = await import("@zesdex/api"); - const baseDir = - process.env.ZESDEX_STORE ?? path.join(os.homedir(), ".local", "share", "zesdex"); - const state = newApiState( - baseDir, - process.env.ZESDEX_JWT_SECRET ?? "dev-secret-change-me", - process.env.OPENAI_API_KEY ?? "", - process.env.ZESDEX_MODEL ?? "claude-opus-5", - process.env.OPENAI_API_BASE ?? undefined, - ); - startApiServer(state, opts.apiPort); - await new Promise(() => {}); - } else if (opts.ws) { - const { startWsServer } = await import("@zesdex/ws"); - startWsServer(opts.wsPort); - await new Promise(() => {}); - } else if (opts.grpc) { - const { startGrpcServer } = await import("@zesdex/grpc"); - startGrpcServer(opts.grpcPort); - await new Promise(() => {}); - } else if (opts.web) { - const { startWebServer } = await import("@zesdex/web"); - startWebServer(opts.webPort, process.env.ZESDEX_WEB_DIR ?? undefined); - await new Promise(() => {}); - } else { - // Default: single-process REPL. - await runRepl(); - } + // Default: boot the interactive TUI. + const { runTui } = await import("@zesdex/tui"); + await runTui(); + process.exit(0); } main().catch((e) => { diff --git a/src/interfaces/daemon/client.ts b/src/interfaces/daemon/client.ts deleted file mode 100644 index 2855faf..0000000 --- a/src/interfaces/daemon/client.ts +++ /dev/null @@ -1,61 +0,0 @@ -/** - * Daemon client — connects to a running daemon over its Unix socket and - * exchanges framed IPC messages. Mirrors `apps/interfaces/daemon/src/client.rs`. - */ -import { IpcClient } from "@zesdex/infrastructure"; -import type { DaemonFrame, ClientRequest } from "@zesdex/infrastructure"; - -/** - * Connect to a running daemon socket and send one or more requests, - * printing the frames received back. Returns when the daemon closes or - * after a timeout. - */ -export async function runAttach( - socketPath: string, - initialText?: string, - timeoutMs = 15_000, -): Promise { - const client = await IpcClient.connectUnix(socketPath); - - // Send a submit if provided, so the daemon starts a turn. - if (initialText) { - const req: ClientRequest = { kind: "Submit", text: initialText }; - client.send(req); - } - - // Read frames until closed or timeout. - const deadline = Date.now() + timeoutMs; - while (Date.now() < deadline) { - const frame = await client.receive(); - if (frame === null) break; // daemon closed connection - switch (frame.kind) { - case "StreamToken": - process.stdout.write(frame.token); - break; - case "StateUpdate": - process.stdout.write(`\n[state: ${frame.payload.message_count} msgs]\n`); - break; - case "Closed": - process.stdout.write("\n[closed]\n"); - client.close(); - return; - default: - break; - } - } - - client.close(); -} - -/** Connect and send a Close request, for shutdown testing. */ -export async function sendClose(socketPath: string): Promise { - const client = await IpcClient.connectUnix(socketPath); - const req: ClientRequest = { kind: "Close" }; - client.send(req); - // Give the daemon a moment to reply Closed. - const frame = await client.receive(); - if (frame && frame.kind === "Closed") { - process.stdout.write("[closed]\n"); - } - client.close(); -} diff --git a/src/interfaces/daemon/index.ts b/src/interfaces/daemon/index.ts deleted file mode 100644 index a2cf072..0000000 --- a/src/interfaces/daemon/index.ts +++ /dev/null @@ -1,7 +0,0 @@ -/** - * Zesdex daemon interface package. - * Mirrors `apps/interfaces/daemon/`. - */ -export { runDaemon } from "./server.ts"; -export type { DaemonHandle } from "./server.ts"; -export { runAttach, sendClose } from "./client.ts"; diff --git a/src/interfaces/daemon/main.ts b/src/interfaces/daemon/main.ts deleted file mode 100644 index 2c8327a..0000000 --- a/src/interfaces/daemon/main.ts +++ /dev/null @@ -1,20 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex daemon — standalone entry point. - * Mirrors the `--daemon` mode of the Rust gateway. - */ -import { runDaemon } from "./index.ts"; - -const handle = await runDaemon(); -console.log("daemon running — Ctrl+C to stop"); - -// Keep alive until SIGINT/SIGTERM. -let stopping = false; -async function shutdown(): Promise { - if (stopping) return; - stopping = true; - await handle.stop(); - process.exit(0); -} -process.on("SIGINT", () => void shutdown()); -process.on("SIGTERM", () => void shutdown()); diff --git a/src/interfaces/daemon/server.ts b/src/interfaces/daemon/server.ts deleted file mode 100644 index ac01d4d..0000000 --- a/src/interfaces/daemon/server.ts +++ /dev/null @@ -1,249 +0,0 @@ -/** - * Daemon server — owns the agent state, listens on a per-session Unix socket, - * and drives requests from an attached client. - * Mirrors `apps/interfaces/daemon/src/server.rs` (simplified: focuses on the - * IPC agent-driver loop; full TUI key-handling lives in the Fase 4 TUI). - * - * Flow: `runDaemon()` creates a session → binds a Unix socket under - * `/run/.sock` → accepts a client → loops reading - * `ClientRequest`s. On `Submit` it runs an agent turn and streams - * `StreamToken`/`StateUpdate` frames back. On `Close`/disconnect it cleans - * up the socket file. - */ -import * as fs from "node:fs"; -import * as path from "node:path"; -import { randomUUID } from "node:crypto"; -import { IpcServer } from "@zesdex/infrastructure"; -import { - JsonSettingsRepository, - JsonAppConfigRepository, - LlmClient, - resolveApiKey, - InfrastructureToolExecutor, - allTools, - toolDefs, -} from "@zesdex/infrastructure"; -import type { ToolCtx } from "@zesdex/infrastructure"; -import { AgentTurnServiceImpl } from "@zesdex/application"; -import { - newStore, - ensureStoreDirs, - newSessionId, - resolveEffectiveModel, - type TurnEvent, - type TurnEventSink, - type AgentTurnParams, - userMessage, -} from "@zesdex/domain"; -import type { ClientRequest, DaemonFrame, MessageEntry, StatePayload } from "@zesdex/infrastructure"; - -/** A running daemon handle. */ -export interface DaemonHandle { - stop(): Promise; -} - -/** Array-backed turn event sink (push + drain). */ -export class ArraySink implements TurnEventSink { - constructor(public events: TurnEvent[] = []) {} - push(event: TurnEvent): void { - this.events.push(event); - } - drain(): TurnEvent[] { - const out = this.events; - this.events = []; - return out; - } -} - -/** Wire a fresh agent runtime (mirrors the CLI composition root). */ -async function buildRuntime() { - const store = newStore(); - await ensureStoreDirs(store); - - const settingsRepo = new JsonSettingsRepository(); - const appConfigRepo = new JsonAppConfigRepository(); - const settings = await settingsRepo.load(store.base_dir); - const appConfig = await appConfigRepo.load(store.base_dir); - - const provider = settings.provider; - const model = resolveEffectiveModel(settings, appConfig); - const baseUrl = appConfig.providers[provider]?.api_base ?? undefined; - const apiKey = resolveApiKey(settings, appConfig); - - const llmClient = new LlmClient(apiKey, model, baseUrl); - const toolCtx = { - sessionDir: store.base_dir, - workspaces: [process.cwd()], - turnEvents: new ArraySink(), - workflowFindings: [], - } as unknown as ToolCtx; - const executor = new InfrastructureToolExecutor(toolCtx); - const turnService = new AgentTurnServiceImpl(llmClient, executor, toolDefs(allTools()) as never); - - const buildParams = (message: string, sink: TurnEventSink): AgentTurnParams => ({ - messages: [userMessage(message)], - session_dir: store.base_dir, - workspace_roots: [process.cwd()], - turn_events: sink, - in_flight: { value: false }, - abort: new AbortController(), - api_key: apiKey, - model, - api_base: baseUrl, - }); - - return { store, turnService, buildParams }; -} - -/** Build a StatePayload snapshot from the accumulated messages. */ -function snapshot(sessionId: string, messages: MessageEntry[], dirty: boolean): StatePayload { - return { - session_id: sessionId, - messages, - edit_count: 0, - message_count: messages.length, - overlay: null, - toasts: [], - dirty, - input_buffer: "", - input_cursor: 0, - }; -} - -/** Handle one attached client: process its requests, driving agent turns. */ -async function handleDaemonClient( - conn: { send(m: T): void; receiveJson(): Promise; close(): void }, - sessionId: string, -): Promise { - const runtime = await buildRuntime(); - const messages: MessageEntry[] = []; - - const pushState = (dirty: boolean) => - conn.send({ - kind: "StateUpdate", - payload: snapshot(sessionId, messages, dirty), - }); - - pushState(false); - - let running = true; - while (running) { - let req: ClientRequest | null = null; - try { - req = await conn.receiveJson(); - } catch { - break; - } - if (req === null) break; // client closed - - switch (req.kind) { - case "Submit": { - const text = req.text; - if (text.trim() === "") break; - messages.push({ role: "user", content: text, timestamp: Date.now() }); - pushState(true); - - // Run the turn and stream tokens as they arrive. - const sink = new ArraySink(); - const params = runtime.buildParams(text, sink); - void runTurnAndStream(conn, sink, runtime.turnService, params).then(() => { - messages.push({ - role: "assistant", - content: sink.events - .filter((e) => e.kind === "stream_token") - .map((e) => (e as { content: string }).content) - .join(""), - timestamp: Date.now(), - }); - pushState(false); - }); - break; - } - case "Close": - running = false; - conn.send({ kind: "Closed" }); - break; - default: - // KeyPress/Tick/etc. no-op in this simplified daemon. - break; - } - } - conn.close(); -} - -/** Run a turn, streaming token frames to the client as they are produced. */ -async function runTurnAndStream( - conn: { send(m: T): void }, - sink: ArraySink, - turnService: AgentTurnServiceImpl, - params: AgentTurnParams, -): Promise { - const poll = setInterval(() => { - for (const ev of sink.drain()) { - if (ev.kind === "stream_token") { - conn.send({ kind: "StreamToken", token: (ev as { content: string }).content }); - } - } - }, 50); - - try { - await turnService.runTurn(params); - } finally { - clearInterval(poll); - // Drain any remaining tokens after the turn completes. - for (const ev of sink.drain()) { - if (ev.kind === "stream_token") { - conn.send({ kind: "StreamToken", token: (ev as { content: string }).content }); - } - } - } -} - -/** Run zesdex as a background daemon owning agent state over a Unix socket. */ -export async function runDaemon(): Promise { - const store = newStore(); - await ensureStoreDirs(store); - - const idRes = newSessionId(randomUUID()); - const sessionId = idRes.ok ? idRes.value : `sess-${Date.now()}`; - fs.mkdirSync(path.join(store.base_dir, "sessions", sessionId), { recursive: true }); - - const runDir = path.join(store.base_dir, "run"); - fs.mkdirSync(runDir, { recursive: true }); - const socketPath = path.join(runDir, `${sessionId}.sock`); - - const server = await IpcServer.bindUnix(socketPath); - console.log(`zesdex-daemon listening on ${socketPath}`); - - // Accept clients serially. - void (async () => { - for (;;) { - try { - const conn = await server.accept(); - console.log("daemon client connected"); - try { - await handleDaemonClient( - { send: (m) => conn.send(m), receiveJson: () => conn.receiveJson(), close: () => conn.close() }, - sessionId, - ); - } catch (e) { - console.error("daemon client error:", e); - } - console.log("daemon client disconnected"); - } catch { - break; // server closed - } - } - })(); - - return { - async stop() { - await server.close(); - try { - fs.unlinkSync(socketPath); - } catch { - /* ignore */ - } - }, - }; -} diff --git a/src/interfaces/grpc/index.ts b/src/interfaces/grpc/index.ts deleted file mode 100644 index 4099b5b..0000000 --- a/src/interfaces/grpc/index.ts +++ /dev/null @@ -1,6 +0,0 @@ -/** - * Zesdex gRPC interface package. - * Mirrors `apps/interfaces/grpc/`. - */ -export { startGrpcServer, handleGrpcRequest } from "./server.ts"; -export type { GrpcState } from "./server.ts"; diff --git a/src/interfaces/grpc/main.ts b/src/interfaces/grpc/main.ts deleted file mode 100644 index 87f1463..0000000 --- a/src/interfaces/grpc/main.ts +++ /dev/null @@ -1,9 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex gRPC server — standalone entry point (HTTP health stub). - * Mirrors the `--grpc` mode of the Rust gateway. - */ -import { startGrpcServer } from "./index.ts"; - -const port = Number.parseInt(process.env.ZESDEX_GRPC_PORT ?? "50051", 10); -startGrpcServer(port); diff --git a/src/interfaces/grpc/server.ts b/src/interfaces/grpc/server.ts deleted file mode 100644 index 6edff94..0000000 --- a/src/interfaces/grpc/server.ts +++ /dev/null @@ -1,41 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex gRPC interface — high-performance RPC with protobuf. - * Mirrors `apps/interfaces/grpc/src/lib.rs`. - * - * For now this crate provides a minimal HTTP health-check endpoint so - * consumers can verify the gRPC server is reachable. Full gRPC would use - * tonic+prost (Rust) — in Bun/TS a real gRPC server would use `@grpc/grpc-js` - * + generated protobuf types. This is intentionally a health stub. - */ - -/** gRPC server state (minimal for health checks). */ -export interface GrpcState { - version: string; -} - -/** Build the gRPC health router. */ -export function handleGrpcRequest(_state: GrpcState, req: Request): Response { - const url = new URL(req.url); - if (url.pathname === "/grpc/health") { - return new Response("gRPC server is running", { - status: 200, - headers: { "Content-Type": "text/plain; charset=utf-8" }, - }); - } - return new Response("Not found", { status: 404 }); -} - -/** Start the gRPC server (HTTP health stub). */ -export function startGrpcServer(port: number): { stop: () => void } { - const state: GrpcState = { version: "1.21.2" }; - const server = Bun.serve({ - port, - fetch(req) { - return handleGrpcRequest(state, req); - }, - }); - console.log(`zesdex-grpc listening on http://0.0.0.0:${server.port}`); - console.log("Note: gRPC currently runs HTTP health endpoint. Add @grpc/grpc-js for full gRPC."); - return { stop: () => server.stop(true) }; -} diff --git a/src/interfaces/tui/action.ts b/src/interfaces/tui/action.ts index 4dfe0ae..65c64b2 100644 --- a/src/interfaces/tui/action.ts +++ b/src/interfaces/tui/action.ts @@ -40,7 +40,6 @@ export type Action = | { tag: "LessonAccept"; name: string } | { tag: "LessonReject"; name: string } | { tag: "LessonDelete"; name: string } - | { tag: "StartOAuth"; provider: string } | { tag: "OpenEditor"; path: string } | { tag: "McpAdd"; name: string; command: string } | { tag: "ModelList" } @@ -154,16 +153,13 @@ export function applyAction(state: AppStateRest, action: Action): void { toastInfo(state, `Lesson deleted: ${action.name}`); break; - case "StartOAuth": - toastInfo(state, `OAuth login started for ${action.provider}`); - break; - case "McpAdd": toastInfo(state, `MCP server added: ${action.name} (${action.command})`); break; case "AbortTurn": state.abortFlag.value = true; + state.turnAbort?.abort(); toastInfo(state, "Aborting current turn..."); break; diff --git a/src/interfaces/tui/command.ts b/src/interfaces/tui/command.ts index 3c28a69..d0a95bb 100644 --- a/src/interfaces/tui/command.ts +++ b/src/interfaces/tui/command.ts @@ -13,7 +13,6 @@ export type Command = | { tag: "McpOpen" } | { tag: "Clear" } | { tag: "ClearConfirm" } - | { tag: "Login"; provider: string } | { tag: "Edit"; path: string } | { tag: "McpAdd"; name: string; command: string } | { tag: "ModelList" } @@ -41,8 +40,6 @@ export function parseCommand(text: string): Command { return { tag: "Quit" }; case "/clear": return arg1 === "" ? { tag: "ClearConfirm" } : { tag: "Clear" }; - case "/login": - return arg1 === "" ? { tag: "Login", provider: "" } : { tag: "Login", provider: arg1 }; case "/edit": return { tag: "Edit", path: arg1 === "" ? "." : arg1 }; case "/mcp": @@ -86,8 +83,6 @@ export function applyCommand(cmd: Command): Action[] { return [{ tag: "SystemNote", kind: "clear", message: "Transcript cleared." }]; case "ClearConfirm": return [{ tag: "OpenOverlay", overlay: "clear_confirm" }]; - case "Login": - return [{ tag: "StartOAuth", provider: cmd.provider }]; case "Edit": return [{ tag: "OpenEditor", path: cmd.path }]; case "McpAdd": diff --git a/src/interfaces/tui/main.ts b/src/interfaces/tui/main.ts deleted file mode 100644 index 130a03e..0000000 --- a/src/interfaces/tui/main.ts +++ /dev/null @@ -1,8 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex TUI — standalone entry point. - * Mirrors the default (TUI) mode of the Rust gateway. - */ -import { runTui } from "./run.ts"; - -void runTui(); diff --git a/src/interfaces/tui/run.ts b/src/interfaces/tui/run.ts index 2352d8f..51617b8 100644 --- a/src/interfaces/tui/run.ts +++ b/src/interfaces/tui/run.ts @@ -1,22 +1,65 @@ #!/usr/bin/env bun /** - * Zesdex TUI — OpenTUI entry point. + * Zesdex TUI — bootstrapper / composition root. * - * Boots the OpenTUI CLI renderer and mounts `ZesdexApp`, which routes input - * through `handleKey` → `applyAction`. Replaced the earlier raw-mode stdin - * event loop with OpenTUI's native renderer/keyboard handling. + * Wires the real agent runtime (via `runSingleProcess` from the CLI compose + * root) into the OpenTUI app, and installs a `turnRunner` that executes real + * turns against the LLM + tool executor, streaming events into the transcript. */ import * as fs from "node:fs"; -import { createTuiState } from "./state.ts"; +import { newStore, userMessage, type TurnEventSink } from "@zesdex/domain"; +import { runSingleProcess, type WiredRuntime } from "../cli/compose.ts"; +import { createTuiState, type AppStateRest, type TurnRunner } from "./state.ts"; import { bootTui } from "./ui.tsx"; -/** Run the interactive OpenTUI until the user quits. */ +/** Build the real turn runner from the wired runtime. */ +function makeTurnRunner(runtime: WiredRuntime): TurnRunner { + return async (text, onEvent) => { + const abort = new AbortController(); + // Live sink: forward every event to the UI immediately as it happens. + const sink: TurnEventSink = { + push(ev) { + onEvent(ev); + }, + drain() { + return []; + }, + }; + const messages = [userMessage(text)]; + await runtime.turnService.runTurn({ + messages, + session_dir: runtime.store.base_dir, + workspace_roots: [process.cwd()], + turn_events: sink, + in_flight: { value: false }, + abort, + api_key: runtime.apiKey, + model: runtime.model, + api_base: runtime.apiBase, + }); + }; +} + +/** Run the interactive OpenTUI with the real agent wired in. */ export async function runTui(): Promise { const cwd = process.cwd(); - const dataDir = `${process.env.HOME ?? ""}/.local/share/zesdex`; - const sessionDir = `${dataDir}/sessions/${Date.now()}`; - fs.mkdirSync(sessionDir, { recursive: true }); - const state = createTuiState([cwd], sessionDir, `${dataDir}/memory`); + // Build the real runtime once (LLM client + tool executor + turn service). + const runtime = await runSingleProcess(); + const store = newStore(); + void store; + const dataDir = `${process.env.HOME ?? ""}/.local/share/zesdex`; + fs.mkdirSync(dataDir, { recursive: true }); + const sessionLabel = new Date().toISOString().replace(/[:T]/g, "-").slice(0, 19); + const sessionDir = `${dataDir}/sessions/${sessionLabel}`; + fs.mkdirSync(sessionDir, { recursive: true }); + const memoryDir = `${dataDir}/memory`; + + const state: AppStateRest = createTuiState([cwd], sessionDir, memoryDir); + state.turnRunner = makeTurnRunner(runtime); + state.misc.apiConnected = true; + // Surface the real effective model/provider in the header. + state.settings = { provider: runtime.model } as Record; + await bootTui(state); } diff --git a/src/interfaces/tui/state.ts b/src/interfaces/tui/state.ts index 84fdefd..c2dbb97 100644 --- a/src/interfaces/tui/state.ts +++ b/src/interfaces/tui/state.ts @@ -7,7 +7,18 @@ * unit-testable in Bun. */ -import { Roles, type Role } from "@zesdex/domain"; +import { Roles, type Role, type TurnEvent } from "@zesdex/domain"; + +/** + * Driver for executing a single agent turn in the TUI. Implemented by the + * composition root (`run.ts` → `runSingleProcess`), keeping the TUI logic + * layer free of direct wiring. Returns a promise that resolves when the turn + * (including all tool calls) finishes or is aborted. + */ +export type TurnRunner = ( + text: string, + onEvent: (ev: TurnEvent) => void, +) => Promise; /* ── Input state ──────────────────────────────────────────────────── */ @@ -23,9 +34,6 @@ export const COMMANDS: string[] = [ "/help", "/quit", "/clear", - "/login", - "/login zen", - "/login openai", "/edit", "/mcp add", "/model", @@ -390,6 +398,10 @@ export interface AppStateRest { renderWidthAtCache: number; mentionFiles: string[]; pendingSubmit: string | null; // set when a turn should be spawned + /** Real agent turn driver (wired by the composition root). */ + turnRunner: TurnRunner | null; + /** Abort controller for the in-flight turn (driven by AbortTurn). */ + turnAbort: AbortController | null; } export const DEFAULT_HELP_TEXT = ` Zesdex TUI — Keyboard Shortcuts @@ -459,6 +471,8 @@ export function createTuiState( renderWidthAtCache: 0, mentionFiles: [], pendingSubmit: null, + turnRunner: null, + turnAbort: null, }; } diff --git a/src/interfaces/tui/ui.tsx b/src/interfaces/tui/ui.tsx index 623540e..22d8c8e 100644 --- a/src/interfaces/tui/ui.tsx +++ b/src/interfaces/tui/ui.tsx @@ -119,27 +119,78 @@ export function ZesdexApp(props: { state: AppStateRest; onQuit: () => void }): R const dimensions = useTerminalDimensionsHook(); - // Local-only "turn drain": when a turn is submitted without a live LLM, - // simulate a short thinking delay then append a note. Real wiring happens - // in the CLI. Replaced the old raw-loop poll. + // Turn driver: when a turn is submitted, run it against the real agent and + // stream every event into the transcript as it happens. useEffect(() => { - const timer = setInterval(() => { - if (state.pendingSubmit != null) { - const text = state.pendingSubmit; + let turnRunning = false; + const poll = async () => { + if (turnRunning) return; + const text = state.pendingSubmit; + if (text == null) return; + if (state.turnRunner == null) { + appendToLastTranscript( + state, + "\n(no agent runtime — run via `zesdex` on the CLI / ensure settings)", + false, + ); state.pendingSubmit = null; - state.turnInFlightFlag.value = true; - setTimeout(() => { - state.turnInFlightFlag.value = false; - appendToLastTranscript(state, `(no LLM configured — processed: "${text}")`, false); - setTick((x: number) => x + 1); - }, 50); + return; } - setTick((x: number) => x + 1); - }, 120); + turnRunning = true; + const captured = state.pendingSubmit; + state.pendingSubmit = null; + state.turnInFlightFlag.value = true; + const abort = new AbortController(); + state.turnAbort = abort; + + // Open an assistant message we stream into. + pushTranscript(state, makeChatMessage(Roles.Assistant, "")); + + try { + await state.turnRunner(captured ?? text, (ev) => { + switch (ev.kind) { + case "stream_token": + appendToLastTranscript(state, ev.content, false); + break; + case "stream_reasoning": + appendToLastTranscript(state, ev.content, true); + break; + case "system_note": + pushTranscript( + state, + makeChatMessage(Roles.System, `[${ev.systemKind}] ${ev.message}`), + ); + break; + case "error": + pushTranscript(state, makeChatMessage(Roles.System, `error: ${ev.message}`)); + break; + case "tool_result": + pushTranscript( + state, + makeChatMessage(Roles.Tool, `⇄ ${ev.tool_name} · ${ev.is_error ? "err" : "ok"}`), + ); + break; + default: + break; + } + }); + } finally { + state.turnInFlightFlag.value = false; + state.turnAbort = null; + state.abortFlag.value = false; + // Re-prime the next poll so queued turns run. + setTimeout(() => { + turnRunning = false; + setTick((x: number) => x + 1); + }, 0); + } + }; + const timer = setInterval(() => poll(), 120); return () => clearInterval(timer); }, [state, setTick]); - const header = ` zesdex · ${state.settings.provider ?? "?"} `; + const model = String(state.settings.provider ?? "?"); + const header = ` zesdex · ${model} `; const messages = state.transcriptCache.messages; const scrollFromEnd = messages.length - state.scroll.offset; const visible = messages.slice(Math.max(0, scrollFromEnd - 40), scrollFromEnd); @@ -159,6 +210,19 @@ export function ZesdexApp(props: { state: AppStateRest; onQuit: () => void }): R {oneLine(m.content)} ))} + {/* Reasoning block for the last assistant message (dimmed). */} + {messages.length > 0 && + (() => { + const last = messages[messages.length - 1]; + if (last && last.role === Roles.Assistant && last.reasoning.length > 0) { + return ( + + {`reasoning ${oneLine(last.reasoning)}`} + + ); + } + return null; + })()} {/* Overlay or input area */} diff --git a/src/interfaces/web/index.ts b/src/interfaces/web/index.ts deleted file mode 100644 index f01cc2d..0000000 --- a/src/interfaces/web/index.ts +++ /dev/null @@ -1,6 +0,0 @@ -/** - * Zesdex web frontend interface package. - * Mirrors `apps/interfaces/web/`. - */ -export { startWebServer, handleWebRequest } from "./server.ts"; -export type { WebState } from "./server.ts"; diff --git a/src/interfaces/web/main.ts b/src/interfaces/web/main.ts deleted file mode 100644 index acdcda6..0000000 --- a/src/interfaces/web/main.ts +++ /dev/null @@ -1,10 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex web server — standalone entry point. - * Mirrors the `--web` mode of the Rust gateway. - */ -import { startWebServer } from "./index.ts"; - -const port = Number.parseInt(process.env.ZESDEX_WEB_PORT ?? "3000", 10); -const staticDir = process.env.ZESDEX_WEB_DIR ?? undefined; -startWebServer(port, staticDir); diff --git a/src/interfaces/web/server.ts b/src/interfaces/web/server.ts deleted file mode 100644 index b08c51d..0000000 --- a/src/interfaces/web/server.ts +++ /dev/null @@ -1,121 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex web frontend interface — serves static browser assets. - * Mirrors `apps/interfaces/web/src/lib.rs`. - * - * Serves files from a `dist/`/static directory with path-traversal protection: - * every requested path is canonicalized and must resolve inside `static_dir`. - * Falls back to a minimal placeholder page when no frontend is built. - */ -import * as fs from "node:fs"; -import * as path from "node:path"; - -/** Web server state. */ -export interface WebState { - static_dir: string; -} - -const MIME: Record = { - ".html": "text/html; charset=utf-8", - ".js": "text/javascript; charset=utf-8", - ".mjs": "text/javascript; charset=utf-8", - ".css": "text/css; charset=utf-8", - ".json": "application/json", - ".png": "image/png", - ".jpg": "image/jpeg", - ".jpeg": "image/jpeg", - ".gif": "image/gif", - ".svg": "image/svg+xml", - ".ico": "image/x-icon", - ".webp": "image/webp", - ".woff": "font/woff", - ".woff2": "font/woff2", - ".ttf": "font/ttf", - ".map": "application/json", - ".txt": "text/plain; charset=utf-8", -}; - -function mimeFor(file: string): string { - return MIME[path.extname(file).toLowerCase()] ?? "application/octet-stream"; -} - -/** Placeholder page when no frontend is built. */ -const PLACEHOLDER = ` -Zesdex Web - - - -

Zesdex Web

-

Web interface is ready.

-

To connect the frontend:

-
    -
  1. Build the frontend: cd apps/interfaces/web && npm install && npm run build
  2. -
  3. Restart with --web-dir apps/interfaces/web/dist
  4. -
-`; - -/** Resolve a requested path safely inside `staticDir`; returns null on traversal/not-found. */ -function safeResolve(staticDir: string, rel: string): string | null { - // Strip any leading slash so a relative path is resolved against staticDir - // (path.resolve would otherwise ignore staticDir for absolute-ish inputs). - const clean = rel.replace(/^\/+/, ""); - const candidate = path.resolve(staticDir, clean || "index.html"); - const canonStatic = path.resolve(staticDir); - // Prevent directory traversal: resolved path must stay inside static dir. - if (!candidate.startsWith(canonStatic + path.sep) && candidate !== canonStatic) { - return null; - } - // Only serve regular files; reject if the target is a directory or missing. - try { - const st = fs.statSync(candidate); - return st.isFile() ? candidate : null; - } catch { - return null; - } -} - -/** Build the web router (static file serving). */ -export function handleWebRequest(state: WebState, req: Request): Response { - const url = new URL(req.url); - let rel = decodeURIComponent(url.pathname); - if (rel.endsWith("/")) rel += "index.html"; - - const file = safeResolve(state.static_dir, rel); - if (!file) { - // Root always returns the placeholder index (even without a built frontend). - if (rel === "/index.html") { - return new Response(PLACEHOLDER, { - status: 200, - headers: { "Content-Type": "text/html; charset=utf-8" }, - }); - } - return new Response("Not found", { status: 404 }); - } - - const data = fs.readFileSync(file); - return new Response(data, { - status: 200, - headers: { "Content-Type": mimeFor(file) }, - }); -} - -/** Start the web frontend server. */ -export function startWebServer(port: number, staticDir?: string): { stop: () => void } { - const dir = - staticDir && staticDir !== "" - ? staticDir - : (() => { - const p = path.resolve("apps/interfaces/web/dist"); - return fs.existsSync(p) ? p : process.cwd(); - })(); - const state: WebState = { static_dir: dir }; - const server = Bun.serve({ - port, - fetch(req) { - return handleWebRequest(state, req); - }, - }); - console.log(`zesdex-web listening on http://0.0.0.0:${server.port} (dir: ${state.static_dir})`); - return { stop: () => server.stop(true) }; -} diff --git a/src/interfaces/ws/index.ts b/src/interfaces/ws/index.ts deleted file mode 100644 index efd0d57..0000000 --- a/src/interfaces/ws/index.ts +++ /dev/null @@ -1,6 +0,0 @@ -/** - * Zesdex WebSocket interface package. - * Mirrors `apps/interfaces/ws/`. - */ -export { startWsServer } from "./server.ts"; -export type { WsState, SocketSend } from "./server.ts"; diff --git a/src/interfaces/ws/main.ts b/src/interfaces/ws/main.ts deleted file mode 100644 index 481ffa5..0000000 --- a/src/interfaces/ws/main.ts +++ /dev/null @@ -1,9 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex WS server — standalone entry point. - * Mirrors the `--ws` mode of the Rust gateway. - */ -import { startWsServer } from "./index.ts"; - -const port = Number.parseInt(process.env.ZESDEX_WS_PORT ?? "8081", 10); -startWsServer(port); diff --git a/src/interfaces/ws/server.ts b/src/interfaces/ws/server.ts deleted file mode 100644 index 278961d..0000000 --- a/src/interfaces/ws/server.ts +++ /dev/null @@ -1,220 +0,0 @@ -#!/usr/bin/env bun -/** - * Zesdex WebSocket interface — real-time bidirectional communication. - * Mirrors `apps/interfaces/ws/src/lib.rs`. - * - * Security: accepts an optional `?token=` query param. When `ZESDEX_WS_TOKEN` - * env is set, connections MUST present a matching token, else rejected — - * prevents the endpoint from being used as an open LLM proxy. - * - * Protocol (JSON text frames): - * client → { "type": "prompt", "message": "...", "model"?: "..." } - * server → { "type": "connected", "session": ..., "message": "..." } (on connect) - * server → { "type": "token", "content": "..." } (streaming) - * server → { "type": "done" } | { "type": "error", "message": "..." } (terminal) - * server → { "type": "echo", "data": "..." } (fallback) - */ -import { - JsonSettingsRepository, - JsonAppConfigRepository, - LlmClient, - resolveApiKey, - InfrastructureToolExecutor, - allTools, - toolDefs, -} from "@zesdex/infrastructure"; -import type { ToolCtx } from "@zesdex/infrastructure"; -import { AgentTurnServiceImpl } from "@zesdex/application"; -import { - newStore, - ensureStoreDirs, - type TurnEvent, - type TurnEventSink, - type AgentTurnParams, - userMessage, - resolveEffectiveModel, -} from "@zesdex/domain"; - -/** Shared application state for the WS server. */ -export interface WsState { - store_base_dir: string; - session_id: string | null; -} - -/** Minimal send interface used by streaming helpers. */ -export interface SocketSend { - send(json: string): void; -} - -/** Array-backed turn event sink (push + drain). */ -class ArraySink implements TurnEventSink { - constructor(public events: TurnEvent[] = []) {} - push(event: TurnEvent): void { - this.events.push(event); - } - drain(): TurnEvent[] { - const out = this.events; - this.events = []; - return out; - } -} - -/** Validate the token against ZESDEX_WS_TOKEN (when configured). */ -function tokenAllowed(url: URL): boolean { - const configured = process.env.ZESDEX_WS_TOKEN; - const t = url.searchParams.get("token"); - if (configured) { - return t === configured; - } - return true; -} - -/** Wire a turn service for one prompt (matches the CLI composition root). */ -async function buildSingleShot( - sessionDir: string, - eventSink: TurnEventSink, -): Promise<{ turnService: AgentTurnServiceImpl; params: AgentTurnParams }> { - const store = newStore(); - await ensureStoreDirs(store); - - const settingsRepo = new JsonSettingsRepository(); - const appConfigRepo = new JsonAppConfigRepository(); - const settings = await settingsRepo.load(store.base_dir); - const appConfig = await appConfigRepo.load(store.base_dir); - - const provider = settings.provider; - const model = resolveEffectiveModel(settings, appConfig); - const baseUrl = appConfig.providers[provider]?.api_base ?? undefined; - const apiKey = resolveApiKey(settings, appConfig); - - const llmClient = new LlmClient(apiKey, model, baseUrl); - const toolCtx = { - sessionDir, - workspaces: [sessionDir], - turnEvents: eventSink, - workflowFindings: [], - } as unknown as ToolCtx; - const executor = new InfrastructureToolExecutor(toolCtx); - const turnService = new AgentTurnServiceImpl(llmClient, executor, toolDefs(allTools()) as never); - - const params: AgentTurnParams = { - messages: [userMessage("")], - session_dir: sessionDir, - workspace_roots: [sessionDir], - turn_events: eventSink, - in_flight: { value: false }, - abort: new AbortController(), - api_key: apiKey, - model, - api_base: baseUrl, - }; - - return { turnService, params }; -} - -/** - * Start the WS server. `onPrompt` is invoked with the client's prompt payload, - * and `forward` sends raw strings to the client. - */ -export function startWsServer(port: number): { stop: () => void } { - const state: WsState = { store_base_dir: ".", session_id: null }; - - const server = Bun.serve({ - port, - fetch(req, server) { - const url = new URL(req.url); - if (url.pathname === "/ws") { - if (!tokenAllowed(url)) { - return new Response("missing or invalid token", { status: 401 }); - } - if (server.upgrade(req, {})) { - return undefined; - } - return new Response("upgrade failed", { status: 400 }); - } - return new Response("Not found", { status: 404 }); - }, - websocket: { - open(ws) { - const sender: SocketSend = { send: (s) => ws.send(s) }; - // Welcome message on connect. - sender.send( - JSON.stringify({ - type: "connected", - session: state.session_id, - message: "Connected to Zesdex WebSocket server", - }), - ); - }, - message(ws, msg) { - const text = typeof msg === "string" ? msg : JSON.stringify(msg); - if (!text.startsWith("{")) { - ws.send(JSON.stringify({ type: "echo", data: text })); - return; - } - let val: Record; - try { - val = JSON.parse(text) as Record; - } catch { - ws.send(JSON.stringify({ type: "echo", data: text })); - return; - } - if (val.type === "prompt" && typeof val.message === "string") { - const sender: SocketSend = { send: (s) => ws.send(s) }; - void runPrompt(sender, val.message, typeof val.model === "string" ? val.model : undefined); - } else { - ws.send(JSON.stringify({ type: "echo", data: text })); - } - }, - close() {}, - }, - }); - console.log(`zesdex-ws listening on ws://0.0.0.0:${server.port}`); - return { stop: () => server.stop(true) }; -} - -/** Run a single prompt and stream events to the client. */ -async function runPrompt(ws: SocketSend, message: string, modelOverride?: string): Promise { - const sessionDir = process.cwd(); - const sink = new ArraySink(); - try { - const { turnService, params } = await buildSingleShot(sessionDir, sink); - if (modelOverride) params.model = modelOverride; - params.messages = [userMessage(message)]; - - // Poll the sink and forward events to the client. - const poll = setInterval(() => { - const evs = sink.drain(); - for (const ev of evs) { - if (ev.kind === "stream_token") { - ws.send(JSON.stringify({ type: "token", content: ev.content })); - } else if (ev.kind === "done") { - ws.send(JSON.stringify({ type: "done" })); - clearInterval(poll); - } else if (ev.kind === "error") { - ws.send(JSON.stringify({ type: "error", message: ev.message })); - clearInterval(poll); - } - } - }, 50); - - try { - await turnService.runTurn(params); - // Ensure a terminal done/error is sent even if none auto-emitted - // while the poll stopped early. - const remaining = sink.drain(); - for (const ev of remaining) { - if (ev.kind === "stream_token") ws.send(JSON.stringify({ type: "token", content: ev.content })); - else if (ev.kind === "done") ws.send(JSON.stringify({ type: "done" })); - else if (ev.kind === "error") ws.send(JSON.stringify({ type: "error", message: ev.message })); - } - if (!remaining.some((e) => e.kind === "done" || e.kind === "error")) { - ws.send(JSON.stringify({ type: "done" })); - } - } finally { - clearInterval(poll); - } - } catch (e) { - ws.send(JSON.stringify({ type: "error", message: (e as Error).message })); - } -} diff --git a/tsconfig.json b/tsconfig.json index bde085e..d636d0b 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -25,11 +25,6 @@ "@zesdex/domain": ["./src/domain/index.ts"], "@zesdex/application": ["./src/application/index.ts"], "@zesdex/infrastructure": ["./src/infrastructure/index.ts"], - "@zesdex/api": ["./src/interfaces/api/index.ts"], - "@zesdex/ws": ["./src/interfaces/ws/index.ts"], - "@zesdex/grpc": ["./src/interfaces/grpc/index.ts"], - "@zesdex/web": ["./src/interfaces/web/index.ts"], - "@zesdex/daemon": ["./src/interfaces/daemon/index.ts"], "@zesdex/tui": ["./src/interfaces/tui/index.ts"] } },