302 lines
16 KiB
TypeScript
302 lines
16 KiB
TypeScript
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<typeof teacherHistoryMessages> }
|
|
interface Manifest {
|
|
topic: Omit<TeacherTopic, 'requests'>;
|
|
turns: TurnIndex[];
|
|
/** Retained verbatim when present in older manifests; no longer imported or interpreted. */
|
|
importedDiscussions?: unknown;
|
|
importedGlobal?: boolean;
|
|
}
|
|
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<void>;
|
|
private writes: Promise<unknown> = Promise.resolve();
|
|
private readonly importedProjects = new Set<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'); }
|
|
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
|
|
|| 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.
|
|
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<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: this.projectId, sourceConversationId: 'project', definition: structuredClone(definition), version,
|
|
createdAt: now, updatedAt: now, requests: [],
|
|
conversation: { agentId: this.agentId, segmentTurns: 0 },
|
|
};
|
|
// 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 {
|
|
// 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<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)));
|
|
// Newly imported turns can still be live before they have passed through turn().
|
|
// Repair the returned view too, keeping archived bytes unchanged.
|
|
const projection = structuredClone(requests);
|
|
projection.forEach(repairCompletedReply);
|
|
return { requests: projection, 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,
|
|
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,
|
|
// 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) { /* One conversation per project and 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, '');
|
|
}
|
|
/** 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);
|
|
}
|
|
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);
|
|
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);
|
|
}
|
|
// Retired component state stays in the untouched legacy topic. Import
|
|
// its conversation turns without creating active per-project components.
|
|
}
|
|
}
|
|
if (this.manifest) await this.recent();
|
|
this.importedProjects.add(project.id);
|
|
}
|
|
}
|