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 { repairCompletedReply, teacherTopicId } from './store'; interface TurnIndex { replyProjectionVersion?: 1; id: string; createdAt: string; origin?: string; messages: ReturnType } interface Manifest { topic: Omit; turns: TurnIndex[]; importedDiscussions?: Record; } 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 => { const projection = structuredClone(turn); repairCompletedReply(projection); return { id: turn.id, createdAt: turn.createdAt, origin: originKey(turn), replyProjectionVersion: 1, messages: teacherHistoryMessages([projection]).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; private writes: Promise = Promise.resolve(); private readonly importedProjects = new Set(); 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); const index = manifest.turns.findIndex(item => item.id === turn.id); if (index >= 0) manifest.turns[index] = indexTurn(turn); recovered = true; } if (recovered) { this.live.revision++; // Only restart metadata changes here; completed display repairs stay in memory. 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', '智能体消息无法读取,请保留本机记录。'); repairCompletedReply(turn); const index = this.manifest?.turns.findIndex(item => item.id === turn.id) ?? -1; if (this.manifest && index >= 0) this.manifest.turns[index] = indexTurn(turn); return turn; } async exists() { await this.load(); return Boolean(this.manifest); } async ensure(definition: TeacherDefinition, version: number): Promise { 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).requests; await atomicWriteJson(this.manifestPath(), manifest); this.manifest = manifest; this.live = topic; return topic; } async read(id?: string): Promise { 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 { 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 { // Older indexes may contain the unsafe quoted prefix. Keep them lazy: // expose an ID for reading, not an unverified assistant excerpt. messages: entries.flatMap(turn => turn.messages.map(message => turn.replyProjectionVersion === 1 || message.role !== 'assistant' ? message : { ...message, text: '(按消息 ID 读取完整回复)' })), 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 { // 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); } }