refactor: rombak jadi full CLI/TUI only — hapus semua interface server (api/ws/grpc/web/daemon) + wire TUI ke agent beneran
- 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 '<prompt>' 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'
This commit is contained in:
@@ -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 "<prompt>"` — 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 "<prompt>"` — one-shot agent turn.
|
||||
- `bun run bootstrap` — seed default settings/app_config.
|
||||
|
||||
+2
-4
@@ -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:
|
||||
|
||||
+4
-8
@@ -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",
|
||||
|
||||
@@ -1,3 +0,0 @@
|
||||
/** Auth application module — OAuth PKCE + session management use-cases. */
|
||||
export * from "./oauth_service.ts";
|
||||
export * from "./session_service.ts";
|
||||
@@ -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<void> {
|
||||
this.verifier = v;
|
||||
this.state = s;
|
||||
}
|
||||
async loadVerifier(): Promise<string> {
|
||||
return this.verifier;
|
||||
}
|
||||
async loadState(): Promise<string> {
|
||||
return this.state;
|
||||
}
|
||||
async clear(): Promise<void> {
|
||||
this.verifier = "";
|
||||
this.state = "";
|
||||
}
|
||||
}
|
||||
|
||||
function makeRepo() {
|
||||
let token: OAuthToken | null = null;
|
||||
return {
|
||||
repo: {
|
||||
async saveToken(_path: string, t: OAuthToken): Promise<void> {
|
||||
token = t;
|
||||
},
|
||||
async loadToken(): Promise<OAuthToken | null> {
|
||||
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<OAuthToken> {
|
||||
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<OAuthToken> {
|
||||
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<OAuthToken> {
|
||||
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");
|
||||
});
|
||||
});
|
||||
@@ -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<void>;
|
||||
loadVerifier(): Promise<string>;
|
||||
loadState(): Promise<string>;
|
||||
clear(): Promise<void>;
|
||||
}
|
||||
|
||||
/** 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<OAuthToken>;
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* 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<OAuthToken> {
|
||||
// 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<OAuthToken | null> {
|
||||
return this.tokenRepo.loadToken(this.tokenPath);
|
||||
}
|
||||
}
|
||||
@@ -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<R, L>`
|
||||
private lockRepo: SessionLockRepository,
|
||||
private baseDir: string,
|
||||
) {}
|
||||
|
||||
async createSession(title: string): Promise<Session> {
|
||||
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<Session[]> {
|
||||
return this.sessionRepo.listSessions(this.baseDir);
|
||||
}
|
||||
|
||||
async archiveSession(id: SessionId): Promise<void> {
|
||||
const session = await this.sessionRepo.loadSession(this.baseDir, id);
|
||||
session.archived = true;
|
||||
session.updated_at = Date.now();
|
||||
await this.sessionRepo.saveSession(this.baseDir, session);
|
||||
}
|
||||
}
|
||||
@@ -5,5 +5,4 @@
|
||||
*/
|
||||
export * from "./ports/index.ts";
|
||||
export * from "./agent/index.ts";
|
||||
export * from "./auth/index.ts";
|
||||
export * from "./cms/index.ts";
|
||||
@@ -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 };
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -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";
|
||||
@@ -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],
|
||||
};
|
||||
}
|
||||
@@ -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 `<base_dir>/sessions/`. */
|
||||
listSessions(baseDir: string): Promise<Session[]> | Session[];
|
||||
/** Load a single session by id. */
|
||||
loadSession(baseDir: string, id: SessionId): Promise<Session> | Session;
|
||||
/** Save a session's metadata to disk. */
|
||||
saveSession(baseDir: string, session: Session): Promise<void> | void;
|
||||
/** Delete a session directory and all its contents. */
|
||||
deleteSession(baseDir: string, id: SessionId): Promise<void> | 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> | boolean;
|
||||
/** Release the lock. */
|
||||
unlock(sessionDir: string): Promise<void> | 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> | void;
|
||||
/** Load an OAuth token, returning `null` if the file does not exist. */
|
||||
loadToken(path: string): Promise<OAuthToken | null> | OAuthToken | 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> | Session;
|
||||
/** List all available sessions. */
|
||||
listAll(): Promise<Session[]> | Session[];
|
||||
/** Archive a session by id (sets `archived = true`). */
|
||||
archiveSession(id: SessionId): Promise<void> | 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> | OAuthToken;
|
||||
/** Retrieve the currently stored OAuth token (if any). */
|
||||
getToken(): Promise<OAuthToken | null> | OAuthToken | 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 `<base_dir>/sessions/<id>`. */
|
||||
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");
|
||||
}
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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/<pid>/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 (`<session_dir>/.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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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";
|
||||
|
||||
@@ -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";
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -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:<ephemeral>`, 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<void> {
|
||||
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<string> {
|
||||
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<void> {
|
||||
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"));
|
||||
}
|
||||
@@ -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<string> {
|
||||
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<boolean> {
|
||||
return verify(hashStr, password);
|
||||
}
|
||||
}
|
||||
@@ -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";
|
||||
|
||||
@@ -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<IpcClient> {
|
||||
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<T>(msg: T): void {
|
||||
this.conn.send(msg);
|
||||
}
|
||||
|
||||
/** Receive the next message and JSON-decode it (null on close). */
|
||||
async receive<T>(): Promise<T | null> {
|
||||
return this.conn.receiveJson<T>();
|
||||
}
|
||||
|
||||
/** Close the connection. */
|
||||
close(): void {
|
||||
this.conn.close();
|
||||
}
|
||||
}
|
||||
@@ -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<T>(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<Buffer | null> {
|
||||
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<T>(): Promise<T | null> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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]);
|
||||
}
|
||||
@@ -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";
|
||||
@@ -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" };
|
||||
@@ -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<IpcServer> {
|
||||
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<Connection> {
|
||||
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<void> {
|
||||
return new Promise((resolve) => {
|
||||
this.server.close(() => {
|
||||
try {
|
||||
fs.unlinkSync(this.path);
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -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";
|
||||
@@ -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<void> {
|
||||
fs.mkdirSync(path.dirname(file), { recursive: true });
|
||||
writeJsonAtomic(file, token, 0o600);
|
||||
}
|
||||
|
||||
async loadToken(file: string): Promise<OAuthToken | null> {
|
||||
if (!fs.existsSync(file)) return null;
|
||||
const data = fs.readFileSync(file, "utf8");
|
||||
return JSON.parse(data) as OAuthToken;
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -1,61 +0,0 @@
|
||||
/**
|
||||
* Filesystem-backed `SessionRepository`. Each session is
|
||||
* `<base_dir>/sessions/<id>/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<Session[]> {
|
||||
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<Session> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
const dir = path.join(baseDir, "sessions", id);
|
||||
if (fs.existsSync(dir)) fs.rmSync(dir, { recursive: true, force: true });
|
||||
}
|
||||
}
|
||||
@@ -1,3 +1,2 @@
|
||||
/** Persistence layer — CMS + IAM file repositories. */
|
||||
/** Persistence layer — CMS file repositories. */
|
||||
export * from "./cms/index.ts";
|
||||
export * from "./iam/index.ts";
|
||||
@@ -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 };
|
||||
}
|
||||
@@ -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<ApiErrorKind, number> = {
|
||||
BadRequest: 400,
|
||||
Unauthorized: 401,
|
||||
NotFound: 404,
|
||||
Conflict: 409,
|
||||
TooManyRequests: 429,
|
||||
Internal: 500,
|
||||
ChatProxy: 502,
|
||||
};
|
||||
|
||||
export class ApiError extends Error {
|
||||
constructor(public kind: ApiErrorKind, message: string) {
|
||||
super(message);
|
||||
this.name = "ApiError";
|
||||
}
|
||||
|
||||
get status(): number {
|
||||
return HTTP_STATUS[this.kind];
|
||||
}
|
||||
|
||||
static badRequest(msg: string): ApiError {
|
||||
return new ApiError("BadRequest", msg);
|
||||
}
|
||||
static unauthorized(msg: string): ApiError {
|
||||
return new ApiError("Unauthorized", msg);
|
||||
}
|
||||
static notFound(msg: string): ApiError {
|
||||
return new ApiError("NotFound", msg);
|
||||
}
|
||||
static conflict(msg: string): ApiError {
|
||||
return new ApiError("Conflict", msg);
|
||||
}
|
||||
static tooManyRequests(msg: string): ApiError {
|
||||
return new ApiError("TooManyRequests", msg);
|
||||
}
|
||||
static internal(msg: string): ApiError {
|
||||
return new ApiError("Internal", msg);
|
||||
}
|
||||
static chatProxy(msg: string): ApiError {
|
||||
return new ApiError("ChatProxy", msg);
|
||||
}
|
||||
|
||||
/** Whether an internal error's details should be hidden from the client. */
|
||||
get exposeDetails(): boolean {
|
||||
// Internal errors hide details; all others surface the message.
|
||||
return this.kind !== "Internal";
|
||||
}
|
||||
}
|
||||
@@ -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<AuthResponse> {
|
||||
enforceRateLimit(state.auth_rate_limiter, headers);
|
||||
requireFields(body, false);
|
||||
|
||||
const users = loadUsers(state.store_base_dir);
|
||||
if (!users) {
|
||||
throw ApiError.unauthorized("Invalid username or password");
|
||||
}
|
||||
const storedHash = users[body.username];
|
||||
if (!storedHash) {
|
||||
throw ApiError.unauthorized("Invalid username or password");
|
||||
}
|
||||
|
||||
const valid = await state.password_service.verify(body.password, storedHash);
|
||||
if (!valid) {
|
||||
throw ApiError.unauthorized("Invalid username or password");
|
||||
}
|
||||
|
||||
return authResponse(state, body.username);
|
||||
}
|
||||
|
||||
/** POST /auth/register — create a new user account and issue a JWT pair. */
|
||||
export async function registerHandler(
|
||||
state: ApiState,
|
||||
headers: Headers,
|
||||
body: RegisterRequest,
|
||||
): Promise<AuthResponse> {
|
||||
enforceRateLimit(state.auth_rate_limiter, headers);
|
||||
requireFields(body, true);
|
||||
|
||||
// Hash first (async), then serialize the read-modify-write of users.json.
|
||||
const hash = await state.password_service.hash(body.password);
|
||||
|
||||
const unlock = state.users_lock.lock();
|
||||
try {
|
||||
const users = loadUsers(state.store_base_dir) ?? {};
|
||||
if (body.username in users) {
|
||||
throw ApiError.conflict("Username already exists");
|
||||
}
|
||||
users[body.username] = hash;
|
||||
saveUsers(state.store_base_dir, users);
|
||||
} finally {
|
||||
unlock();
|
||||
}
|
||||
|
||||
return authResponse(state, body.username);
|
||||
}
|
||||
|
||||
/** POST /auth/refresh — exchange a refresh token for a fresh pair. */
|
||||
export async function refreshHandler(
|
||||
state: ApiState,
|
||||
headers: Headers,
|
||||
body: RefreshRequest,
|
||||
): Promise<AuthResponse> {
|
||||
enforceRateLimit(state.auth_rate_limiter, headers);
|
||||
if (body.refresh_token === "") {
|
||||
throw ApiError.badRequest("Refresh token is required");
|
||||
}
|
||||
|
||||
let sub: string;
|
||||
try {
|
||||
sub = state.token_service.verifyRefreshToken(body.refresh_token);
|
||||
} catch {
|
||||
throw ApiError.unauthorized("Invalid or expired refresh token");
|
||||
}
|
||||
|
||||
return authResponse(state, sub);
|
||||
}
|
||||
@@ -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<ChatCompletionResponse> {
|
||||
if (req.session_id === "") throw ApiError.badRequest("session_id is required");
|
||||
if (req.message === "") throw ApiError.badRequest("message is required");
|
||||
|
||||
// Load or create the conversation for this session.
|
||||
let conversation;
|
||||
try {
|
||||
conversation = await state.conversation_service.loadConversation(req.session_id);
|
||||
} catch {
|
||||
conversation = newConversation("", req.session_id);
|
||||
conversation.model = req.model ?? state.llm_client.model;
|
||||
if (req.max_tokens !== undefined) conversation.max_tokens = req.max_tokens;
|
||||
if (req.temperature !== undefined) conversation.temperature = req.temperature;
|
||||
}
|
||||
|
||||
if (req.model !== undefined) conversation.model = req.model;
|
||||
|
||||
const userMsg: ChatMessage = { role: Roles.User, content: req.message };
|
||||
conversation.messages.push(userMsg);
|
||||
|
||||
// Call the LLM (non-streaming), using the conversation's resolved model.
|
||||
let response: ChatMessage;
|
||||
let usage: [number, number] | null = null;
|
||||
try {
|
||||
const result = await state.llm_client.chat(
|
||||
conversation.messages,
|
||||
undefined,
|
||||
conversation.max_tokens,
|
||||
conversation.temperature,
|
||||
);
|
||||
response = result.message;
|
||||
usage = result.usage;
|
||||
} catch (e) {
|
||||
throw ApiError.chatProxy(`LLM request failed: ${(e as Error).message}`);
|
||||
}
|
||||
|
||||
const [promptTokens, completionTokens] = usage ?? [0, 0];
|
||||
const replyText = response.content ?? "";
|
||||
|
||||
const assistantMsg: ChatMessage = { role: Roles.Assistant, content: replyText };
|
||||
await state.conversation_service.addMessage(conversation, assistantMsg);
|
||||
|
||||
return { reply: replyText, prompt_tokens: promptTokens, completion_tokens: completionTokens };
|
||||
}
|
||||
@@ -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<ConversationResponse> {
|
||||
if (id === "") throw ApiError.badRequest("Session ID is required");
|
||||
const conversation = await state.conversation_service.loadConversation(id);
|
||||
return toConversationResponse(conversation);
|
||||
}
|
||||
|
||||
/** POST /sessions/:id/conversations — append a message to a conversation. */
|
||||
export async function addMessageHandler(
|
||||
state: ApiState,
|
||||
id: string,
|
||||
req: AddMessageRequest,
|
||||
): Promise<{ status: number; body: ConversationResponse }> {
|
||||
if (id === "") throw ApiError.badRequest("Session ID is required");
|
||||
if (req.content === "") throw ApiError.badRequest("Message content is required");
|
||||
|
||||
const role = req.role.toLowerCase();
|
||||
if (role !== "user" && role !== "assistant") {
|
||||
throw ApiError.badRequest(`Invalid role: ${req.role}`);
|
||||
}
|
||||
|
||||
const msg: ChatMessage = {
|
||||
role: role === "user" ? Roles.User : Roles.Assistant,
|
||||
content: req.content,
|
||||
};
|
||||
|
||||
const conversation = await state.conversation_service.loadConversation(id);
|
||||
await state.conversation_service.addMessage(conversation, msg);
|
||||
return { status: 200, body: toConversationResponse(conversation) };
|
||||
}
|
||||
|
||||
/** DELETE /sessions/:id/conversations/:cid — delete a message by index. */
|
||||
export async function deleteMessageHandler(
|
||||
state: ApiState,
|
||||
id: string,
|
||||
cid: string,
|
||||
): Promise<{ status: number }> {
|
||||
if (id === "") throw ApiError.badRequest("Session ID is required");
|
||||
|
||||
const index = Number.parseInt(cid, 10);
|
||||
if (Number.isNaN(index)) {
|
||||
throw ApiError.badRequest(`Invalid message index: ${cid}`);
|
||||
}
|
||||
|
||||
const conversation = await state.conversation_service.loadConversation(id);
|
||||
if (index >= conversation.messages.length) {
|
||||
throw ApiError.notFound(
|
||||
`Message index ${index} out of bounds (max: ${conversation.messages.length - 1})`,
|
||||
);
|
||||
}
|
||||
|
||||
conversation.messages.splice(index, 1);
|
||||
await state.conversation_service.saveConversation(conversation);
|
||||
return { status: 204 };
|
||||
}
|
||||
@@ -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<SessionListResponse> {
|
||||
const sessions = await state.session_service.listAll();
|
||||
const sessionResponses = sessions.map(toSessionResponse);
|
||||
return { sessions: sessionResponses, total: sessionResponses.length };
|
||||
}
|
||||
|
||||
/** POST /sessions — create a new session. */
|
||||
export async function createSessionHandler(
|
||||
state: ApiState,
|
||||
req: CreateSessionRequest,
|
||||
): Promise<{ status: number; body: SessionResponse }> {
|
||||
if (req.title.trim() === "") {
|
||||
throw ApiError.badRequest("Session title is required");
|
||||
}
|
||||
const session = await state.session_service.createSession(req.title);
|
||||
return { status: 201, body: toSessionResponse(session) };
|
||||
}
|
||||
|
||||
/** DELETE /sessions/:id — archive/close a session (returns 204). */
|
||||
export async function deleteSessionHandler(
|
||||
state: ApiState,
|
||||
id: string,
|
||||
): Promise<{ status: number }> {
|
||||
if (id === "") {
|
||||
throw ApiError.badRequest("Session ID is required");
|
||||
}
|
||||
const parsed = newSessionId(id);
|
||||
if (!parsed.ok) {
|
||||
throw ApiError.badRequest(`Invalid session ID: ${parsed.error}`);
|
||||
}
|
||||
await state.session_service.archiveSession(parsed.value);
|
||||
return { status: 204 };
|
||||
}
|
||||
@@ -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";
|
||||
@@ -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);
|
||||
@@ -1,51 +0,0 @@
|
||||
/**
|
||||
* JWT authentication middleware for the API.
|
||||
* Mirrors `apps/interfaces/api/src/middleware/auth.rs`.
|
||||
*
|
||||
* Validates the `Authorization: Bearer <token>` header and injects the
|
||||
* validated subject into `req.user`. Refresh tokens are rejected on
|
||||
* protected routes — only access tokens are accepted.
|
||||
*/
|
||||
import { verifyToken } from "@zesdex/infrastructure";
|
||||
import type { ApiState } from "../state.ts";
|
||||
|
||||
/** Claims decoded from a valid JWT, attached to the request. */
|
||||
export interface JwtClaims {
|
||||
sub: string;
|
||||
}
|
||||
|
||||
/** Result of authentication: either the subject or an error status/detail. */
|
||||
export type AuthResult = { ok: true; sub: string } | { ok: false; status: number; detail: string };
|
||||
|
||||
/**
|
||||
* Authenticate a request against the API state's JWT secret.
|
||||
* Reads the `Authorization: Bearer <token>` header, verifies the signature
|
||||
* and token type (must be `access`), and returns the subject on success.
|
||||
*/
|
||||
export function authenticateRequest(state: ApiState, headers: Headers): AuthResult {
|
||||
const authHeader = headers.get("Authorization");
|
||||
if (!authHeader) {
|
||||
return { ok: false, status: 401, detail: "Missing or invalid Authorization header" };
|
||||
}
|
||||
const prefix = "Bearer ";
|
||||
if (!authHeader.startsWith(prefix)) {
|
||||
return { ok: false, status: 401, detail: "Missing or invalid Authorization header" };
|
||||
}
|
||||
const token = authHeader.slice(prefix.length).trim();
|
||||
if (!token) {
|
||||
return { ok: false, status: 401, detail: "Missing or invalid Authorization header" };
|
||||
}
|
||||
try {
|
||||
const claims = verifyToken(state.jwt_secret, token);
|
||||
if (claims.typ !== "access") {
|
||||
return {
|
||||
ok: false,
|
||||
status: 401,
|
||||
detail: "refresh tokens are not accepted on protected routes",
|
||||
};
|
||||
}
|
||||
return { ok: true, sub: claims.sub };
|
||||
} catch (e) {
|
||||
return { ok: false, status: 401, detail: `Invalid token: ${(e as Error).message}` };
|
||||
}
|
||||
}
|
||||
@@ -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<string, string> = {
|
||||
"Access-Control-Allow-Origin": "*",
|
||||
"Access-Control-Allow-Methods": "GET, POST, DELETE, OPTIONS",
|
||||
"Access-Control-Allow-Headers": "Content-Type, Authorization",
|
||||
};
|
||||
|
||||
/** A handler's response: JSON body plus optional status code. */
|
||||
interface Result {
|
||||
status: number;
|
||||
body: unknown;
|
||||
}
|
||||
|
||||
/** Per-request context passed to handlers. */
|
||||
interface Ctx {
|
||||
state: ApiState;
|
||||
params: string[];
|
||||
body: unknown;
|
||||
headers: Headers;
|
||||
}
|
||||
|
||||
/** Route mapping: regex matches URL path after `/api/v1`. */
|
||||
interface Route {
|
||||
method: string;
|
||||
pattern: RegExp;
|
||||
protected: boolean;
|
||||
handle: (ctx: Ctx) => Promise<Result>;
|
||||
}
|
||||
|
||||
function result(body: unknown, status = 200): Result {
|
||||
return { status, body };
|
||||
}
|
||||
|
||||
function json(res: Result): Response {
|
||||
return new Response(JSON.stringify(res.body), {
|
||||
status: res.status,
|
||||
headers: { "Content-Type": "application/json", ...CORS_HEADERS },
|
||||
});
|
||||
}
|
||||
|
||||
function noContent(): Response {
|
||||
return new Response(null, { status: 204, headers: CORS_HEADERS });
|
||||
}
|
||||
|
||||
/** Build the route table for the API. Matches against paths without the /api/v1 prefix. */
|
||||
function buildRoutes(): Route[] {
|
||||
return [
|
||||
// Auth — public (rate-limited inside handlers)
|
||||
{ method: "POST", pattern: /^\/auth\/login$/, protected: false, handle: (c) => loginHandler(c.state, c.headers, c.body as never).then((r) => result(r)) },
|
||||
{ method: "POST", pattern: /^\/auth\/register$/, protected: false, handle: (c) => registerHandler(c.state, c.headers, c.body as never).then((r) => result(r)) },
|
||||
{ method: "POST", pattern: /^\/auth\/refresh$/, protected: false, handle: (c) => refreshHandler(c.state, c.headers, c.body as never).then((r) => result(r)) },
|
||||
|
||||
// Health — public
|
||||
{ method: "GET", pattern: /^\/health$/, protected: false, handle: () => Promise.resolve(result({ status: "ok" })) },
|
||||
|
||||
// Sessions — protected
|
||||
{ method: "GET", pattern: /^\/sessions$/, protected: true, handle: (c) => listSessionsHandler(c.state).then((s) => result(s)) },
|
||||
{ method: "POST", pattern: /^\/sessions$/, protected: true, handle: async (c) => {
|
||||
const r = await createSessionHandler(c.state, c.body as { title: string });
|
||||
return result(r.body, r.status);
|
||||
} },
|
||||
{ method: "DELETE", pattern: /^\/sessions\/([^/]+)$/, protected: true, handle: async (c) => {
|
||||
await deleteSessionHandler(c.state, c.params[0]!);
|
||||
return { status: 204, body: null };
|
||||
} },
|
||||
|
||||
// Conversations (sub-resource of sessions) — protected
|
||||
{ method: "GET", pattern: /^\/sessions\/([^/]+)\/conversations$/, protected: true, handle: async (c) => {
|
||||
const conv = await getConversationHandler(c.state, c.params[0]!);
|
||||
return result(conv);
|
||||
} },
|
||||
{ method: "POST", pattern: /^\/sessions\/([^/]+)\/conversations$/, protected: true, handle: async (c) => {
|
||||
const r = await addMessageHandler(c.state, c.params[0]!, c.body as { role: string; content: string });
|
||||
return result(r.body, r.status);
|
||||
} },
|
||||
{ method: "DELETE", pattern: /^\/sessions\/([^/]+)\/conversations\/([^/]+)$/, protected: true, handle: async (c) => {
|
||||
await deleteMessageHandler(c.state, c.params[0]!, c.params[1]!);
|
||||
return { status: 204, body: null };
|
||||
} },
|
||||
|
||||
// Chat — protected
|
||||
{ method: "POST", pattern: /^\/chat\/completions$/, protected: true, handle: (c) => chatCompletionsHandler(c.state, c.body as never).then((s) => result(s)) },
|
||||
];
|
||||
}
|
||||
|
||||
/** Handle a request against the router; maps ApiError → HTTP response. */
|
||||
async function handleRequest(state: ApiState, routes: Route[], req: Request): Promise<Response> {
|
||||
const url = new URL(req.url);
|
||||
// Strip the /api/v1 version prefix so route patterns match cleanly.
|
||||
const fullPath = url.pathname;
|
||||
const path = fullPath.startsWith("/api/v1") ? fullPath.slice("/api/v1".length) : fullPath;
|
||||
const method = req.method;
|
||||
|
||||
// CORS preflight
|
||||
if (method === "OPTIONS") {
|
||||
return new Response(null, { status: 204, headers: CORS_HEADERS });
|
||||
}
|
||||
|
||||
const route = routes.find((r) => r.method === method && r.pattern.test(path));
|
||||
if (!route) {
|
||||
return json(result(errorBody("Not found", 404), 404));
|
||||
}
|
||||
|
||||
// JWT auth for protected routes (subject is validated but not currently exposed).
|
||||
if (route.protected) {
|
||||
const auth = authenticateRequest(state, req.headers);
|
||||
if (!auth.ok) {
|
||||
return json(result({ error: "Unauthorized", detail: auth.detail }, auth.status));
|
||||
}
|
||||
}
|
||||
|
||||
// Parse JSON body if present; read once and reuse.
|
||||
let body: unknown = undefined;
|
||||
const raw = await req.text();
|
||||
if (raw) {
|
||||
try {
|
||||
body = JSON.parse(raw);
|
||||
} catch {
|
||||
return json(result(errorBody("Bad request: malformed JSON body", 400), 400));
|
||||
}
|
||||
}
|
||||
|
||||
const match = path.match(route.pattern);
|
||||
const ctx: Ctx = {
|
||||
state,
|
||||
params: match ? match.slice(1) : [],
|
||||
body,
|
||||
headers: req.headers,
|
||||
};
|
||||
|
||||
try {
|
||||
const res = await route.handle(ctx);
|
||||
if (res.status === 204) return noContent();
|
||||
return json(res);
|
||||
} catch (e) {
|
||||
if (e instanceof ApiError) {
|
||||
const msg = e.exposeDetails ? e.message : "An internal error occurred";
|
||||
return json(result(errorBody(msg, e.status), e.status));
|
||||
}
|
||||
console.error("Unhandled error:", e);
|
||||
return json(result(errorBody("An internal error occurred", 500)));
|
||||
}
|
||||
}
|
||||
|
||||
/** Start the API server on the given port. Returns a handle. */
|
||||
export function startApiServer(state: ApiState, port: number): { stop: () => void; port: number } {
|
||||
const routes = buildRoutes();
|
||||
const server = Bun.serve({
|
||||
port,
|
||||
async fetch(req) {
|
||||
return handleRequest(state, routes, req);
|
||||
},
|
||||
});
|
||||
console.log(`zesdex-api listening on http://0.0.0.0:${server.port}`);
|
||||
return { stop: () => server.stop(true), port: server.port ?? port };
|
||||
}
|
||||
@@ -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<string, number[]>();
|
||||
|
||||
/** Check whether `key` is within `max` requests per `windowSecs`. Returns true if allowed. */
|
||||
check(key: string, max: number, windowSecs: number): boolean {
|
||||
const now = Date.now();
|
||||
const cutoff = now - windowSecs * 1000;
|
||||
const list = (this.hits.get(key) ?? []).filter((t) => t > cutoff);
|
||||
if (list.length >= max) {
|
||||
this.hits.set(key, list);
|
||||
return false;
|
||||
}
|
||||
list.push(now);
|
||||
this.hits.set(key, list);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
/** Concrete session repository wired from infrastructure. */
|
||||
const sessionRepo = new FileSystemSessionRepository();
|
||||
const sessionLockRepo = new FileSystemSessionLockRepository();
|
||||
const conversationRepo = new JsonConversationRepository();
|
||||
const settingsRepo = new JsonSettingsRepository();
|
||||
const appConfigRepo = new JsonAppConfigRepository();
|
||||
const memoryRepo = new MarkdownMemoryRepository();
|
||||
|
||||
/** Rooted, long-lived shared API state. */
|
||||
export interface ApiState {
|
||||
store_base_dir: string;
|
||||
jwt_secret: string;
|
||||
session_service: SessionServiceImpl;
|
||||
conversation_service: ConversationServiceImpl;
|
||||
settings_service: SettingsServiceImpl;
|
||||
memory_service: MemoryServiceImpl;
|
||||
password_service: Argon2PasswordService;
|
||||
token_service: Hs256TokenService;
|
||||
auth_rate_limiter: RateLimiter;
|
||||
/** Serializes read-modify-write of `users.json` (TOCTOU race guard). */
|
||||
users_lock: { lock: () => () => void };
|
||||
llm_client: LlmClient;
|
||||
}
|
||||
|
||||
/**
|
||||
* A no-op lock that returns an unlock function.
|
||||
* Bun is single-threaded per process, so the JS event loop already
|
||||
* serialises synchronous read-modify-write of users.json — the lock is
|
||||
* retained for structural parity with the Rust mutex guard.
|
||||
*/
|
||||
function noopLock(): { lock: () => () => void } {
|
||||
return { lock: () => () => {} };
|
||||
}
|
||||
|
||||
/** Construct a new API state with all services wired to their defaults. */
|
||||
export function newApiState(
|
||||
baseDir: string,
|
||||
jwtSecret: string,
|
||||
llmApiKey: string,
|
||||
llmModel: string,
|
||||
llmBaseUrl?: string,
|
||||
): ApiState {
|
||||
const sessionsDir = path.join(baseDir, "sessions");
|
||||
const memoryDir = path.join(baseDir, "memories");
|
||||
|
||||
const session_service = new SessionServiceImpl(sessionRepo, sessionLockRepo, baseDir);
|
||||
const conversation_service = new ConversationServiceImpl(conversationRepo, sessionsDir);
|
||||
const settings_service = new SettingsServiceImpl(settingsRepo, appConfigRepo, baseDir);
|
||||
const memory_service = new MemoryServiceImpl(memoryRepo, memoryDir);
|
||||
|
||||
const token_service = new Hs256TokenService(jwtSecret);
|
||||
const llm = new LlmClient(llmApiKey, llmModel, llmBaseUrl);
|
||||
|
||||
return {
|
||||
store_base_dir: baseDir,
|
||||
jwt_secret: jwtSecret,
|
||||
session_service,
|
||||
conversation_service,
|
||||
settings_service,
|
||||
memory_service,
|
||||
password_service: new Argon2PasswordService(),
|
||||
token_service,
|
||||
auth_rate_limiter: new RateLimiter(),
|
||||
users_lock: noopLock(),
|
||||
llm_client: llm,
|
||||
};
|
||||
}
|
||||
|
||||
/** Load the `users.json` map (username → password hash), or null if absent. */
|
||||
export function loadUsers(baseDir: string): Record<string, string> | null {
|
||||
const p = path.join(baseDir, "users.json");
|
||||
try {
|
||||
return JSON.parse(fs.readFileSync(p, "utf8")) as Record<string, string>;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** Persist the users map to `users.json` (pretty-printed). */
|
||||
export function saveUsers(baseDir: string, users: Record<string, string>): void {
|
||||
const p = path.join(baseDir, "users.json");
|
||||
fs.mkdirSync(path.dirname(p), { recursive: true });
|
||||
fs.writeFileSync(p, JSON.stringify(users, null, 2));
|
||||
}
|
||||
+18
-102
@@ -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 <session>, --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 <prompt>` 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 "<prompt>" Run a single one-shot agent turn and exit
|
||||
|
||||
Options:
|
||||
--daemon Run as background daemon with IPC socket
|
||||
--attach <session> Attach TUI client to a running daemon session
|
||||
--api Run REST API server
|
||||
--api-port <PORT> REST API port [default: 8080]
|
||||
--ws Run WebSocket server
|
||||
--ws-port <PORT> WebSocket port [default: 8081]
|
||||
--grpc Run gRPC server
|
||||
--grpc-port <PORT> gRPC port [default: 50051]
|
||||
--web Serve web frontend
|
||||
--web-port <PORT> Web frontend port [default: 3000]
|
||||
-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<string, unknown>)[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}`);
|
||||
if (arg === "--headless") {
|
||||
const next = argv[i + 1];
|
||||
if (next === undefined) throw new Error("flag '--headless' requires a value");
|
||||
opts.headless = next;
|
||||
i++;
|
||||
continue;
|
||||
}
|
||||
(opts as unknown as Record<string, unknown>)[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}`);
|
||||
}
|
||||
}
|
||||
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;
|
||||
}
|
||||
|
||||
+24
-106
@@ -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 "<p>"` → 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 <data-dir>/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<void> {
|
||||
/** Run a single one-shot agent turn non-interactively, printing the outcome. */
|
||||
async function runHeadless(prompt: string): Promise<void> {
|
||||
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 });
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
const events: import("@zesdex/domain").TurnEvent[] = [];
|
||||
const abort = new AbortController();
|
||||
const params = buildTurnParamsFrom(runtime, text, events, abort);
|
||||
const params = buildTurnParams(runtime, prompt, events, abort);
|
||||
|
||||
process.stdout.write("…");
|
||||
try {
|
||||
await runtime.turnService.runTurn(params);
|
||||
// Print the final assistant content from events.
|
||||
const lastTokens: string[] = [];
|
||||
|
||||
const tokens: string[] = [];
|
||||
for (const ev of events) {
|
||||
if (ev.kind === "stream_token") lastTokens.push(ev.content);
|
||||
if (ev.kind === "stream_token") tokens.push(ev.content);
|
||||
}
|
||||
const full = lastTokens.join("");
|
||||
const full = tokens.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();
|
||||
}
|
||||
|
||||
// Imported lazily to avoid a circular type-only dependency.
|
||||
import { buildTurnParams } from "./compose.ts";
|
||||
function buildTurnParamsFrom(
|
||||
runtime: Parameters<typeof buildTurnParams>[0],
|
||||
text: string,
|
||||
events: import("@zesdex/domain").TurnEvent[],
|
||||
abort: AbortController,
|
||||
) {
|
||||
return buildTurnParams(runtime, text, events, abort);
|
||||
}
|
||||
|
||||
/** Main entry. */
|
||||
/** Main entry. `zesdex` → TUI; `--headless <prompt>` → one-shot turn. */
|
||||
async function main(): Promise<void> {
|
||||
initLogging();
|
||||
|
||||
@@ -119,57 +79,15 @@ async function main(): Promise<void> {
|
||||
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) => {
|
||||
|
||||
@@ -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<void> {
|
||||
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<DaemonFrame>();
|
||||
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<void> {
|
||||
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<DaemonFrame>();
|
||||
if (frame && frame.kind === "Closed") {
|
||||
process.stdout.write("[closed]\n");
|
||||
}
|
||||
client.close();
|
||||
}
|
||||
@@ -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";
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user