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,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;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user