Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,8 @@ it('keeps a Space that has a directory on the file stores, whatever the structur
).ok,
).toBe(true);
const root = mkdtempSync(path.join(tmpdir(), 'huabu-agenetes-history-'));
const namespace: Namespace = { name: CANVAS_ID, storage: { root } };
const namespace = canvasAcpNamespace(CANVAS_ID);
namespace.storage = { root };
disposals.push(async () => {
// The file turn store holds an open database under `root`; the handle
// arbitration point is how a Space's owners are asked to let go.
Expand Down
103 changes: 85 additions & 18 deletions apps/server/src/modules/agent/agenetes/conversation-stores.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,12 @@ import {
getStructuredStore,
registerSpaceDirHandleOwner,
} from '../../storage/index.js';
import { resolveCanvasAcpNamespace } from '../../workspace/paths.js';
import { acquireWorkspaceOperationLease } from '../../workspace.js';
import {
appendRecentConversationTurn,
preserveRecentConversationBeforeHistoryDelete,
} from '../recent-conversation.js';

import type {
EventLogRecord,
Expand Down Expand Up @@ -122,46 +128,107 @@ function backingFor(namespace: Namespace): Backing {
throw new Error('A named Disk conversation requires a storage root');
}

async function inNamespace<T>(
namespace: Namespace,
action: (current: Namespace, backing: Backing) => T | Promise<T>,
): Promise<T> {
if (!namespace.name) return action(namespace, memory);
const lease = acquireWorkspaceOperationLease();
try {
const current = resolveCanvasAcpNamespace(namespace);
return await action(current, backingFor(current));
} finally {
lease.release();
}
}

export const conversationThreadStore: ThreadStore = {
upsert: (namespace, threadId, record: ThreadRecord) =>
backingFor(namespace).threads.upsert(namespace, threadId, record),
inNamespace(namespace, (current, backing) =>
backing.threads.upsert(
current,
threadId,
current.name
? { ...record, spec: { ...record.spec, namespace: current } }
: record,
),
),
get: (namespace, threadId) =>
backingFor(namespace).threads.get(namespace, threadId),
list: (namespace) => backingFor(namespace).threads.list(namespace),
inNamespace(namespace, (current, backing) =>
backing.threads.get(current, threadId),
),
list: (namespace) =>
inNamespace(namespace, (current, backing) => backing.threads.list(current)),
delete: (namespace, threadId) =>
backingFor(namespace).threads.delete(namespace, threadId),
inNamespace(namespace, (current, backing) =>
backing.threads.delete(current, threadId),
),
};

export const conversationEventLogStore: EventLogStore = {
appendTurnStart: (namespace, threadId, request: AgentSubmission | null) =>
backingFor(namespace).events.appendTurnStart(namespace, threadId, request),
inNamespace(namespace, (current, backing) =>
appendRecentConversationTurn(current, threadId, request, backing.events),
),
append: (namespace, threadId, event) =>
backingFor(namespace).events.append(namespace, threadId, event),
inNamespace(namespace, (current, backing) =>
backing.events.append(current, threadId, event),
),
read: (namespace, threadId, sinceSeq) =>
backingFor(namespace).events.read(namespace, threadId, sinceSeq),
inNamespace(namespace, (current, backing) =>
backing.events.read(current, threadId, sinceSeq),
),
readRecords: (namespace, threadId, sinceSeq) =>
backingFor(namespace).events.readRecords(namespace, threadId, sinceSeq),
inNamespace(namespace, (current, backing) =>
backing.events.readRecords(current, threadId, sinceSeq),
),
maxSeq: (namespace, threadId) =>
backingFor(namespace).events.maxSeq(namespace, threadId),
inNamespace(namespace, (current, backing) =>
backing.events.maxSeq(current, threadId),
),
replace: (namespace, threadId, records: readonly EventLogRecord[]) =>
backingFor(namespace).events.replace(namespace, threadId, records),
inNamespace(namespace, (current, backing) =>
backing.events.replace(current, threadId, records),
),
delete: (namespace, threadId) =>
backingFor(namespace).events.delete(namespace, threadId),
inNamespace(namespace, async (current, backing) => {
const events = backing.events;
await preserveRecentConversationBeforeHistoryDelete(
current,
threadId,
events,
);
await events.delete(current, threadId);
}),
};

export const conversationTurnStore: TurnStore = {
append: (namespace, threadId, persisted: PersistedTurn) =>
backingFor(namespace).turns.append(namespace, threadId, persisted),
inNamespace(namespace, (current, backing) =>
backing.turns.append(current, threadId, persisted),
),
list: (namespace, threadId) =>
backingFor(namespace).turns.list(namespace, threadId),
inNamespace(namespace, (current, backing) =>
backing.turns.list(current, threadId),
),
page: (namespace, threadId, options: TurnStorePageOptions) =>
backingFor(namespace).turns.page(namespace, threadId, options),
inNamespace(namespace, (current, backing) =>
backing.turns.page(current, threadId, options),
),
count: (namespace, threadId) =>
backingFor(namespace).turns.count(namespace, threadId),
inNamespace(namespace, (current, backing) =>
backing.turns.count(current, threadId),
),
fence: (namespace, threadId) =>
backingFor(namespace).turns.fence(namespace, threadId),
inNamespace(namespace, (current, backing) =>
backing.turns.fence(current, threadId),
),
replace: (namespace, threadId, persisted: readonly PersistedTurn[]) =>
backingFor(namespace).turns.replace(namespace, threadId, persisted),
inNamespace(namespace, (current, backing) =>
backing.turns.replace(current, threadId, persisted),
),
delete: (namespace, threadId) =>
backingFor(namespace).turns.delete(namespace, threadId),
inNamespace(namespace, (current, backing) =>
backing.turns.delete(current, threadId),
),
};
16 changes: 11 additions & 5 deletions apps/server/src/modules/agent/conversation/ink-ocr.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -274,6 +274,9 @@ describe('Ink OCR deadline and cancellation', () => {
{ durationMs: 1499, slow: false },
{ durationMs: 1500, slow: true },
{ durationMs: 1999, slow: true },
{ durationMs: 2000, slow: true },
{ durationMs: 2500, slow: true },
{ durationMs: 2999, slow: true },
])(
'uses 1500ms only for telemetry, allowing success at $durationMs ms',
async ({ durationMs, slow }) => {
Expand All @@ -292,17 +295,20 @@ describe('Ink OCR deadline and cancellation', () => {
);

it.each(['resolve', 'reject'] as const)(
'settles at exactly 2000ms and discards a late %s',
'settles at exactly 3000ms and discards a late %s',
async (completion) => {
const pending = deferred<Response>();
fetchMock.mockReturnValue(pending.promise);
const result = recognizeInk(params);
await vi.advanceTimersByTimeAsync(1999);
await vi.advanceTimersByTimeAsync(2999);
expect(debug).not.toHaveBeenCalled();
expect(warn).not.toHaveBeenCalled();
expect(fetchMock.mock.calls[0]?.[1]?.signal?.aborted).toBe(false);
await vi.advanceTimersByTimeAsync(1);
expect(await result).toBeUndefined();
expect(lastOutcome()).toMatchObject({
outcome: 'timeout',
durationMs: 2000,
durationMs: 3000,
});
expect(fetchMock.mock.calls[0]?.[1]?.signal?.aborted).toBe(true);
if (completion === 'resolve') pending.resolve(response());
Expand Down Expand Up @@ -358,11 +364,11 @@ describe('Ink OCR deadline and cancellation', () => {
});

const result = recognizeInk(params);
await vi.advanceTimersByTimeAsync(2000);
await vi.advanceTimersByTimeAsync(3000);
expect(await result).toBeUndefined();
expect(lastOutcome()).toMatchObject({
outcome: 'timeout',
durationMs: 2000,
durationMs: 3000,
});
expect(cancel).toHaveBeenCalled();
pendingRead.resolve({ done: true, value: undefined });
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/modules/agent/conversation/ink-ocr.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import type { InkRecognition } from '@huabu/shared';
import type { FastifyBaseLogger } from 'fastify';

const API_VERSION = '2024-02-01';
const DEADLINE_MS = 2000;
const DEADLINE_MS = 3000;
const SLOW_MS = 1500;
const MAX_RESPONSE_BYTES = 1024 * 1024;

Expand Down
Loading
Loading