diff --git a/.gitignore b/.gitignore index c4df330..a1e26a7 100644 --- a/.gitignore +++ b/.gitignore @@ -2,7 +2,6 @@ target/ .env .claude/settings.local.json node_modules/ -package.json package-lock.json .superpowers/ docs/lesson/ diff --git a/apps/packages/domain/package.json b/apps/packages/domain/package.json new file mode 100644 index 0000000..cae1842 --- /dev/null +++ b/apps/packages/domain/package.json @@ -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" + } +} \ No newline at end of file diff --git a/apps/packages/domain/src/agent/defaults.ts b/apps/packages/domain/src/agent/defaults.ts new file mode 100644 index 0000000..4a2f6a1 --- /dev/null +++ b/apps/packages/domain/src/agent/defaults.ts @@ -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; \ No newline at end of file diff --git a/apps/packages/domain/src/agent/index.ts b/apps/packages/domain/src/agent/index.ts new file mode 100644 index 0000000..607c8b5 --- /dev/null +++ b/apps/packages/domain/src/agent/index.ts @@ -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"; \ No newline at end of file diff --git a/apps/packages/domain/src/agent/mod.ts b/apps/packages/domain/src/agent/mod.ts new file mode 100644 index 0000000..817e6b4 --- /dev/null +++ b/apps/packages/domain/src/agent/mod.ts @@ -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; + +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>`). */ +export interface TurnEventSink { + /** Append one event. */ + push(event: TurnEvent): void; + /** Drain all pending events (FIFO) and return them. */ + drain(): TurnEvent[]; +} \ No newline at end of file diff --git a/apps/packages/domain/src/agent/progress.ts b/apps/packages/domain/src/agent/progress.ts new file mode 100644 index 0000000..7e5d6f2 --- /dev/null +++ b/apps/packages/domain/src/agent/progress.ts @@ -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 | "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 }; +} \ No newline at end of file diff --git a/apps/packages/domain/src/agent/prompt.ts b/apps/packages/domain/src/agent/prompt.ts new file mode 100644 index 0000000..27b3914 --- /dev/null +++ b/apps/packages/domain/src/agent/prompt.ts @@ -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.`; +} \ No newline at end of file diff --git a/apps/packages/domain/src/auth/commands.ts b/apps/packages/domain/src/auth/commands.ts new file mode 100644 index 0000000..079bc28 --- /dev/null +++ b/apps/packages/domain/src/auth/commands.ts @@ -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 }; +} \ No newline at end of file diff --git a/apps/packages/domain/src/auth/error.ts b/apps/packages/domain/src/auth/error.ts new file mode 100644 index 0000000..648a060 --- /dev/null +++ b/apps/packages/domain/src/auth/error.ts @@ -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; + } +} \ No newline at end of file diff --git a/apps/packages/domain/src/auth/index.ts b/apps/packages/domain/src/auth/index.ts new file mode 100644 index 0000000..4e85570 --- /dev/null +++ b/apps/packages/domain/src/auth/index.ts @@ -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"; \ No newline at end of file diff --git a/apps/packages/domain/src/auth/oauth.ts b/apps/packages/domain/src/auth/oauth.ts new file mode 100644 index 0000000..7b07b17 --- /dev/null +++ b/apps/packages/domain/src/auth/oauth.ts @@ -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], + }; +} \ No newline at end of file diff --git a/apps/packages/domain/src/auth/repository.ts b/apps/packages/domain/src/auth/repository.ts new file mode 100644 index 0000000..0ce4b52 --- /dev/null +++ b/apps/packages/domain/src/auth/repository.ts @@ -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 `/sessions/`. */ + listSessions(baseDir: string): Promise | Session[]; + /** Load a single session by id. */ + loadSession(baseDir: string, id: SessionId): Promise | Session; + /** Save a session's metadata to disk. */ + saveSession(baseDir: string, session: Session): Promise | void; + /** Delete a session directory and all its contents. */ + deleteSession(baseDir: string, id: SessionId): Promise | void; +} + +/** Repository for per-session PID-file advisory locks. */ +export interface SessionLockRepository { + /** Try to acquire the lock. `true` if acquired, `false` if a live process holds it. */ + tryLock(sessionDir: string): Promise | boolean; + /** Release the lock. */ + unlock(sessionDir: string): Promise | void; + /** Check whether a process with the given PID is alive. */ + isAlive(pid: number): boolean; +} + +/** Repository for persisting and loading OAuth tokens. */ +export interface OAuthRepository { + /** Persist an OAuth token to a JSON file. */ + saveToken(path: string, token: OAuthToken): Promise | void; + /** Load an OAuth token, returning `null` if the file does not exist. */ + loadToken(path: string): Promise | OAuthToken | null; +} + diff --git a/apps/packages/domain/src/auth/service.ts b/apps/packages/domain/src/auth/service.ts new file mode 100644 index 0000000..7a3a37d --- /dev/null +++ b/apps/packages/domain/src/auth/service.ts @@ -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; + /** List all available sessions. */ + listAll(): Promise | Session[]; + /** Archive a session by id (sets `archived = true`). */ + archiveSession(id: SessionId): Promise | void; +} + +/** OAuth flow use-case boundary. */ +export interface OAuthService { + /** + * Start an OAuth authorization-code + PKCE flow. Returns `{ authUrl, state }`: + * the URL to send the user to, and the CSRF state token that must be passed + * back into `completeFlow` unchanged. + */ + startFlow(config: OAuthConfig, redirectUri: string): Promise<{ authUrl: string; state: string }> | { authUrl: string; state: string }; + /** + * Complete the flow: validate `state`, then exchange `code` for a token. + * Validates `state` against the persisted value (CSRF check). + */ + completeFlow(config: OAuthConfig, redirectUri: string, code: string, state: string): Promise | OAuthToken; + /** Retrieve the currently stored OAuth token (if any). */ + getToken(): Promise | OAuthToken | null; +} \ No newline at end of file diff --git a/apps/packages/domain/src/auth/session.ts b/apps/packages/domain/src/auth/session.ts new file mode 100644 index 0000000..a39cd6f --- /dev/null +++ b/apps/packages/domain/src/auth/session.ts @@ -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 `/sessions/`. */ +export function sessionDir(session: Session, baseDir: string): string { + return path.join(baseDir, "sessions", session.id); +} + +/** Compute this session's `conversation.json` path. */ +export function conversationPath(session: Session, baseDir: string): string { + return path.join(sessionDir(session, baseDir), "conversation.json"); +} \ No newline at end of file diff --git a/apps/packages/domain/src/auth/session_id.test.ts b/apps/packages/domain/src/auth/session_id.test.ts new file mode 100644 index 0000000..9ce3a96 --- /dev/null +++ b/apps/packages/domain/src/auth/session_id.test.ts @@ -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"); + }); +}); \ No newline at end of file diff --git a/apps/packages/domain/src/auth/session_id.ts b/apps/packages/domain/src/auth/session_id.ts new file mode 100644 index 0000000..0ede37e --- /dev/null +++ b/apps/packages/domain/src/auth/session_id.ts @@ -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); +} \ No newline at end of file diff --git a/apps/packages/domain/src/auth/session_lock.ts b/apps/packages/domain/src/auth/session_lock.ts new file mode 100644 index 0000000..4e7cbdf --- /dev/null +++ b/apps/packages/domain/src/auth/session_lock.ts @@ -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//exe` and compare to self, plus + * `kill(pid, 0)`) → stale locks overwritten atomically (temp + rename + fsync). + */ +import * as fs from "node:fs"; +import * as path from "node:path"; + +/** A PID-file lock (`/.lock`) tied to the current process. */ +export class SessionLock { + readonly path: string; + readonly pid: number; + + constructor(sessionDir: string) { + this.path = path.join(sessionDir, ".lock"); + this.pid = process.pid; + } + + /** Attempt to acquire the lock. Returns `true` on success, `false` if held by a live process. */ + tryLock(): boolean { + // Phase 1: atomic create (O_CREAT|O_EXCL equivalent via flag 'wx'). + try { + const fd = fs.openSync(this.path, "wx"); + try { + fs.writeFileSync(fd, String(this.pid)); + fs.fsyncSync(fd); + } finally { + fs.closeSync(fd); + } + return true; + } catch (e) { + const code = (e as NodeJS.ErrnoException).code; + if (code !== "EEXIST") throw e; + // Lock file exists — check staleness. + } + + // Phase 2: read owning PID and check liveness. + const content = fs.readFileSync(this.path, "utf8").trim(); + if (content !== "") { + const pid = Number(content); + if (Number.isFinite(pid) && pid > 0) { + if (pidIsAlive(pid)) return false; + } + } + + // Phase 3: stale lock — overwrite atomically. + const tmp = this.path + ".tmp"; + const fd = fs.openSync(tmp, "w"); + try { + fs.writeFileSync(fd, String(this.pid)); + fs.fsyncSync(fd); + } finally { + fs.closeSync(fd); + } + fs.renameSync(tmp, this.path); + // Best-effort parent dir fsync. + const parent = path.dirname(this.path); + try { + const dirFd = fs.openSync(parent, "r"); + try { + fs.fsyncSync(dirFd); + } finally { + fs.closeSync(dirFd); + } + } catch { + /* best-effort */ + } + return true; + } + + /** Explicitly release the lock by removing the lock file. */ + unlock(): void { + try { + fs.unlinkSync(this.path); + } catch { + /* ignore */ + } + } +} + +/** + * Check whether a process with the given PID is alive. Conservative on + * non-Unix: returns true. + */ +export function pidIsAlive(pid: number): boolean { + if (process.platform === "win32") return true; + try { + const selfExe = fs.readlinkSync("/proc/self/exe"); + let target: string; + try { + target = fs.readlinkSync(`/proc/${pid}/exe`); + } catch { + return false; + } + if (target !== selfExe) return false; + // kill(pid, 0): throws if the process is gone / not permitted. + try { + process.kill(pid, 0); + } catch { + return false; + } + // Re-check to close the TOCTOU window. + try { + return fs.readlinkSync(`/proc/${pid}/exe`) === selfExe; + } catch { + return false; + } + } catch { + // No /proc (macOS) — fall back to kill(pid,0) only. + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } + } +} \ No newline at end of file diff --git a/apps/packages/domain/src/cms/app_config.ts b/apps/packages/domain/src/cms/app_config.ts new file mode 100644 index 0000000..aada52d --- /dev/null +++ b/apps/packages/domain/src/cms/app_config.ts @@ -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; + model_roles: Record; + 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, + }; +} \ No newline at end of file diff --git a/apps/packages/domain/src/cms/commands.ts b/apps/packages/domain/src/cms/commands.ts new file mode 100644 index 0000000..bdd04d7 --- /dev/null +++ b/apps/packages/domain/src/cms/commands.ts @@ -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; + 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; + 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[]; +} \ No newline at end of file diff --git a/apps/packages/domain/src/cms/conversation.ts b/apps/packages/domain/src/cms/conversation.ts new file mode 100644 index 0000000..e79a3c9 --- /dev/null +++ b/apps/packages/domain/src/cms/conversation.ts @@ -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"; \ No newline at end of file diff --git a/apps/packages/domain/src/cms/edit_log.ts b/apps/packages/domain/src/cms/edit_log.ts new file mode 100644 index 0000000..e262b91 --- /dev/null +++ b/apps/packages/domain/src/cms/edit_log.ts @@ -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(); + } + } +} \ No newline at end of file diff --git a/apps/packages/domain/src/cms/error.ts b/apps/packages/domain/src/cms/error.ts new file mode 100644 index 0000000..c81d982 --- /dev/null +++ b/apps/packages/domain/src/cms/error.ts @@ -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; + } +} \ No newline at end of file diff --git a/apps/packages/domain/src/cms/index.ts b/apps/packages/domain/src/cms/index.ts new file mode 100644 index 0000000..9f78c43 --- /dev/null +++ b/apps/packages/domain/src/cms/index.ts @@ -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"; \ No newline at end of file diff --git a/apps/packages/domain/src/cms/memory.ts b/apps/packages/domain/src/cms/memory.ts new file mode 100644 index 0000000..49f6b81 --- /dev/null +++ b/apps/packages/domain/src/cms/memory.ts @@ -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); +} \ No newline at end of file diff --git a/apps/packages/domain/src/cms/repository.ts b/apps/packages/domain/src/cms/repository.ts new file mode 100644 index 0000000..bba103d --- /dev/null +++ b/apps/packages/domain/src/cms/repository.ts @@ -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; + save(baseDir: string, settings: Settings): Promise | void; +} + +/** Persistence contract for `AppConfig`. */ +export interface AppConfigRepository { + load(baseDir: string): Promise | AppConfig; + save(baseDir: string, config: AppConfig): Promise | void; +} + +/** Persistence contract for `Conversation`. */ +export interface ConversationRepository { + load(sessionDir: string): Promise | Conversation; + save(sessionDir: string, conversation: Conversation): Promise | void; +} + +/** Persistence contract for `Memory`. */ +export interface MemoryRepository { + list(memoryDir: string): Promise | string[]; + load(memoryDir: string, name: string): Promise | Memory; + save(memoryDir: string, memory: Memory): Promise | void; + delete(memoryDir: string, name: string): Promise | void; +} + +/** Persistence contract for rewind-snapshot binary blobs. */ +export interface RewindBlobRepository { + storeBlob(sessionDir: string, blobKey: string, data: Uint8Array, mimeType?: string): Promise | void; + retrieveBlob(sessionDir: string, blobKey: string): Promise | Uint8Array | null; + listBlobKeys(sessionDir: string): Promise | string[]; +} + +/** Persistence contract for `EditLog`. */ +export interface EditLogRepository { + open(sessionDir: string): Promise | EditLog; + append(sessionDir: string, log: EditLog, entry: EditLogEntry): Promise | void; + entries(log: EditLog): EditLogEntry[]; +} + diff --git a/apps/packages/domain/src/cms/service.ts b/apps/packages/domain/src/cms/service.ts new file mode 100644 index 0000000..1c0de05 --- /dev/null +++ b/apps/packages/domain/src/cms/service.ts @@ -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; + saveSettings(settings: Settings): Promise | void; + updateProvider(name: string, config: ProviderConfig): Promise | void; +} + +/** Use-cases for conversation (session message) management. */ +export interface ConversationService { + loadConversation(sessionId: string): Promise | Conversation; + saveConversation(conv: Conversation): Promise | void; + addMessage(conv: Conversation, msg: ChatMessage): Promise | void; +} + +/** Use-cases for long-term memory management. */ +export interface MemoryService { + listMemories(): Promise | string[]; + saveMemory(memory: Memory): Promise | void; + deleteMemory(name: string): Promise | void; +} \ No newline at end of file diff --git a/apps/packages/domain/src/cms/settings.test.ts b/apps/packages/domain/src/cms/settings.test.ts new file mode 100644 index 0000000..42d9493 --- /dev/null +++ b/apps/packages/domain/src/cms/settings.test.ts @@ -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); + }); +}); \ No newline at end of file diff --git a/apps/packages/domain/src/cms/settings.ts b/apps/packages/domain/src/cms/settings.ts new file mode 100644 index 0000000..c07c3b8 --- /dev/null +++ b/apps/packages/domain/src/cms/settings.ts @@ -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; + +/** 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; + 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; +} \ No newline at end of file diff --git a/apps/packages/domain/src/core/conversation.ts b/apps/packages/domain/src/core/conversation.ts new file mode 100644 index 0000000..168a410 --- /dev/null +++ b/apps/packages/domain/src/core/conversation.ts @@ -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 }; \ No newline at end of file diff --git a/apps/packages/domain/src/core/error.ts b/apps/packages/domain/src/core/error.ts new file mode 100644 index 0000000..bdf6ef4 --- /dev/null +++ b/apps/packages/domain/src/core/error.ts @@ -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 + ) + ); +} \ No newline at end of file diff --git a/apps/packages/domain/src/core/index.ts b/apps/packages/domain/src/core/index.ts new file mode 100644 index 0000000..d0e1bf5 --- /dev/null +++ b/apps/packages/domain/src/core/index.ts @@ -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"; \ No newline at end of file diff --git a/apps/packages/domain/src/core/message.ts b/apps/packages/domain/src/core/message.ts new file mode 100644 index 0000000..b946761 --- /dev/null +++ b/apps/packages/domain/src/core/message.ts @@ -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; + +/** 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; +} \ No newline at end of file diff --git a/apps/packages/domain/src/core/provider.test.ts b/apps/packages/domain/src/core/provider.test.ts new file mode 100644 index 0000000..408bd5b --- /dev/null +++ b/apps/packages/domain/src/core/provider.test.ts @@ -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 => e.kind === "tool_call_delta"); + const secondDeltas = second.filter((e): e is Extract => 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); + }); +}); \ No newline at end of file diff --git a/apps/packages/domain/src/core/provider.ts b/apps/packages/domain/src/core/provider.ts new file mode 100644 index 0000000..da7af34 --- /dev/null +++ b/apps/packages/domain/src/core/provider.ts @@ -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)[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; +} \ No newline at end of file diff --git a/apps/packages/domain/src/core/store.ts b/apps/packages/domain/src/core/store.ts new file mode 100644 index 0000000..39cbc5f --- /dev/null +++ b/apps/packages/domain/src/core/store.ts @@ -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 { + 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; +} \ No newline at end of file diff --git a/apps/packages/domain/src/core/tool_call.test.ts b/apps/packages/domain/src/core/tool_call.test.ts new file mode 100644 index 0000000..8c1090e --- /dev/null +++ b/apps/packages/domain/src/core/tool_call.test.ts @@ -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"); + }); +}); \ No newline at end of file diff --git a/apps/packages/domain/src/core/tool_call.ts b/apps/packages/domain/src/core/tool_call.ts new file mode 100644 index 0000000..2664a2e --- /dev/null +++ b/apps/packages/domain/src/core/tool_call.ts @@ -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" }; +} \ No newline at end of file diff --git a/apps/packages/domain/src/core/tool_result.ts b/apps/packages/domain/src/core/tool_result.ts new file mode 100644 index 0000000..622d3d1 --- /dev/null +++ b/apps/packages/domain/src/core/tool_result.ts @@ -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 }; +} \ No newline at end of file diff --git a/apps/packages/domain/src/core/usage.ts b/apps/packages/domain/src/core/usage.ts new file mode 100644 index 0000000..e6b574b --- /dev/null +++ b/apps/packages/domain/src/core/usage.ts @@ -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, + }; +} \ No newline at end of file diff --git a/apps/packages/domain/src/index.ts b/apps/packages/domain/src/index.ts new file mode 100644 index 0000000..35d85de --- /dev/null +++ b/apps/packages/domain/src/index.ts @@ -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"; \ No newline at end of file diff --git a/apps/packages/domain/src/subagent/mod.ts b/apps/packages/domain/src/subagent/mod.ts new file mode 100644 index 0000000..996d569 --- /dev/null +++ b/apps/packages/domain/src/subagent/mod.ts @@ -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"; \ No newline at end of file diff --git a/apps/packages/domain/src/workflow/mod.ts b/apps/packages/domain/src/workflow/mod.ts new file mode 100644 index 0000000..1d60cb8 --- /dev/null +++ b/apps/packages/domain/src/workflow/mod.ts @@ -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; +} \ No newline at end of file diff --git a/bun.lock b/bun.lock new file mode 100644 index 0000000..a527ba0 --- /dev/null +++ b/bun.lock @@ -0,0 +1,34 @@ +{ + "lockfileVersion": 1, + "configVersion": 1, + "workspaces": { + "": { + "name": "zesdex", + "devDependencies": { + "@types/bun": "^1.2.0", + "typescript": "^5.7.0", + }, + }, + "apps/packages/domain": { + "name": "@zesdex/domain", + "version": "1.21.2", + "devDependencies": { + "@types/bun": "^1.2.0", + "typescript": "^5.7.0", + }, + }, + }, + "packages": { + "@types/bun": ["@types/bun@1.4.0", "", { "dependencies": { "bun-types": "1.4.0" } }, "sha512-K+lZULY23vRgK/CfTjFIV+tyifaNdSMlPh9j+6mQ/cLfpOznLyAuzgV/JQysyECpkBQLVMSyvjlr2fBUSA9wFQ=="], + + "@types/node": ["@types/node@26.4.1", "", { "dependencies": { "undici-types": "~8.3.0" } }, "sha512-k97ENvZWtvA6yqz5/FS6a7duDgOPEeOQOc2iKS/nY6mX6qJUKtLnWzQS+Xj6tXweyj6ZcTAK2Qecetnvi9nCLA=="], + + "@zesdex/domain": ["@zesdex/domain@workspace:apps/packages/domain"], + + "bun-types": ["bun-types@1.4.0", "", { "dependencies": { "@types/node": "*" } }, "sha512-iIKw23BspnQQYd3prITOBxeUsxBHnwzX6YJfGMuNOZzeNcMmVqzIIVGRm1l69ogaPQmb4wB6BN8mA5bE9YuC5Q=="], + + "typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="], + + "undici-types": ["undici-types@8.3.0", "", {}, "sha512-j375ScV60dom+YkPFIfTLcOiPxkN/buHz5GobjLhixFuANaNs3C9l4GmrWqejgXWJ7BbJcFYpTEUkS1Ge8bpZQ=="], + } +} diff --git a/package.json b/package.json new file mode 100644 index 0000000..121e8b1 --- /dev/null +++ b/package.json @@ -0,0 +1,25 @@ +{ + "name": "zesdex", + "version": "1.21.2", + "private": true, + "type": "module", + "packageManager": "bun@1.3.14", + "description": "Zesdex — AI-native autonomous coding agent (TypeScript/Bun rewrite of the Rust project)", + "license": "MIT", + "workspaces": [ + "apps/packages/*", + "apps/interfaces/*" + ], + "scripts": { + "build": "bun run --cwd apps/interfaces/cli build", + "check": "tsc --noEmit", + "test": "bun test", + "format": "bun format", + "cli": "bun run --cwd apps/interfaces/cli cli", + "bootstrap": "bun run --cwd apps/interfaces/cli bootstrap" + }, + "devDependencies": { + "@types/bun": "^1.2.0", + "typescript": "^5.7.0" + } +} \ No newline at end of file diff --git a/tsconfig.json b/tsconfig.json new file mode 100644 index 0000000..668e5a1 --- /dev/null +++ b/tsconfig.json @@ -0,0 +1,24 @@ +{ + "compilerOptions": { + "target": "ESNext", + "module": "ESNext", + "moduleResolution": "bundler", + "moduleDetection": "force", + "lib": ["ESNext", "DOM"], + "types": ["bun"], + "strict": true, + "noEmit": true, + "allowImportingTsExtensions": true, + "verbatimModuleSyntax": true, + "skipLibCheck": true, + "resolveJsonModule": true, + "esModuleInterop": true, + "isolatedModules": true, + "forceConsistentCasingInFileNames": true, + "noUncheckedIndexedAccess": true, + "noUnusedLocals": true, + "noUnusedParameters": true, + "noFallthroughCasesInSwitch": true + }, + "include": ["apps/packages/*/src", "apps/interfaces/*/src", "apps/interfaces/*/bin"] +} \ No newline at end of file