feat(coding): keep one ongoing conversation per delivered agent
This commit is contained in:
@@ -17,6 +17,8 @@ export async function handleCodingTeacherRoutes(
|
||||
url: URL,
|
||||
ctx: HostApiContext
|
||||
): Promise<boolean> {
|
||||
const legacy = url.pathname.match(/^\/api\/coding\/projects\/([^/]+)\/agent-history(?:\/([^/]+))?$/);
|
||||
const conversation = url.pathname.match(/^\/api\/coding\/agent-conversations\/([^/]+)(?:\/(messages|events|save|seen|discussion|requests\/([^/]+)\/cancel))?$/);
|
||||
const source = url.pathname.match(
|
||||
/^\/api\/coding\/projects\/([^/]+)\/conversations\/([^/]+)\/teacher-topics(?:\/([^/]+))?(?:\/(messages|events|save|discussion|requests\/([^/]+)\/cancel))?$/
|
||||
);
|
||||
@@ -32,7 +34,7 @@ export async function handleCodingTeacherRoutes(
|
||||
const catalog = url.pathname === '/api/coding/teacher/teachers';
|
||||
const draft = url.pathname === '/api/coding/teacher-preview';
|
||||
const pending = url.pathname === '/api/coding/teacher-preview/pending-link';
|
||||
if (!source && !preview && !projectTopics && !checkIn && !config && !catalog && !draft && !pending) return false;
|
||||
if (!legacy && !conversation && !source && !preview && !projectTopics && !checkIn && !config && !catalog && !draft && !pending) return false;
|
||||
if ((config || catalog || draft || pending) && req.method !== 'GET') {
|
||||
sendJson(res, 405, { error: '不支持此操作。' });
|
||||
return true;
|
||||
@@ -47,6 +49,27 @@ export async function handleCodingTeacherRoutes(
|
||||
return true;
|
||||
}
|
||||
try {
|
||||
if (legacy) {
|
||||
if (req.method !== 'GET') sendJson(res, 405, { error: '旧记录仅供查看。' });
|
||||
else sendJson(res, 200, await service.legacyHistory(decodeURIComponent(legacy[1]), legacy[2] ? {
|
||||
id: decodeURIComponent(legacy[2]), sourceId: url.searchParams.get('source') ?? 'project', friend: url.searchParams.get('friend') === 'true',
|
||||
} : undefined));
|
||||
return true;
|
||||
}
|
||||
const agentId = conversation ? decodeURIComponent(conversation[1]) : undefined;
|
||||
if (agentId && !conversation![2] && req.method === 'GET') {
|
||||
sendJson(res, 200, await service.conversation(agentId, url.searchParams.get('before') ?? undefined));
|
||||
return true;
|
||||
}
|
||||
if (agentId && conversation![2] === 'messages' && req.method === 'POST') {
|
||||
sendJson(res, 202, await service.sendConversation(agentId, await parseJsonBody<TeacherSend>(req)));
|
||||
return true;
|
||||
}
|
||||
if (agentId && conversation![2] === 'seen' && req.method === 'POST') {
|
||||
const input = await parseJsonBody<{ requestId: string }>(req);
|
||||
sendJson(res, 200, await service.markConversationSeen(agentId, input.requestId));
|
||||
return true;
|
||||
}
|
||||
if (checkIn) {
|
||||
if (req.method !== 'POST') sendJson(res, 405, { error: '不支持此操作。' });
|
||||
else sendJson(res, 200, await service.checkIn(
|
||||
@@ -70,14 +93,16 @@ export async function handleCodingTeacherRoutes(
|
||||
sendJson(res, 200, await service.previewDefinition(revision));
|
||||
return true;
|
||||
}
|
||||
const scope: TeacherScope = source
|
||||
const continuous = agentId ? await service.conversation(agentId) : undefined;
|
||||
if (agentId && !continuous?.topic) throw new TeacherError(404, 'teacher_topic_not_found', '智能体聊天尚未开始。');
|
||||
const scope: TeacherScope = agentId ? { projectId: '', sourceId: 'project', agentId } : source
|
||||
? { projectId: decodeURIComponent(source[1]), sourceId: decodeURIComponent(source[2]) }
|
||||
: projectTopics
|
||||
? { projectId: decodeURIComponent(projectTopics[1]), sourceId: 'project', role }
|
||||
: { projectId: 'preview', sourceId: 'preview' };
|
||||
const id = source?.[3] ?? projectTopics?.[3] ?? preview?.[1],
|
||||
action = source?.[4] ?? projectTopics?.[4] ?? preview?.[2],
|
||||
requestId = source?.[5] ?? projectTopics?.[5] ?? preview?.[3];
|
||||
const id = continuous?.topic?.id ?? source?.[3] ?? projectTopics?.[3] ?? preview?.[1],
|
||||
action = conversation?.[2] ?? source?.[4] ?? projectTopics?.[4] ?? preview?.[2],
|
||||
requestId = conversation?.[3] ?? source?.[5] ?? projectTopics?.[5] ?? preview?.[3];
|
||||
if (!id && req.method === 'GET') {
|
||||
sendJson(res, 200, await service.list(scope));
|
||||
return true;
|
||||
@@ -112,6 +137,9 @@ export async function handleCodingTeacherRoutes(
|
||||
const buffered: unknown[] = [];
|
||||
let started = false;
|
||||
const close = await service.subscribe(scope, id, (topic) => {
|
||||
// GET owns paged history. Stream only the current turn and metadata;
|
||||
// Renderer merges by request ID/revision without replacing older pages.
|
||||
if (conversation) topic = { ...topic, requests: topic.requests.slice(-1) };
|
||||
if (!started) {
|
||||
buffered.push(topic);
|
||||
return;
|
||||
|
||||
@@ -212,7 +212,7 @@ export function prepareCloudTeacher(
|
||||
id: requestId,
|
||||
read_protocol: TEACHER_READ_PROTOCOL,
|
||||
scope: {
|
||||
project_id: topic.projectId,
|
||||
project_id: currentRequest?.projectId ?? topic.projectId,
|
||||
source_session_id: currentRequest?.sourceConversationId ?? topic.sourceConversationId,
|
||||
},
|
||||
tools: tools?.definitions.map((item) => item.function.name) ?? [],
|
||||
@@ -238,7 +238,7 @@ export function prepareCloudTeacher(
|
||||
if (previous && previous.status !== 'completed') {
|
||||
try {
|
||||
await transport.json(
|
||||
'/questions/' + encodeURIComponent(previous.id) + '/cancel',
|
||||
'/questions/' + encodeURIComponent(previous.cloudRequestId ?? previous.id) + '/cancel',
|
||||
{},
|
||||
bounded
|
||||
);
|
||||
@@ -249,8 +249,8 @@ export function prepareCloudTeacher(
|
||||
let queued = await transport.json(
|
||||
'/questions',
|
||||
{
|
||||
teacher_version: topic.version,
|
||||
thread_id: topic.id,
|
||||
teacher_version: currentRequest?.teacherVersion ?? topic.version,
|
||||
thread_id: currentRequest?.runtimeThreadId ?? topic.id,
|
||||
request_id: requestId,
|
||||
query,
|
||||
local_context: localContext,
|
||||
|
||||
@@ -60,7 +60,8 @@ export function teacherHistoryMessages(history: TeacherRequest[]): TeacherSource
|
||||
return history.filter(request => request.status === 'completed').flatMap((request): TeacherSourceMessage[] => [
|
||||
...(request.intent === 'check-in' ? [] : [{
|
||||
id: 'teacher:' + request.id + ':user', role: 'user' as const,
|
||||
text: [...request.references.map(ref => '明确引用:\n' + ref.text), request.text].join('\n\n'),
|
||||
text: [...(request.projectId ? ['当时的项目:' + (request.projectName ?? request.projectId)] : []),
|
||||
...request.references.map(ref => '明确引用:\n' + ref.text), request.text].join('\n\n'),
|
||||
}]),
|
||||
{
|
||||
id: 'teacher:' + request.id + ':assistant', role: 'assistant' as const,
|
||||
|
||||
235
electron/coding-teacher/conversation-store.ts
Normal file
235
electron/coding-teacher/conversation-store.ts
Normal file
@@ -0,0 +1,235 @@
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import { readdir } from 'node:fs/promises';
|
||||
import path from 'node:path';
|
||||
import type { TeacherDefinition, TeacherHistoryPage, TeacherRequest, TeacherTopic } from '../../shared/coding-teacher';
|
||||
import { excerptTeacherText, teacherHistoryMessages } from './context';
|
||||
import { atomicWriteJson, readJsonFile } from '../coding-projects/atomic-json';
|
||||
import { TeacherError } from './config-client';
|
||||
import { teacherTopicId } from './store';
|
||||
|
||||
interface TurnIndex { id: string; createdAt: string; origin?: string; messages: ReturnType<typeof teacherHistoryMessages> }
|
||||
interface Manifest {
|
||||
topic: Omit<TeacherTopic, 'requests'>;
|
||||
turns: TurnIndex[];
|
||||
importedDiscussions?: Record<string, string>;
|
||||
}
|
||||
const PAGE_SIZE = 50;
|
||||
const isMissing = (error: unknown) => (error as NodeJS.ErrnoException)?.code === 'ENOENT';
|
||||
const originKey = (request: TeacherRequest) => request.origin
|
||||
? [request.origin.projectId, request.origin.topicId, request.origin.requestId].join(':') : undefined;
|
||||
const order = (a: TurnIndex, b: TurnIndex) => a.createdAt.localeCompare(b.createdAt) || (a.origin ?? a.id).localeCompare(b.origin ?? b.id);
|
||||
const indexTurn = (turn: TeacherRequest): TurnIndex => ({ id: turn.id, createdAt: turn.createdAt, origin: originKey(turn),
|
||||
messages: teacherHistoryMessages([turn]).map(message => ({ ...message, text: excerptTeacherText(message.text, 180).replaceAll('\n', ' ') })) });
|
||||
|
||||
/** A small index and one atomic file per turn. Only the latest page is held live. */
|
||||
export class TeacherConversationStore {
|
||||
private manifest?: Manifest;
|
||||
private live?: TeacherTopic;
|
||||
private loaded?: Promise<void>;
|
||||
private writes: Promise<unknown> = Promise.resolve();
|
||||
private readonly importedProjects = new Set<string>();
|
||||
constructor(readonly directory: string, readonly accountId: string, readonly agentId: string) {}
|
||||
|
||||
private manifestPath() { return path.join(this.directory, 'conversation.json'); }
|
||||
private turnPath(id: string) { return path.join(this.directory, 'turns', teacherTopicId(id) + '.json'); }
|
||||
private async load() {
|
||||
this.loaded ??= this.loadFiles().catch(error => { this.loaded = undefined; throw error; });
|
||||
return this.loaded;
|
||||
}
|
||||
private async loadFiles() {
|
||||
try { this.manifest = await readJsonFile(this.manifestPath()) as Manifest; }
|
||||
catch (error) { if (isMissing(error)) return; throw error; }
|
||||
const manifest = this.manifest;
|
||||
if (manifest.topic.accountId !== this.accountId || manifest.topic.conversation?.agentId !== this.agentId
|
||||
|| !Array.isArray(manifest.turns))
|
||||
throw new TeacherError(409, 'teacher_history_invalid', '智能体历史无法读取,请保留本机记录。');
|
||||
// A turn can finish writing just before the index write is interrupted.
|
||||
// Recover only unindexed files, without reading every historical body.
|
||||
let names: string[] = [];
|
||||
try { names = await readdir(path.join(this.directory, 'turns')); }
|
||||
catch (error) { if (!isMissing(error)) throw error; }
|
||||
const indexed = new Set(manifest.turns.map(turn => turn.id));
|
||||
let recovered = false;
|
||||
for (const name of names.filter(name => name.endsWith('.json'))) {
|
||||
const id = name.slice(0, -5);
|
||||
if (indexed.has(id)) continue;
|
||||
const turn = await this.turn(id);
|
||||
manifest.turns.push(indexTurn(turn));
|
||||
recovered = true;
|
||||
}
|
||||
manifest.turns.sort(order);
|
||||
this.live = { ...manifest.topic, requests: (await this.page()).requests };
|
||||
for (const turn of this.live.requests) {
|
||||
if (turn.status !== 'preparing' && turn.status !== 'running') continue;
|
||||
turn.status = 'interrupted';
|
||||
turn.error = '应用已重启,本次回复中断。';
|
||||
await atomicWriteJson(this.turnPath(turn.id), turn);
|
||||
recovered = true;
|
||||
}
|
||||
if (recovered) {
|
||||
this.live.revision++;
|
||||
await this.save(this.live);
|
||||
}
|
||||
}
|
||||
private async turn(id: string) {
|
||||
const turn = await readJsonFile(this.turnPath(id)) as TeacherRequest;
|
||||
if (turn.id !== id || !Array.isArray(turn.references))
|
||||
throw new TeacherError(409, 'teacher_history_invalid', '智能体消息无法读取,请保留本机记录。');
|
||||
return turn;
|
||||
}
|
||||
async exists() { await this.load(); return Boolean(this.manifest); }
|
||||
async ensure(definition: TeacherDefinition, version: number): Promise<TeacherTopic> {
|
||||
await this.load();
|
||||
if (this.live) return this.live;
|
||||
const now = new Date().toISOString();
|
||||
const topic: TeacherTopic = {
|
||||
schemaVersion: 1, revision: 0, id: randomUUID(), accountId: this.accountId,
|
||||
projectId: '', sourceConversationId: 'project', definition: structuredClone(definition), version,
|
||||
createdAt: now, updatedAt: now, requests: [],
|
||||
conversation: { agentId: this.agentId, segmentTurns: 0, discussions: {} },
|
||||
};
|
||||
// The manifest is durable before any turn file, so orphan recovery knows its owner.
|
||||
const manifest: Manifest = { topic: { ...topic }, turns: [] };
|
||||
delete (manifest.topic as Partial<TeacherTopic>).requests;
|
||||
await atomicWriteJson(this.manifestPath(), manifest);
|
||||
this.manifest = manifest;
|
||||
this.live = topic;
|
||||
return topic;
|
||||
}
|
||||
async read(id?: string): Promise<TeacherTopic> {
|
||||
await this.load();
|
||||
if (!this.live || (id && id !== this.live.id))
|
||||
throw new TeacherError(404, 'teacher_topic_not_found', '智能体聊天尚未开始。');
|
||||
return this.live;
|
||||
}
|
||||
async findRequest(id: string): Promise<TeacherRequest | undefined> {
|
||||
await this.load();
|
||||
const current = this.live?.requests.find(turn => turn.id === id);
|
||||
if (current) return current;
|
||||
return this.manifest?.turns.some(turn => turn.id === id) ? this.turn(id) : undefined;
|
||||
}
|
||||
async history() {
|
||||
await this.load();
|
||||
const entries = structuredClone(this.manifest?.turns ?? []);
|
||||
const owners = new Map(entries.flatMap(turn => turn.messages.map(message => [message.id, turn.id] as const)));
|
||||
return {
|
||||
messages: entries.flatMap(turn => turn.messages),
|
||||
read: async (id: string) => {
|
||||
const owner = owners.get(id);
|
||||
return owner ? teacherHistoryMessages([await this.turn(owner)]).find(message => message.id === id) : undefined;
|
||||
},
|
||||
};
|
||||
}
|
||||
async page(before?: string, limit = PAGE_SIZE): Promise<TeacherHistoryPage> {
|
||||
// loadFiles calls page only after the manifest is available.
|
||||
if (!this.manifest) {
|
||||
if (!this.loaded) await this.load();
|
||||
if (!this.manifest) return { requests: [], before: null };
|
||||
}
|
||||
const turns = this.manifest.turns;
|
||||
const end = before ? turns.findIndex(turn => turn.id === before) : turns.length;
|
||||
if (end < 0) throw new TeacherError(400, 'teacher_cursor_invalid', '历史位置已变化,请重新打开聊天。');
|
||||
const start = Math.max(0, end - Math.min(PAGE_SIZE, Math.max(1, limit)));
|
||||
const requests = await Promise.all(turns.slice(start, end).map(turn =>
|
||||
this.live?.requests.find(live => live.id === turn.id) ?? this.turn(turn.id)));
|
||||
return { requests: structuredClone(requests), before: start > 0 ? turns[start].id : null };
|
||||
}
|
||||
async save(topic: TeacherTopic, requestId?: string) {
|
||||
const next = this.writes.catch(() => undefined).then(() => this.persist(topic, requestId));
|
||||
this.writes = next;
|
||||
await next;
|
||||
}
|
||||
private async persist(topic: TeacherTopic, requestId?: string) {
|
||||
if (!this.manifest) throw new Error('Conversation not initialized');
|
||||
const requests = requestId !== undefined ? topic.requests.filter(turn => turn.id === requestId) : topic.requests;
|
||||
const turns = [...this.manifest.turns];
|
||||
for (const turn of requests) {
|
||||
await atomicWriteJson(this.turnPath(turn.id), turn);
|
||||
const next = indexTurn(turn);
|
||||
const index = turns.findIndex(item => item.id === turn.id);
|
||||
if (index < 0) turns.push(next); else turns[index] = next;
|
||||
}
|
||||
turns.sort(order);
|
||||
const { requests: _requests, unsaved: _unsaved, ...metadata } = topic;
|
||||
const next = { topic: metadata, turns, importedDiscussions: this.manifest.importedDiscussions };
|
||||
try { await atomicWriteJson(this.manifestPath(), next); }
|
||||
catch (error) {
|
||||
// The body may already be durable. Reload and recover it before a retry,
|
||||
// including its original legacy identity, instead of importing twice.
|
||||
this.loaded = undefined;
|
||||
throw error;
|
||||
}
|
||||
this.manifest = next;
|
||||
topic.unsaved = false;
|
||||
// Keep references used by an active run; trim only at the next read/accept boundary.
|
||||
this.live = topic;
|
||||
}
|
||||
async recent() {
|
||||
const topic = await this.read();
|
||||
topic.requests = (await this.page()).requests;
|
||||
return topic;
|
||||
}
|
||||
async select(_id: string) { /* There is only one conversation per agent. */ }
|
||||
async list() {
|
||||
await this.load();
|
||||
const topic = this.live;
|
||||
return { items: topic ? [{ id: topic.id, title: topic.definition.name, updatedAt: topic.updatedAt,
|
||||
version: topic.version, teacherId: this.agentId }] : [], lastSelectedTopicId: topic?.id ?? null };
|
||||
}
|
||||
async markSeen(requestId: string) {
|
||||
const topic = await this.read();
|
||||
const turns = this.manifest!.turns;
|
||||
const index = turns.findIndex(turn => turn.id === requestId);
|
||||
if (index < 0 || index <= turns.findIndex(turn => turn.id === topic.conversation!.seenThrough)) return;
|
||||
topic.conversation!.seenThrough = requestId;
|
||||
topic.revision++;
|
||||
await this.save(topic, '');
|
||||
}
|
||||
/** Old files remain untouched. Only a proven account + config ID can be imported. */
|
||||
async importProject(project: { id: string; path: string; name: string }) {
|
||||
if (this.importedProjects.has(project.id)) return;
|
||||
await this.load();
|
||||
const root = path.join(project.path, '.makelore', 'teacher-conversations', this.accountId);
|
||||
let scopes: string[];
|
||||
try { scopes = await readdir(root); }
|
||||
catch (error) { if (isMissing(error)) return; throw error; }
|
||||
for (const scope of scopes.sort()) {
|
||||
let names: string[];
|
||||
try { names = await readdir(path.join(root, scope)); }
|
||||
catch (error) { if (isMissing(error) || (error as NodeJS.ErrnoException).code === 'ENOTDIR') continue; throw error; }
|
||||
for (const name of names.filter(name => /^[0-9a-f-]{36}\.json$/i.test(name)).sort()) {
|
||||
const old = await readJsonFile(path.join(root, scope, name)) as TeacherTopic;
|
||||
if (old.accountId !== this.accountId || old.projectId !== project.id
|
||||
|| old.definition?.config_id !== this.agentId || old.role === 'friend') continue;
|
||||
const topic = await this.ensure(old.definition, old.version);
|
||||
const origins = new Set(this.manifest!.turns.map(turn => turn.origin).filter(Boolean));
|
||||
for (const request of old.requests) {
|
||||
const origin = { projectId: project.id, topicId: old.id, requestId: request.id };
|
||||
if (origins.has([origin.projectId, origin.topicId, origin.requestId].join(':'))) continue;
|
||||
const turn: TeacherRequest = {
|
||||
...structuredClone(request), id: randomUUID(), origin, projectId: project.id, projectName: project.name,
|
||||
sourceConversationId: request.sourceConversationId ?? (old.sourceConversationId !== 'project' ? old.sourceConversationId : undefined),
|
||||
teacherVersion: old.version, runtimeThreadId: old.id,
|
||||
...(old.definition.runtime === 'yuxi' ? { cloudRequestId: request.cloudRequestId ?? request.id } : {}),
|
||||
};
|
||||
if (turn.status === 'running' || turn.status === 'preparing') {
|
||||
turn.status = 'interrupted'; turn.error = '旧话题中的回复已中断。';
|
||||
}
|
||||
topic.requests.push(turn);
|
||||
topic.updatedAt = [topic.updatedAt, old.updatedAt].sort().at(-1)!;
|
||||
topic.revision++;
|
||||
await this.save(topic, turn.id);
|
||||
topic.requests = topic.requests.slice(-PAGE_SIZE);
|
||||
}
|
||||
if (old.discussion && (!topic.conversation!.discussions[project.id]
|
||||
|| old.updatedAt > (this.manifest!.importedDiscussions?.[project.id] ?? ''))) {
|
||||
(this.manifest!.importedDiscussions ??= {})[project.id] = old.updatedAt;
|
||||
topic.conversation!.discussions[project.id] = structuredClone(old.discussion);
|
||||
await this.save(topic, '');
|
||||
}
|
||||
}
|
||||
}
|
||||
if (this.manifest) await this.recent();
|
||||
this.importedProjects.add(project.id);
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
import path from 'node:path';
|
||||
import type { TeacherRequest, TeacherSourceContext } from '../../shared/coding-teacher';
|
||||
import type { TeacherRequest, TeacherSourceContext, TeacherSourceMessage } from '../../shared/coding-teacher';
|
||||
import { CodingProjectFileService } from '../coding-projects/project-files';
|
||||
import { excerptTeacherText, teacherHistoryMessages } from './context';
|
||||
import { readInteger, readPage, textChunks } from './read-page';
|
||||
@@ -14,6 +14,7 @@ export interface TeacherReadAccess {
|
||||
projectPath: string;
|
||||
source: TeacherSourceContext;
|
||||
history?: TeacherRequest[];
|
||||
archive?: { messages: TeacherSourceMessage[]; read(id: string): Promise<TeacherSourceMessage | undefined> };
|
||||
assertCurrent(): void;
|
||||
}
|
||||
|
||||
@@ -32,7 +33,7 @@ export const teacherReadToolDefinitions = [
|
||||
parameters: { type: 'object', properties: { path: { type: 'string' }, ...lineParameters }, required: ['path'], additionalProperties: false },
|
||||
} },
|
||||
{ type: 'function', function: {
|
||||
name: 'read_conversation', description: 'Read the captured active coding conversation and current teacher topic. Omit message_id to list messages; supply it to read continuous original text. Results give start_line/start_column and next; do not reread unchanged ranges.',
|
||||
name: 'read_conversation', description: 'Read the captured active coding conversation and this agent chat history. Historical messages can belong to earlier projects and are labelled as such; files remain limited to the current project. Omit message_id to list messages; supply it to read continuous original text. Results give start_line/start_column and next; do not reread unchanged ranges.',
|
||||
parameters: { type: 'object', properties: { message_id: { type: 'string' }, ...lineParameters }, additionalProperties: false },
|
||||
} },
|
||||
];
|
||||
@@ -50,6 +51,7 @@ function projectPath(value: unknown): string {
|
||||
export function createTeacherReadTools(access: TeacherReadAccess) {
|
||||
const files = new CodingProjectFileService();
|
||||
const messages = [...access.source.messages, ...teacherHistoryMessages(access.history ?? [])];
|
||||
const index = [...access.source.messages, ...(access.archive?.messages ?? teacherHistoryMessages(access.history ?? []))];
|
||||
return {
|
||||
definitions: teacherReadToolDefinitions,
|
||||
async execute(name: string, rawArguments: string, signal: AbortSignal, maxBytes = TEACHER_READ_PAGE_BYTES): Promise<string> {
|
||||
@@ -96,9 +98,10 @@ export function createTeacherReadTools(access: TeacherReadAccess) {
|
||||
break;
|
||||
}
|
||||
case 'read_conversation': {
|
||||
const message = args.message_id === undefined ? undefined : messages.find(item => item.id === args.message_id);
|
||||
const message = args.message_id === undefined ? undefined : messages.find(item => item.id === args.message_id)
|
||||
?? (typeof args.message_id === 'string' ? await access.archive?.read(args.message_id) : undefined);
|
||||
if (args.message_id !== undefined && !message) throw new Error('Message is not in the current conversation.');
|
||||
const text = message?.text ?? (messages.map(item =>
|
||||
const text = message?.text ?? (index.map(item =>
|
||||
`${item.id} ${item.role}: ${excerptTeacherText(item.text, 180).replaceAll('\n', ' ')}`).join('\n')
|
||||
|| '(no completed text messages in this conversation)');
|
||||
({ content: result, truncated } = await readPage(message ? message.id + ' ' + message.role : 'conversation message index (snippets)',
|
||||
|
||||
@@ -11,9 +11,11 @@ import type {
|
||||
TeacherDefinition,
|
||||
TeacherDiscussionAction,
|
||||
TeacherReference,
|
||||
TeacherRequest,
|
||||
TeacherSend,
|
||||
TeacherSourceContext,
|
||||
TeacherTopic,
|
||||
TeacherTopicList,
|
||||
} from '../../shared/coding-teacher';
|
||||
import { TEACHER_CHECK_IN_INTERVAL_MS, TEACHER_UNCHANGED_CHECK_IN_INTERVAL_MS } from '../../shared/coding-teacher';
|
||||
import {
|
||||
@@ -27,6 +29,8 @@ import {
|
||||
type TeacherAccount,
|
||||
} from './config-client';
|
||||
import { TeacherTopicStore, teacherTopicId } from './store';
|
||||
import { TeacherConversationStore } from './conversation-store';
|
||||
import { readJsonFile } from '../coding-projects/atomic-json';
|
||||
import { compileTeacherContext } from './context';
|
||||
import { prepareTeacherModel } from './model-runner';
|
||||
import { prepareCloudTeacher } from './cloud-runner';
|
||||
@@ -36,6 +40,7 @@ import { applyDiscussionReply, discussionInstructions, editDiscussion, validateD
|
||||
import { subscribeWorksSquareSession } from '../services/works-square-session';
|
||||
|
||||
export interface TeacherScope {
|
||||
agentId?: string;
|
||||
projectId: string;
|
||||
sourceId: string;
|
||||
role?: LegacyConsultationRole;
|
||||
@@ -59,6 +64,7 @@ export interface TeacherServiceOptions {
|
||||
readSource?(scope: TeacherScope): Promise<TeacherSourceContext>;
|
||||
}
|
||||
export class CodingTeacherService {
|
||||
private readonly conversations = new Map<string, TeacherConversationStore>();
|
||||
private readonly stores = new Map<string, TeacherTopicStore>();
|
||||
private readonly tails = new Map<string, Promise<unknown>>();
|
||||
private readonly active = new Map<
|
||||
@@ -99,10 +105,103 @@ export class CodingTeacherService {
|
||||
async catalog() {
|
||||
return (this.options.catalog ?? teacherCatalog)(await this.account());
|
||||
}
|
||||
async legacyHistory(projectId: string, selection?: { id: string; sourceId: string; friend: boolean }) {
|
||||
const account = await this.account();
|
||||
const project = await this.options.projects.getProject(projectId);
|
||||
const directory = (friend: boolean) => path.join(project.path, '.makelore', friend ? 'friend-conversations' : 'teacher-conversations', account.id);
|
||||
const read = async (friend: boolean, source: string, id: string) => {
|
||||
const topic = await readJsonFile(path.join(directory(friend), source === 'project' ? source : teacherTopicId(source), teacherTopicId(id) + '.json')) as TeacherTopic;
|
||||
if (topic.accountId !== account.id || topic.projectId !== projectId)
|
||||
throw new TeacherError(404, 'teacher_topic_not_found', '旧记录不存在。');
|
||||
this.assertAccount(account);
|
||||
return topic;
|
||||
};
|
||||
if (selection) return read(selection.friend, selection.sourceId, selection.id);
|
||||
const items: TeacherTopicList['items'] = [];
|
||||
for (const friend of [false, true]) {
|
||||
let scopes: string[];
|
||||
try { scopes = await readdir(directory(friend)); }
|
||||
catch (error) { if ((error as NodeJS.ErrnoException).code === 'ENOENT') continue; throw error; }
|
||||
for (const source of scopes.filter(value => value === 'project' || /^[0-9a-f-]{36}$/i.test(value))) {
|
||||
for (const name of await readdir(path.join(directory(friend), source))) {
|
||||
if (!/^[0-9a-f-]{36}\.json$/i.test(name)) continue;
|
||||
const topic = await read(friend, source, name.slice(0, -5));
|
||||
items.push({ id: topic.id, title: topic.requests.find(turn => turn.text)?.text.slice(0, 48) || topic.definition.name,
|
||||
updatedAt: topic.updatedAt, version: topic.version, teacherId: topic.definition.config_id,
|
||||
legacySourceId: source, ...(friend ? { legacyRole: 'friend' as const } : {}) });
|
||||
}
|
||||
}
|
||||
}
|
||||
this.assertAccount(account);
|
||||
return { items: items.sort((a, b) => b.updatedAt.localeCompare(a.updatedAt)), lastSelectedTopicId: null };
|
||||
}
|
||||
private conversationStore(account: TeacherAccount, agentId: string) {
|
||||
if (!agentId || agentId.length > 128 || !/^[a-zA-Z0-9_-]+$/.test(agentId))
|
||||
throw new TeacherError(400, 'teacher_agent_invalid', '智能体标识无效。');
|
||||
const key = account.id + ':' + agentId;
|
||||
let store = this.conversations.get(key);
|
||||
if (!store) {
|
||||
store = new TeacherConversationStore(path.join(this.options.userDataDir, 'agent-conversations', account.id, agentId), account.id, agentId);
|
||||
this.conversations.set(key, store);
|
||||
}
|
||||
return store;
|
||||
}
|
||||
private async importConversations(account: TeacherAccount, agentId: string) {
|
||||
const store = this.conversationStore(account, agentId);
|
||||
if (![...this.active.keys()].some(key => key.startsWith(account.id + ':agent:' + agentId + ':'))) {
|
||||
for (const project of await this.options.projects.listProjects()) {
|
||||
this.assertAccount(account);
|
||||
await store.importProject(project);
|
||||
}
|
||||
}
|
||||
return store;
|
||||
}
|
||||
async conversation(agentId: string, before?: string) {
|
||||
const account = await this.account();
|
||||
return this.serialize(account.id + ':agent:' + agentId + ':acceptance', async () => {
|
||||
const store = await this.importConversations(account, agentId);
|
||||
if (!await store.exists()) return { topic: null, before: null };
|
||||
const topic = await store.read();
|
||||
const page = await store.page(before);
|
||||
this.assertAccount(account);
|
||||
return { topic: { ...structuredClone(topic), requests: page.requests }, before: page.before };
|
||||
});
|
||||
}
|
||||
async sendConversation(agentId: string, input: TeacherSend) {
|
||||
if (!input.projectId) throw new TeacherError(422, 'teacher_project_required', '请先选择本次提问的项目。');
|
||||
const account = await this.account();
|
||||
const scope: TeacherScope = { projectId: input.projectId, sourceId: 'project', agentId };
|
||||
return this.serialize(this.acceptanceKey(account, scope), async () => {
|
||||
const store = await this.importConversations(account, agentId);
|
||||
if (!await store.exists()) {
|
||||
const catalog = await (this.options.catalog ?? teacherCatalog)(account);
|
||||
const selected = catalog.items.find(item => item.teacher_id === agentId);
|
||||
if (!selected) throw new TeacherError(409, 'teacher_disabled', '智能体已停用,历史仍可查看。');
|
||||
this.assertAccount(account);
|
||||
await store.ensure(selected.definition, selected.version);
|
||||
}
|
||||
const topic = await store.read();
|
||||
if (!this.active.has(this.key(account, scope, topic.id))) await store.recent();
|
||||
return this.sendRequest(account, scope, topic.id, input);
|
||||
});
|
||||
}
|
||||
async markConversationSeen(agentId: string, requestId: string) {
|
||||
const account = await this.account();
|
||||
const scope = { agentId, projectId: '', sourceId: 'project' };
|
||||
return this.serialize(this.acceptanceKey(account, scope), async () => {
|
||||
const store = this.conversationStore(account, agentId);
|
||||
this.assertAccount(account);
|
||||
await store.markSeen(requestId);
|
||||
const topic = await store.read();
|
||||
this.events.emit(this.key(account, scope, topic.id), structuredClone(topic));
|
||||
return structuredClone(topic);
|
||||
});
|
||||
}
|
||||
private async scopedStore(
|
||||
account: TeacherAccount,
|
||||
scope: TeacherScope
|
||||
): Promise<TeacherTopicStore> {
|
||||
): Promise<TeacherTopicStore | TeacherConversationStore> {
|
||||
if (scope.agentId) return this.conversationStore(account, scope.agentId);
|
||||
if (this.deletingSources.has(scope.projectId + ':' + scope.sourceId))
|
||||
throw new TeacherError(404, 'teacher_source_not_found', '来源会话已删除。');
|
||||
let directory: string;
|
||||
@@ -128,9 +227,11 @@ export class CodingTeacherService {
|
||||
return store;
|
||||
}
|
||||
private key(account: TeacherAccount, scope: TeacherScope, id: string) {
|
||||
if (scope.agentId) return account.id + ':agent:' + scope.agentId + ':' + id;
|
||||
return account.id + ':' + scope.projectId + ':' + scope.sourceId + ':' + (scope.role ?? 'teacher') + ':' + id;
|
||||
}
|
||||
private acceptanceKey(account: TeacherAccount, scope: TeacherScope) {
|
||||
if (scope.agentId) return account.id + ':agent:' + scope.agentId + ':acceptance';
|
||||
return account.id + ':' + scope.projectId + ':teacher-acceptance';
|
||||
}
|
||||
private async serialize<T>(key: string, operation: () => Promise<T>): Promise<T> {
|
||||
@@ -146,7 +247,7 @@ export class CodingTeacherService {
|
||||
async list(scope: TeacherScope) {
|
||||
const account = await this.account();
|
||||
const current = await (await this.scopedStore(account, scope)).list();
|
||||
if (scope.sourceId !== 'project' || scope.role === 'friend') return current;
|
||||
if (scope.agentId || scope.sourceId !== 'project' || scope.role === 'friend') return current;
|
||||
const legacy = await (await this.scopedStore(account, { ...scope, role: 'friend' })).list();
|
||||
return { ...current, items: [...current.items, ...legacy.items.map(item => ({ ...item, legacyRole: 'friend' as const }))]
|
||||
.sort((a, b) => b.updatedAt.localeCompare(a.updatedAt)) };
|
||||
@@ -216,8 +317,8 @@ export class CodingTeacherService {
|
||||
const topic = (await store.read(id)) as PreviewTopic;
|
||||
if (
|
||||
topic.accountId !== account.id ||
|
||||
topic.projectId !== scope.projectId ||
|
||||
topic.sourceConversationId !== scope.sourceId ||
|
||||
(scope.agentId ? topic.conversation?.agentId !== scope.agentId :
|
||||
topic.projectId !== scope.projectId || topic.sourceConversationId !== scope.sourceId) ||
|
||||
(topic.role ?? 'teacher') !== (scope.role ?? 'teacher')
|
||||
)
|
||||
throw new TeacherError(404, 'teacher_topic_not_found', '智能体话题不存在。');
|
||||
@@ -347,12 +448,17 @@ export class CodingTeacherService {
|
||||
throw new TeacherError(422, 'teacher_reference_invalid', '引用内容无效或超过 12000 字。');
|
||||
const key = this.key(account, scope, id);
|
||||
return await this.serialize(key, async () => {
|
||||
const { store, topic } = await this.readOwned(account, scope, id);
|
||||
const { store, topic: stored } = await this.readOwned(account, scope, id);
|
||||
if (this.active.has(key) && !stored.requests.some(request => ['preparing', 'running'].includes(request.status)))
|
||||
await this.finishes.get(key);
|
||||
const topic = scope.agentId ? structuredClone(stored) : stored;
|
||||
this.assertWritable(scope, topic);
|
||||
const existing = topic.requests.find((request) => request.id === input.requestId);
|
||||
const existing = store instanceof TeacherConversationStore ? await store.findRequest(input.requestId)
|
||||
: topic.requests.find((request) => request.id === input.requestId);
|
||||
if (existing) {
|
||||
if (
|
||||
existing.text !== input.text ||
|
||||
(scope.agentId && existing.projectId !== scope.projectId) ||
|
||||
(existing.intent ?? 'question') !== intent ||
|
||||
JSON.stringify(existing.references) !== JSON.stringify(refs) ||
|
||||
(existing.sourceConversationId ?? undefined) !== (input.sourceConversationId ?? undefined) ||
|
||||
@@ -369,6 +475,27 @@ export class CodingTeacherService {
|
||||
)
|
||||
)
|
||||
throw new TeacherError(409, 'teacher_topic_busy', '请等待当前回复完成,或先停止。');
|
||||
let newConversationSegment = false;
|
||||
if (scope.agentId) {
|
||||
const selected = (await (this.options.catalog ?? teacherCatalog)(account)).items.find(item => item.teacher_id === scope.agentId);
|
||||
if (!selected) throw new TeacherError(409, 'teacher_disabled', '智能体已停用,历史仍可查看。');
|
||||
this.assertAccount(account);
|
||||
topic.definition = structuredClone(selected.definition);
|
||||
topic.version = selected.version;
|
||||
topic.projectId = scope.projectId;
|
||||
topic.sourceConversationId = 'project';
|
||||
topic.discussion = topic.conversation!.discussions[scope.projectId];
|
||||
const previous = topic.requests.at(-1);
|
||||
const reuse = previous && !previous.origin && previous.projectId === scope.projectId
|
||||
&& previous.sourceConversationId === input.sourceConversationId
|
||||
&& previous.teacherVersion === selected.version && previous.runtimeThreadId;
|
||||
newConversationSegment = !reuse;
|
||||
// Yuxi binds both project and Pi source to a thread. Preserve its native
|
||||
// summary middleware within a segment; explicitly carry bounded public
|
||||
// history only when a new scope/release starts another segment.
|
||||
topic.conversation!.runtimeThreadId = reuse || randomUUID();
|
||||
topic.conversation!.segmentTurns = reuse ? topic.conversation!.segmentTurns + 1 : 1;
|
||||
}
|
||||
const discussionContext = validateDiscussionContext(topic, input.discussion);
|
||||
if (topic.draftRevision) {
|
||||
await (this.options.preview ?? teacherPreview)(account, topic.draftRevision);
|
||||
@@ -432,20 +559,21 @@ export class CodingTeacherService {
|
||||
projectPath: (await this.options.projects.getProject(scope.projectId)).path,
|
||||
source,
|
||||
history: topic.requests,
|
||||
...(store instanceof TeacherConversationStore ? { archive: await store.history() } : {}),
|
||||
assertCurrent: () => this.assertAccount(account),
|
||||
};
|
||||
const model = isCloud && access
|
||||
? (this.options.prepareCloud ?? prepareCloudTeacher)(account, topic, input.requestId, access,
|
||||
(progress) => {
|
||||
const current = topic.requests.at(-1)!;
|
||||
const current = topic.requests.find(request => request.id === input.requestId)!;
|
||||
if (current.progress === progress) return;
|
||||
current.progress = progress; topic.revision++;
|
||||
this.events.emit(key, structuredClone(topic));
|
||||
}, async cloudRequestId => {
|
||||
topic.requests.at(-1)!.cloudRequestId = cloudRequestId;
|
||||
await store.save(topic);
|
||||
topic.requests.find(request => request.id === input.requestId)!.cloudRequestId = cloudRequestId;
|
||||
await store.save(topic, input.requestId);
|
||||
}, undefined, activity => {
|
||||
const current = topic.requests.at(-1)!;
|
||||
const current = topic.requests.find(request => request.id === input.requestId)!;
|
||||
const activities = current.toolActivity ??= [];
|
||||
const index = activities.findIndex(item => item.id === activity.id);
|
||||
if (index < 0) activities.push(activity);
|
||||
@@ -460,7 +588,7 @@ export class CodingTeacherService {
|
||||
const compiled = compileTeacherContext(
|
||||
topic.definition,
|
||||
source,
|
||||
isCloud ? [] : topic.requests,
|
||||
isCloud && (!scope.agentId || !newConversationSegment) ? [] : topic.requests,
|
||||
input.text,
|
||||
references,
|
||||
model.inputLimit,
|
||||
@@ -482,7 +610,10 @@ export class CodingTeacherService {
|
||||
throw new TeacherError(409, 'teacher_source_archived', '来源会话已归档。');
|
||||
this.assertAccount(account);
|
||||
}
|
||||
const request = {
|
||||
const request: TeacherRequest = {
|
||||
...(scope.agentId ? { projectId: scope.projectId,
|
||||
projectName: (await this.options.projects.getProject(scope.projectId)).name,
|
||||
teacherVersion: topic.version, runtimeThreadId: topic.conversation!.runtimeThreadId } : {}),
|
||||
id: input.requestId,
|
||||
intent,
|
||||
...(input.presentation ? { presentation: input.presentation } : {}),
|
||||
@@ -504,7 +635,7 @@ export class CodingTeacherService {
|
||||
topic.updatedAt = request.createdAt;
|
||||
topic.revision++;
|
||||
try {
|
||||
await store.save(topic);
|
||||
await store.save(topic, request.id);
|
||||
} catch (error) {
|
||||
topic.requests.pop();
|
||||
throw error;
|
||||
@@ -517,7 +648,7 @@ export class CodingTeacherService {
|
||||
this.active.set(key, { account, controller, sourceId, projectId: scope.projectId });
|
||||
const release = this.options.acquireLease?.(key) ?? (() => undefined);
|
||||
const finish = async () => {
|
||||
const current = topic.requests.at(-1)!;
|
||||
const current = topic.requests.find(request => request.id === input.requestId)!;
|
||||
try {
|
||||
this.assertAccount(account);
|
||||
if (controller.signal.aborted) throw controller.signal.reason;
|
||||
@@ -554,7 +685,8 @@ export class CodingTeacherService {
|
||||
topic.updatedAt = new Date().toISOString();
|
||||
topic.revision++;
|
||||
try {
|
||||
await store.save(topic);
|
||||
if (topic.conversation && topic.discussion) topic.conversation.discussions[scope.projectId] = topic.discussion;
|
||||
await store.save(topic, current.id);
|
||||
} catch {
|
||||
topic.unsaved = true;
|
||||
}
|
||||
@@ -578,11 +710,17 @@ export class CodingTeacherService {
|
||||
if (this.active.has(key) || topic.requests.some(request => ['preparing', 'running'].includes(request.status)))
|
||||
throw new TeacherError(409, 'teacher_topic_busy', '请等待智能体回复,或先停止。');
|
||||
const next = structuredClone(topic);
|
||||
if (next.conversation) {
|
||||
if (!input.projectId) throw new TeacherError(422, 'teacher_project_required', '请选择讨论所属的项目。');
|
||||
await this.options.projects.getProject(input.projectId);
|
||||
next.discussion = next.conversation.discussions[input.projectId];
|
||||
}
|
||||
editDiscussion(next, input);
|
||||
if (next.conversation && next.discussion) next.conversation.discussions[input.projectId!] = next.discussion;
|
||||
next.revision++;
|
||||
next.updatedAt = new Date().toISOString();
|
||||
this.assertAccount(account);
|
||||
await store.save(next);
|
||||
await store.save(next, scope.agentId ? '' : undefined);
|
||||
this.events.emit(key, structuredClone(next));
|
||||
return structuredClone(next);
|
||||
});
|
||||
@@ -590,7 +728,7 @@ export class CodingTeacherService {
|
||||
async cancel(scope: TeacherScope, id: string, requestId: string) {
|
||||
const account = await this.account();
|
||||
const { topic } = await this.readOwned(account, scope, id);
|
||||
if (topic.requests.at(-1)?.id === requestId)
|
||||
if (topic.requests.some(request => request.id === requestId && ['preparing', 'running'].includes(request.status)))
|
||||
this.active.get(this.key(account, scope, id))?.controller.abort();
|
||||
return structuredClone(topic);
|
||||
}
|
||||
|
||||
@@ -101,7 +101,7 @@ export class TeacherTopicStore {
|
||||
}
|
||||
return await pending;
|
||||
}
|
||||
async save(topic: TeacherTopic) {
|
||||
async save(topic: TeacherTopic, _requestId?: string) {
|
||||
await mkdir(this.directory, { recursive: true });
|
||||
const saved = { ...topic };
|
||||
delete saved.unsaved;
|
||||
|
||||
Reference in New Issue
Block a user