feat(rewrite): tambah application layer — port traits & use cases TypeScript
Paket @zesdex/application (apps/packages/application): - ports: ProviderService (chat/chatStream+abort), PasswordService, TokenService, AuthService - agent: AgentTurnServiceImpl (loop 50 iterasi, auto-compact 60k chars, adaptive max-tokens 800/1600/4096, temp 0.2/0.7, ErrorTracker, eksekusi tool read-only paralel terbatas mempertahankan urutan) + compact_messages_with_ai - auth: OAuthUseCase (PKCE S256 + CSRF state), SessionServiceImpl - cms: ConversationServiceImpl, MemoryServiceImpl, SettingsServiceImpl - 11 unit test Bun (PKCE, OAuth CSRF, turn_service helper) - tsconfig paths untuk workspace @zesdex/* Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
07e84e43b3
commit
03538a455a
@@ -0,0 +1,21 @@
|
||||
{
|
||||
"name": "@zesdex/application",
|
||||
"version": "1.21.2",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"description": "Zesdex application layer — port traits, use cases, turn service (depends only on @zesdex/domain)",
|
||||
"exports": {
|
||||
".": "./src/index.ts"
|
||||
},
|
||||
"scripts": {
|
||||
"test": "bun test",
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@zesdex/domain": "workspace:*"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/bun": "^1.2.0",
|
||||
"typescript": "^5.7.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
/**
|
||||
* Agent application module — ToolExecutor, AgentTurnService, and the
|
||||
* AgentTurnServiceImpl turn loop + compaction. Mirrors `agent/` in Rust.
|
||||
*/
|
||||
import type { JsonValue } from "@zesdex/domain";
|
||||
|
||||
/** Interface for dispatching tool calls to their concrete implementations. */
|
||||
export interface ToolExecutor {
|
||||
/** Execute a tool call asynchronously. */
|
||||
execute(toolName: string, args: JsonValue): Promise<string>;
|
||||
/** Whether a tool is read-only / safe to run concurrently. Default `false`. */
|
||||
isParallelSafe?(toolName: string): boolean;
|
||||
}
|
||||
|
||||
/** Service for running agent turns asynchronously. */
|
||||
export interface AgentTurnService {
|
||||
runTurn(params: import("@zesdex/domain").AgentTurnParams): Promise<void>;
|
||||
}
|
||||
|
||||
export * from "./turn_service.ts";
|
||||
@@ -0,0 +1,59 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import {
|
||||
adaptiveMaxTokens,
|
||||
conversationChars,
|
||||
ErrorTracker,
|
||||
truncateToolOutput,
|
||||
} from "./turn_service.ts";
|
||||
import { type ChatMessage, newConversation, systemMessage, userMessage, toolResultMessage } from "@zesdex/domain";
|
||||
|
||||
describe("truncateToolOutput", () => {
|
||||
it("short output is unchanged", () => {
|
||||
expect(truncateToolOutput("short")).toBe("short");
|
||||
});
|
||||
|
||||
it("long output preserves head and marks cut", () => {
|
||||
const long = "x".repeat(12_000 + 500);
|
||||
const truncated = truncateToolOutput(long);
|
||||
expect(truncated.length).toBeLessThan(long.length);
|
||||
expect(truncated).toContain("...[truncated");
|
||||
expect(truncated.startsWith("xxx")).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
describe("adaptiveMaxTokens", () => {
|
||||
it("scales with request length", () => {
|
||||
expect(adaptiveMaxTokens(10)).toBe(800);
|
||||
expect(adaptiveMaxTokens(200)).toBe(1600);
|
||||
expect(adaptiveMaxTokens(5000)).toBe(4096);
|
||||
});
|
||||
});
|
||||
|
||||
describe("ErrorTracker", () => {
|
||||
it("injects recovery note after repeated errors", () => {
|
||||
const tracker = new ErrorTracker();
|
||||
const messages: ChatMessage[] = [];
|
||||
tracker.record("read", "Error: File not found", messages);
|
||||
tracker.record("read", "Error: File not found", messages);
|
||||
expect(tracker.shouldStop()).toBe(false);
|
||||
tracker.record("read", "Error: File not found", messages);
|
||||
expect(messages.some((m) => m.content?.includes("[System note]"))).toBe(true);
|
||||
});
|
||||
|
||||
it("stops after too many errors", () => {
|
||||
const tracker = new ErrorTracker();
|
||||
const messages: ChatMessage[] = [];
|
||||
for (let i = 0; i < 8; i++) {
|
||||
tracker.record("bash", `Error: boom ${i}`, messages);
|
||||
}
|
||||
expect(tracker.shouldStop()).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
describe("conversationChars", () => {
|
||||
it("sums content only", () => {
|
||||
const conv = newConversation("sys", "s1");
|
||||
conv.messages.push(systemMessage("sys"), userMessage("hello world"), toolResultMessage("id", "output"));
|
||||
expect(conversationChars(conv.messages)).toBe(3 + 11 + 6);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,358 @@
|
||||
/**
|
||||
* Agent turn service — the core adaptive turn loop. Mirrors
|
||||
* `apps/application/src/agent/turn_service.rs`.
|
||||
*/
|
||||
import {
|
||||
type AgentTurnParams,
|
||||
type ChatMessage,
|
||||
type StreamEvent,
|
||||
type ToolCall,
|
||||
type ToolDef,
|
||||
type TurnEventSink,
|
||||
type JsonValue,
|
||||
systemMessage,
|
||||
userMessage,
|
||||
assistantMessage,
|
||||
toolResultMessage,
|
||||
mainAgentPromptWithProjectContext,
|
||||
compactionPrompt,
|
||||
errorRecoveryNote,
|
||||
sanitizeToolArguments,
|
||||
} from "@zesdex/domain";
|
||||
import type { ProviderService } from "../ports/index.ts";
|
||||
import type { ToolExecutor } from "./index.ts";
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* Constants (mirrors Rust) */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
const MAX_TURN_ITERATIONS = 50;
|
||||
const MAX_CONSECUTIVE_TOOL_ERRORS = 3;
|
||||
const MAX_TOTAL_TOOL_ERRORS = 8;
|
||||
const TOOL_OUTPUT_MAX_CHARS = 12_000;
|
||||
const AUTO_COMPACT_CHARS = 60_000;
|
||||
const MAX_PARALLEL_TOOLS = 8;
|
||||
const PROJECT_CONTEXT_MAX_CHARS = 12_000;
|
||||
const RULE_FILENAMES = ["AGENTS.md", "agent.md", "CLAUDE.md", "claude.md", ".cursorrules", ".zesdexrules"];
|
||||
const COMPACT_KEEP_TAIL = 6;
|
||||
|
||||
/** Whether the output string denotes a tool error. */
|
||||
function isErrorOutput(output: string): boolean {
|
||||
return output.startsWith("Error:");
|
||||
}
|
||||
|
||||
/** Truncate a long tool output, preserving the head + truncation marker. */
|
||||
export function truncateToolOutput(output: string): string {
|
||||
if (output.length <= TOOL_OUTPUT_MAX_CHARS) return output;
|
||||
const head = output.slice(0, TOOL_OUTPUT_MAX_CHARS);
|
||||
return `${head}\n...[truncated ${output.length - TOOL_OUTPUT_MAX_CHARS} chars]`;
|
||||
}
|
||||
|
||||
/** Pick a `max_tokens` budget based on the user's request length. */
|
||||
export function adaptiveMaxTokens(requestLen: number): number {
|
||||
if (requestLen <= 80) return 800;
|
||||
if (requestLen <= 400) return 1600;
|
||||
return 4096;
|
||||
}
|
||||
|
||||
/** Sum character length of message content as a context-size proxy. */
|
||||
export function conversationChars(messages: ChatMessage[]): number {
|
||||
return messages.reduce((acc, m) => acc + (m.content?.length ?? 0), 0);
|
||||
}
|
||||
|
||||
/** Best-effort build of project context from convention rule files. */
|
||||
export function buildProjectContext(root: string): string {
|
||||
let ctx = "";
|
||||
for (const file of RULE_FILENAMES) {
|
||||
try {
|
||||
const content = requireNodeFsReadFile(root, file);
|
||||
ctx += `\n### ${file}\n\`\`\`\n${content.trim()}\n\`\`\``;
|
||||
} catch {
|
||||
/* file missing — skip */
|
||||
}
|
||||
}
|
||||
const context = ctx.trim();
|
||||
if (context.length <= PROJECT_CONTEXT_MAX_CHARS) return context;
|
||||
return `${context.slice(0, PROJECT_CONTEXT_MAX_CHARS)}\n...[project context truncated]`;
|
||||
}
|
||||
|
||||
/** Read a repo rule file synchronously (Bun-compatible). */
|
||||
function requireNodeFsReadFile(root: string, file: string): string {
|
||||
const fs = require("node:fs");
|
||||
return fs.readFileSync(`${root}/${file}`, "utf8");
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* ErrorTracker */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
/** Track repeated tool-call errors so the loop can recover. */
|
||||
export class ErrorTracker {
|
||||
consecutive = 0;
|
||||
total = 0;
|
||||
lastTool: string | null = null;
|
||||
lastError = "";
|
||||
|
||||
record(toolName: string, error: string, messages: ChatMessage[]): void {
|
||||
if (this.lastTool === toolName) {
|
||||
this.consecutive += 1;
|
||||
} else {
|
||||
this.consecutive = 1;
|
||||
}
|
||||
this.lastTool = toolName;
|
||||
this.lastError = error;
|
||||
this.total += 1;
|
||||
|
||||
const sysNoteInContext = messages.some((m) => m.content?.includes("[System note]"));
|
||||
if (this.consecutive >= MAX_CONSECUTIVE_TOOL_ERRORS && !sysNoteInContext) {
|
||||
messages.push(systemMessage(errorRecoveryNote(toolName, error)));
|
||||
}
|
||||
}
|
||||
|
||||
shouldStop(): boolean {
|
||||
return this.consecutive >= MAX_CONSECUTIVE_TOOL_ERRORS * 2 || this.total >= MAX_TOTAL_TOOL_ERRORS;
|
||||
}
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* Tool execution */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
/** Execute one tool call, push events, return the result string. */
|
||||
async function executeToolCall(
|
||||
executor: ToolExecutor,
|
||||
sink: TurnEventSink,
|
||||
tc: ToolCall,
|
||||
): Promise<string> {
|
||||
const name = tc.function.name;
|
||||
const args = sanitizeToolArguments(tc.function.arguments);
|
||||
|
||||
let output: string;
|
||||
try {
|
||||
output = await executor.execute(name, args as JsonValue);
|
||||
} catch (e) {
|
||||
output = `Error: ${(e as Error).message}`;
|
||||
}
|
||||
|
||||
const isError = isErrorOutput(output);
|
||||
const truncated = truncateToolOutput(output);
|
||||
|
||||
sink.push({
|
||||
kind: "tool_result",
|
||||
tool_call_id: tc.id,
|
||||
tool_name: name,
|
||||
output: truncated,
|
||||
is_error: isError,
|
||||
path: null,
|
||||
});
|
||||
|
||||
return truncated;
|
||||
}
|
||||
|
||||
/** Execute a batch of read-only tool calls concurrently (bounded) in original order. */
|
||||
async function executeToolCallsInParallel(
|
||||
executor: ToolExecutor,
|
||||
sink: TurnEventSink,
|
||||
toolCalls: ToolCall[],
|
||||
): Promise<string[]> {
|
||||
// Simple bounded concurrency preserving input order.
|
||||
const results: string[] = new Array(toolCalls.length);
|
||||
let next = 0;
|
||||
|
||||
async function worker() {
|
||||
while (true) {
|
||||
const idx = next++;
|
||||
if (idx >= toolCalls.length) return;
|
||||
results[idx] = await executeToolCall(executor, sink, toolCalls[idx]!);
|
||||
}
|
||||
}
|
||||
|
||||
const workers = Array.from({ length: Math.min(MAX_PARALLEL_TOOLS, toolCalls.length) }, () => worker());
|
||||
await Promise.all(workers);
|
||||
return results;
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* Compaction */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
/**
|
||||
* Compact oversized conversation history using AI summarisation. At most once
|
||||
* per turn. Keeps the last COMPACT_KEEP_TAIL messages.
|
||||
*/
|
||||
export async function compactMessagesWithAi(
|
||||
messages: ChatMessage[],
|
||||
provider: ProviderService,
|
||||
): Promise<void> {
|
||||
if (messages.length <= COMPACT_KEEP_TAIL + 2) return;
|
||||
|
||||
const splitIdx = messages.length - COMPACT_KEEP_TAIL;
|
||||
const evicted = messages.splice(0, splitIdx);
|
||||
|
||||
const summaryPrompt: ChatMessage[] = [systemMessage(compactionPrompt()), ...evicted, userMessage("Please summarise our previous conversation above for context continuity.")];
|
||||
|
||||
try {
|
||||
const { message } = await provider.chat(summaryPrompt, undefined, 1024, 0.3);
|
||||
const summaryText = message.content ?? "Previous context summarised.";
|
||||
messages.unshift(systemMessage(`[AI Summary of Previous Conversation]\n${summaryText.trim()}`));
|
||||
} catch {
|
||||
messages.unshift(systemMessage("[Earlier conversation messages compacted to save context window]"));
|
||||
}
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* AgentTurnServiceImpl */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
/** Service implementation for executing an agent turn asynchronously. */
|
||||
export class AgentTurnServiceImpl {
|
||||
private provider: ProviderService;
|
||||
private toolExecutor: ToolExecutor;
|
||||
private toolDefs: ToolDef[];
|
||||
|
||||
constructor(provider: ProviderService, toolExecutor: ToolExecutor, toolDefs: ToolDef[]) {
|
||||
this.provider = provider;
|
||||
this.toolExecutor = toolExecutor;
|
||||
this.toolDefs = toolDefs;
|
||||
}
|
||||
|
||||
/** Emit a TurnEvent onto the sink (no-op if the sink is missing). */
|
||||
private push(sink: TurnEventSink, event: Parameters<TurnEventSink["push"]>[0]): void {
|
||||
sink.push(event);
|
||||
}
|
||||
|
||||
/** Execute a single LLM stream call, forwarding tokens and checking abort. */
|
||||
private async callLlm(
|
||||
messages: ChatMessage[],
|
||||
abort: AbortController,
|
||||
sink: TurnEventSink,
|
||||
maxTokens: number,
|
||||
temperature: number,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }> {
|
||||
const onEvent = (event: StreamEvent): boolean => {
|
||||
if (abort.signal.aborted) return false;
|
||||
if (event.kind === "token") this.push(sink, { kind: "stream_token", content: event.content });
|
||||
else if (event.kind === "reasoning") this.push(sink, { kind: "stream_reasoning", content: event.content });
|
||||
return true;
|
||||
};
|
||||
|
||||
try {
|
||||
return await this.provider.chatStream(messages, this.toolDefs, maxTokens, temperature, onEvent, abort.signal);
|
||||
} catch (e) {
|
||||
throw new Error(`LLM error: ${(e as Error).message}`);
|
||||
}
|
||||
}
|
||||
|
||||
/** Auto-compact history in place if it exceeds the threshold. */
|
||||
private async autoCompactIfNeeded(messages: ChatMessage[]): Promise<void> {
|
||||
if (conversationChars(messages) <= AUTO_COMPACT_CHARS) return;
|
||||
const sys = messages[0];
|
||||
if (!sys) return;
|
||||
const rest = messages.splice(1);
|
||||
const before = rest.length;
|
||||
try {
|
||||
await compactMessagesWithAi(rest, this.provider);
|
||||
} catch (e) {
|
||||
console.warn(`auto-compact failed (non-fatal): ${(e as Error).message}`);
|
||||
}
|
||||
messages.length = 0;
|
||||
messages.push(sys, ...rest);
|
||||
console.info(`auto-compacted history: ${before} messages -> ${rest.length}`);
|
||||
}
|
||||
|
||||
/** Run the full agent turn loop. */
|
||||
async runTurn(params: AgentTurnParams): Promise<void> {
|
||||
const sink = params.turn_events;
|
||||
const abort = params.abort;
|
||||
const { in_flight } = params;
|
||||
|
||||
// Insert system prompt at index 0 with repo conventions loaded.
|
||||
const projectContext = buildProjectContext(params.workspace_roots[0] ?? ".");
|
||||
const systemPrompt = mainAgentPromptWithProjectContext(projectContext);
|
||||
params.messages.unshift(systemMessage(systemPrompt));
|
||||
const originalCount = params.messages.length;
|
||||
|
||||
// Estimate request complexity from the last user message.
|
||||
const last = params.messages[params.messages.length - 1];
|
||||
const requestLen = last?.content?.length ?? 0;
|
||||
|
||||
const errors = new ErrorTracker();
|
||||
let sawToolCalls = false;
|
||||
|
||||
for (let iteration = 0; iteration < MAX_TURN_ITERATIONS; iteration++) {
|
||||
// Check abort flag.
|
||||
if (abort.signal.aborted) {
|
||||
this.push(sink, { kind: "system_note", systemKind: "info", message: "Turn aborted by user" });
|
||||
break;
|
||||
}
|
||||
|
||||
if (errors.shouldStop()) {
|
||||
this.push(sink, { kind: "system_note", systemKind: "warn", message: "Stopping: repeated tool errors without progress" });
|
||||
break;
|
||||
}
|
||||
|
||||
// Auto-compact oversized history before the LLM call.
|
||||
await this.autoCompactIfNeeded(params.messages);
|
||||
|
||||
// Adaptive generation parameters.
|
||||
const maxTokens = adaptiveMaxTokens(requestLen);
|
||||
const temperature = sawToolCalls ? 0.2 : 0.7;
|
||||
|
||||
this.push(sink, { kind: "stream_start" });
|
||||
|
||||
let result;
|
||||
try {
|
||||
result = await this.callLlm(params.messages, abort, sink, maxTokens, temperature);
|
||||
} catch (e) {
|
||||
const msg = (e as Error).message;
|
||||
console.warn(msg);
|
||||
this.push(sink, { kind: "error", message: msg });
|
||||
break;
|
||||
}
|
||||
|
||||
const { message: assistantMsg, usage } = result;
|
||||
const content = assistantMsg.content ?? "";
|
||||
const toolCalls = assistantMsg.tool_calls ?? [];
|
||||
|
||||
this.push(sink, { kind: "stream_done", message: assistantMsg });
|
||||
if (usage) this.push(sink, { kind: "usage", tokens_in: usage[0], tokens_out: usage[1] });
|
||||
|
||||
// No tool calls → assistant is done.
|
||||
if (toolCalls.length === 0) {
|
||||
params.messages.push(assistantMessage(content));
|
||||
break;
|
||||
}
|
||||
|
||||
sawToolCalls = true;
|
||||
params.messages.push(assistantMsg);
|
||||
|
||||
// Execute tool calls — parallel when all read-only, else sequential.
|
||||
const isParallelSafe = (tc: ToolCall): boolean =>
|
||||
typeof this.toolExecutor.isParallelSafe === "function"
|
||||
? this.toolExecutor.isParallelSafe(tc.function.name)
|
||||
: false;
|
||||
const parallel = toolCalls.length > 1 && toolCalls.every(isParallelSafe);
|
||||
|
||||
const outputs = parallel
|
||||
? await executeToolCallsInParallel(this.toolExecutor, sink, toolCalls)
|
||||
: await (async () => {
|
||||
const seq: string[] = [];
|
||||
for (const tc of toolCalls) seq.push(await executeToolCall(this.toolExecutor, sink, tc));
|
||||
return seq;
|
||||
})();
|
||||
|
||||
for (let i = 0; i < toolCalls.length; i++) {
|
||||
const tc = toolCalls[i]!;
|
||||
const output = outputs[i]!;
|
||||
if (isErrorOutput(output)) errors.record(tc.function.name, output, params.messages);
|
||||
params.messages.push(toolResultMessage(tc.id, output));
|
||||
}
|
||||
}
|
||||
|
||||
// Remove the synthetic sys_msg before emitting to the transcript.
|
||||
const compacted = params.messages.splice(originalCount - 1);
|
||||
this.push(sink, { kind: "compacted", messages: compacted });
|
||||
this.push(sink, { kind: "done" });
|
||||
in_flight.value = false;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
/** Auth application module — OAuth PKCE + session management use-cases. */
|
||||
export * from "./oauth_service.ts";
|
||||
export * from "./session_service.ts";
|
||||
@@ -0,0 +1,141 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { generatePkcePair, OAuthUseCase } from "./oauth_service.ts";
|
||||
import type { OAuthFlowStore } from "./oauth_service.ts";
|
||||
import { createHash } from "node:crypto";
|
||||
import type { OAuthToken } from "@zesdex/domain";
|
||||
|
||||
describe("generatePkcePair", () => {
|
||||
it("produces a verifier ≥43 chars and a valid S256 challenge", () => {
|
||||
const [verifier, challenge] = generatePkcePair();
|
||||
expect(verifier.length).toBeGreaterThanOrEqual(43);
|
||||
// Challenge = base64url(sha256(verifier)).
|
||||
const expected = Buffer.from(createHash("sha256").update(verifier, "utf8").digest()).toString("base64url");
|
||||
expect(challenge).toBe(expected);
|
||||
});
|
||||
|
||||
it("is unique across calls", () => {
|
||||
const [a] = generatePkcePair();
|
||||
const [b] = generatePkcePair();
|
||||
expect(a).not.toBe(b);
|
||||
});
|
||||
});
|
||||
|
||||
class MemoryFlowStore implements OAuthFlowStore {
|
||||
private verifier = "";
|
||||
private state = "";
|
||||
async saveFlowState(v: string, s: string): Promise<void> {
|
||||
this.verifier = v;
|
||||
this.state = s;
|
||||
}
|
||||
async loadVerifier(): Promise<string> {
|
||||
return this.verifier;
|
||||
}
|
||||
async loadState(): Promise<string> {
|
||||
return this.state;
|
||||
}
|
||||
async clear(): Promise<void> {
|
||||
this.verifier = "";
|
||||
this.state = "";
|
||||
}
|
||||
}
|
||||
|
||||
function makeRepo() {
|
||||
let token: OAuthToken | null = null;
|
||||
return {
|
||||
repo: {
|
||||
async saveToken(_path: string, t: OAuthToken): Promise<void> {
|
||||
token = t;
|
||||
},
|
||||
async loadToken(): Promise<OAuthToken | null> {
|
||||
return token;
|
||||
},
|
||||
},
|
||||
getToken: () => token,
|
||||
};
|
||||
}
|
||||
|
||||
describe("OAuthUseCase", () => {
|
||||
it("builds an auth URL with PKCE params and persists flow state", async () => {
|
||||
const { repo } = makeRepo();
|
||||
const store = new MemoryFlowStore();
|
||||
const exchanger = {
|
||||
async exchangeCode(): Promise<OAuthToken> {
|
||||
return { access_token: "at", refresh_token: "rt", expires_at: 9999, token_type: "Bearer" };
|
||||
},
|
||||
};
|
||||
const useCase = new OAuthUseCase(repo, store, exchanger, "/tmp/token.json");
|
||||
|
||||
const { authUrl, state } = await useCase.startFlow(
|
||||
{ auth_url: "https://provider.example/oauth/authorize", token_url: "https://provider.example/oauth/token", client_id: "cid", scopes: ["openid", "profile"] },
|
||||
"http://localhost:9999/callback",
|
||||
);
|
||||
|
||||
const url = new URL(authUrl);
|
||||
expect(url.searchParams.get("response_type")).toBe("code");
|
||||
expect(url.searchParams.get("client_id")).toBe("cid");
|
||||
expect(url.searchParams.get("scope")).toBe("openid profile");
|
||||
expect(url.searchParams.get("code_challenge_method")).toBe("S256");
|
||||
expect(url.searchParams.get("state")).toBe(state);
|
||||
expect(url.searchParams.get("code_challenge")).toBeTruthy();
|
||||
// Flow state persisted.
|
||||
expect(await store.loadState()).toBe(state);
|
||||
expect(await store.loadVerifier()).toBeTruthy();
|
||||
});
|
||||
|
||||
it("enforces CSRF state match in completeFlow", async () => {
|
||||
const { repo } = makeRepo();
|
||||
const store = new MemoryFlowStore();
|
||||
const exchanger = {
|
||||
async exchangeCode(): Promise<OAuthToken> {
|
||||
return { access_token: "at", refresh_token: "rt", expires_at: 9999, token_type: "Bearer" };
|
||||
},
|
||||
};
|
||||
const useCase = new OAuthUseCase(repo, store, exchanger, "/tmp/token.json");
|
||||
await useCase.startFlow(
|
||||
{ auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] },
|
||||
"http://localhost:9999/callback",
|
||||
);
|
||||
|
||||
await expect(
|
||||
useCase.completeFlow(
|
||||
{ auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] },
|
||||
"http://localhost:9999/callback",
|
||||
"code123",
|
||||
"wrong-state",
|
||||
),
|
||||
).rejects.toMatchObject({ kind: "state_mismatch" });
|
||||
});
|
||||
|
||||
it("completes flow successfully with matching state", async () => {
|
||||
const { repo, getToken } = makeRepo();
|
||||
const store = new MemoryFlowStore();
|
||||
let exchangedVerifier = "";
|
||||
const exchanger = {
|
||||
async exchangeCode(_tokenUrl: string, _clientId: string, _secret: string | null, _redirectUri: string, code: string, codeVerifier: string): Promise<OAuthToken> {
|
||||
exchangedVerifier = codeVerifier;
|
||||
return { access_token: `at-${code}`, refresh_token: "rt", expires_at: 9999, token_type: "Bearer" };
|
||||
},
|
||||
};
|
||||
const useCase = new OAuthUseCase(repo, store, exchanger, "/tmp/token.json");
|
||||
const { state } = await useCase.startFlow(
|
||||
{ auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] },
|
||||
"http://localhost:9999/callback",
|
||||
);
|
||||
const savedVerifier = await store.loadVerifier();
|
||||
|
||||
const token = await useCase.completeFlow(
|
||||
{ auth_url: "https://provider.example/authorize", token_url: "https://provider.example/token", client_id: "cid", scopes: ["openid"] },
|
||||
"http://localhost:9999/callback",
|
||||
"code123",
|
||||
state,
|
||||
);
|
||||
|
||||
expect(token.access_token).toBe("at-code123");
|
||||
// The code verifier passed to the exchanger is the one saved at start.
|
||||
expect(exchangedVerifier).toBe(savedVerifier);
|
||||
// Flow state cleared after completion.
|
||||
expect(await store.loadState()).toBe("");
|
||||
// Token persisted via repo.
|
||||
expect(getToken()?.access_token).toBe("at-code123");
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,121 @@
|
||||
/**
|
||||
* OAuth 2.0 authorization-code + PKCE flow use-case. Mirrors
|
||||
* `apps/application/src/auth/oauth_service.rs`.
|
||||
*/
|
||||
import {
|
||||
type OAuthConfig,
|
||||
type OAuthRepository,
|
||||
type OAuthToken,
|
||||
authInvalidConfig,
|
||||
authStateMismatch,
|
||||
} from "@zesdex/domain";
|
||||
import { createHash, randomUUID } from "node:crypto";
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* Port traits */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
/** Persistence contract for ephemeral OAuth flow state. */
|
||||
export interface OAuthFlowStore {
|
||||
saveFlowState(verifier: string, state: string): Promise<void>;
|
||||
loadVerifier(): Promise<string>;
|
||||
loadState(): Promise<string>;
|
||||
clear(): Promise<void>;
|
||||
}
|
||||
|
||||
/** Abstraction for exchanging an authorization code for tokens. */
|
||||
export interface TokenExchanger {
|
||||
exchangeCode(
|
||||
tokenUrl: string,
|
||||
clientId: string,
|
||||
clientSecret: string | null,
|
||||
redirectUri: string,
|
||||
code: string,
|
||||
codeVerifier: string,
|
||||
): Promise<OAuthToken>;
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* PKCE helpers */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
function base64url(input: string | Uint8Array): string {
|
||||
const buf = typeof input === "string" ? new TextEncoder().encode(input) : Buffer.from(input as Uint8Array);
|
||||
return Buffer.from(buf).toString("base64url");
|
||||
}
|
||||
|
||||
/** Generate a PKCE code-verifier and its S256 code-challenge. */
|
||||
export function generatePkcePair(): [string, string] {
|
||||
const bytes = new Uint8Array(32);
|
||||
crypto.getRandomValues(bytes);
|
||||
const verifier = base64url(bytes);
|
||||
const challenge = base64url(createHash("sha256").update(verifier, "utf8").digest());
|
||||
return [verifier, challenge];
|
||||
}
|
||||
|
||||
/** Generate a random CSRF state token (UUID-based). */
|
||||
export function generateStateToken(): string {
|
||||
return randomUUID();
|
||||
}
|
||||
|
||||
/* -------------------------------------------------------------------------- */
|
||||
/* OAuthUseCase */
|
||||
/* -------------------------------------------------------------------------- */
|
||||
|
||||
/** Concrete OAuth flow use-case over injected repositories. */
|
||||
export class OAuthUseCase {
|
||||
constructor(
|
||||
private tokenRepo: OAuthRepository,
|
||||
private flowStore: OAuthFlowStore,
|
||||
private tokenExchanger: TokenExchanger,
|
||||
private tokenPath: string,
|
||||
) {}
|
||||
|
||||
async startFlow(config: OAuthConfig, redirectUri: string): Promise<{ authUrl: string; state: string }> {
|
||||
if (config.auth_url === "") {
|
||||
throw authInvalidConfig("OAuth auth_url is empty");
|
||||
}
|
||||
|
||||
const [verifier, challenge] = generatePkcePair();
|
||||
const state = generateStateToken();
|
||||
|
||||
await this.flowStore.saveFlowState(verifier, state);
|
||||
|
||||
const url = new URL(config.auth_url);
|
||||
url.searchParams.set("response_type", "code");
|
||||
url.searchParams.set("client_id", config.client_id);
|
||||
url.searchParams.set("redirect_uri", redirectUri);
|
||||
url.searchParams.set("scope", config.scopes.join(" "));
|
||||
url.searchParams.set("state", state);
|
||||
url.searchParams.set("code_challenge_method", "S256");
|
||||
url.searchParams.set("code_challenge", challenge);
|
||||
|
||||
return { authUrl: url.toString(), state };
|
||||
}
|
||||
|
||||
async completeFlow(config: OAuthConfig, redirectUri: string, code: string, state: string): Promise<OAuthToken> {
|
||||
// CSRF check.
|
||||
const expectedState = await this.flowStore.loadState();
|
||||
if (expectedState !== state) throw authStateMismatch();
|
||||
|
||||
// Read the PKCE verifier saved in start_flow.
|
||||
const verifier = await this.flowStore.loadVerifier();
|
||||
|
||||
const token = await this.tokenExchanger.exchangeCode(
|
||||
config.token_url,
|
||||
config.client_id,
|
||||
config.client_secret ?? null,
|
||||
redirectUri,
|
||||
code,
|
||||
verifier,
|
||||
);
|
||||
|
||||
await this.tokenRepo.saveToken(this.tokenPath, token);
|
||||
await this.flowStore.clear();
|
||||
return token;
|
||||
}
|
||||
|
||||
async getToken(): Promise<OAuthToken | null> {
|
||||
return this.tokenRepo.loadToken(this.tokenPath);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,44 @@
|
||||
/**
|
||||
* Session management use-case. Mirrors `apps/application/src/auth/session_service.rs`.
|
||||
*/
|
||||
import { randomUUID } from "node:crypto";
|
||||
import {
|
||||
type Session,
|
||||
type SessionId,
|
||||
type SessionLockRepository,
|
||||
type SessionRepository,
|
||||
newSession,
|
||||
newSessionId,
|
||||
authOtherError as otherErr,
|
||||
} from "@zesdex/domain";
|
||||
|
||||
/** Concrete session service backed by injected repositories. */
|
||||
export class SessionServiceImpl {
|
||||
constructor(
|
||||
private sessionRepo: SessionRepository,
|
||||
/** Lock repository is injected for future lock acquire/release (matches Rust contract). */
|
||||
// @ts-expect-error -- kept for structural parity with the Rust `SessionServiceImpl<R, L>`
|
||||
private lockRepo: SessionLockRepository,
|
||||
private baseDir: string,
|
||||
) {}
|
||||
|
||||
async createSession(title: string): Promise<Session> {
|
||||
const idRes = newSessionId(randomUUID());
|
||||
if (!idRes.ok) throw otherErr(idRes.error);
|
||||
const titleOwned = title === "" ? "New Session" : title;
|
||||
const session = newSession(idRes.value, titleOwned);
|
||||
await this.sessionRepo.saveSession(this.baseDir, session);
|
||||
return session;
|
||||
}
|
||||
|
||||
async listAll(): Promise<Session[]> {
|
||||
return this.sessionRepo.listSessions(this.baseDir);
|
||||
}
|
||||
|
||||
async archiveSession(id: SessionId): Promise<void> {
|
||||
const session = await this.sessionRepo.loadSession(this.baseDir, id);
|
||||
session.archived = true;
|
||||
session.updated_at = Date.now();
|
||||
await this.sessionRepo.saveSession(this.baseDir, session);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
/**
|
||||
* Conversation use-case. Mirrors `apps/application/src/cms/conversation_service.rs`.
|
||||
*/
|
||||
import * as path from "node:path";
|
||||
import { type Conversation, type ConversationRepository, type ChatMessage, pushMessage } from "@zesdex/domain";
|
||||
|
||||
/** Service implementation for conversation CRUD operations. */
|
||||
export class ConversationServiceImpl {
|
||||
constructor(
|
||||
private repo: ConversationRepository,
|
||||
private sessionsDir: string,
|
||||
) {}
|
||||
|
||||
private sessionDir(sessionId: string): string {
|
||||
return path.join(this.sessionsDir, sessionId);
|
||||
}
|
||||
|
||||
async loadConversation(sessionId: string): Promise<Conversation> {
|
||||
const dir = this.sessionDir(sessionId);
|
||||
return this.repo.load(dir);
|
||||
}
|
||||
|
||||
async saveConversation(conv: Conversation): Promise<void> {
|
||||
const dir = this.sessionDir(conv.session_id);
|
||||
await this.repo.save(dir, conv);
|
||||
}
|
||||
|
||||
async addMessage(conv: Conversation, msg: ChatMessage): Promise<void> {
|
||||
pushMessage(conv, msg);
|
||||
const dir = this.sessionDir(conv.session_id);
|
||||
await this.repo.save(dir, conv);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,4 @@
|
||||
/** CMS application module — conversation, memory, and settings use-cases. */
|
||||
export * from "./conversation_service.ts";
|
||||
export * from "./memory_service.ts";
|
||||
export * from "./settings_service.ts";
|
||||
@@ -0,0 +1,24 @@
|
||||
/**
|
||||
* Memory use-case. Mirrors `apps/application/src/cms/memory_service.rs`.
|
||||
*/
|
||||
import { type Memory, type MemoryRepository } from "@zesdex/domain";
|
||||
|
||||
/** Service implementation for memory CRUD operations. */
|
||||
export class MemoryServiceImpl {
|
||||
constructor(
|
||||
private repo: MemoryRepository,
|
||||
private memoryDir: string,
|
||||
) {}
|
||||
|
||||
async listMemories(): Promise<string[]> {
|
||||
return this.repo.list(this.memoryDir);
|
||||
}
|
||||
|
||||
async saveMemory(memory: Memory): Promise<void> {
|
||||
await this.repo.save(this.memoryDir, memory);
|
||||
}
|
||||
|
||||
async deleteMemory(name: string): Promise<void> {
|
||||
await this.repo.delete(this.memoryDir, name);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
/**
|
||||
* Settings + app-config use-case. Mirrors `apps/application/src/cms/settings_service.rs`.
|
||||
*/
|
||||
import {
|
||||
type AppConfig,
|
||||
type AppConfigRepository,
|
||||
type ProviderConfig,
|
||||
type Settings,
|
||||
type SettingsRepository,
|
||||
} from "@zesdex/domain";
|
||||
|
||||
/** Service implementation for settings and app-config operations. */
|
||||
export class SettingsServiceImpl {
|
||||
constructor(
|
||||
private settingsRepo: SettingsRepository,
|
||||
private appConfigRepo: AppConfigRepository,
|
||||
private baseDir: string,
|
||||
) {}
|
||||
|
||||
async loadSettings(): Promise<Settings> {
|
||||
return this.settingsRepo.load(this.baseDir);
|
||||
}
|
||||
|
||||
async saveSettings(settings: Settings): Promise<void> {
|
||||
await this.settingsRepo.save(this.baseDir, settings);
|
||||
}
|
||||
|
||||
async updateProvider(name: string, config: ProviderConfig): Promise<void> {
|
||||
const appConfig: AppConfig = await this.appConfigRepo.load(this.baseDir);
|
||||
appConfig.providers[name] = config;
|
||||
await this.appConfigRepo.save(this.baseDir, appConfig);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
/**
|
||||
* Zesdex Application Layer — port traits, use cases, turn service.
|
||||
* Depends only on @zesdex/domain. Higher-level ports implemented by
|
||||
* infrastructure adapters.
|
||||
*/
|
||||
export * from "./ports/index.ts";
|
||||
export * from "./agent/index.ts";
|
||||
export * from "./auth/index.ts";
|
||||
export * from "./cms/index.ts";
|
||||
@@ -0,0 +1,53 @@
|
||||
/**
|
||||
* Port traits (interfaces) to external services. Mirrors `apps/application/src/ports/`.
|
||||
* Concrete implementations live in the infrastructure layer.
|
||||
*/
|
||||
import type { ChatMessage, StreamEvent, ToolDef } from "@zesdex/domain";
|
||||
|
||||
/** Abstraction for an LLM provider chat-completion service. */
|
||||
export interface ProviderService {
|
||||
/**
|
||||
* Send a non-streaming chat completion request.
|
||||
* Returns the assistant's `ChatMessage` and optional `[prompt, completion]` token usage.
|
||||
*/
|
||||
chat(
|
||||
messages: ChatMessage[],
|
||||
tools?: ToolDef[],
|
||||
maxTokens?: number,
|
||||
temperature?: number,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }>;
|
||||
|
||||
/**
|
||||
* Send a streaming request. `onEvent` is called per parsed SSE event and
|
||||
* returns `false` to abort. Returns the fully assembled assistant message
|
||||
* and optional usage once the stream completes.
|
||||
*/
|
||||
chatStream(
|
||||
messages: ChatMessage[],
|
||||
tools: ToolDef[],
|
||||
maxTokens: number,
|
||||
temperature: number,
|
||||
onEvent: (event: StreamEvent) => boolean,
|
||||
signal?: AbortSignal,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }>;
|
||||
}
|
||||
|
||||
/** Abstraction for password hashing and verification. */
|
||||
export interface PasswordService {
|
||||
hash(password: string): Promise<string>;
|
||||
verify(password: string, hash: string): Promise<boolean>;
|
||||
}
|
||||
|
||||
/** Abstraction for JWT-based token generation and verification. */
|
||||
export interface TokenService {
|
||||
/** Generate an `[access, refresh]` token pair for the given subject. */
|
||||
generateTokens(sub: string): [string, string];
|
||||
verifyAccessToken(token: string): string;
|
||||
verifyRefreshToken(token: string): string;
|
||||
}
|
||||
|
||||
/** High-level authentication service combining password + token issuance. */
|
||||
export interface AuthService {
|
||||
authenticate(password: string, hash: string): Promise<boolean>;
|
||||
issueTokens(sub: string): [string, string];
|
||||
}
|
||||
@@ -9,6 +9,17 @@
|
||||
"typescript": "^5.7.0",
|
||||
},
|
||||
},
|
||||
"apps/packages/application": {
|
||||
"name": "@zesdex/application",
|
||||
"version": "1.21.2",
|
||||
"dependencies": {
|
||||
"@zesdex/domain": "workspace:*",
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/bun": "^1.2.0",
|
||||
"typescript": "^5.7.0",
|
||||
},
|
||||
},
|
||||
"apps/packages/domain": {
|
||||
"name": "@zesdex/domain",
|
||||
"version": "1.21.2",
|
||||
@@ -23,6 +34,8 @@
|
||||
|
||||
"@types/node": ["@types/node@26.4.1", "", { "dependencies": { "undici-types": "~8.3.0" } }, "sha512-k97ENvZWtvA6yqz5/FS6a7duDgOPEeOQOc2iKS/nY6mX6qJUKtLnWzQS+Xj6tXweyj6ZcTAK2Qecetnvi9nCLA=="],
|
||||
|
||||
"@zesdex/application": ["@zesdex/application@workspace:apps/packages/application"],
|
||||
|
||||
"@zesdex/domain": ["@zesdex/domain@workspace:apps/packages/domain"],
|
||||
|
||||
"bun-types": ["bun-types@1.4.0", "", { "dependencies": { "@types/node": "*" } }, "sha512-iIKw23BspnQQYd3prITOBxeUsxBHnwzX6YJfGMuNOZzeNcMmVqzIIVGRm1l69ogaPQmb4wB6BN8mA5bE9YuC5Q=="],
|
||||
|
||||
+6
-1
@@ -18,7 +18,12 @@
|
||||
"noUncheckedIndexedAccess": true,
|
||||
"noUnusedLocals": true,
|
||||
"noUnusedParameters": true,
|
||||
"noFallthroughCasesInSwitch": true
|
||||
"noFallthroughCasesInSwitch": true,
|
||||
"paths": {
|
||||
"@zesdex/domain": ["./apps/packages/domain/src/index.ts"],
|
||||
"@zesdex/application": ["./apps/packages/application/src/index.ts"],
|
||||
"@zesdex/infrastructure": ["./apps/packages/infrastructure/src/index.ts"]
|
||||
}
|
||||
},
|
||||
"include": ["apps/packages/*/src", "apps/interfaces/*/src", "apps/interfaces/*/bin"]
|
||||
}
|
||||
Reference in New Issue
Block a user