337 lines
24 KiB
TypeScript
337 lines
24 KiB
TypeScript
// @vitest-environment node
|
|
import { randomUUID } from 'node:crypto';
|
|
import { mkdtemp, readFile, rm } from 'node:fs/promises';
|
|
import { tmpdir } from 'node:os';
|
|
import path from 'node:path';
|
|
import { afterEach, describe, expect, it, vi } from 'vitest';
|
|
import { TeacherConversationStore } from '../../electron/coding-teacher/conversation-store';
|
|
import { TeacherTopicStore } from '../../electron/coding-teacher/store';
|
|
import { CodingTeacherService } from '../../electron/coding-teacher/service';
|
|
import { CodingProjectService } from '../../electron/coding-projects/project-service';
|
|
import { createCodingProjectStore, createMemoryCodingProjectStorage } from '../../electron/coding-projects/project-store';
|
|
import { InMemoryConversationRuntime } from '../../electron/coding-runtime/in-memory-conversation-runtime';
|
|
import { atomicWriteJson } from '../../electron/coding-projects/atomic-json';
|
|
import * as atomic from '../../electron/coding-projects/atomic-json';
|
|
import { createTeacherReadTools } from '../../electron/coding-teacher/read-tools';
|
|
import { createServer } from 'node:http';
|
|
import { handleCodingTeacherRoutes } from '../../electron/api/routes/coding-teacher';
|
|
import type { HostApiContext } from '../../electron/api/context';
|
|
import type { TeacherDefinition, TeacherRequest, TeacherTopic } from '../../shared/coding-teacher';
|
|
|
|
const definition: TeacherDefinition = { schema_version: 1, config_id: 'agent-a', teacher_id: 'teacher',
|
|
runtime: 'yuxi', name: '设计老师', description: '', avatar_id: 'avatar-01', system_prompt: '', skills: [],
|
|
welcome_message: '', suggested_questions: [], model: { model_id: 'model', reasoning_choice: { mode: 'default' } },
|
|
limits: { max_input_tokens: 8000, max_output_tokens: 1000 } };
|
|
const roots: string[] = [];
|
|
const services: CodingTeacherService[] = [];
|
|
afterEach(async () => {
|
|
vi.restoreAllMocks();
|
|
await Promise.all(services.splice(0).map(service => service.dispose()));
|
|
await Promise.all(roots.splice(0).map(root => rm(root, { recursive: true, force: true })));
|
|
});
|
|
async function root() { const dir = await mkdtemp(path.join(tmpdir(), 'single-agent-chat-')); roots.push(dir); return dir; }
|
|
function turn(index: number): TeacherRequest {
|
|
return { id: randomUUID(), text: '重复的问题', references: [], createdAt: new Date(1000 * index).toISOString(),
|
|
sourceCursor: { seq: 0, workerGeneration: 0 }, sourceCapturedAt: 'now', includedSourceMessageIds: [],
|
|
omittedMessages: 0, status: 'completed', response: '回答 ' + index };
|
|
}
|
|
describe('continuous conversation persistence', () => {
|
|
it('splits mixed global history by recorded project, preserving originals, seen positions and legacy identities', async () => {
|
|
const dir = await root(), global = path.join(dir, 'global');
|
|
const project = { id: randomUUID(), path: path.join(dir, 'project'), name: '天气' };
|
|
const other = randomUUID(), oldId = randomUUID();
|
|
const legacyTurn = turn(0);
|
|
const imported = { ...legacyTurn, id: randomUUID(), projectId: project.id,
|
|
origin: { projectId: project.id, topicId: oldId, requestId: legacyTurn.id } };
|
|
const turns = [imported, ...Array.from({ length: 110 }, (_, i) => ({
|
|
...turn(i + 1), projectId: i % 2 ? other : project.id, runtimeThreadId: 'mixed-cloud-thread',
|
|
})), turn(112)]; // This old turn has no proven project and must remain only in the original archive.
|
|
const metadata = { schemaVersion: 1, revision: 4, accountId: 'account', id: randomUUID(),
|
|
projectId: other, sourceConversationId: 'project', definition, version: 1,
|
|
createdAt: 'now', updatedAt: 'now', conversation: { agentId: 'agent-a', segmentTurns: 0, discussions: {},
|
|
seenThrough: turns[42].id } };
|
|
const manifest = { topic: metadata, turns: turns.map(({ id, createdAt }) => ({ id, createdAt, messages: [] })) };
|
|
await atomicWriteJson(path.join(global, 'conversation.json'), manifest);
|
|
for (const request of turns) await atomicWriteJson(path.join(global, 'turns', request.id + '.json'), request);
|
|
const orphan = { ...turn(113), projectId: other, status: 'running' };
|
|
await atomicWriteJson(path.join(global, 'turns', orphan.id + '.json'), orphan);
|
|
const original = await readFile(path.join(global, 'conversation.json'), 'utf8');
|
|
const originalTurn = await readFile(path.join(global, 'turns', orphan.id + '.json'), 'utf8');
|
|
const target = path.join(dir, 'a');
|
|
const a = new TeacherConversationStore(target, 'account', 'agent-a', project.id);
|
|
await a.importGlobal(global);
|
|
const b = new TeacherConversationStore(path.join(dir, 'b'), 'account', 'agent-a', other);
|
|
await b.importGlobal(global);
|
|
expect((await a.read()).id).not.toBe((await b.read()).id);
|
|
expect((await a.history()).messages).toHaveLength(56 * 2);
|
|
expect((await b.history()).messages).toHaveLength(55 * 2); // Interrupted turns stay visible but are not model context.
|
|
expect((await a.page()).requests.every(request => request.projectId === project.id)).toBe(true);
|
|
expect((await a.read()).conversation?.seenThrough).toBe(turns[41].id);
|
|
expect((await b.read()).conversation?.seenThrough).toBe(turns[42].id);
|
|
expect((await b.findRequest(orphan.id))?.status).toBe('interrupted');
|
|
expect((await a.findRequest(turns[1].id))?.origin?.topicId).toBe(metadata.id);
|
|
expect((await a.findRequest(turns[1].id))?.cloudRequestId).toBe(turns[1].id);
|
|
const legacy = new TeacherTopicStore(path.join(project.path, '.makelore', 'teacher-conversations', 'account', 'project'));
|
|
await legacy.save({ ...metadata, id: oldId, projectId: project.id, conversation: undefined, requests: [legacyTurn] });
|
|
await a.importProject(project);
|
|
const resumed = new TeacherConversationStore(target, 'account', 'agent-a', project.id);
|
|
await resumed.importGlobal(global); await resumed.importProject(project);
|
|
expect((await resumed.history()).messages).toHaveLength(56 * 2);
|
|
expect(await readFile(path.join(global, 'conversation.json'), 'utf8')).toBe(original);
|
|
expect(await readFile(path.join(global, 'turns', orphan.id + '.json'), 'utf8')).toBe(originalTurn);
|
|
const foreign = new TeacherConversationStore(path.join(dir, 'foreign'), 'other-account', 'agent-a', project.id);
|
|
await foreign.importGlobal(global);
|
|
expect(await foreign.exists()).toBe(false);
|
|
});
|
|
it('resumes a partial global split after an index write failure without duplicate messages', async () => {
|
|
const dir = await root(), old = path.join(dir, 'global'), dest = path.join(dir, 'project');
|
|
const request = { ...turn(1), projectId: 'project-a' };
|
|
await atomicWriteJson(path.join(old, 'conversation.json'), { topic: {
|
|
accountId: 'account', id: randomUUID(), definition, version: 1,
|
|
conversation: { agentId: 'agent-a', discussions: {} },
|
|
}, turns: [{ id: request.id, createdAt: request.createdAt, messages: [] }] });
|
|
await atomicWriteJson(path.join(old, 'turns', request.id + '.json'), request);
|
|
const store = new TeacherConversationStore(dest, 'account', 'agent-a', 'project-a');
|
|
await store.ensure(definition, 1);
|
|
const write = atomic.atomicWriteJson;
|
|
let fail = true;
|
|
vi.spyOn(atomic, 'atomicWriteJson').mockImplementation(async (file, value) => {
|
|
if (fail && file === path.join(dest, 'conversation.json')) { fail = false; throw new Error('disk unavailable'); }
|
|
return write(file, value);
|
|
});
|
|
await expect(store.importGlobal(old)).rejects.toThrow('disk unavailable');
|
|
await store.importGlobal(old);
|
|
expect((await store.read()).requests.map(turn => turn.id)).toEqual([request.id]);
|
|
await new TeacherConversationStore(dest, 'account', 'agent-a', 'project-a').importGlobal(old);
|
|
expect((await store.history()).messages).toHaveLength(2);
|
|
});
|
|
it('pages thousands of turns, restores the latest page and preserves repeated questions', async () => {
|
|
const dir = await root();
|
|
const store = new TeacherConversationStore(dir, 'account', 'agent-a', 'project-a');
|
|
const topic = await store.ensure(definition, 1);
|
|
topic.requests = Array.from({ length: 1001 }, (_, i) => turn(i));
|
|
await store.save(topic);
|
|
const resumed = new TeacherConversationStore(dir, 'account', 'agent-a', 'project-a');
|
|
expect((await resumed.read()).requests).toHaveLength(50);
|
|
const latest = await resumed.page();
|
|
expect(latest.requests.at(-1)?.response).toBe('回答 1000');
|
|
const older = await resumed.page(latest.before!);
|
|
expect(older.requests.at(-1)?.response).toBe('回答 950');
|
|
expect(new Set([...older.requests, ...latest.requests].map(turn => turn.id)).size).toBe(100);
|
|
expect((await resumed.findRequest(topic.requests[0].id))?.response).toBe('回答 0');
|
|
const archive = await resumed.history();
|
|
topic.requests.push(turn(1002)); await store.save(topic, topic.requests.at(-1)!.id);
|
|
const tools = createTeacherReadTools({ projectPath: dir, history: (await resumed.read()).requests, archive,
|
|
source: { messages: [], cursor: { seq: 0, workerGeneration: 0 }, capturedAt: 'now' }, assertCurrent() {} });
|
|
const signal = new AbortController().signal;
|
|
expect(await tools.execute('read_conversation', JSON.stringify({message_id: `teacher:${topic.requests[0].id}:assistant`}), signal)).toContain('回答 0');
|
|
expect((await tools.executeResult('read_conversation', JSON.stringify({message_id: `teacher:${topic.requests.at(-1)!.id}:assistant`}), signal)).status).toBe('error');
|
|
await resumed.markSeen(latest.requests.at(-1)!.id);
|
|
await resumed.markSeen(older.requests[0].id);
|
|
expect((await new TeacherConversationStore(dir, 'account', 'agent-a', 'project-a').read()).conversation?.seenThrough).toBe(latest.requests.at(-1)!.id);
|
|
});
|
|
it('recovers a turn written before its index and marks only unfinished work interrupted', async () => {
|
|
const dir = await root();
|
|
const store = new TeacherConversationStore(dir, 'account', 'agent-a', 'project-a');
|
|
await store.ensure(definition, 1);
|
|
const pending = { ...turn(1), status: 'running' as const };
|
|
await atomicWriteJson(path.join(dir, 'turns', pending.id + '.json'), pending);
|
|
const resumed = new TeacherConversationStore(dir, 'account', 'agent-a', 'project-a');
|
|
expect((await resumed.read()).requests[0]).toMatchObject({ id: pending.id, status: 'interrupted' });
|
|
expect((await resumed.page()).requests).toHaveLength(1);
|
|
});
|
|
it('imports separate topics idempotently without changing old files or merging repeated text', async () => {
|
|
const dir = await root(), project = { id: randomUUID(), path: path.join(dir, 'project'), name: '旧项目' };
|
|
const oldStore = new TeacherTopicStore(path.join(project.path, '.makelore', 'teacher-conversations', 'account', 'project'));
|
|
const old: TeacherTopic = { schemaVersion: 1, revision: 1, accountId: 'account', id: randomUUID(),
|
|
projectId: project.id, sourceConversationId: 'project', definition, version: 1, createdAt: 'now', updatedAt: 'now',
|
|
requests: [turn(1), turn(2)] };
|
|
await oldStore.save(old);
|
|
const file = path.join(oldStore.directory, old.id + '.json'), before = await readFile(file, 'utf8');
|
|
const target = path.join(dir, 'new');
|
|
const store = new TeacherConversationStore(target, 'account', 'agent-a', project.id);
|
|
await store.importProject(project);
|
|
const resumed = new TeacherConversationStore(target, 'account', 'agent-a', project.id);
|
|
await resumed.importProject(project);
|
|
expect((await resumed.read()).requests).toHaveLength(2);
|
|
expect((await resumed.read()).requests[0].origin?.topicId).toBe(old.id);
|
|
expect(await readFile(file, 'utf8')).toBe(before);
|
|
const other = new TeacherConversationStore(path.join(dir, 'other'), 'another-account', 'agent-a', project.id);
|
|
await other.importProject(project);
|
|
expect(await other.exists()).toBe(false);
|
|
});
|
|
it('recovers an import index failure before retrying in the same process', async () => {
|
|
const dir = await root(), project = { id: randomUUID(), path: path.join(dir, 'project'), name: '旧项目' };
|
|
const oldStore = new TeacherTopicStore(path.join(project.path, '.makelore', 'teacher-conversations', 'account', 'project'));
|
|
const old: TeacherTopic = { schemaVersion: 1, revision: 1, accountId: 'account', id: randomUUID(),
|
|
projectId: project.id, sourceConversationId: 'project', definition, version: 1, createdAt: 'now', updatedAt: 'now', requests: [turn(1)] };
|
|
await oldStore.save(old);
|
|
const store = new TeacherConversationStore(path.join(dir, 'chat'), 'account', 'agent-a', project.id);
|
|
await store.ensure(definition, 1);
|
|
const write = atomic.atomicWriteJson;
|
|
let fail = true;
|
|
vi.spyOn(atomic, 'atomicWriteJson').mockImplementation(async (file, value) => {
|
|
if (fail && file.endsWith('conversation.json')) { fail = false; throw new Error('disk unavailable'); }
|
|
return write(file, value);
|
|
});
|
|
await expect(store.importProject(project)).rejects.toThrow('disk unavailable');
|
|
await store.importProject(project);
|
|
expect((await store.page()).requests).toHaveLength(1);
|
|
expect((await store.page()).requests[0].origin?.requestId).toBe(old.requests[0].id);
|
|
});
|
|
});
|
|
|
|
it('serves paged account chats and compact SSE, preserves legacy source records and never dispatches on reads', async () => {
|
|
const f = await fixture();
|
|
const oldStore = new TeacherTopicStore(path.join(f.a.project.path, '.makelore', 'teacher-conversations', 'account', f.aSource.id));
|
|
const old: TeacherTopic = { schemaVersion: 1, revision: 1, accountId: 'account', id: randomUUID(),
|
|
projectId: f.a.project.id, sourceConversationId: f.aSource.id, definition, version: 1,
|
|
createdAt: 'now', updatedAt: 'now', requests: Array.from({length: 65}, (_, i) => turn(i)) };
|
|
await oldStore.save(old);
|
|
const server = createServer((req, res) => { void handleCodingTeacherRoutes(req, res, new URL(req.url!, 'http://localhost'),
|
|
{codingProducts: {teacher: f.service}} as HostApiContext); });
|
|
await new Promise<void>(resolve => server.listen(0, '127.0.0.1', resolve));
|
|
const address = server.address(); if (!address || typeof address === 'string') throw new Error('no address');
|
|
const origin = 'http://127.0.0.1:' + address.port;
|
|
const base = origin + '/api/coding/projects/' + f.a.project.id + '/agent-conversations/agent-a';
|
|
const otherBase = origin + '/api/coding/projects/' + f.b.project.id + '/agent-conversations/agent-a';
|
|
const post = (url: string, input: unknown) => fetch(url, {method: 'POST', headers: {'Content-Type': 'application/json'}, body: JSON.stringify(input)});
|
|
try {
|
|
const page = await (await fetch(base)).json();
|
|
expect(page.topic.requests).toHaveLength(50);
|
|
expect((await (await fetch(base + '?before=' + page.before)).json()).topic.requests).toHaveLength(15);
|
|
const archive = origin + '/api/coding/projects/' + f.a.project.id + '/agent-history';
|
|
expect((await (await fetch(archive)).json()).items[0]).toMatchObject({id: old.id, legacySourceId: f.aSource.id});
|
|
expect((await (await fetch(archive + '/' + old.id + '?source=' + f.aSource.id)).json()).requests).toHaveLength(65);
|
|
expect((await post(archive, {})).status).toBe(405);
|
|
expect(f.run).not.toHaveBeenCalled();
|
|
const input = {projectId: f.b.project.id, sourceConversationId: f.bSource.id, requestId: randomUUID(), text: '继续'};
|
|
expect((await (await fetch(otherBase)).json()).topic).toBeNull();
|
|
expect((await post(base + '/messages', input)).status).toBe(422);
|
|
expect((await post(otherBase + '/messages', input)).status).toBe(202);
|
|
expect((await post(otherBase + '/messages', input)).status).toBe(202);
|
|
await f.settled(f.b); expect(f.run).toHaveBeenCalledTimes(1);
|
|
const reader = (await fetch(otherBase + '/events')).body!.getReader();
|
|
const chunk = new TextDecoder().decode((await reader.read()).value);
|
|
const snapshot = JSON.parse(chunk.split('data: ')[1].split('\n')[0]);
|
|
expect(snapshot.requests).toHaveLength(1); expect(snapshot.requests[0].id).toBe(input.requestId);
|
|
await reader.cancel();
|
|
expect((await (await post(otherBase + '/seen', {requestId: input.requestId})).json()).conversation.seenThrough).toBe(input.requestId);
|
|
expect((await (await post(base + '/seen', {requestId: input.requestId})).json()).conversation.seenThrough).toBeUndefined();
|
|
expect((await (await fetch(base)).json()).topic.requests).toHaveLength(50);
|
|
expect((await post(otherBase + '/discussion', {projectId: f.a.project.id, action: 'archive'})).status).toBe(422);
|
|
expect((await fetch(base + '?before=missing')).status).toBe(400);
|
|
} finally { server.closeAllConnections(); await new Promise<void>(resolve => server.close(() => resolve())); }
|
|
});
|
|
|
|
async function fixture() {
|
|
const dir = await root();
|
|
const projects = new CodingProjectService(createCodingProjectStore(createMemoryCodingProjectStorage()));
|
|
const a = await projects.createProject({ projectPath: path.join(dir, 'weather'), identity: { kind: 'create' } });
|
|
const b = await projects.createProject({ projectPath: path.join(dir, 'game'), identity: { kind: 'create' } });
|
|
const source = async (project: typeof a) => projects.conversationStore(project.project.path)
|
|
.create({ agentId: project.config.defaultAgentId!, title: '操作对话', model: null, modelResolution: 'required' });
|
|
const aSource = await source(a), aOther = await source(a), bSource = await source(b);
|
|
let release: (() => void) | undefined;
|
|
let wait = false, version = 1, enabled = true, accountId = 'account';
|
|
const run = vi.fn(async (messages, _signal: AbortSignal, onText: (text: string) => void) => {
|
|
onText('已查看');
|
|
if (wait) { wait = false; await new Promise<void>(resolve => { release = resolve; }); }
|
|
return { inputTokens: 10, outputTokens: 5 };
|
|
});
|
|
const prepareCloud = vi.fn((_account, _topic, _id, _access) => ({ inputLimit: 8000, run }));
|
|
const service = new CodingTeacherService({ projects, runtime: new InMemoryConversationRuntime(), userDataDir: dir,
|
|
account: async () => ({ id: accountId, binding: { accountKey: accountId, epoch: 1 } }), assertAccount: () => undefined,
|
|
catalog: async () => ({ items: enabled ? [{ teacher_id: 'agent-a', version, definition, is_default: true }] : [] }),
|
|
prepareCloud,
|
|
readSource: async () => ({ messages: [{ id: 'pi', role: 'user', text: '项目上下文' }],
|
|
cursor: { seq: 0, workerGeneration: 0 }, capturedAt: 'now' }),
|
|
});
|
|
services.push(service);
|
|
const send = (project = a, session = aSource, requestId = randomUUID()) => service.sendConversation('agent-a', {
|
|
projectId: project.project.id, sourceConversationId: session.id, requestId, text: '帮我看看',
|
|
});
|
|
const settled = async (project = a) => {
|
|
await vi.waitFor(async () => expect((await service.conversation(project.project.id, 'agent-a')).topic?.requests.at(-1)?.status).toBe('completed'));
|
|
return (await service.conversation(project.project.id, 'agent-a')).topic!;
|
|
};
|
|
return { dir, service, a, b, aSource, aOther, bSource, prepareCloud, run, send, settled,
|
|
wait: () => { wait = true; }, release: () => { wait = false; release?.(); },
|
|
version: (v: number) => { version = v; }, enabled: (v: boolean) => { enabled = v; },
|
|
account: (id: string) => { accountId = id; } };
|
|
}
|
|
describe('one account-project-agent conversation', () => {
|
|
it('starts a clean cloud segment after splitting mixed history and supplies only the selected project messages', async () => {
|
|
const f = await fixture(), global = path.join(f.dir, 'agent-conversations', 'account', 'agent-a');
|
|
const turns = [
|
|
{ ...turn(1), projectId: f.a.project.id, sourceConversationId: f.aSource.id, runtimeThreadId: 'old-mixed-thread', teacherVersion: 1, text: '天气专属历史' },
|
|
{ ...turn(2), projectId: f.b.project.id, sourceConversationId: f.bSource.id, runtimeThreadId: 'old-mixed-thread', teacherVersion: 1, text: '游戏专属历史' },
|
|
];
|
|
await atomicWriteJson(path.join(global, 'conversation.json'), { topic: {
|
|
id: randomUUID(), accountId: 'account', definition, version: 1, projectId: f.b.project.id,
|
|
conversation: { agentId: 'agent-a', runtimeThreadId: 'old-mixed-thread', discussions: {} },
|
|
}, turns: turns.map(({ id, createdAt }) => ({ id, createdAt, messages: [] })) });
|
|
for (const turn of turns) await atomicWriteJson(path.join(global, 'turns', turn.id + '.json'), turn);
|
|
await f.send(); const current = await f.settled();
|
|
expect(current.requests).toHaveLength(2);
|
|
expect(current.requests.at(-1)?.runtimeThreadId).not.toBe('old-mixed-thread');
|
|
const messages = JSON.stringify(f.run.mock.calls[0][0]);
|
|
expect(messages).toContain('天气专属历史');
|
|
expect(messages).not.toContain('游戏专属历史');
|
|
});
|
|
it('isolates the same agent between projects and restores each project chat', async () => {
|
|
const f = await fixture();
|
|
const a = await f.send(); await f.settled();
|
|
const b = await f.send(f.b, f.bSource); await f.settled(f.b);
|
|
expect(b.id).not.toBe(a.id);
|
|
expect(b.requests).toHaveLength(1);
|
|
expect(JSON.stringify(f.run.mock.calls[1][0])).not.toContain('weather');
|
|
const back = await f.send(); await f.settled();
|
|
expect(back.id).toBe(a.id);
|
|
expect(back.requests).toHaveLength(2);
|
|
expect(back.requests.every(request => request.projectId === f.a.project.id)).toBe(true);
|
|
});
|
|
it('keeps one visible chat per project across versions while freezing each execution scope', async () => {
|
|
const f = await fixture();
|
|
await f.send(); const first = await f.settled();
|
|
await f.send(); const same = await f.settled();
|
|
expect(same.id).toBe(first.id);
|
|
expect(same.requests[1].runtimeThreadId).toBe(first.requests[0].runtimeThreadId);
|
|
// The same cloud checkpoint already owns the previous exchange.
|
|
expect(f.run.mock.calls[1][0].filter((m: { content: string }) => m.content.includes('已查看'))).toHaveLength(0);
|
|
await f.send(f.b, f.bSource); const second = await f.settled(f.b);
|
|
expect(second.id).not.toBe(first.id);
|
|
expect(second.requests[0].runtimeThreadId).not.toBe(first.requests[0].runtimeThreadId);
|
|
expect(f.prepareCloud.mock.calls[2][3].projectPath).toBe(f.b.project.path);
|
|
expect(JSON.stringify(f.run.mock.calls[2][0])).not.toContain('weather');
|
|
f.version(2); await f.send(f.b, f.bSource); const updated = await f.settled(f.b);
|
|
expect(updated.id).toBe(second.id);
|
|
expect(updated.requests.map(request => request.teacherVersion)).toEqual([1, 2]);
|
|
expect(updated.requests[1].runtimeThreadId).not.toBe(updated.requests[0].runtimeThreadId);
|
|
f.enabled(false);
|
|
await expect(f.send()).rejects.toMatchObject({ code: 'teacher_disabled' });
|
|
expect((await f.service.conversation(f.a.project.id, 'agent-a')).topic?.requests).toHaveLength(2);
|
|
});
|
|
it('deduplicates concurrent submissions, rejects a second active turn and rotates on Pi source change', async () => {
|
|
const f = await fixture(); f.wait();
|
|
const requestId = randomUUID();
|
|
const [a, b] = await Promise.all([f.send(f.a, f.aSource, requestId), f.send(f.a, f.aSource, requestId)]);
|
|
expect(a.id).toBe(b.id);
|
|
expect(f.run).toHaveBeenCalledTimes(1);
|
|
await expect(f.send()).rejects.toMatchObject({ code: 'teacher_topic_busy' });
|
|
const other = await f.send(f.b, f.bSource); await f.settled(f.b);
|
|
expect((await f.service.conversation(f.a.project.id, 'agent-a')).topic?.requests[0].status).toBe('running');
|
|
const otherScope = { projectId: f.b.project.id, sourceId: 'project', agentId: 'agent-a' };
|
|
await f.service.cancel(otherScope, other.id, requestId);
|
|
expect(f.run.mock.calls[0][1].aborted).toBe(false);
|
|
await expect(f.service.read(otherScope, a.id)).rejects.toMatchObject({ code: 'teacher_topic_not_found' });
|
|
expect(f.prepareCloud.mock.calls[0][3].projectPath).toBe(f.a.project.path);
|
|
f.release(); await f.settled();
|
|
await f.send(f.a, f.aOther); const next = await f.settled();
|
|
expect(next.requests[1].runtimeThreadId).not.toBe(next.requests[0].runtimeThreadId);
|
|
f.account('other');
|
|
expect((await f.service.conversation(f.a.project.id, 'agent-a')).topic).toBe(null);
|
|
await f.send(); expect((await f.settled()).id).not.toBe(a.id);
|
|
});
|
|
});
|