feat(rewrite): bootstrap monorepo Bun + port domain layer ke TypeScript

Mulai migrasi Zesdex dari Rust ke TypeScript/Bun.

- Root workspace package.json, tsconfig strict, bun.lock, gitignore fix
- Paket @zesdex/domain (apps/packages/domain):
  - core: ChatMessage/Conversation/ToolCall/SseParser/repairJson/
    sanitizeToolArguments/UsageStats/Store/DomainError
  - auth: Session/SessionId/SessionLock/OAuth + repository & service interfaces
  - cms: AppConfig/Settings/SettingsPatch/Memory/EditLog + repository &
    service interfaces
  - agent: TurnEvent/SessionRuntime/AgentTurnParams/AgentProgress/prompt
  - subagent & workflow models
- Wire-shape JSON disamakan dengan Rust (serde rename/flatten/skip_if)
- 27 unit test Bun port dari test Rust (SseParser, repair_json, resolve_effective_model, session id)

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
asepharyana
2026-09-02 22:50:38 +07:00
co-authored by Claude Opus 5
parent 9d1544d799
commit 07e84e43b3
45 changed files with 2315 additions and 1 deletions
+18
View File
@@ -0,0 +1,18 @@
{
"name": "@zesdex/domain",
"version": "1.21.2",
"private": true,
"type": "module",
"description": "Zesdex domain layer — pure types, value objects, and port interfaces (zero I/O)",
"exports": {
".": "./src/index.ts"
},
"scripts": {
"test": "bun test",
"typecheck": "tsc --noEmit"
},
"devDependencies": {
"@types/bun": "^1.2.0",
"typescript": "^5.7.0"
}
}
@@ -0,0 +1,28 @@
/** Shared default constants. Mirrors `agent/defaults.rs`. */
/** Default LLM provider API base URL. */
export const DEFAULT_API_BASE = "https://opencode.ai/zen/v1";
/** Default LLM model identifier. */
export const DEFAULT_MODEL = "deepseek-v4-flash-free";
/** Fallback JWT secret used only when `JWT_SECRET` env var is unset. */
export const FALLBACK_JWT_SECRET = "dev-secret";
/** Default context window size (256k tokens). */
export const DEFAULT_CONTEXT_WINDOW = 256_000;
/** Maximum tool-call iterations per agent turn. */
export const MAX_TOOL_ITERATIONS = 50;
/** Maximum subagent tool-call iterations. */
export const MAX_SUBAGENT_ITERATIONS = 25;
/** Default LLM request max tokens. */
export const DEFAULT_MAX_TOKENS = 4096;
/** Default temperature for the main agent. */
export const DEFAULT_TEMPERATURE = 0.7;
/** Default temperature for compaction / summary calls. */
export const DEFAULT_COMPACT_TEMPERATURE = 0.3;
+5
View File
@@ -0,0 +1,5 @@
/** Agent domain module — turn events, runtime, progress, prompts, constants. */
export * from "./mod.ts";
export * from "./defaults.ts";
export * from "./progress.ts";
export * from "./prompt.ts";
+194
View File
@@ -0,0 +1,194 @@
/**
* Domain types for agent lifecycle: turn events, session runtime, progress.
* Mirrors `apps/domain/src/agent/mod.rs`.
*/
import type { ChatMessage, JsonValue, ToolCallResult, UsageStats } from "../core/index.ts";
import type { AgentProgress } from "./progress.ts";
/** Which kind of caller (main agent vs subagent vs reviewer) is invoking a tool. */
export type Origin = "Main" | "SubAgent" | "Reviewer";
export const OriginTag = {
Main: "main",
SubAgent: "subagent",
Reviewer: "reviewer",
} as const satisfies Record<Origin, string>;
export function originTag(o: Origin): string {
return OriginTag[o];
}
/** Severity/category of a toast notification. */
export type ToastKind = "Info" | "Success" | "Warning" | "Error" | "Lesson";
/** A transient status message shown in the TUI, auto-dismissed after lifetime_ms. */
export interface Toast {
kind: ToastKind;
message: string;
created_at: number;
lifetime_ms: number;
}
export const DEFAULT_TOAST_LIFETIME_MS = 5000;
export function newToast(kind: ToastKind, message: string): Toast {
return { kind, message, created_at: Date.now(), lifetime_ms: DEFAULT_TOAST_LIFETIME_MS };
}
export function toastExpired(toast: Toast, nowMs: number): boolean {
return nowMs - toast.created_at > toast.lifetime_ms;
}
/**
* Agent status for workflow progress tracking. `Failed` carries a message;
* we represent it as a string union plus an optional error on failures.
*/
export type AgentStatus = "Pending" | "Running" | "Completed" | "Failed" | "Cancelled";
export function agentStatusDisplay(s: AgentStatus, error?: string): string {
switch (s) {
case "Pending":
return "pending";
case "Running":
return "running";
case "Completed":
return "completed";
case "Failed":
return error ? `failed: ${error}` : "failed";
case "Cancelled":
return "cancelled";
}
}
/** Events emitted onto the turn-event queue while an agent turn runs. */
export type TurnEvent =
| { kind: "assistant_message"; message: ChatMessage }
| { kind: "tool_result"; tool_call_id: string; tool_name: string; output: string; is_error: boolean; path: string | null }
| { kind: "system_note"; systemKind: string; message: string }
| { kind: "stream_start" }
| { kind: "stream_token"; content: string }
| { kind: "stream_reasoning"; content: string }
| { kind: "stream_done"; message: ChatMessage }
| { kind: "usage"; tokens_in: number; tokens_out: number }
| { kind: "review_usage"; tokens_in: number; tokens_out: number }
| { kind: "compacted"; messages: ChatMessage[] }
| { kind: "error"; message: string }
| { kind: "done" }
| { kind: "workflow_agent_update"; agent_id: string; agent_name: string; status: AgentStatus; error?: string }
| { kind: "todo_update"; content: string }
| { kind: "plan_update"; content: string }
| { kind: "agent_progress"; progress: AgentProgress };
/** How a pending tool call should be executed when the turn resumes. */
export type ExecutionModel = "Inline" | "Deferred" | "AsyncTokio";
/** A tool call awaiting execution. */
export interface PendingTool {
tool_name: string;
args: JsonValue;
execution_model: ExecutionModel;
}
/** Reference to a background bash job tracked in session state. */
export interface BashJobRef {
id: string;
command: string;
started_at: number;
running: boolean;
}
/** Tracks counts of learned patterns by outcome and lifecycle stage. */
export interface LessonStats {
total: number;
user: number;
feedback: number;
project: number;
reference: number;
active: number;
stale: number;
contradicted: number;
human: number;
verified: number;
unverified: number;
}
export function newLessonStats(): LessonStats {
return {
total: 0, user: 0, feedback: 0, project: 0, reference: 0,
active: 0, stale: 0, contradicted: 0, human: 0, verified: 0, unverified: 0,
};
}
/** Per-session runtime state: message history, pending tool queue, jobs, counters. */
export interface SessionRuntime {
messages: ChatMessage[];
tool_call_results: ToolCallResult[];
pending_tool_queue: PendingTool[];
bash_jobs: BashJobRef[];
subagent_queue: number;
edit_count: number;
consecutive_empty_reviews: number;
session_start: number;
lessons: LessonStats;
review_count: number;
session_dir: string;
usage: UsageStats;
hive_mind_converged: boolean;
}
export function newSessionRuntime(sessionDir: string): SessionRuntime {
return {
messages: [],
tool_call_results: [],
pending_tool_queue: [],
bash_jobs: [],
subagent_queue: 0,
edit_count: 0,
consecutive_empty_reviews: 0,
session_start: Date.now(),
lessons: newLessonStats(),
review_count: 0,
session_dir: sessionDir,
usage: {
tokens_in: 0, tokens_out: 0, last_tokens_in: 0, last_tokens_out: 0,
api_calls: 0, review_tokens: 0, total_ms: 0,
},
hive_mind_converged: false,
};
}
export function runtimePushMessage(rt: SessionRuntime, msg: ChatMessage): void {
rt.messages.push(msg);
}
/** Simple ASCII progress display for a long-running operation. */
export interface ProgressState {
current: number;
total: number;
message: string;
start_time: number;
}
/** Owned parameters required to spawn and execute an agent turn. */
export interface AgentTurnParams {
messages: ChatMessage[];
session_dir: string;
workspace_roots: string[];
/** Event sink — array or queue implementing push/drain semantics. */
turn_events: TurnEventSink;
/** Whether a turn is currently in flight. */
in_flight: { value: boolean };
/** Abort control. */
abort: AbortController;
api_key: string;
model: string;
api_base?: string;
}
/** Minimal event-sink abstraction (backs the Rust `Arc<Mutex<VecDeque>>`). */
export interface TurnEventSink {
/** Append one event. */
push(event: TurnEvent): void;
/** Drain all pending events (FIFO) and return them. */
drain(): TurnEvent[];
}
@@ -0,0 +1,44 @@
/** Progress reporting types for agent/subagent operations. Mirrors `progress.rs`. */
import type { AgentStatus } from "./mod.ts";
/**
* A failure status carries the error text alongside the `Failed` kind;
* success statuses have no error.
*/
export type AgentStatusWithError = Exclude<AgentStatus, "Failed"> | "Failed";
/** Describes progress within a single subagent or workflow-node execution. */
export interface AgentProgress {
/** Unique identifier for this agent. */
agent_id: string;
/** Human-readable display name shown in the TUI sidebar. */
agent_name: string;
/** Current lifecycle status. */
status: AgentStatus;
/** Error message when status === "Failed". */
error?: string;
/** Optional current tool / step being executed. `null` when idle. */
current_tool: string | null;
/** Optional progress range: [completed, total]. `null` = indeterminate. */
steps: [number, number] | null;
}
/** Mark this agent as running with an optional tool name. */
export function agentProgressRunning(agentId: string, agentName: string, currentTool: string | null): AgentProgress {
return { agent_id: agentId, agent_name: agentName, status: "Running", current_tool: currentTool, steps: null };
}
/** Mark this agent as pending (queued but not yet started). */
export function agentProgressPending(agentId: string, agentName: string): AgentProgress {
return { agent_id: agentId, agent_name: agentName, status: "Pending", current_tool: null, steps: null };
}
/** Mark this agent as completed successfully. */
export function agentProgressCompleted(agentId: string, agentName: string): AgentProgress {
return { agent_id: agentId, agent_name: agentName, status: "Completed", current_tool: null, steps: null };
}
/** Mark this agent as failed with an error message. */
export function agentProgressFailed(agentId: string, agentName: string, error: string): AgentProgress {
return { agent_id: agentId, agent_name: agentName, status: "Failed", error, current_tool: null, steps: null };
}
+68
View File
@@ -0,0 +1,68 @@
/** System prompts and directive templates. Mirrors `agent/prompt.rs`. */
/** Build the main-agent system prompt. */
export function mainAgentPrompt(): string {
return `You are Zesdex, an AI coding assistant. You have access to various tools via native function calling to help the user.
TOKEN BUDGET — BE EFFICIENT:
- For simple/factual questions, answer directly. Do NOT call tools.
- For complex or unfamiliar code tasks, call \`explore_codebase\` ONCE at the start to locate relevant code, then work from that context.
- Keep tool usage minimal: prefer \`grep\`/\`glob\`/\`read\` for targeted lookups; avoid re-reading files you already have in context.
- Keep responses concise; do not repeat tool output verbatim.
CRITICAL DIRECTIVES & PRIORITY HIERARCHY:
1. WORKFLOW FIRST: For any multi-step, complex, or non-trivial task, you MUST prioritise using \`workflow_run\` (to construct and execute a multi-phase YAML workflow) or \`hive_mind\` (to orchestrate parallel autonomous agents). Workflows are your primary strategy.
2. PLANNING & TODOs: Use \`plan_enter\` to establish high-level architectural plans and \`todowrite\` to maintain granular task checklists.
3. REASONING: Use \`seq_think\` for deep step-by-step analysis.
4. TOOL EXECUTION: Execute individual tools (file edits, terminal commands) within or guided by your workflows. If an error occurs, analyse and fix it.
VERIFY AFTER EDIT (CLAUDE-CODE STYLE):
- After modifying code (edit/write), run the repo's check command via \`bash\` before ending the turn: \`cargo check\` / \`cargo clippy\` / \`cargo test\` for Rust, or the equivalent lint/test (\`bun run lint && bun run test\`, \`npm test\`, etc.) for other stacks. Pick the project's actual verify command (see PROJECT CONTEXT / AGENTS.md when present).
- If the check fails, fix the errors you can see and re-run; only end the turn after the check passes or you cannot resolve a failure yourself (then report it explicitly).
- Do NOT claim code compiles or works without running a real check.
Respond conversationally, concisely, and helpfully.`;
}
/** Main-agent prompt with an injected `## PROJECT CONTEXT` block. Empty context → base prompt. */
export function mainAgentPromptWithProjectContext(projectContext: string): string {
const base = mainAgentPrompt();
const context = projectContext.trim();
if (context === "") return base;
return `${base}
## PROJECT CONTEXT (repo rules — follow these conventions)
${context}`;
}
/** Build a subagent directive prompt. */
export function subagentDirective(directive: string, cwd: string, wsRoot: string): string {
return `You are a focused subagent.
Current directory (PWD): ${cwd}
Workspace root: ${wsRoot}
Your directive:
${directive}
Complete the directive autonomously using the tools available to you. Return your final answer when done.`;
}
/** Build a conversation-compaction prompt. */
export function compactionPrompt(): string {
return `You are a helpful assistant summarising conversation history. Provide a concise summary of the key user requests, decisions, tools executed, and modified files. Format as a clear bulleted list.`;
}
/** Directive for a lightweight context-scout subagent. */
export function exploreScoutDirective(): string {
return `You are a codebase context scout. Given the workspace root, quickly locate the code that is most relevant to the user's request:
1. Run semantic_search once with the user's key terms.
2. Read up to the 3 most relevant files (use grep for symbols if needed).
3. Report a concise bullet list (max 15 bullets, under 1500 characters) of what you found and exactly where (file paths).
Do NOT rebuild the index. Do NOT enumerate unrelated files. Be brief.`;
}
/** System note injected after repeated tool errors. */
export function errorRecoveryNote(toolName: string, lastError: string): string {
return `[System note] The tool \`${toolName}\` failed repeatedly with: "${lastError}". Try an alternative approach (verify paths, correct arguments, use a different tool, or finish without this tool). Do NOT retry the same call.`;
}
+11
View File
@@ -0,0 +1,11 @@
/**
* 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 };
}
+47
View File
@@ -0,0 +1,47 @@
/**
* 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;
}
}
+20
View File
@@ -0,0 +1,20 @@
/** 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";
+40
View File
@@ -0,0 +1,40 @@
/**
* 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],
};
}
@@ -0,0 +1,39 @@
/**
* 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;
}
+35
View File
@@ -0,0 +1,35 @@
/**
* 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;
}
+56
View File
@@ -0,0 +1,56 @@
/**
* 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");
}
@@ -0,0 +1,25 @@
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");
});
});
@@ -0,0 +1,41 @@
/**
* 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);
}
@@ -0,0 +1,120 @@
/**
* 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;
}
}
}
@@ -0,0 +1,62 @@
/**
* Application configuration entities. Mirrors `app_config.rs`.
*/
import { DEFAULT_CONTEXT_WINDOW } from "../agent/defaults.ts";
/** Top-level application configuration. */
export interface AppConfig {
providers: Record<string, ProviderConfig>;
model_roles: Record<string, ModelRole>;
default_provider: string;
default_model: string;
default_context_window: number;
}
/** Connection details for a single LLM provider endpoint. */
export interface ProviderConfig {
api_base: string;
api_key_env?: string;
default_model?: string;
default_api_key?: string;
}
/** A named model role mapping to a provider/model with parameters. */
export interface ModelRole {
provider: string;
model: string;
max_tokens?: number;
context_window?: number;
temperature?: number;
}
/** Returns the default AppConfig with built-in "zen" and "router" providers. */
export function newAppConfig(): AppConfig {
return {
providers: {
zen: {
api_base: "https://opencode.ai/zen/v1",
api_key_env: "API_KEY",
default_model: "deepseek-v4-flash-free",
default_api_key: undefined,
},
router: {
api_base: "https://9router.asepharyana.my.id/v1",
api_key_env: "ROUTER_API_KEY",
default_model: "claude-opus-5",
default_api_key: undefined,
},
},
model_roles: {
default: {
provider: "zen",
model: "deepseek-v4-flash-free",
max_tokens: undefined,
context_window: undefined,
temperature: 0.7,
},
},
default_provider: "zen",
default_model: "deepseek-v4-flash-free",
default_context_window: DEFAULT_CONTEXT_WINDOW,
};
}
+81
View File
@@ -0,0 +1,81 @@
/**
* Command types for CMS domain operations. Mirrors `commands.rs`.
*/
import type { InternetMode } from "./settings.ts";
const VALID_INTERNET_MODES: InternetMode[] = ["Off", "ReadOnly", "Full"];
/** Partial update command for `Settings` — only non-null fields are applied. */
export interface SettingsPatch {
internet_mode?: string;
provider?: string;
model?: string;
api_keys?: Record<string, string>;
max_tokens?: number | null;
temperature?: number | null;
review_max_lessons_per_run?: number;
adaptive_review_max_skip?: number;
verify_command?: string | null;
verify_timeout_ms?: number;
workflow_max_concurrency?: number;
review_enabled?: boolean;
session_archive_enabled?: boolean;
hive_mind_node_timeout_ms?: number;
}
/** The `apply` fields a Settings-like target must expose (subset of Settings). */
export interface SettingsPatchTarget {
internet_mode: InternetMode;
provider: string;
model: string;
api_keys: Record<string, string>;
max_tokens?: number | null;
temperature?: number | null;
review_max_lessons_per_run: number;
adaptive_review_max_skip: number;
verify_command?: string | null;
verify_timeout_ms: number;
workflow_max_concurrency: number;
review_enabled: boolean;
session_archive_enabled: boolean;
hive_mind_node_timeout_ms: number;
}
/** Merge a patch into a settings-like target; only defined fields are applied. */
export function applySettingsPatch(patch: SettingsPatch, settings: SettingsPatchTarget): string | null {
if (patch.internet_mode !== undefined) {
const mode = patch.internet_mode;
if (!VALID_INTERNET_MODES.includes(mode as InternetMode)) {
return `invalid internet_mode '${mode}'; expected Off, ReadOnly, or Full`;
}
settings.internet_mode = mode as InternetMode;
}
if (patch.provider !== undefined) settings.provider = patch.provider;
if (patch.model !== undefined) settings.model = patch.model;
if (patch.api_keys !== undefined) settings.api_keys = patch.api_keys;
if (patch.max_tokens !== undefined) settings.max_tokens = patch.max_tokens;
if (patch.temperature !== undefined) settings.temperature = patch.temperature;
if (patch.review_max_lessons_per_run !== undefined) settings.review_max_lessons_per_run = patch.review_max_lessons_per_run;
if (patch.adaptive_review_max_skip !== undefined) settings.adaptive_review_max_skip = patch.adaptive_review_max_skip;
if (patch.verify_command !== undefined) settings.verify_command = patch.verify_command;
if (patch.verify_timeout_ms !== undefined) settings.verify_timeout_ms = patch.verify_timeout_ms;
if (patch.workflow_max_concurrency !== undefined) settings.workflow_max_concurrency = patch.workflow_max_concurrency;
if (patch.review_enabled !== undefined) settings.review_enabled = patch.review_enabled;
if (patch.session_archive_enabled !== undefined) settings.session_archive_enabled = patch.session_archive_enabled;
if (patch.hive_mind_node_timeout_ms !== undefined) settings.hive_mind_node_timeout_ms = patch.hive_mind_node_timeout_ms;
return null;
}
/** Command to create a new memory entry. */
export interface NewMemory {
name: string;
description: string;
content: string;
kind?: string;
outcome?: string;
lifecycle?: string;
scope?: string;
before_snippet?: string;
after_snippet?: string;
provenances?: string[];
}
@@ -0,0 +1,7 @@
/**
* Re-export of conversation types from core (matches Rust `cms/conversation.rs`).
*/
export type { Conversation } from "../core/conversation.ts";
export type { ChatMessage, Role } from "../core/message.ts";
export { newConversation, pushMessage, rebuildSystem, toApiMessages } from "../core/conversation.ts";
export { Roles } from "../core/message.ts";
+53
View File
@@ -0,0 +1,53 @@
/**
* Edit log — append-only log of file mutations. Mirrors `edit_log.rs`.
*/
/** A single recorded file edit event. */
export interface EditLogEntry {
/** Unix timestamp (ms) when the edit occurred. */
ts: number;
/** Name of the tool that performed the edit. */
tool: string;
/** Absolute file path that was modified. */
path: string;
/** Human-readable explanation of why the edit was made. */
reason: string;
/** SHA-256 hex digest of the content after the edit. */
content_sha256: string;
/** Signed byte count change (+added, -removed). */
bytes_delta: number;
/** Origin identifier (which agent / session context). */
origin: string;
/** Session in which this edit was performed. */
session_id: string;
}
/** Maximum number of edit entries held in memory at once. */
export const MAX_MEMORY_ENTRIES = 10_000;
/** In-memory view of a session's edit log. */
export class EditLog {
entries: EditLogEntry[];
constructor() {
this.entries = [];
}
/** Number of in-memory entries. */
get length(): number {
return this.entries.length;
}
/** Whether the log contains no entries. */
get isEmpty(): boolean {
return this.entries.length === 0;
}
/** Append an entry, evicting oldest beyond MAX_MEMORY_ENTRIES. */
push(entry: EditLogEntry): void {
this.entries.push(entry);
if (this.entries.length > MAX_MEMORY_ENTRIES) {
this.entries.shift();
}
}
}
+34
View File
@@ -0,0 +1,34 @@
/**
* Domain error types for the CMS module. Mirrors `cms/error.rs`.
*/
import type { DomainError } from "../core/error.ts";
/** Shared repository error type for CMS persistence. */
export type RepositoryError = DomainError;
/** Errors from service / use-case operations in the CMS domain. */
export type ServiceError =
| { kind: "repository"; error: DomainError }
| { kind: "invalid_input"; message: string }
| { kind: "other"; message: string };
export function cmsRepositoryErr(err: DomainError): ServiceError {
return { kind: "repository", error: err };
}
export function invalidInput(message: string): ServiceError {
return { kind: "invalid_input", message };
}
export function cmsOther(message: string): ServiceError {
return { kind: "other", message };
}
export function cmsServiceErrorToString(e: ServiceError): string {
switch (e.kind) {
case "repository":
return `repository error: ${e.error.message}`;
case "invalid_input":
return `invalid input: ${e.message}`;
case "other":
return e.message;
}
}
+19
View File
@@ -0,0 +1,19 @@
/** CMS domain module — settings, app config, memory, edit log, services. */
export * from "./app_config.ts";
export * from "./settings.ts";
export * from "./commands.ts";
export * from "./memory.ts";
export * from "./edit_log.ts";
export * from "./repository.ts";
export * from "./service.ts";
export {
type ServiceError as CmsServiceError,
type RepositoryError as CmsRepositoryError,
} from "./error.ts";
export {
cmsRepositoryErr,
invalidInput,
cmsOther,
cmsServiceErrorToString,
} from "./error.ts";
export * from "./conversation.ts";
+44
View File
@@ -0,0 +1,44 @@
/**
* Long-term agent memory entity. Mirrors `memory.rs`.
*/
import * as path from "node:path";
/** A single memory entry with frontmatter metadata and markdown content. */
export interface Memory {
name: string;
description: string;
content: string;
kind: string;
created_at: number;
updated_at: number;
outcome?: string;
lifecycle: string;
scope?: string;
before_snippet?: string;
after_snippet?: string;
provenances: string[];
}
/** Convert an arbitrary string into a filesystem-safe slug. */
export function slugify(s: string): string | null {
// Phase 1: replace every non-alphanumeric char with '-'
let slug = s.toLowerCase().replace(/[^a-z0-9]/g, "-");
// Phase 2: collapse consecutive '-'
slug = slug.split("-").filter((seg) => seg !== "").join("-");
if (slug === "" || slug.length > 80) return null;
return slug;
}
/**
* Compute the on-disk path for a memory of the given name.
* Falls back to `"memory.md"` when the slug is empty/invalid.
*/
export function memoryPath(memoryDir: string, name: string): string {
const slug = slugify(name) ?? "memory";
const clean = `${slug}.md`
.split("")
.map((c) => (/[a-z0-9.-]/.test(c) ? c : "-"))
.join("");
const trimmed = clean.replace(/^\.+/, "");
return path.join(memoryDir, trimmed === "" ? "memory.md" : trimmed);
}
@@ -0,0 +1,51 @@
/**
* Repository trait definitions (interfaces) for persistence.
* Mirrors `apps/domain/src/cms/repository.rs`. Infrastructure adapters
* implement these.
*/
import type { AppConfig } from "./app_config.ts";
import type { EditLog, EditLogEntry } from "./edit_log.ts";
import type { Memory } from "./memory.ts";
import type { Settings } from "./settings.ts";
import type { Conversation } from "../core/conversation.ts";
/** Persistence contract for `Settings`. */
export interface SettingsRepository {
load(baseDir: string): Promise<Settings> | Settings;
save(baseDir: string, settings: Settings): Promise<void> | void;
}
/** Persistence contract for `AppConfig`. */
export interface AppConfigRepository {
load(baseDir: string): Promise<AppConfig> | AppConfig;
save(baseDir: string, config: AppConfig): Promise<void> | void;
}
/** Persistence contract for `Conversation`. */
export interface ConversationRepository {
load(sessionDir: string): Promise<Conversation> | Conversation;
save(sessionDir: string, conversation: Conversation): Promise<void> | void;
}
/** Persistence contract for `Memory`. */
export interface MemoryRepository {
list(memoryDir: string): Promise<string[]> | string[];
load(memoryDir: string, name: string): Promise<Memory> | Memory;
save(memoryDir: string, memory: Memory): Promise<void> | void;
delete(memoryDir: string, name: string): Promise<void> | void;
}
/** Persistence contract for rewind-snapshot binary blobs. */
export interface RewindBlobRepository {
storeBlob(sessionDir: string, blobKey: string, data: Uint8Array, mimeType?: string): Promise<void> | void;
retrieveBlob(sessionDir: string, blobKey: string): Promise<Uint8Array | null> | Uint8Array | null;
listBlobKeys(sessionDir: string): Promise<string[]> | string[];
}
/** Persistence contract for `EditLog`. */
export interface EditLogRepository {
open(sessionDir: string): Promise<EditLog> | EditLog;
append(sessionDir: string, log: EditLog, entry: EditLogEntry): Promise<void> | void;
entries(log: EditLog): EditLogEntry[];
}
+30
View File
@@ -0,0 +1,30 @@
/**
* Service trait definitions — use-case boundaries for CMS operations.
* Mirrors `apps/domain/src/cms/service.rs`. Implemented by the application layer.
*/
import type { Conversation } from "../core/conversation.ts";
import type { ChatMessage } from "../core/message.ts";
import type { Memory } from "./memory.ts";
import type { Settings } from "./settings.ts";
import type { ProviderConfig } from "./app_config.ts";
/** Use-cases for application settings. */
export interface SettingsService {
loadSettings(): Promise<Settings> | Settings;
saveSettings(settings: Settings): Promise<void> | void;
updateProvider(name: string, config: ProviderConfig): Promise<void> | void;
}
/** Use-cases for conversation (session message) management. */
export interface ConversationService {
loadConversation(sessionId: string): Promise<Conversation> | Conversation;
saveConversation(conv: Conversation): Promise<void> | void;
addMessage(conv: Conversation, msg: ChatMessage): Promise<void> | void;
}
/** Use-cases for long-term memory management. */
export interface MemoryService {
listMemories(): Promise<string[]> | string[];
saveMemory(memory: Memory): Promise<void> | void;
deleteMemory(name: string): Promise<void> | void;
}
@@ -0,0 +1,50 @@
import { describe, expect, it } from "bun:test";
import { newSettings, resolveEffectiveModel } from "./settings.ts";
import { newAppConfig } from "./app_config.ts";
function claudeAppConfig() {
const cfg = newAppConfig();
cfg.providers.claude = {
api_base: "https://9router.example/v1",
api_key_env: "ANTHROPIC_API_KEY",
default_model: "claude-opus-5",
default_api_key: "sk-test",
};
cfg.default_provider = "claude";
cfg.default_model = "claude-opus-5";
return cfg;
}
describe("resolveEffectiveModel", () => {
it("claude provider uses opus model over stale settings model", () => {
const settings = { ...newSettings(), provider: "claude", model: "deepseek-v4-flash-free" };
expect(resolveEffectiveModel(settings, claudeAppConfig())).toBe("claude-opus-5");
});
it("non-claude provider uses settings model", () => {
const settings = { ...newSettings(), provider: "zen", model: "my-model" };
expect(resolveEffectiveModel(settings, newAppConfig())).toBe("my-model");
});
it("claude falls back to app default", () => {
const settings = { ...newSettings(), provider: "claude", model: "" };
const cfg = newAppConfig();
expect(resolveEffectiveModel(settings, cfg)).toBe(cfg.default_model);
});
});
describe("newSettings defaults", () => {
it("matches Rust defaults", () => {
const s = newSettings();
expect(s.internet_mode).toBe("Off");
expect(s.provider).toBe("zen");
expect(s.model).toBe("deepseek-v4-flash-free");
expect(s.review_max_lessons_per_run).toBe(5);
expect(s.adaptive_review_max_skip).toBe(3);
expect(s.verify_timeout_ms).toBe(30_000);
expect(s.workflow_max_concurrency).toBe(5);
expect(s.review_enabled).toBe(true);
expect(s.session_archive_enabled).toBe(true);
expect(s.hive_mind_node_timeout_ms).toBe(600_000);
});
});
+78
View File
@@ -0,0 +1,78 @@
/**
* Application settings domain entity. Mirrors `settings.rs`.
*/
import type { AppConfig } from "./app_config.ts";
/** Controls how much network access the agent is permitted. */
export type InternetMode = "Off" | "ReadOnly" | "Full";
export const InternetModeLiteral = {
Off: "Off" as const,
ReadOnly: "ReadOnly" as const,
Full: "Full" as const,
} satisfies Record<string, InternetMode>;
/** Grouped boolean feature toggles. */
export interface SettingsFlags {
review_enabled: boolean;
session_archive_enabled: boolean;
}
export function newSettingsFlags(): SettingsFlags {
return { review_enabled: true, session_archive_enabled: true };
}
const DEFAULT_HIVE_MIND_NODE_TIMEOUT_MS = 600_000;
/** Top-level application settings model (serialized to `settings.json`). */
export interface Settings {
internet_mode: InternetMode;
provider: string;
model: string;
api_keys: Record<string, string>;
max_tokens?: number;
temperature?: number;
review_max_lessons_per_run: number;
adaptive_review_max_skip: number;
verify_command?: string;
verify_timeout_ms: number;
workflow_max_concurrency: number;
review_enabled: boolean;
session_archive_enabled: boolean;
hive_mind_node_timeout_ms: number;
}
export function newSettings(): Settings {
return {
internet_mode: "Off",
provider: "zen",
model: "deepseek-v4-flash-free",
api_keys: {},
max_tokens: undefined,
temperature: undefined,
review_max_lessons_per_run: 5,
adaptive_review_max_skip: 3,
verify_command: undefined,
verify_timeout_ms: 30_000,
workflow_max_concurrency: 5,
review_enabled: true,
session_archive_enabled: true,
hive_mind_node_timeout_ms: DEFAULT_HIVE_MIND_NODE_TIMEOUT_MS,
};
}
/**
* Pick the effective model name for the main agent.
*
* When `settings.provider === "claude"`, the provider's `default_model`
* (or the app-level `default_model`) wins over a possibly-stale persisted
* `settings.model`. Otherwise the user's explicit `settings.model` is used.
*/
export function resolveEffectiveModel(settings: Settings, appConfig: AppConfig): string {
if (settings.provider === "claude") {
const m = appConfig.providers["claude"]?.default_model;
if (m) return m;
return appConfig.default_model;
}
return settings.model;
}
@@ -0,0 +1,74 @@
/**
* In-memory conversation state: message history plus the system prompt and
* model parameters used to drive the LLM. Mirrors `conversation.rs`.
*/
import type { ChatMessage, Role } from "./message.ts";
import { Roles, systemMessage, isSystemMessage } from "./message.ts";
/** A single conversation's message history and generation settings. */
export interface Conversation {
/** Ordered list of chat messages. */
messages: ChatMessage[];
/** System prompt prepended at request time (see `toApiMessages`). */
system_prompt: string;
/** Foreign key referencing the owning session. */
session_id: string;
/** Model identifier string, e.g. `"anthropic/claude-opus-4-8"`. */
model: string;
/** Optional cap on output tokens. */
max_tokens?: number;
/** Optional temperature (0.0 – 2.0). */
temperature?: number;
}
/** Default model used when a conversation is created. */
export const DEFAULT_CONVERSATION_MODEL = "anthropic/claude-opus-4-8";
/** Create an empty conversation with the given system prompt and session id. */
export function newConversation(systemPrompt: string, sessionId: string): Conversation {
return {
messages: [],
system_prompt: systemPrompt,
session_id: sessionId,
model: DEFAULT_CONVERSATION_MODEL,
max_tokens: undefined,
temperature: undefined,
};
}
/** Append a message to the conversation history. */
export function pushMessage(conv: Conversation, msg: ChatMessage): void {
conv.messages.push(msg);
}
/**
* Replace the system prompt and strip any prior `System`-role messages from
* history. The system prompt is re-injected fresh at request time via
* `toApiMessages`, so stale `System` messages would be redundant.
*/
export function rebuildSystem(conv: Conversation, newPrompt: string): void {
conv.system_prompt = newPrompt;
conv.messages = conv.messages.filter((m) => !isSystemMessage(m));
}
/**
* Build the message list to send to the LLM API, with the system prompt
* prepended at index 0.
*/
export function toApiMessages(conv: Conversation): ChatMessage[] {
return [systemMessage(conv.system_prompt), ...conv.messages];
}
/** Number of messages in history (excluding the synthesized system message). */
export function conversationLen(conv: Conversation): number {
return conv.messages.length;
}
/** Whether the conversation has no messages. */
export function isConversationEmpty(conv: Conversation): boolean {
return conv.messages.length === 0;
}
/** Re-export role bits for convenience. */
export type { ChatMessage, Role };
export { Roles, isSystemMessage };
+46
View File
@@ -0,0 +1,46 @@
/**
* Unified domain error types. Mirrors `apps/domain/src/error.rs`.
*
* Infrastructure adapters convert native errors into `DomainError`.
* Domain service layers wrap `DomainError` in their own `ServiceError`.
*/
export type DomainError = {
kind: "not_found" | "conflict" | "io" | "serde" | "invalid_id" | "other";
message: string;
source?: Error;
};
export function notFound(msg: string): DomainError {
return { kind: "not_found", message: msg };
}
export function conflict(msg: string): DomainError {
return { kind: "conflict", message: msg };
}
export function ioError(err: Error): DomainError {
return { kind: "io", message: `I/O error: ${err.message}`, source: err };
}
export function serdeError(msg: string): DomainError {
return { kind: "serde", message: `serialization error: ${msg}` };
}
export function invalidId(msg: string): DomainError {
return { kind: "invalid_id", message: `invalid id: ${msg}` };
}
export function otherError(msg: string): DomainError {
return { kind: "other", message: msg };
}
export function domainErrorToString(e: DomainError): string {
return e.message;
}
export function isDomainError(v: unknown): v is DomainError {
return (
typeof v === "object" &&
v !== null &&
"kind" in v &&
"message" in v &&
["not_found", "conflict", "io", "serde", "invalid_id", "other"].includes(
(v as DomainError).kind
)
);
}
+10
View File
@@ -0,0 +1,10 @@
/** Core domain module — shared entities, value objects, and provider types. */
export * from "./error.ts";
export * from "./message.ts";
export * from "./conversation.ts";
export * from "./tool_call.ts";
export * from "./provider.ts";
export * from "./tool_result.ts";
export * from "./usage.ts";
export * from "./store.ts";
export * from "./error.ts";
+103
View File
@@ -0,0 +1,103 @@
/**
* Chat message types shared across the entity layer.
*
* Provides `Role` (conversation participant), `ChatMessage` (a single
* message with optional tool-call metadata), and `ToolCall`/`ToolFunction`
* DTOs. Includes convenience constructors for each role.
*
* Mirrors Rust `apps/domain/src/core/message.rs` + `tool_call.rs` types.
*/
/* -------------------------------------------------------------------------- */
/* ToolCall / ToolFunction — defined here to avoid circular imports */
/* -------------------------------------------------------------------------- */
/** A single tool-call request emitted by the model in an assistant message. */
export interface ToolCall {
/** Unique identifier for this tool call (referenced by tool results). */
id: string;
/** Discriminator, e.g. `"function"`. Serialized as `type` on wire. */
type: string;
/** The function to invoke (name + arguments). */
function: ToolFunction;
}
/** The function name and raw arguments payload for a `ToolCall`. */
export interface ToolFunction {
/** The function/tool name to dispatch against. */
name: string;
/** Arguments as a JSON value. */
arguments: JsonValue;
}
/** JSON value type (subset matching serde_json::Value). */
export type JsonValue =
| null
| boolean
| number
| string
| JsonValue[]
| { [key: string]: JsonValue };
/* -------------------------------------------------------------------------- */
/* Role */
/* -------------------------------------------------------------------------- */
/** The conversation participant who authored a message. */
export type Role = "user" | "assistant" | "system" | "tool";
/** Canonical role string constants (lowercase, wire format). */
export const Roles = {
User: "user" as const,
Assistant: "assistant" as const,
System: "system" as const,
Tool: "tool" as const,
} satisfies Record<string, Role>;
/** Return the role as its lowercase wire string. */
export function roleAsString(role: Role): string {
return role;
}
/* -------------------------------------------------------------------------- */
/* ChatMessage */
/* -------------------------------------------------------------------------- */
/** A single message in a conversation, OpenAI/Anthropic chat-completion shaped. */
export interface ChatMessage {
/** Who sent this message (user, assistant, system, tool). */
role: Role;
/** The message text content. `null` for assistant messages that only carry tool calls. */
content: string | null;
/** Tool-call requests attached to an assistant message. Omitted on wire when absent. */
tool_calls?: ToolCall[];
/** For tool-role messages: the `id` of the `ToolCall` being responded to. */
tool_call_id?: string;
/** Optional function name for the tool invocation. */
name?: string;
}
/** Build a user-role message with text content. */
export function userMessage(content: string): ChatMessage {
return { role: Roles.User, content };
}
/** Build an assistant-role message with optional text response. */
export function assistantMessage(content: string | null): ChatMessage {
return { role: Roles.Assistant, content };
}
/** Build a system-role message with instruction text. */
export function systemMessage(content: string): ChatMessage {
return { role: Roles.System, content };
}
/** Build a tool-role result message referencing a prior tool call. */
export function toolResultMessage(toolCallId: string, content: string): ChatMessage {
return { role: Roles.Tool, content, tool_call_id: toolCallId };
}
/** Type guard: is this message a system-role message? */
export function isSystemMessage(m: ChatMessage): boolean {
return m.role === Roles.System;
}
@@ -0,0 +1,74 @@
import { describe, expect, it } from "bun:test";
import { SseParser } from "./provider.ts";
describe("SseParser", () => {
it("parses OpenAI-style content delta", () => {
const p = new SseParser();
const events = p.feed(`data: {"choices":[{"delta":{"content":"Hello"}}]}\n\n`);
expect(events).toContainEqual({ kind: "token", content: "Hello" });
});
it("parses Anthropic-style top-level delta content", () => {
const p = new SseParser();
const events = p.feed(`data: {"delta":{"content":"Hello"}}\n\ndata: [DONE]\n\n`);
expect(events).toContainEqual({ kind: "token", content: "Hello" });
expect(events).toContainEqual({ kind: "done" });
});
it("handles [DONE] sentinel", () => {
const p = new SseParser();
const events = p.feed(`data: [DONE]\n\n`);
expect(events).toEqual([{ kind: "done" }]);
});
it("captures tool-call deltas across multiple chunks", () => {
const p = new SseParser();
const first = p.feed(`data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","function":{"name":"read","arguments":"{\\"path\\""}}]}}]}\n\n`);
const second = p.feed(`data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":": \\"a.ts\\"}"}}]}}]}\n\n`);
const firstDeltas = first.filter((e): e is Extract<typeof e, { kind: "tool_call_delta" }> => e.kind === "tool_call_delta");
const secondDeltas = second.filter((e): e is Extract<typeof e, { kind: "tool_call_delta" }> => e.kind === "tool_call_delta");
expect(firstDeltas).toHaveLength(1);
expect(firstDeltas[0]?.id).toBe("call_1");
expect(firstDeltas[0]?.name).toBe("read");
expect(firstDeltas[0]?.arguments_delta).toBe('{"path"');
expect(secondDeltas).toHaveLength(1);
expect(secondDeltas[0]?.index).toBe(0);
expect(secondDeltas[0]?.arguments_delta).toBe(': "a.ts"}');
});
it("parses usage chunk", () => {
const p = new SseParser();
const events = p.feed(`data: {"usage":{"prompt_tokens":10,"completion_tokens":5,"total_tokens":15}}\n\n`);
expect(events).toContainEqual({ kind: "usage", prompt_tokens: 10, completion_tokens: 5, total_tokens: 15 });
});
it("handles message.stop as done", () => {
const p = new SseParser();
const events = p.feed(`event: message.stop\ndata: {"message":{"stop_reason":"end_turn"}}\n\n`);
expect(events).toContainEqual({ kind: "done" });
});
it("handles multi-line data frames by joining with newline", () => {
const p = new SseParser();
const events = p.feed(`data: {"delta":{"content":"line1"}}\ndata: line2\n\n`);
// Should parse the JSON only if valid; here line2 breaks it → no events, no throw.
expect(Array.isArray(events)).toBe(true);
});
it("keeps partial lines in buffer for next feed", () => {
const p = new SseParser();
p.feed(`data: {"choices":[{"delta":{"content":"Hel`);
const events = p.feed(`lo"}}]}\n\n`);
expect(events).toContainEqual({ kind: "token", content: "Hello" });
});
it("caps buffer growth to ~1MB by resetting oversized chunks", () => {
const p = new SseParser();
const big = "x".repeat(1_050_000);
const events = p.feed(`data: ${big}\n\n`);
expect(Array.isArray(events)).toBe(true);
});
});
+297
View File
@@ -0,0 +1,297 @@
/**
* Provider-facing DTOs: chat completion request, response, streaming types,
* and the SSE stream parser. Mirrors `apps/domain/src/core/provider.rs`.
*/
import type { ChatMessage, JsonValue, Role, ToolCall } from "./message.ts";
/* -------------------------------------------------------------------------- */
/* Chat request / response */
/* -------------------------------------------------------------------------- */
/** Outbound chat completion request body (OpenAI/Anthropic-compatible). */
export interface ChatRequest {
model: string;
messages: ChatMessage[];
max_tokens?: number;
temperature?: number;
tools?: ToolDef[];
tool_choice?: JsonValue;
stream?: boolean;
top_p?: number;
stop?: string[];
stream_options?: StreamOptions;
}
/** Streaming options; `include_usage` asks for a final usage chunk. */
export interface StreamOptions {
include_usage: boolean;
}
/** Wire format for a single tool definition sent to the provider. */
export interface ToolDef {
type: string;
function: ToolFunctionDef;
}
/** Name, description, and JSON schema parameters for a tool definition. */
export interface ToolFunctionDef {
name: string;
description: string;
parameters: JsonValue;
}
/** Non-streaming chat completion response. */
export interface ChatResponse {
id: string;
object?: string;
model: string;
choices: Choice[];
usage?: TokenUsage;
created?: number;
}
/** One completion candidate. */
export interface Choice {
index: number;
message?: ChatMessage;
delta?: Delta;
finish_reason?: string;
}
/** Incremental delta in a streaming SSE chunk. */
export interface Delta {
role?: Role;
content?: string;
tool_calls?: ToolCall[];
}
/** Token counts and optional cost breakdown. */
export interface TokenUsage {
prompt_tokens: number;
completion_tokens: number;
total_tokens: number;
prompt_tokens_cost?: number;
completion_tokens_cost?: number;
}
/* -------------------------------------------------------------------------- */
/* StreamEvent */
/* -------------------------------------------------------------------------- */
/** One atomic event from an LLM streaming response. */
export type StreamEvent =
| { kind: "token"; content: string }
| { kind: "reasoning"; content: string }
| {
kind: "tool_call_delta";
index: number;
id: string | null;
name: string | null;
arguments_delta: string;
}
| {
kind: "usage";
prompt_tokens: number;
completion_tokens: number;
total_tokens: number;
}
| { kind: "done" }
| { kind: "error"; message: string };
/* -------------------------------------------------------------------------- */
/* SSE Parser */
/* -------------------------------------------------------------------------- */
const MAX_BUFFER_SIZE = 1_048_576; // 1 MB
const MAX_TOOL_CALLS = 64;
/**
* Buffered SSE frame parser. Feed raw chunks → get `StreamEvent[]`.
*
* Handles both OpenAI-style (`choices[].delta`) and Anthropic-style
* (top-level `delta`, `content_block_delta`) formats.
*/
export class SseParser {
private buffer = "";
private eventType: string | null = null;
private dataLines: string[] = [];
/** Create a new parser with an empty buffer. */
constructor() {}
/**
* Feed a raw SSE chunk and return any completed events.
* Normalizes `\r\n`/`\r` → `\n`; caps buffer at 1 MB.
*/
feed(chunk: string): StreamEvent[] {
const normalized = chunk.replace(/\r\n/g, "\n").replace(/\r/g, "\n");
// Prevent unbounded buffer growth
if (this.buffer.length + normalized.length > MAX_BUFFER_SIZE) {
this.buffer = "";
this.eventType = null;
this.dataLines = [];
}
this.buffer += normalized;
const events: StreamEvent[] = [];
let idx: number;
while ((idx = this.buffer.indexOf("\n")) !== -1) {
const line = this.buffer.slice(0, idx).replace(/\r$/, "");
this.buffer = this.buffer.slice(idx + 1);
if (line === "") {
events.push(...this.flushEvent());
} else if (line.startsWith("event:")) {
// Handles both "event:foo" and "event: foo"
this.eventType = line.slice(6).trim();
} else if (line.startsWith("data:")) {
this.dataLines.push(line.slice(5).trimStart());
}
}
return events;
}
/**
* Flush the current buffered `data:` lines as one or more StreamEvents.
*/
private flushEvent(): StreamEvent[] {
const data = this.dataLines.join("\n");
this.dataLines = [];
const eventType = this.eventType ?? "";
this.eventType = null;
if (data === "") return [];
if (data === "[DONE]") return [{ kind: "done" }];
let value: JsonValue;
try {
value = JSON.parse(data) as JsonValue;
} catch {
return [];
}
if (typeof value !== "object" || value === null) return [];
const events: StreamEvent[] = [];
// Handle top-level usage chunk
const usage = getNested(value, "usage");
if (usage && typeof usage === "object" && usage !== null) {
const prompt = getNum(usage, "prompt_tokens") ?? 0;
const completion = getNum(usage, "completion_tokens") ?? 0;
const total = getNum(usage, "total_tokens") ?? prompt + completion;
events.push({ kind: "usage", prompt_tokens: prompt, completion_tokens: completion, total_tokens: total });
}
// Dispatch on event_type
if (eventType === "message.stop") {
events.push({ kind: "done" });
} else if (eventType === "message.delta" || eventType === "") {
events.push(...parseDeltaEvent(value));
} else if (eventType === "content_block_delta") {
events.push(...parseContentBlockDelta(value));
}
return events;
}
}
/* -------------------------------------------------------------------------- */
/* Delta parsing helpers */
/* -------------------------------------------------------------------------- */
function parseDeltaEvent(value: JsonValue): StreamEvent[] {
const events: StreamEvent[] = [];
// Try OpenAI-style: value.choices[0].delta
const choices = getNested(value, "choices");
if (Array.isArray(choices) && choices.length > 0) {
const choice = choices[0];
if (typeof choice === "object" && choice !== null) {
const delta = getNested(choice, "delta");
if (delta && typeof delta === "object" && delta !== null) {
// Content token
const content = getStr(delta, "content");
if (content !== null) events.push({ kind: "token", content });
// Reasoning token
const reasoning = getStr(delta, "reasoning_content");
if (reasoning !== null) events.push({ kind: "reasoning", content: reasoning });
// Tool calls — iterate ALL entries
const toolCalls = getNested(delta, "tool_calls");
if (Array.isArray(toolCalls)) {
for (const tc of toolCalls) {
if (typeof tc !== "object" || tc === null) continue;
const rawIndex = getNum(tc, "index") ?? 0;
const index = Math.min(Math.max(rawIndex, 0), MAX_TOOL_CALLS - 1);
events.push({
kind: "tool_call_delta",
index,
id: getStr(tc, "id"),
name: (() => {
const fn = getNested(tc, "function");
return fn && typeof fn === "object" ? getStr(fn, "name") : null;
})(),
arguments_delta: (() => {
const fn = getNested(tc, "function");
return (fn && typeof fn === "object" ? getStr(fn, "arguments") : null) ?? "";
})(),
});
}
}
// Finish reason
const finishReason = getStr(choice, "finish_reason");
if (finishReason === "stop" || finishReason === "tool_calls") {
events.push({ kind: "done" });
}
}
}
return events;
}
// Try Anthropic-style: top-level value.delta.content
const delta = getNested(value, "delta");
if (delta && typeof delta === "object" && delta !== null) {
const content = getStr(delta, "content");
if (content !== null) events.push({ kind: "token", content });
}
return events;
}
function parseContentBlockDelta(value: JsonValue): StreamEvent[] {
const events: StreamEvent[] = [];
const delta = getNested(value, "delta");
if (delta && typeof delta === "object" && delta !== null) {
const text = getStr(delta, "text");
if (text !== null) events.push({ kind: "token", content: text });
const reasoning = getStr(delta, "reasoning_content");
if (reasoning !== null) events.push({ kind: "reasoning", content: reasoning });
}
return events;
}
/* -------------------------------------------------------------------------- */
/* JSON helpers (typed access on JsonValue) */
/* -------------------------------------------------------------------------- */
function getNested(obj: JsonValue, key: string): JsonValue | undefined {
if (typeof obj === "object" && obj !== null && !Array.isArray(obj)) {
return (obj as Record<string, JsonValue>)[key];
}
return undefined;
}
function getStr(obj: JsonValue, key: string): string | null {
const v = getNested(obj, key);
return typeof v === "string" ? v : null;
}
function getNum(obj: JsonValue, key: string): number | null {
const v = getNested(obj, key);
return typeof v === "number" ? v : null;
}
+63
View File
@@ -0,0 +1,63 @@
/**
* Filesystem layout for zesdex's persistent and scratch data directories.
* Mirrors `apps/domain/src/core/store.rs`.
*/
import * as os from "node:os";
import * as path from "node:path";
/** Resolved paths for all data directories zesdex reads from and writes to. */
export interface Store {
base_dir: string;
scratch_root: string;
memory_dir: string;
session_images_dir: string;
download_dir: string;
}
function env(name: string): string | undefined {
if (typeof process !== "undefined" && process.env) return process.env[name];
return undefined;
}
/** Compute the standard set of zesdex data directory paths. */
export function newStore(): Store {
let base: string;
const dataHome = env("XDG_DATA_HOME");
const home = env("HOME");
if (dataHome) {
base = path.join(dataHome, "zesdex");
} else if (home) {
base = path.join(home, ".local", "share", "zesdex");
} else {
base = path.join(".local", "share", "zesdex");
}
const scratch = path.join(os.tmpdir(), "zesdex-scratch");
return {
base_dir: base,
scratch_root: scratch,
memory_dir: path.join(base, "memory"),
session_images_dir: path.join(base, "session-images"),
download_dir: path.join(base, "downloads"),
};
}
/**
* Create all store directories if missing.
* Return: `null` on success, or the error message on the first failure.
*/
export async function ensureStoreDirs(store: Store): Promise<string | null> {
for (const dir of [
store.base_dir,
store.memory_dir,
store.scratch_root,
store.session_images_dir,
store.download_dir,
]) {
try {
await import("node:fs/promises").then((fs) => fs.mkdir(dir, { recursive: true }));
} catch (err) {
return (err as Error).message;
}
}
return null;
}
@@ -0,0 +1,53 @@
import { describe, expect, it } from "bun:test";
import { repairJson, sanitizeToolArguments } from "./tool_call.ts";
describe("repairJson", () => {
it("closes an unclosed object", () => {
expect(repairJson('{"a": 1')).toBe('{"a": 1}');
});
it("closes unclosed object and array in reverse nesting order", () => {
expect(repairJson('{"a": [1, 2')).toBe('{"a": [1, 2]}');
});
it("closes an unclosed string", () => {
expect(repairJson('{"a": "hello')).toBe('{"a": "hello"}');
});
it("handles unclosed trailing escape by popping it", () => {
expect(repairJson('{"a": "text\\')).toBe('{"a": "text"}');
});
it("leaves already-valid JSON unchanged", () => {
expect(repairJson('{"a": [1,2], "b": {"c": true}}')).toBe('{"a": [1,2], "b": {"c": true}}');
});
});
describe("sanitizeToolArguments", () => {
it("passes objects through unchanged", () => {
const obj = { path: "src/main.ts", mode: "append" };
expect(sanitizeToolArguments(obj)).toBe(obj);
});
it("parses a string-encoded JSON object", () => {
expect(sanitizeToolArguments('{"path": "a.ts"}')).toEqual({ path: "a.ts" });
});
it("strips control chars and reparses", () => {
expect(sanitizeToolArguments('{\n "path": "a.ts",\n "x": 1\n}')).toEqual({ path: "a.ts", x: 1 });
});
it("repairs truncated string JSON", () => {
expect(sanitizeToolArguments('{"path": "a.ts", "content": "partial')).toEqual({
path: "a.ts",
content: "partial",
});
});
it("wraps unparseable strings in _raw on last resort", () => {
const out = sanitizeToolArguments("not json at all ]}}");
// repair attempts can't fix this; must produce an object with _raw.
expect(typeof out).toBe("object");
expect(out).toHaveProperty("_raw");
});
});
@@ -0,0 +1,99 @@
/**
* Tool-call argument utilities.
*
* `sanitize_tool_arguments` and `repair_json` — ported from
* `apps/domain/src/core/tool_call.rs`. ToolCall/ToolFunction types live
* in `message.ts` to avoid circular imports.
*/
import type { JsonValue } from "./message.ts";
/* -------------------------------------------------------------------------- */
/* JSON truncation repair */
/* -------------------------------------------------------------------------- */
/**
* Repair truncated JSON by closing open strings, braces and brackets.
*
* Single-pass character scan tracking string/escape state with a LIFO stack
* for `{`/`[` → append missing `"`, `]`, `}` in reverse nesting order.
*/
export function repairJson(s: string): string {
const stack: string[] = [];
let inString = false;
let prevWasBackslash = false;
let endsWithUnclosedEscape = false;
for (const c of s) {
if (prevWasBackslash) {
prevWasBackslash = false;
endsWithUnclosedEscape = false;
continue;
}
if (c === "\\" && inString) {
prevWasBackslash = true;
endsWithUnclosedEscape = true;
continue;
}
endsWithUnclosedEscape = false;
if (c === '"') {
inString = !inString;
continue;
}
if (inString) continue;
if (c === "{" || c === "[") stack.push(c);
else if (c === "}" || c === "]") stack.pop();
}
let result = s;
if (endsWithUnclosedEscape) result = result.slice(0, -1);
if (inString) result += '"';
for (let i = stack.length - 1; i >= 0; i--) {
if (stack[i] === "{") result += "}";
else if (stack[i] === "[") result += "]";
}
return result;
}
/* -------------------------------------------------------------------------- */
/* Argument sanitization */
/* -------------------------------------------------------------------------- */
function tryParseJson(s: string): JsonValue | null {
try { return JSON.parse(s) as JsonValue; } catch { return null; }
}
function isControl(c: string): boolean {
const code = c.codePointAt(0) ?? 0;
return code >= 0 && code <= 31;
}
/**
* Normalize tool-call arguments into a JSON value.
*
* Handles: string-encoded JSON → parsed object; control character stripping;
* truncated JSON repair; total failure → `{ _raw, _parse_error }` wrapper.
*/
export function sanitizeToolArguments(args: JsonValue): JsonValue {
if (typeof args !== "string") return args;
const s = args;
// Attempt 1: direct parse.
const direct = tryParseJson(s);
if (direct !== null) return direct;
// Attempt 2: strip control chars (0x00-0x1F except \t, \n, \r).
const cleaned = [...s].filter((c) => !isControl(c) || c === "\t" || c === "\n" || c === "\r").join("");
if (cleaned.length !== s.length) {
const p = tryParseJson(cleaned);
if (p !== null) return p;
}
// Attempt 3: repair truncated JSON and retry.
const input = cleaned.length === s.length ? s : cleaned;
const repaired = repairJson(input);
const rp = tryParseJson(repaired);
if (rp !== null) return rp;
// Fallback: wrap raw in object with parse error.
return { _raw: s, _parse_error: "failed to parse tool argument string" };
}
@@ -0,0 +1,26 @@
/**
* Record of one completed tool invocation. Mirrors `tool_result.rs`.
*/
export interface ToolCallResult {
/** The `id` of the `ToolCall` this result responds to. */
tool_call_id: string;
/** The name of the tool that was invoked. */
tool_name: string;
/** The text output produced by the tool (or error message). */
output: string;
/** Whether the tool exited with an error. */
is_error: boolean;
/** Wall-clock execution duration in milliseconds. */
duration_ms: number;
}
/** Create a new tool call result. */
export function newToolCallResult(
tool_call_id: string,
tool_name: string,
output: string,
is_error: boolean,
duration_ms: number,
): ToolCallResult {
return { tool_call_id, tool_name, output, is_error, duration_ms };
}
+33
View File
@@ -0,0 +1,33 @@
/**
* Token usage accounting shared by streaming and non-streaming responses.
* Mirrors `apps/domain/src/core/usage.rs`.
*/
export interface UsageStats {
/** Total tokens consumed as input (prompt). */
tokens_in: number;
/** Total tokens generated as output (completion). */
tokens_out: number;
/** Most recent call's input tokens (for live display). */
last_tokens_in: number;
/** Most recent call's output tokens (for live display). */
last_tokens_out: number;
/** Total number of LLM API calls made this session. */
api_calls: number;
/** Tokens consumed by auto-review subagent calls. */
review_tokens: number;
/** Total wall-clock time spent on LLM API calls (milliseconds). */
total_ms: number;
}
/** Create a new `UsageStats` with all counters zeroed. */
export function newUsageStats(): UsageStats {
return {
tokens_in: 0,
tokens_out: 0,
last_tokens_in: 0,
last_tokens_out: 0,
api_calls: 0,
review_tokens: 0,
total_ms: 0,
};
}
+10
View File
@@ -0,0 +1,10 @@
/**
* Zesdex Domain Layer — pure types, value objects, and port interfaces.
* 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";
export * from "./workflow/mod.ts";
+7
View File
@@ -0,0 +1,7 @@
/** Subagent domain models. Mirrors `subagent/mod.rs`. */
/**
* Access tier for subagent tool permissions. Cumulative: Write includes Read,
* Full includes Write.
*/
export type AccessTier = "Read" | "Write" | "Full";
+37
View File
@@ -0,0 +1,37 @@
/** Workflow and Hive-mind domain models. Mirrors `workflow/mod.rs`. */
/** A single phase in a parsed workflow script. */
export interface WorkflowPhase {
name: string;
directive: string;
}
/** A parsed workflow script with named phases. */
export interface WorkflowScript {
name: string;
phases: WorkflowPhase[];
}
/** A directive for a single processing node in the hive mind. */
export interface NodeDirective {
directive: string;
access_tier: string;
}
/** A cognitive cycle plan — ordered cycles of parallel node directives. */
export interface CognitiveCyclePlan {
cycles: NodeDirective[][];
}
/** A single cycle in a cognitive cycle plan. */
export interface CognitiveCycle {
index: number;
directives: NodeDirective[];
}
/** Output from a single hive-mind node after a cycle completes. */
export interface NodeOutput {
id: string;
directive: string;
output: string;
}