fix: isolate delivered agent chats by project

This commit is contained in:
2026-09-28 16:21:14 +08:00
parent a21a1f077c
commit 4b41c23a3a
15 changed files with 518 additions and 165 deletions

View File

@@ -18,7 +18,7 @@ export async function handleCodingTeacherRoutes(
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 conversation = url.pathname.match(/^\/api\/coding\/projects\/([^/]+)\/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))?$/
);
@@ -56,18 +56,22 @@ export async function handleCodingTeacherRoutes(
} : 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));
const projectId = conversation ? decodeURIComponent(conversation[1]) : undefined;
const agentId = conversation ? decodeURIComponent(conversation[2]) : undefined;
if (agentId && !conversation![3] && req.method === 'GET') {
sendJson(res, 200, await service.conversation(projectId!, 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)));
if (agentId && conversation![3] === 'messages' && req.method === 'POST') {
const input = await parseJsonBody<TeacherSend>(req);
if (input.projectId && input.projectId !== projectId)
throw new TeacherError(422, 'teacher_project_mismatch', '提问不属于当前项目。');
sendJson(res, 202, await service.sendConversation(agentId, { ...input, projectId }));
return true;
}
if (agentId && conversation![2] === 'seen' && req.method === 'POST') {
if (agentId && conversation![3] === 'seen' && req.method === 'POST') {
const input = await parseJsonBody<{ requestId: string }>(req);
sendJson(res, 200, await service.markConversationSeen(agentId, input.requestId));
sendJson(res, 200, await service.markConversationSeen(projectId!, agentId, input.requestId));
return true;
}
if (checkIn) {
@@ -93,16 +97,16 @@ export async function handleCodingTeacherRoutes(
sendJson(res, 200, await service.previewDefinition(revision));
return true;
}
const continuous = agentId ? await service.conversation(agentId) : undefined;
const continuous = agentId ? await service.conversation(projectId!, agentId) : undefined;
if (agentId && !continuous?.topic) throw new TeacherError(404, 'teacher_topic_not_found', '智能体聊天尚未开始。');
const scope: TeacherScope = agentId ? { projectId: '', sourceId: 'project', agentId } : source
const scope: TeacherScope = agentId ? { projectId: 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 = 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];
action = conversation?.[3] ?? source?.[4] ?? projectTopics?.[4] ?? preview?.[2],
requestId = conversation?.[4] ?? source?.[5] ?? projectTopics?.[5] ?? preview?.[3];
if (!id && req.method === 'GET') {
sendJson(res, 200, await service.list(scope));
return true;

View File

@@ -12,6 +12,7 @@ interface Manifest {
topic: Omit<TeacherTopic, 'requests'>;
turns: TurnIndex[];
importedDiscussions?: Record<string, string>;
importedGlobal?: boolean;
}
const PAGE_SIZE = 50;
const isMissing = (error: unknown) => (error as NodeJS.ErrnoException)?.code === 'ENOENT';
@@ -28,7 +29,8 @@ export class TeacherConversationStore {
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 checkedGlobal = false;
constructor(readonly directory: string, readonly accountId: string, readonly agentId: string, readonly projectId: string) {}
private manifestPath() { return path.join(this.directory, 'conversation.json'); }
private turnPath(id: string) { return path.join(this.directory, 'turns', teacherTopicId(id) + '.json'); }
@@ -41,7 +43,7 @@ export class TeacherConversationStore {
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))
|| manifest.topic.projectId !== this.projectId || !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.
@@ -84,7 +86,7 @@ export class TeacherConversationStore {
const now = new Date().toISOString();
const topic: TeacherTopic = {
schemaVersion: 1, revision: 0, id: randomUUID(), accountId: this.accountId,
projectId: '', sourceConversationId: 'project', definition: structuredClone(definition), version,
projectId: this.projectId, sourceConversationId: 'project', definition: structuredClone(definition), version,
createdAt: now, updatedAt: now, requests: [],
conversation: { agentId: this.agentId, segmentTurns: 0, discussions: {} },
};
@@ -151,7 +153,8 @@ export class TeacherConversationStore {
}
turns.sort(order);
const { requests: _requests, unsaved: _unsaved, ...metadata } = topic;
const next = { topic: metadata, turns, importedDiscussions: this.manifest.importedDiscussions };
const next = { topic: metadata, turns, importedDiscussions: this.manifest.importedDiscussions,
importedGlobal: this.manifest.importedGlobal };
try { await atomicWriteJson(this.manifestPath(), next); }
catch (error) {
// The body may already be durable. Reload and recover it before a retry,
@@ -169,7 +172,7 @@ export class TeacherConversationStore {
topic.requests = (await this.page()).requests;
return topic;
}
async select(_id: string) { /* There is only one conversation per agent. */ }
async select(_id: string) { /* One conversation per project and agent. */ }
async list() {
await this.load();
const topic = this.live;
@@ -185,8 +188,64 @@ export class TeacherConversationStore {
topic.revision++;
await this.save(topic, '');
}
/** Old files remain untouched. Only a proven account + config ID can be imported. */
/** Split the former account-wide chat using recorded turn ownership; never mutate its files. */
async importGlobal(directory: string) {
await this.load();
if (this.checkedGlobal || this.manifest?.importedGlobal) return;
let old: Manifest;
try { old = await readJsonFile(path.join(directory, 'conversation.json')) as Manifest; }
catch (error) { if (isMissing(error)) { this.checkedGlobal = true; return; } throw error; }
if (old.topic.accountId !== this.accountId || old.topic.conversation?.agentId !== this.agentId) {
this.checkedGlobal = true; return;
}
const ids = new Set(old.turns.map(turn => turn.id));
// Include a body whose original index write was interrupted.
try {
for (const name of await readdir(path.join(directory, 'turns')))
if (/^[0-9a-f-]{36}\.json$/i.test(name)) ids.add(name.slice(0, -5));
} catch (error) { if (!isMissing(error)) throw error; }
const seenIndex = old.turns.findIndex(turn => turn.id === old.topic.conversation?.seenThrough);
const seenIds = new Set(old.turns.slice(0, seenIndex + 1).map(turn => turn.id));
let seen: string | undefined;
for (const id of ids) {
const request = await readJsonFile(path.join(directory, 'turns', teacherTopicId(id) + '.json')) as TeacherRequest;
if ((request.projectId ?? request.origin?.projectId) !== this.projectId) continue;
const topic = await this.ensure(old.topic.definition, old.topic.version);
if (seenIds.has(id)) seen = id;
if (this.manifest!.turns.some(turn => turn.id === id)) continue;
const turn: TeacherRequest = { ...structuredClone(request), projectId: this.projectId,
// A split chat must start a new cloud checkpoint, since the old one may contain other projects.
origin: request.origin ?? { projectId: this.projectId, topicId: old.topic.id, requestId: id },
...(old.topic.definition.runtime === 'yuxi' ? { cloudRequestId: request.cloudRequestId ?? id } : {}) };
if (turn.status === 'running' || turn.status === 'preparing') {
turn.status = 'interrupted'; turn.error = '项目聊天已分开,本次回复中断。';
}
topic.requests.push(turn);
topic.updatedAt = [topic.updatedAt, turn.createdAt].sort().at(-1)!;
topic.revision++;
await this.save(topic, turn.id);
topic.requests = topic.requests.slice(-PAGE_SIZE);
}
// Retain an empty project's discussion too, without assigning unowned turns to it.
const discussion = old.topic.conversation.discussions[this.projectId];
if (discussion) {
const topic = await this.ensure(old.topic.definition, old.topic.version);
if (!topic.conversation!.discussions[this.projectId]) {
topic.conversation!.discussions[this.projectId] = structuredClone(discussion);
(this.manifest!.importedDiscussions ??= {})[this.projectId] = old.topic.updatedAt;
}
}
if (this.manifest) {
if (seen) await this.markSeen(seen);
this.manifest.importedGlobal = true;
await this.save(await this.read(), '');
await this.recent();
}
this.checkedGlobal = true;
}
/** Old files remain untouched. Only a proven account + project + config ID can be imported. */
async importProject(project: { id: string; path: string; name: string }) {
if (project.id !== this.projectId) throw new Error('Conversation project mismatch');
if (this.importedProjects.has(project.id)) return;
await this.load();
const root = path.join(project.path, '.makelore', 'teacher-conversations', this.accountId);

View File

@@ -135,31 +135,32 @@ export class CodingTeacherService {
this.assertAccount(account);
return { items: items.sort((a, b) => b.updatedAt.localeCompare(a.updatedAt)), lastSelectedTopicId: null };
}
private conversationStore(account: TeacherAccount, agentId: string) {
private conversationStore(account: TeacherAccount, projectId: string, 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;
teacherTopicId(projectId);
const key = account.id + ':' + projectId + ':' + 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);
store = new TeacherConversationStore(path.join(this.options.userDataDir, 'agent-conversations', account.id, agentId, 'projects', projectId), account.id, agentId, projectId);
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);
}
private async importConversations(account: TeacherAccount, projectId: string, agentId: string) {
const project = await this.options.projects.getProject(projectId);
const store = this.conversationStore(account, projectId, agentId);
if (![...this.active.keys()].some(key => key.startsWith(account.id + ':agent:' + projectId + ':' + agentId + ':'))) {
this.assertAccount(account);
await store.importGlobal(path.join(this.options.userDataDir, 'agent-conversations', account.id, agentId));
await store.importProject(project);
}
return store;
}
async conversation(agentId: string, before?: string) {
async conversation(projectId: string, 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);
return this.serialize(this.acceptanceKey(account, { projectId, agentId, sourceId: 'project' }), async () => {
const store = await this.importConversations(account, projectId, agentId);
if (!await store.exists()) return { topic: null, before: null };
const topic = await store.read();
const page = await store.page(before);
@@ -172,7 +173,7 @@ export class CodingTeacherService {
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);
const store = await this.importConversations(account, scope.projectId, agentId);
if (!await store.exists()) {
const catalog = await (this.options.catalog ?? teacherCatalog)(account);
const selected = catalog.items.find(item => item.teacher_id === agentId);
@@ -185,11 +186,11 @@ export class CodingTeacherService {
return this.sendRequest(account, scope, topic.id, input);
});
}
async markConversationSeen(agentId: string, requestId: string) {
async markConversationSeen(projectId: string, agentId: string, requestId: string) {
const account = await this.account();
const scope = { agentId, projectId: '', sourceId: 'project' };
const scope = { agentId, projectId, sourceId: 'project' };
return this.serialize(this.acceptanceKey(account, scope), async () => {
const store = this.conversationStore(account, agentId);
const store = this.conversationStore(account, projectId, agentId);
this.assertAccount(account);
await store.markSeen(requestId);
const topic = await store.read();
@@ -201,7 +202,7 @@ export class CodingTeacherService {
account: TeacherAccount,
scope: TeacherScope
): Promise<TeacherTopicStore | TeacherConversationStore> {
if (scope.agentId) return this.conversationStore(account, scope.agentId);
if (scope.agentId) return this.conversationStore(account, scope.projectId, scope.agentId);
if (this.deletingSources.has(scope.projectId + ':' + scope.sourceId))
throw new TeacherError(404, 'teacher_source_not_found', '来源会话已删除。');
let directory: string;
@@ -227,11 +228,11 @@ export class CodingTeacherService {
return store;
}
private key(account: TeacherAccount, scope: TeacherScope, id: string) {
if (scope.agentId) return account.id + ':agent:' + scope.agentId + ':' + id;
if (scope.agentId) return account.id + ':agent:' + scope.projectId + ':' + 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';
if (scope.agentId) return this.key(account, scope, 'acceptance');
return account.id + ':' + scope.projectId + ':teacher-acceptance';
}
private async serialize<T>(key: string, operation: () => Promise<T>): Promise<T> {
@@ -317,8 +318,8 @@ export class CodingTeacherService {
const topic = (await store.read(id)) as PreviewTopic;
if (
topic.accountId !== account.id ||
(scope.agentId ? topic.conversation?.agentId !== scope.agentId :
topic.projectId !== scope.projectId || topic.sourceConversationId !== scope.sourceId) ||
topic.projectId !== scope.projectId ||
(scope.agentId ? topic.conversation?.agentId !== scope.agentId : topic.sourceConversationId !== scope.sourceId) ||
(topic.role ?? 'teacher') !== (scope.role ?? 'teacher')
)
throw new TeacherError(404, 'teacher_topic_not_found', '智能体话题不存在。');
@@ -712,6 +713,7 @@ export class CodingTeacherService {
const next = structuredClone(topic);
if (next.conversation) {
if (!input.projectId) throw new TeacherError(422, 'teacher_project_required', '请选择讨论所属的项目。');
if (input.projectId !== scope.projectId) throw new TeacherError(422, 'teacher_project_mismatch', '讨论不属于当前项目。');
await this.options.projects.getProject(input.projectId);
next.discussion = next.conversation.discussions[input.projectId];
}