diff --git a/.ai-project/state.yaml b/.ai-project/state.yaml index 5e36c2f..abe99a4 100644 --- a/.ai-project/state.yaml +++ b/.ai-project/state.yaml @@ -26,10 +26,10 @@ 功能: - 编号: "F-001" 名称: "应用外壳、首次状态与普通会话" - 状态: "待用户查看" + 状态: "已完成" - 编号: "F-002" 名称: "Agent Core 与 DeepSeek 流式回复" - 状态: "待开始" + 状态: "待用户查看" - 编号: "F-003" 名称: "项目管理与项目聊天" 状态: "待开始" diff --git a/apps/web-server/src/composition/create-app.ts b/apps/web-server/src/composition/create-app.ts index 910ff58..c119b6b 100644 --- a/apps/web-server/src/composition/create-app.ts +++ b/apps/web-server/src/composition/create-app.ts @@ -4,16 +4,25 @@ import { createErrorHandler } from "../http/error-handler"; import { requestId } from "../http/request-id"; import { createHealthRoutes } from "../routes/health"; import { createConversationRoutes } from "../routes/conversations"; -import type { ConversationService } from "@great-agent/agent-core"; +import type { + AgentRunService, + ConversationService, +} from "@great-agent/agent-core"; +import type { RunRegistry } from "./run-registry"; +import { createRunRoutes } from "../routes/runs"; export function createApp( logger: Logger, conversations?: ConversationService, + runs?: AgentRunService, + runRegistry?: RunRegistry, ): Hono { const app = new Hono(); app.use("*", requestId); app.onError(createErrorHandler(logger)); app.route("/api", createHealthRoutes()); if (conversations) app.route("/api", createConversationRoutes(conversations)); + if (runs && runRegistry) + app.route("/api", createRunRoutes(runs, runRegistry)); return app; } diff --git a/apps/web-server/src/composition/run-registry.ts b/apps/web-server/src/composition/run-registry.ts new file mode 100644 index 0000000..241f3fa --- /dev/null +++ b/apps/web-server/src/composition/run-registry.ts @@ -0,0 +1,36 @@ +import type { RunEvent } from "@great-agent/agent-core"; + +type Listener = (event: RunEvent) => Promise; + +export class RunRegistry { + private readonly history = new Map(); + private readonly listeners = new Map>(); + private readonly finished = new Set(); + + async publish(event: RunEvent): Promise { + const events = this.history.get(event.runId) ?? []; + events.push(event); + this.history.set(event.runId, events); + if (event.type === "run.completed" || event.type === "run.failed") + this.finished.add(event.runId); + await Promise.all( + [...(this.listeners.get(event.runId) ?? [])].map((listener) => + listener(event), + ), + ); + } + + events(runId: string): readonly RunEvent[] { + return this.history.get(runId) ?? []; + } + isFinished(runId: string): boolean { + return this.finished.has(runId); + } + + subscribe(runId: string, listener: Listener): () => void { + const values = this.listeners.get(runId) ?? new Set(); + values.add(listener); + this.listeners.set(runId, values); + return () => values.delete(listener); + } +} diff --git a/apps/web-server/src/index.ts b/apps/web-server/src/index.ts index ed3fe1f..637660c 100644 --- a/apps/web-server/src/index.ts +++ b/apps/web-server/src/index.ts @@ -1,14 +1,22 @@ import { resolve } from "node:path"; import { FileConversationRepository, + FileRunRepository, createDataLayout, ensureDataLayout, } from "@great-agent/local-data"; -import { ConversationService } from "@great-agent/agent-core"; +import { + AgentRunService, + ConversationService, + CoreError, + type ModelPort, +} from "@great-agent/agent-core"; +import { DeepSeekModelAdapter } from "@great-agent/model-deepseek"; import { createApp } from "./composition/create-app"; import { createLogger } from "./composition/create-logger"; import { mountStaticWeb } from "./composition/static-web"; import { loadEnvironment } from "./config/environment"; +import { RunRegistry } from "./composition/run-registry"; const environment = loadEnvironment(); const logger = createLogger(environment.logLevel); @@ -20,7 +28,31 @@ const conversations = new ConversationService({ clock: { now: () => new Date() }, ids: { create: () => crypto.randomUUID() }, }); -const app = createApp(logger, conversations); +const clock = { now: () => new Date() }; +const ids = { create: () => crypto.randomUUID() }; +const model: ModelPort = environment.deepSeekApiKey + ? new DeepSeekModelAdapter({ + apiKey: environment.deepSeekApiKey, + baseURL: environment.deepSeekBaseUrl, + model: environment.deepSeekModel, + }) + : { + async *stream() { + yield { type: "response.completed" as const }; + throw new CoreError( + "MODEL_CONFIG_MISSING", + "尚未配置 DeepSeek API 密钥", + ); + }, + }; +const runs = new AgentRunService( + conversations, + model, + new FileRunRepository(layout), + clock, + ids, +); +const app = createApp(logger, conversations, runs, new RunRegistry()); mountStaticWeb(app, resolve(import.meta.dir, "../../web/dist")); logger.info( @@ -32,8 +64,10 @@ logger.info( "Great Agent 2 服务启动", ); -export default { +const server = Bun.serve({ hostname: environment.host, port: environment.port, fetch: app.fetch, -}; +}); + +logger.info({ url: server.url.href }, "Great Agent 2 已就绪"); diff --git a/apps/web-server/src/routes/conversations.test.ts b/apps/web-server/src/routes/conversations.test.ts index 9521c7e..a3509a1 100644 --- a/apps/web-server/src/routes/conversations.test.ts +++ b/apps/web-server/src/routes/conversations.test.ts @@ -42,29 +42,14 @@ function createTestApp() { } describe("普通会话路由", () => { - test("首条消息创建会话并可读取", async () => { + test("会话路由只提供查询,消息必须通过 Run 创建", async () => { const app = createTestApp(); - const created = await app.request("/api/conversations", { + expect(await (await app.request("/api/conversations")).json()).toEqual([]); + const directWrite = await app.request("/api/conversations", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ content: "你好" }), }); - expect(created.status).toBe(201); - const conversation = (await created.json()) as Conversation; - expect( - (await app.request(`/api/conversations/${conversation.id}`)).status, - ).toBe(200); - expect(await (await app.request("/api/conversations")).json()).toHaveLength( - 1, - ); - }); - - test("拒绝空消息", async () => { - const response = await createTestApp().request("/api/conversations", { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ content: " " }), - }); - expect(response.status).toBe(400); + expect(directWrite.status).toBe(404); }); }); diff --git a/apps/web-server/src/routes/conversations.ts b/apps/web-server/src/routes/conversations.ts index 25fa298..d46f3f9 100644 --- a/apps/web-server/src/routes/conversations.ts +++ b/apps/web-server/src/routes/conversations.ts @@ -1,6 +1,5 @@ import { Hono } from "hono"; import type { ConversationService } from "@great-agent/agent-core"; -import { messageInputSchema } from "@great-agent/web-contracts"; export function createConversationRoutes(service: ConversationService): Hono { const routes = new Hono(); @@ -10,18 +9,5 @@ export function createConversationRoutes(service: ConversationService): Hono { routes.get("/conversations/:id", async (context) => context.json(await service.getConversation(context.req.param("id"))), ); - routes.post("/conversations", async (context) => { - const input = messageInputSchema.parse(await context.req.json()); - return context.json( - await service.createWithFirstMessage(input.content), - 201, - ); - }); - routes.post("/conversations/:id/messages", async (context) => { - const input = messageInputSchema.parse(await context.req.json()); - return context.json( - await service.appendUserMessage(context.req.param("id"), input.content), - ); - }); return routes; } diff --git a/apps/web-server/src/routes/runs.test.ts b/apps/web-server/src/routes/runs.test.ts new file mode 100644 index 0000000..a4927d4 --- /dev/null +++ b/apps/web-server/src/routes/runs.test.ts @@ -0,0 +1,105 @@ +import { describe, expect, test } from "bun:test"; +import pino from "pino"; +import type { + AgentRun, + Conversation, + ConversationRepository, + ModelPort, + RunEvent, + RunRepository, +} from "@great-agent/agent-core"; +import { AgentRunService, ConversationService } from "@great-agent/agent-core"; +import { createApp } from "../composition/create-app"; +import { RunRegistry } from "../composition/run-registry"; + +class Conversations implements ConversationRepository { + value: Conversation | null = null; + async listRecent() { + return this.value ? [this.value] : []; + } + async getById(id: string) { + return this.value?.id === id ? this.value : null; + } + async create(value: Conversation) { + this.value = value; + } + async appendMessage(_id: string, message: Conversation["messages"][number]) { + if (!this.value) throw new Error("not found"); + this.value = { ...this.value, messages: [...this.value.messages, message] }; + return this.value; + } +} + +class Runs implements RunRepository { + value: AgentRun | null = null; + events: RunEvent[] = []; + async create(run: AgentRun) { + this.value = run; + } + async update(run: AgentRun) { + this.value = run; + } + async getById(id: string) { + return this.value?.id === id ? this.value : null; + } + async appendEvent(event: RunEvent) { + this.events.push(event); + } + async listEvents(runId: string) { + return this.events.filter((event) => event.runId === runId); + } +} + +describe("Run HTTP 与 SSE", () => { + test("启动普通 Run 并按顺序输出流式事件", async () => { + let id = 0; + const conversations = new Conversations(); + const service = new ConversationService({ + conversations, + clock: { now: () => new Date("2026-08-12T00:00:00.000Z") }, + ids: { create: () => `id_${++id}` }, + }); + const model: ModelPort = { + async *stream() { + yield { type: "text.delta", delta: "流式" }; + yield { type: "response.completed" }; + }, + }; + const runs = new AgentRunService( + service, + model, + new Runs(), + { now: () => new Date("2026-08-12T00:00:00.000Z") }, + { create: () => `id_${++id}` }, + ); + const app = createApp( + pino({ enabled: false }), + service, + runs, + new RunRegistry(), + ); + const started = await app.request("/api/runs", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ kind: "ordinary", message: "开始" }), + }); + expect(started.status).toBe(202); + const body = (await started.json()) as { + runId: string; + conversationId: string; + }; + const events = await app.request(`/api/runs/${body.runId}/events`); + const text = await events.text(); + expect(events.status).toBe(200); + expect(text).toContain("event: message.delta"); + expect(text).toContain("event: run.completed"); + expect( + (await service.getConversation(body.conversationId)).messages.at(-1) + ?.content, + ).toBe("流式"); + const conversation = await service.getConversation(body.conversationId); + const assistantMessage = conversation.messages.at(-1); + expect(assistantMessage?.runId).toBe(body.runId); + expect(text).toContain(`"messageId":"${assistantMessage?.id}"`); + }); +}); diff --git a/apps/web-server/src/routes/runs.ts b/apps/web-server/src/routes/runs.ts new file mode 100644 index 0000000..c74fe9d --- /dev/null +++ b/apps/web-server/src/routes/runs.ts @@ -0,0 +1,71 @@ +import { Hono } from "hono"; +import { streamSSE } from "hono/streaming"; +import type { AgentRunService, RunEvent } from "@great-agent/agent-core"; +import { startRunRequestSchema } from "@great-agent/web-contracts"; +import type { RunRegistry } from "../composition/run-registry"; + +export function createRunRoutes( + service: AgentRunService, + registry: RunRegistry, +): Hono { + const routes = new Hono(); + routes.post("/runs", async (context) => { + const input = startRunRequestSchema.parse(await context.req.json()); + const controller = new AbortController(); + const started = await service.start(input, controller.signal); + void consume(started.events, registry); + return context.json( + { + conversationId: started.run.conversationId, + runId: started.run.id, + status: "running" as const, + }, + 202, + ); + }); + routes.get("/runs/:runId/events", (context) => + streamSSE(context, async (stream) => { + const runId = context.req.param("runId"); + for (const event of registry.events(runId)) + await writeEvent(stream, event); + if (registry.isFinished(runId)) { + // EventSource needs the terminal event to reach the browser before the + // server closes a replay-only stream. + await stream.sleep(100); + return; + } + await new Promise((resolve) => { + const unsubscribe = registry.subscribe(runId, async (event) => { + await writeEvent(stream, event); + if (event.type === "run.completed" || event.type === "run.failed") { + unsubscribe(); + resolve(); + } + }); + stream.onAbort(() => { + unsubscribe(); + resolve(); + }); + }); + }), + ); + return routes; +} + +async function consume( + events: AsyncIterable, + registry: RunRegistry, +): Promise { + for await (const event of events) await registry.publish(event); +} + +async function writeEvent( + stream: Parameters[1]>[0], + event: RunEvent, +): Promise { + await stream.writeSSE({ + event: event.type, + id: String(event.sequence), + data: JSON.stringify(event), + }); +} diff --git a/apps/web/src/api/conversations.ts b/apps/web/src/api/conversations.ts index fef81b2..f8e2595 100644 --- a/apps/web/src/api/conversations.ts +++ b/apps/web/src/api/conversations.ts @@ -21,29 +21,6 @@ export async function getConversation( ); } -export async function createConversation( - content: string, -): Promise { - return conversationSchema.parse( - await request("/api/conversations", { - method: "POST", - body: JSON.stringify({ content }), - }), - ); -} - -export async function appendMessage( - id: string, - content: string, -): Promise { - return conversationSchema.parse( - await request(`/api/conversations/${encodeURIComponent(id)}/messages`, { - method: "POST", - body: JSON.stringify({ content }), - }), - ); -} - async function request(path: string, init?: RequestInit): Promise { const response = await fetch(path, { ...init, diff --git a/apps/web/src/api/runs.ts b/apps/web/src/api/runs.ts new file mode 100644 index 0000000..e8803eb --- /dev/null +++ b/apps/web/src/api/runs.ts @@ -0,0 +1,77 @@ +import { + runEventSchema, + startedRunSchema, + type RunEventResponse, + type StartedRunResponse, +} from "@great-agent/web-contracts"; + +export type StartRunInput = + | { kind: "ordinary"; message: string } + | { kind: "existing"; conversationId: string; message: string }; + +export async function startRun( + input: StartRunInput, +): Promise { + const response = await fetch("/api/runs", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify(input), + }); + const value = await response.json(); + if (!response.ok) throw new Error(readErrorMessage(value)); + return startedRunSchema.parse(value); +} + +export async function streamRun( + runId: string, + onEvent: (event: RunEventResponse) => void, +): Promise { + const response = await fetch( + `/api/runs/${encodeURIComponent(runId)}/events`, + { headers: { accept: "text/event-stream" } }, + ); + if (!response.ok || !response.body) throw new Error("流式连接中断"); + + const reader = response.body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + while (true) { + const { done, value } = await reader.read(); + buffer += decoder.decode(value, { stream: !done }).replaceAll("\r\n", "\n"); + const blocks = buffer.split("\n\n"); + buffer = blocks.pop() ?? ""; + for (const block of blocks) { + const data = block + .split("\n") + .filter((line) => line.startsWith("data:")) + .map((line) => line.slice(5).trimStart()) + .join("\n"); + if (!data) continue; + const event = runEventSchema.parse(JSON.parse(data)); + onEvent(event); + if (event.type === "run.completed") return; + if (event.type === "run.failed") + throw new Error(readPayloadMessage(event.payload)); + } + if (done) break; + } + throw new Error("流式连接提前结束"); +} + +function readErrorMessage(value: unknown): string { + if (typeof value === "object" && value && "error" in value) { + const error = value.error; + if ( + typeof error === "object" && + error && + "message" in error && + typeof error.message === "string" + ) + return error.message; + } + return "无法启动 Agent 任务"; +} + +function readPayloadMessage(payload: Record): string { + return typeof payload.message === "string" ? payload.message : "模型运行失败"; +} diff --git a/apps/web/src/app/app-controller.ts b/apps/web/src/app/app-controller.ts index ad61e96..b05db1f 100644 --- a/apps/web/src/app/app-controller.ts +++ b/apps/web/src/app/app-controller.ts @@ -3,12 +3,8 @@ import type { ConversationResponse, ConversationSummaryResponse, } from "@great-agent/web-contracts"; -import { - appendMessage, - createConversation, - getConversation, - listConversations, -} from "../api/conversations"; +import { getConversation, listConversations } from "../api/conversations"; +import { startRun, streamRun } from "../api/runs"; export type Selection = | { kind: "none" } @@ -21,6 +17,7 @@ export type AppController = Readonly<{ active: State; loading: State; sending: State; + assistantDraft: State; error: State; initialize(): Promise; startNewTask(): void; @@ -34,6 +31,7 @@ export function createAppController(): AppController { const active = van.state(null); const loading = van.state(false); const sending = van.state(false); + const assistantDraft = van.state(""); const error = van.state(""); async function initialize() { @@ -70,16 +68,28 @@ export function createAppController(): AppController { async function send(content: string) { if (!content.trim() || sending.val) return; sending.val = true; + assistantDraft.val = ""; error.val = ""; try { const current = selection.val; - const conversation = + const started = await startRun( current.kind === "conversation" - ? await appendMessage(current.id, content) - : await createConversation(content); - active.val = conversation; - selection.val = { kind: "conversation", id: conversation.id }; + ? { kind: "existing", conversationId: current.id, message: content } + : { kind: "ordinary", message: content }, + ); + active.val = await getConversation(started.conversationId); + selection.val = { kind: "conversation", id: started.conversationId }; recent.val = await listConversations(); + await streamRun(started.runId, (event) => { + if ( + event.type === "message.delta" && + typeof event.payload.delta === "string" + ) { + assistantDraft.val += event.payload.delta; + } + }); + active.val = await getConversation(started.conversationId); + assistantDraft.val = ""; } catch (cause) { error.val = readMessage(cause); } finally { @@ -93,6 +103,7 @@ export function createAppController(): AppController { active, loading, sending, + assistantDraft, error, initialize, startNewTask, diff --git a/apps/web/src/app/app-shell.css b/apps/web/src/app/app-shell.css index b3f0662..c5fb996 100644 --- a/apps/web/src/app/app-shell.css +++ b/apps/web/src/app/app-shell.css @@ -1,7 +1,9 @@ .app-shell { display: grid; grid-template-columns: var(--sidebar-width) minmax(0, 1fr); - min-height: 100vh; + width: 100vw; + height: 100dvh; + overflow: hidden; background: var(--color-canvas); } diff --git a/apps/web/src/features/conversations/conversation-pane.css b/apps/web/src/features/conversations/conversation-pane.css index e79732a..8d47dee 100644 --- a/apps/web/src/features/conversations/conversation-pane.css +++ b/apps/web/src/features/conversations/conversation-pane.css @@ -1,12 +1,23 @@ .main-pane { position: relative; + display: flex; + flex-direction: column; min-width: 0; - min-height: 100vh; - padding-bottom: 170px; + min-height: 0; + height: 100%; + overflow: hidden; background-color: var(--color-main); background-image: radial-gradient(#2e2e2c 0.65px, transparent 0.65px); background-size: 9px 9px; } +.pane-content { + position: relative; + flex: 1 1 auto; + min-height: 0; + overflow-x: hidden; + overflow-y: auto; + scrollbar-gutter: stable; +} .mobile-header { display: none; } @@ -37,7 +48,7 @@ .conversation-view { width: min(720px, calc(100% - 48px)); margin: 0 auto; - padding: 55px 0 30px; + padding: 55px 0 42px; } .conversation-title { margin: 0 0 34px; @@ -64,6 +75,20 @@ font-size: 11px; font-weight: 700; } +.assistant-role { + color: #ddd8cf; + background: #494641; +} +.message.streaming p::after { + display: inline-block; + width: 5px; + height: 14px; + margin-left: 3px; + background: var(--color-accent); + vertical-align: -2px; + animation: cursor-blink 1s steps(1) infinite; + content: ""; +} .message p { margin: 2px 0 0; color: #e1ded8; @@ -72,12 +97,10 @@ white-space: pre-wrap; } .composer-wrap { - position: fixed; - right: 0; - bottom: 18px; - left: var(--sidebar-width); - width: min(var(--composer-width), calc(100% - 236px)); - margin: 0 auto; + z-index: 1; + flex: 0 0 auto; + width: min(var(--composer-width), calc(100% - 40px)); + margin: 0 auto 18px; } .composer { min-height: 84px; @@ -154,15 +177,21 @@ font-size: 11px; } +@keyframes cursor-blink { + 50% { + opacity: 0; + } +} + @media (max-width: 700px) { .mobile-header { display: block; + flex: 0 0 auto; padding: 14px 16px; border-bottom: 1px solid #323230; font-size: 13px; } .composer-wrap { - left: 0; width: calc(100% - 24px); } .conversation-view { diff --git a/apps/web/src/features/conversations/conversation-pane.ts b/apps/web/src/features/conversations/conversation-pane.ts index 7d3fe98..6eda657 100644 --- a/apps/web/src/features/conversations/conversation-pane.ts +++ b/apps/web/src/features/conversations/conversation-pane.ts @@ -11,13 +11,13 @@ export function ConversationPane(controller: AppController): HTMLElement { const content = draft.val; if (!content.trim()) return; await controller.send(content); - if (!controller.error.val) draft.val = ""; + draft.val = ""; } return main( { class: "main-pane" }, header({ class: "mobile-header" }, "Great Agent 2"), - () => + div({ class: "pane-content" }, () => controller.selection.val.kind === "conversation" && controller.active.val ? div( { class: "conversation-view" }, @@ -27,10 +27,21 @@ export function ConversationPane(controller: AppController): HTMLElement { controller.active.val.messages.map((message) => article( { class: `message ${message.role}` }, - span({ class: "message-role" }, "你"), + span( + { class: `message-role ${message.role}-role` }, + message.role === "user" ? "你" : "G", + ), p(message.content), ), ), + () => + controller.assistantDraft.val + ? article( + { class: "message assistant streaming" }, + span({ class: "message-role assistant-role" }, "G"), + p(controller.assistantDraft), + ) + : null, ), ) : div( @@ -38,6 +49,7 @@ export function ConversationPane(controller: AppController): HTMLElement { h1(span({ class: "spark" }, "✳"), " 今天想完成什么?"), p("从一个问题、想法或具体任务开始。"), ), + ), Composer(controller, draft, submit), ); } @@ -49,10 +61,14 @@ function Composer( ): HTMLElement { return div( { class: "composer-wrap" }, - () => - controller.error.val - ? div({ class: "error-banner", role: "alert" }, controller.error.val) - : null, + div( + { + class: "error-banner", + role: "alert", + hidden: () => !controller.error.val, + }, + () => controller.error.val, + ), div( { class: "composer" }, textarea({ diff --git a/apps/web/src/features/conversations/sidebar.css b/apps/web/src/features/conversations/sidebar.css index bf1d9f9..0fb4d0e 100644 --- a/apps/web/src/features/conversations/sidebar.css +++ b/apps/web/src/features/conversations/sidebar.css @@ -2,11 +2,21 @@ position: relative; display: flex; flex-direction: column; - min-height: 100vh; + min-width: 0; + min-height: 0; + height: 100%; + overflow: hidden; padding: 10px 7px; border-right: 1px solid #323230; background: var(--color-sidebar); } +.sidebar-scroll { + flex: 1 1 auto; + min-height: 0; + overflow-x: hidden; + overflow-y: auto; + scrollbar-gutter: stable; +} .window-controls { display: flex; @@ -106,6 +116,7 @@ } .sidebar-footer { display: flex; + flex: 0 0 auto; gap: 9px; align-items: center; margin-top: auto; diff --git a/apps/web/src/features/conversations/sidebar.ts b/apps/web/src/features/conversations/sidebar.ts index 693af0c..fb74d59 100644 --- a/apps/web/src/features/conversations/sidebar.ts +++ b/apps/web/src/features/conversations/sidebar.ts @@ -19,8 +19,7 @@ export function Sidebar(controller: AppController): HTMLElement { "新任务", ), button({ class: "nav-item", disabled: true }, "▣", " 项目"), - h2("最近"), - () => + div({ class: "sidebar-scroll" }, h2("最近"), () => controller.loading.val && controller.recent.val.length === 0 ? div({ class: "sidebar-empty" }, "正在加载…") : controller.recent.val.length === 0 @@ -36,6 +35,7 @@ export function Sidebar(controller: AppController): HTMLElement { ), ), ), + ), div( { class: "sidebar-footer" }, span({ class: "avatar" }, "G"), diff --git a/apps/web/src/styles/global.css b/apps/web/src/styles/global.css index 899af1d..fafef92 100644 --- a/apps/web/src/styles/global.css +++ b/apps/web/src/styles/global.css @@ -2,11 +2,39 @@ * { box-sizing: border-box; + scrollbar-color: var(--color-scrollbar-thumb) transparent; + scrollbar-width: thin; +} + +*::-webkit-scrollbar { + width: 8px; + height: 8px; +} + +*::-webkit-scrollbar-track { + background: transparent; +} + +*::-webkit-scrollbar-thumb { + border: 2px solid transparent; + border-radius: 999px; + background: var(--color-scrollbar-thumb); + background-clip: padding-box; +} + +*::-webkit-scrollbar-thumb:hover { + background: var(--color-scrollbar-thumb-hover); + background-clip: padding-box; +} + +html, +body { + height: 100%; + overflow: hidden; } body { min-width: 320px; - min-height: 100vh; margin: 0; color: var(--color-text-primary); background: var(--color-canvas); @@ -14,6 +42,10 @@ body { font-synthesis: none; } +#app { + height: 100%; +} + button, textarea { color: inherit; diff --git a/apps/web/src/styles/tokens.css b/apps/web/src/styles/tokens.css index 798f115..26aadf2 100644 --- a/apps/web/src/styles/tokens.css +++ b/apps/web/src/styles/tokens.css @@ -8,6 +8,8 @@ --color-text-primary: #e8e5df; --color-text-secondary: #8a8781; --color-accent: #cf785e; + --color-scrollbar-thumb: #4b4945; + --color-scrollbar-thumb-hover: #66635d; --sidebar-width: 196px; --composer-width: 460px; --font-sans: diff --git a/docs/FEATURES.md b/docs/FEATURES.md index 87cefba..6d3ca9f 100644 --- a/docs/FEATURES.md +++ b/docs/FEATURES.md @@ -11,7 +11,7 @@ ### F-001——应用外壳、首次状态与普通会话 -- 状态:待用户查看(自动化开发验证已完成) +- 状态:已完成(用户于 2026-08-12 审核通过) - 用户可见结果:首次进入看到“最近”和“项目”空状态及右侧对话引导;点击“新任务”不会创建空会话,发送首条有效消息后普通会话出现,选择历史会话可还原消息并继续输入。 - 页面与交互:Claude Desktop 风格三栏外壳、空状态、对话引导、新任务草稿、普通历史会话、输入校验、加载与保存失败状态。 - 服务或接口:普通会话列表、会话详情、首条消息原子创建会话、已有会话追加用户消息。 @@ -28,8 +28,22 @@ ### F-002——Agent Core 与 DeepSeek 流式回复 -- 状态:待开始;F-001 获得用户确认后才能开始 -- 用户可见结果:普通会话获得真实 DeepSeek 流式回复并可连续聊天。 +- 状态:待用户查看 +- 用户可见结果:普通会话可启动 Agent Run,DeepSeek 回复增量显示,完成后的助手消息写入本地会话并可继续聊天;未配置模型密钥时保留用户消息并显示明确错误。 +- Agent Core:新增独立于 HTTP 的 `AgentRunService`、Run/RunEvent 领域对象、`ModelPort` 与 `RunRepository`;模型增量、完成和失败均转为稳定的领域事件。 +- DeepSeek 接入:`model-deepseek` 在适配器内部映射 OpenAI 兼容协议,Core 只依赖自身 Message/ModelEvent 类型;首个模型固定为 DeepSeek。 +- 服务与流式协议:新增 `POST /api/runs` 和 `GET /api/runs/:runId/events`;事件按序包含 Run 开始、消息开始/增量/完成及 Run 完成/失败。前端使用基于 `fetch` 的 SSE 读取,兼容不提供原生 `EventSource` 的 WebView。 +- 数据持久化:每次运行保存 `run.json` 和有序 `events.ndjson`;成功后助手消息写入真实本地 Conversation 文件,失败后保存稳定失败终态和错误代码。 +- 轮次关联:每条用户和助手消息都必须保存 `runId`;Run 使用 `triggerMessageId` 指向触发本轮运行的用户消息,并复用流式 `message.started` 给出的助手 `messageId`,因此运行、触发消息、流式事件和最终消息可以稳定互查。会话写接口只提供查询,新增消息必须通过 Run 用例,避免产生没有 Run 的孤立消息。“第几回合”由用户消息顺序计算,不持久化易失序号;重试链 `retryOfRunId` 仍在 F-005 实现。 +- 页面与交互:发送消息后进入运行状态、实时拼接助手草稿、完成后读取持久化消息;模型配置缺失时显示“尚未配置 DeepSeek API 密钥”,输入恢复可用且不丢失已提交的用户消息。 +- 异常状态:模型密钥缺失保留具体可行动提示;其他模型异常对外收敛为“模型服务暂时不可用”,不泄露供应商原始错误或密钥。 +- 自动化测试:Agent Core 成功/失败和轮次关联、文件 Run Repository、HTTP 与 SSE 消息标识一致性、禁止绕过 Run 直接写消息及原有回归共 13 个测试、43 个断言,全部通过。 +- 验证命令与结果:`bun run format:check`、`bun run lint`、`bun run typecheck`、`bun run check:file-size`、`bun run check:architecture`、`bun test`、`bun run build` 全部通过。 +- 人工验证步骤:以空临时数据目录和未配置 DeepSeek 密钥启动生产构建;发送“最终失败链路验收”;检查用户消息、错误提示、输入恢复以及 Run 事件文件和失败摘要。 +- 人工操作结果:2026-08-12 使用本机应用内浏览器完成;页面准确显示模型密钥缺失提示,输入恢复,用户消息保留;`events.ndjson` 依次写入 `run.started`、`message.started`、`run.failed`,`run.json` 状态为 `failed`。 +- 查看阶段修改:根据用户反馈将应用外壳锁定为浏览器可视区高度;页面根节点不再滚动,侧栏仅“最近”区域内部滚动,ConversationPane 仅消息内容区内部滚动,顶部操作、侧栏底部信息和输入组件保持固定。2026-08-13 在 1280×720 视口验证根页面高度与视口一致且无页面滚动,两个内容区均为独立 `overflow-y: auto` 容器。 +- 视觉参考:失败状态验收截图保存为 `docs/visual-reference/f002-model-config-error.jpg`。 +- 已知限制:真实 DeepSeek 成功调用需要用户在运行环境提供自己的密钥;自动化和本次人工验收未使用或读取真实密钥。停止与重试留在 F-005,服务重启后的事件重放加固留在 F-010。 - 主要验收:AC-004、AC-010、AC-014、AC-016。 ### F-003——项目管理与项目聊天 diff --git a/docs/visual-reference/f002-model-config-error.jpg b/docs/visual-reference/f002-model-config-error.jpg new file mode 100644 index 0000000..ab07a3c Binary files /dev/null and b/docs/visual-reference/f002-model-config-error.jpg differ diff --git a/packages/agent-core/src/domain/agent-run.ts b/packages/agent-core/src/domain/agent-run.ts new file mode 100644 index 0000000..dd789f2 --- /dev/null +++ b/packages/agent-core/src/domain/agent-run.ts @@ -0,0 +1,25 @@ +export type AgentRunStatus = "running" | "completed" | "failed"; + +export type AgentRun = Readonly<{ + id: string; + conversationId: string; + triggerMessageId: string; + status: AgentRunStatus; + createdAt: string; + updatedAt: string; + errorCode?: string; +}>; + +export type RunEvent = Readonly<{ + runId: string; + sequence: number; + timestamp: string; + type: + | "run.started" + | "message.started" + | "message.delta" + | "message.completed" + | "run.completed" + | "run.failed"; + payload: Record; +}>; diff --git a/packages/agent-core/src/domain/conversation.ts b/packages/agent-core/src/domain/conversation.ts index eb8c131..4075e17 100644 --- a/packages/agent-core/src/domain/conversation.ts +++ b/packages/agent-core/src/domain/conversation.ts @@ -5,6 +5,7 @@ export type Message = Readonly<{ role: MessageRole; content: string; createdAt: string; + runId: string; }>; export type Conversation = Readonly<{ diff --git a/packages/agent-core/src/index.ts b/packages/agent-core/src/index.ts index 60d43d9..82e9b61 100644 --- a/packages/agent-core/src/index.ts +++ b/packages/agent-core/src/index.ts @@ -5,8 +5,10 @@ export type { Message, MessageRole, } from "./domain/conversation"; +export type { AgentRun, AgentRunStatus, RunEvent } from "./domain/agent-run"; export type { ConversationRepository } from "./ports/conversation-repository"; export type { ModelEvent, ModelPort, ModelRequest } from "./ports/model-port"; +export type { RunRepository } from "./ports/run-repository"; export type { ClockPort, IdPort, @@ -16,3 +18,8 @@ export { ConversationService, type ConversationServiceDependencies, } from "./use-cases/conversation-service"; +export { + AgentRunService, + type StartedRun, + type StartRunInput, +} from "./use-cases/agent-run-service"; diff --git a/packages/agent-core/src/ports/model-port.ts b/packages/agent-core/src/ports/model-port.ts index b03728c..9a714fe 100644 --- a/packages/agent-core/src/ports/model-port.ts +++ b/packages/agent-core/src/ports/model-port.ts @@ -1,12 +1,13 @@ +import type { Message } from "../domain/conversation"; + export type ModelRequest = Readonly<{ model: string; - messages: readonly unknown[]; + messages: readonly Message[]; }>; -export type ModelEvent = Readonly<{ - type: string; - payload: unknown; -}>; +export type ModelEvent = + | Readonly<{ type: "text.delta"; delta: string }> + | Readonly<{ type: "response.completed" }>; export interface ModelPort { stream(request: ModelRequest, signal: AbortSignal): AsyncIterable; diff --git a/packages/agent-core/src/ports/run-repository.ts b/packages/agent-core/src/ports/run-repository.ts new file mode 100644 index 0000000..f1144d1 --- /dev/null +++ b/packages/agent-core/src/ports/run-repository.ts @@ -0,0 +1,9 @@ +import type { AgentRun, RunEvent } from "../domain/agent-run"; + +export interface RunRepository { + create(run: AgentRun): Promise; + update(run: AgentRun): Promise; + getById(id: string): Promise; + appendEvent(event: RunEvent): Promise; + listEvents(runId: string): Promise; +} diff --git a/packages/agent-core/src/use-cases/agent-run-service.test.ts b/packages/agent-core/src/use-cases/agent-run-service.test.ts new file mode 100644 index 0000000..64c72bd --- /dev/null +++ b/packages/agent-core/src/use-cases/agent-run-service.test.ts @@ -0,0 +1,144 @@ +import { describe, expect, test } from "bun:test"; +import type { + AgentRun, + Conversation, + ConversationRepository, + ModelPort, + RunEvent, + RunRepository, +} from ".."; +import { AgentRunService, ConversationService } from ".."; + +class MemoryConversationRepository implements ConversationRepository { + readonly values = new Map(); + async listRecent() { + return [...this.values.values()]; + } + async getById(id: string) { + return this.values.get(id) ?? null; + } + async create(value: Conversation) { + this.values.set(value.id, value); + } + async appendMessage(id: string, message: Conversation["messages"][number]) { + const current = this.values.get(id); + if (!current) throw new Error("not found"); + const next = { + ...current, + messages: [...current.messages, message], + updatedAt: message.createdAt, + }; + this.values.set(id, next); + return next; + } +} + +class MemoryRunRepository implements RunRepository { + value: AgentRun | null = null; + readonly events: RunEvent[] = []; + async create(run: AgentRun) { + this.value = run; + } + async update(run: AgentRun) { + this.value = run; + } + async getById(id: string) { + return this.value?.id === id ? this.value : null; + } + async appendEvent(event: RunEvent) { + this.events.push(event); + } + async listEvents(runId: string) { + return this.events.filter((event) => event.runId === runId); + } +} + +describe("AgentRunService", () => { + test("不依赖 HTTP 完成一次流式运行并持久化助手消息", async () => { + let nextId = 0; + const conversations = new MemoryConversationRepository(); + const conversationService = new ConversationService({ + conversations, + clock: { now: () => new Date("2026-08-12T00:00:00.000Z") }, + ids: { create: () => `id_${++nextId}` }, + }); + const model: ModelPort = { + async *stream() { + yield { type: "text.delta", delta: "你" }; + yield { type: "text.delta", delta: "好" }; + yield { type: "response.completed" }; + }, + }; + const runs = new MemoryRunRepository(); + const service = new AgentRunService( + conversationService, + model, + runs, + { now: () => new Date("2026-08-12T00:00:00.000Z") }, + { create: () => `id_${++nextId}` }, + ); + const started = await service.start( + { kind: "ordinary", message: "开始" }, + new AbortController().signal, + ); + const events: RunEvent[] = []; + for await (const event of started.events) events.push(event); + expect(events.map((event) => event.type)).toEqual([ + "run.started", + "message.started", + "message.delta", + "message.delta", + "message.completed", + "run.completed", + ]); + expect(runs.value?.status).toBe("completed"); + const conversation = await conversationService.getConversation( + started.run.conversationId, + ); + const assistantMessage = conversation.messages.at(-1); + const startedMessageId = events.find( + (event) => event.type === "message.started", + )?.payload.messageId as string | undefined; + expect(conversation.messages[0]?.id).toBe(started.run.triggerMessageId); + expect(conversation.messages[0]?.runId).toBe(started.run.id); + expect(assistantMessage?.content).toBe("你好"); + expect(assistantMessage?.runId).toBe(started.run.id); + expect(assistantMessage?.id).toBe(startedMessageId); + }); + + test("模型失败形成稳定失败终态", async () => { + const conversations = new MemoryConversationRepository(); + let id = 0; + const conversationService = new ConversationService({ + conversations, + clock: { now: () => new Date() }, + ids: { create: () => `id_${++id}` }, + }); + const model: ModelPort = { + async *stream() { + yield { type: "response.completed" }; + throw new Error("secret provider error"); + }, + }; + const runs = new MemoryRunRepository(); + const service = new AgentRunService( + conversationService, + model, + runs, + { now: () => new Date() }, + { create: () => `id_${++id}` }, + ); + const started = await service.start( + { kind: "ordinary", message: "开始" }, + new AbortController().signal, + ); + const events: RunEvent[] = []; + for await (const event of started.events) events.push(event); + expect(events.at(-1)?.type).toBe("run.failed"); + expect(events.at(-1)?.payload.message).toBe("模型服务暂时不可用"); + expect( + (await conversationService.getConversation(started.run.conversationId)) + .messages[0]?.id, + ).toBe(started.run.triggerMessageId); + }); +}); diff --git a/packages/agent-core/src/use-cases/agent-run-service.ts b/packages/agent-core/src/use-cases/agent-run-service.ts new file mode 100644 index 0000000..c511250 --- /dev/null +++ b/packages/agent-core/src/use-cases/agent-run-service.ts @@ -0,0 +1,124 @@ +import type { AgentRun, RunEvent } from "../domain/agent-run"; +import { CoreError } from "../errors/core-error"; +import type { ModelPort } from "../ports/model-port"; +import type { RunRepository } from "../ports/run-repository"; +import type { ClockPort, IdPort } from "../ports/system-ports"; +import type { ConversationService } from "./conversation-service"; + +export type StartRunInput = + | Readonly<{ kind: "ordinary"; message: string }> + | Readonly<{ kind: "existing"; conversationId: string; message: string }>; + +export type StartedRun = Readonly<{ + run: AgentRun; + events: AsyncIterable; +}>; + +export class AgentRunService { + constructor( + private readonly conversations: ConversationService, + private readonly model: ModelPort, + private readonly runs: RunRepository, + private readonly clock: ClockPort, + private readonly ids: IdPort, + ) {} + + async start(input: StartRunInput, signal: AbortSignal): Promise { + const runId = this.ids.create(); + const conversation = + input.kind === "ordinary" + ? await this.conversations.createWithFirstMessage(input.message, runId) + : await this.conversations.appendUserMessage( + input.conversationId, + input.message, + runId, + ); + const timestamp = this.clock.now().toISOString(); + const triggerMessage = conversation.messages.at(-1); + if (triggerMessage?.role !== "user") + throw new CoreError("RUN_TRIGGER_INVALID", "无法确定本次运行的用户消息"); + const run: AgentRun = { + id: runId, + conversationId: conversation.id, + triggerMessageId: triggerMessage.id, + status: "running", + createdAt: timestamp, + updatedAt: timestamp, + }; + await this.runs.create(run); + return { run, events: this.execute(run, signal) }; + } + + private async *execute( + run: AgentRun, + signal: AbortSignal, + ): AsyncIterable { + let sequence = 0; + const event = async ( + type: RunEvent["type"], + payload: Record, + ) => { + const value: RunEvent = { + runId: run.id, + sequence: ++sequence, + timestamp: this.clock.now().toISOString(), + type, + payload, + }; + await this.runs.appendEvent(value); + return value; + }; + try { + yield await event("run.started", { conversationId: run.conversationId }); + const messageId = this.ids.create(); + yield await event("message.started", { messageId }); + const conversation = await this.conversations.getConversation( + run.conversationId, + ); + let content = ""; + for await (const modelEvent of this.model.stream( + { model: "default", messages: conversation.messages }, + signal, + )) { + if (modelEvent.type !== "text.delta") continue; + content += modelEvent.delta; + yield await event("message.delta", { + messageId, + delta: modelEvent.delta, + }); + } + if (!content) + throw new CoreError("MODEL_RESPONSE_INVALID", "模型没有返回有效内容"); + await this.conversations.appendAssistantMessage( + run.conversationId, + messageId, + run.id, + content, + ); + yield await event("message.completed", { messageId, content }); + await this.runs.update({ + ...run, + status: "completed", + updatedAt: this.clock.now().toISOString(), + }); + yield await event("run.completed", {}); + } catch (cause) { + const code = + cause instanceof CoreError ? cause.code : "MODEL_UNAVAILABLE"; + await this.runs.update({ + ...run, + status: "failed", + errorCode: code, + updatedAt: this.clock.now().toISOString(), + }); + yield await event("run.failed", { + code, + message: safeErrorMessage(cause), + }); + } + } +} + +function safeErrorMessage(cause: unknown): string { + return cause instanceof CoreError ? cause.message : "模型服务暂时不可用"; +} diff --git a/packages/agent-core/src/use-cases/conversation-service.test.ts b/packages/agent-core/src/use-cases/conversation-service.test.ts index 7947905..d111ad2 100644 --- a/packages/agent-core/src/use-cases/conversation-service.test.ts +++ b/packages/agent-core/src/use-cases/conversation-service.test.ts @@ -39,9 +39,13 @@ describe("ConversationService", () => { clock: { now: () => new Date("2026-08-12T00:00:00.000Z") }, ids: { create: () => `id_${++nextId}` }, }); - const created = await service.createWithFirstMessage(" 第一次对话 "); + const created = await service.createWithFirstMessage( + " 第一次对话 ", + "run_1", + ); expect(created.projectId).toBeNull(); expect(created.messages[0]?.content).toBe("第一次对话"); + expect(created.messages[0]?.runId).toBe("run_1"); expect(await service.listRecent()).toHaveLength(1); }); @@ -52,9 +56,9 @@ describe("ConversationService", () => { clock: { now: () => new Date() }, ids: { create: () => "unused" }, }); - expect(service.createWithFirstMessage(" \n ")).rejects.toMatchObject({ - code: "MESSAGE_EMPTY", - }); + expect( + service.createWithFirstMessage(" \n ", "run_1"), + ).rejects.toMatchObject({ code: "MESSAGE_EMPTY" }); expect(conversations.values.size).toBe(0); }); }); diff --git a/packages/agent-core/src/use-cases/conversation-service.ts b/packages/agent-core/src/use-cases/conversation-service.ts index 9221809..58111ec 100644 --- a/packages/agent-core/src/use-cases/conversation-service.ts +++ b/packages/agent-core/src/use-cases/conversation-service.ts @@ -27,14 +27,17 @@ export class ConversationService { return conversation; } - async createWithFirstMessage(content: string): Promise { + async createWithFirstMessage( + content: string, + runId: string, + ): Promise { const normalized = normalizeContent(content); const timestamp = this.dependencies.clock.now().toISOString(); const conversation: Conversation = { id: this.dependencies.ids.create(), projectId: null, title: createTitle(normalized), - messages: [this.createUserMessage(normalized, timestamp)], + messages: [this.createUserMessage(normalized, timestamp, runId)], createdAt: timestamp, updatedAt: timestamp, }; @@ -42,7 +45,11 @@ export class ConversationService { return conversation; } - async appendUserMessage(id: string, content: string): Promise { + async appendUserMessage( + id: string, + content: string, + runId: string, + ): Promise { const normalized = normalizeContent(content); await this.getConversation(id); return this.dependencies.conversations.appendMessage( @@ -50,16 +57,38 @@ export class ConversationService { this.createUserMessage( normalized, this.dependencies.clock.now().toISOString(), + runId, ), ); } - private createUserMessage(content: string, createdAt: string): Message { + async appendAssistantMessage( + id: string, + messageId: string, + runId: string, + content: string, + ): Promise { + await this.getConversation(id); + return this.dependencies.conversations.appendMessage(id, { + id: messageId, + role: "assistant", + content, + createdAt: this.dependencies.clock.now().toISOString(), + runId, + }); + } + + private createUserMessage( + content: string, + createdAt: string, + runId: string, + ): Message { return { id: this.dependencies.ids.create(), role: "user", content, createdAt, + runId, }; } } diff --git a/packages/local-data/src/index.ts b/packages/local-data/src/index.ts index 9809766..2265e0d 100644 --- a/packages/local-data/src/index.ts +++ b/packages/local-data/src/index.ts @@ -7,3 +7,4 @@ export { type DataLayout, } from "./layout/data-layout"; export { FileConversationRepository } from "./repositories/file-conversation-repository"; +export { FileRunRepository } from "./repositories/file-run-repository"; diff --git a/packages/local-data/src/repositories/file-conversation-repository.test.ts b/packages/local-data/src/repositories/file-conversation-repository.test.ts index af67202..6730c53 100644 --- a/packages/local-data/src/repositories/file-conversation-repository.test.ts +++ b/packages/local-data/src/repositories/file-conversation-repository.test.ts @@ -32,6 +32,7 @@ describe("FileConversationRepository", () => { role: "user", content: "你好", createdAt: "2026-08-12T01:00:00.000Z", + runId: "run_1", }); expect((await repository.listRecent())[0]?.updatedAt).toBe( "2026-08-12T01:00:00.000Z", diff --git a/packages/local-data/src/repositories/file-run-repository.test.ts b/packages/local-data/src/repositories/file-run-repository.test.ts new file mode 100644 index 0000000..3382974 --- /dev/null +++ b/packages/local-data/src/repositories/file-run-repository.test.ts @@ -0,0 +1,45 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { mkdtemp, rm } from "node:fs/promises"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import type { AgentRun, RunEvent } from "@great-agent/agent-core"; +import { createDataLayout, ensureDataLayout } from "../layout/data-layout"; +import { FileRunRepository } from "./file-run-repository"; + +let root: string | undefined; +afterEach(async () => { + if (root) await rm(root, { recursive: true, force: true }); + root = undefined; +}); + +describe("FileRunRepository", () => { + test("持久化运行摘要和有序事件", async () => { + root = await mkdtemp(join(tmpdir(), "great-agent2-runs-")); + const layout = createDataLayout(root); + await ensureDataLayout(layout); + const repository = new FileRunRepository(layout); + const run: AgentRun = { + id: "run_1", + conversationId: "conversation_1", + triggerMessageId: "message_1", + status: "running", + createdAt: "2026-08-12T00:00:00.000Z", + updatedAt: "2026-08-12T00:00:00.000Z", + }; + const event: RunEvent = { + runId: run.id, + sequence: 1, + timestamp: run.createdAt, + type: "run.started", + payload: { conversationId: run.conversationId }, + }; + await repository.create(run); + await repository.appendEvent(event); + await repository.update({ ...run, status: "completed" }); + expect((await repository.getById(run.id))?.status).toBe("completed"); + expect((await repository.getById(run.id))?.triggerMessageId).toBe( + "message_1", + ); + expect(await repository.listEvents(run.id)).toEqual([event]); + }); +}); diff --git a/packages/local-data/src/repositories/file-run-repository.ts b/packages/local-data/src/repositories/file-run-repository.ts new file mode 100644 index 0000000..f5f5bc4 --- /dev/null +++ b/packages/local-data/src/repositories/file-run-repository.ts @@ -0,0 +1,64 @@ +import { mkdir, readFile } from "node:fs/promises"; +import { join } from "node:path"; +import type { + AgentRun, + RunEvent, + RunRepository, +} from "@great-agent/agent-core"; +import { appendNdjson } from "../atomic-writes/append-ndjson"; +import { writeJsonAtomically } from "../atomic-writes/write-json-atomically"; +import type { DataLayout } from "../layout/data-layout"; + +export class FileRunRepository implements RunRepository { + constructor(private readonly layout: DataLayout) {} + + async create(run: AgentRun): Promise { + await this.update(run); + } + + async update(run: AgentRun): Promise { + await mkdir(this.runDirectory(run.id), { recursive: true }); + await writeJsonAtomically(this.summaryPath(run.id), run); + } + + async getById(id: string): Promise { + try { + return JSON.parse( + await readFile(this.summaryPath(id), "utf8"), + ) as AgentRun; + } catch (error) { + if (error instanceof Error && "code" in error && error.code === "ENOENT") + return null; + throw error; + } + } + + async appendEvent(event: RunEvent): Promise { + await appendNdjson(this.eventsPath(event.runId), event); + } + + async listEvents(runId: string): Promise { + try { + const text = await readFile(this.eventsPath(runId), "utf8"); + return text + .trim() + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line) as RunEvent); + } catch (error) { + if (error instanceof Error && "code" in error && error.code === "ENOENT") + return []; + throw error; + } + } + + private runDirectory(id: string) { + return join(this.layout.runs, id); + } + private summaryPath(id: string) { + return join(this.runDirectory(id), "run.json"); + } + private eventsPath(id: string) { + return join(this.runDirectory(id), "events.ndjson"); + } +} diff --git a/packages/model-deepseek/src/adapter/deepseek-model-adapter.ts b/packages/model-deepseek/src/adapter/deepseek-model-adapter.ts new file mode 100644 index 0000000..e7cdb16 --- /dev/null +++ b/packages/model-deepseek/src/adapter/deepseek-model-adapter.ts @@ -0,0 +1,49 @@ +import OpenAI from "openai"; +import type { + Message, + ModelEvent, + ModelPort, + ModelRequest, +} from "@great-agent/agent-core"; + +export type DeepSeekModelOptions = Readonly<{ + apiKey: string; + baseURL: string; + model: string; +}>; + +export class DeepSeekModelAdapter implements ModelPort { + private readonly client: OpenAI; + + constructor(private readonly options: DeepSeekModelOptions) { + this.client = new OpenAI({ + apiKey: options.apiKey, + baseURL: options.baseURL, + }); + } + + async *stream( + request: ModelRequest, + signal: AbortSignal, + ): AsyncIterable { + const response = await this.client.chat.completions.create( + { + model: request.model === "default" ? this.options.model : request.model, + messages: request.messages.map(toModelMessage), + stream: true, + }, + { signal }, + ); + for await (const chunk of response) { + const delta = chunk.choices[0]?.delta.content; + if (delta) yield { type: "text.delta", delta }; + } + yield { type: "response.completed" }; + } +} + +function toModelMessage( + message: Message, +): OpenAI.Chat.Completions.ChatCompletionMessageParam { + return { role: message.role, content: message.content }; +} diff --git a/packages/model-deepseek/src/index.ts b/packages/model-deepseek/src/index.ts index 8d61db6..7acc507 100644 --- a/packages/model-deepseek/src/index.ts +++ b/packages/model-deepseek/src/index.ts @@ -1 +1,4 @@ -export const deepSeekAdapterPackage = "@great-agent/model-deepseek"; +export { + DeepSeekModelAdapter, + type DeepSeekModelOptions, +} from "./adapter/deepseek-model-adapter"; diff --git a/packages/web-contracts/src/events/run-event.ts b/packages/web-contracts/src/events/run-event.ts new file mode 100644 index 0000000..af5f06d --- /dev/null +++ b/packages/web-contracts/src/events/run-event.ts @@ -0,0 +1,18 @@ +import { z } from "zod"; + +export const runEventSchema = z.object({ + runId: z.string(), + sequence: z.number().int().positive(), + timestamp: z.string(), + type: z.enum([ + "run.started", + "message.started", + "message.delta", + "message.completed", + "run.completed", + "run.failed", + ]), + payload: z.record(z.string(), z.unknown()), +}); + +export type RunEventResponse = z.infer; diff --git a/packages/web-contracts/src/index.ts b/packages/web-contracts/src/index.ts index 97ae02e..8c43541 100644 --- a/packages/web-contracts/src/index.ts +++ b/packages/web-contracts/src/index.ts @@ -1,5 +1,8 @@ export { healthResponseSchema, type HealthResponse } from "./responses/health"; export { messageInputSchema, type MessageInput } from "./requests/message"; +export { startRunRequestSchema, type StartRunRequest } from "./requests/run"; +export { runEventSchema, type RunEventResponse } from "./events/run-event"; +export { startedRunSchema, type StartedRunResponse } from "./responses/run"; export { conversationSchema, conversationSummarySchema, diff --git a/packages/web-contracts/src/requests/run.ts b/packages/web-contracts/src/requests/run.ts new file mode 100644 index 0000000..3eda103 --- /dev/null +++ b/packages/web-contracts/src/requests/run.ts @@ -0,0 +1,15 @@ +import { z } from "zod"; + +export const startRunRequestSchema = z.discriminatedUnion("kind", [ + z.object({ + kind: z.literal("ordinary"), + message: z.string().trim().min(1).max(32_000), + }), + z.object({ + kind: z.literal("existing"), + conversationId: z.string().min(1), + message: z.string().trim().min(1).max(32_000), + }), +]); + +export type StartRunRequest = z.infer; diff --git a/packages/web-contracts/src/responses/conversation.ts b/packages/web-contracts/src/responses/conversation.ts index 7367aac..8ad7d92 100644 --- a/packages/web-contracts/src/responses/conversation.ts +++ b/packages/web-contracts/src/responses/conversation.ts @@ -5,6 +5,7 @@ export const messageSchema = z.object({ role: z.enum(["user", "assistant"]), content: z.string(), createdAt: z.string(), + runId: z.string(), }); export const conversationSummarySchema = z.object({ diff --git a/packages/web-contracts/src/responses/run.ts b/packages/web-contracts/src/responses/run.ts new file mode 100644 index 0000000..2c9176e --- /dev/null +++ b/packages/web-contracts/src/responses/run.ts @@ -0,0 +1,9 @@ +import { z } from "zod"; + +export const startedRunSchema = z.object({ + conversationId: z.string(), + runId: z.string(), + status: z.literal("running"), +}); + +export type StartedRunResponse = z.infer;