diff --git a/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.client.test.tsx b/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.client.test.tsx index e58e391f2d..8159ecbd39 100644 --- a/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.client.test.tsx +++ b/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.client.test.tsx @@ -832,6 +832,72 @@ describe('FastSessionTranscript', () => { createdAt: new Date(ts), }); + it('shows automatic memory saves and expands their distilled facts', () => { + render( + , + ); + + const summary = screen.getByText('Saved to memory'); + expect(summary).toBeInTheDocument(); + expect(summary.closest('details')).not.toHaveAttribute('open'); + fireEvent.click(summary); + expect(summary.closest('details')).toHaveAttribute('open'); + expect( + screen.getByText('Staging deploys use the release branch.'), + ).toBeInTheDocument(); + }); + + it('renders a persisted memory-save row delivered by the transcript messages stream', () => { + render( + , + ); + + const source = FakeEventSource.instances[0]; + expect(source).toBeDefined(); + act(() => { + source?.emit('messages', { + messages: [ + { + ...textMessage({ + id: 'memory-save-streamed', + role: 'assistant', + text: 'Saved to memory', + ts: 2, + }), + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + role: 'system' as const, + payload: { memories: ['The release branch deploys staging.'] }, + }, + ], + conversationResponding: false, + }); + }); + + expect(screen.getByText('Saved to memory')).toBeInTheDocument(); + }); + it('restores each Session draft and scroll position without focusing after a direct switch', () => { function SessionSwitchHarness() { const [sessionId, setSessionId] = useState('session-a'); diff --git a/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx b/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx index 02107f3e6c..a2a592f64e 100644 --- a/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx +++ b/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx @@ -95,6 +95,7 @@ import { CapabilityOfferCard } from './CapabilityOfferCard'; import { PendingIntegrationKeys } from '@/components/sessions/PendingIntegrationKeys'; import { openIntegrationKeyDialog } from '@/components/sessions/integration-key-dialog'; import { useSessionTitlePropagation } from './use-session-title-propagation'; +import { MemorySavedMessage } from '@/components/ai-elements/MemorySavedMessage'; import { AcpMessageItem, @@ -941,6 +942,9 @@ export function FastSessionTranscript({ ); const renderCapabilityOfferMessage = useCallback( (message: AcpUiMessage) => { + if (message.updateType === ACP_ENVELOPE_EVENT_TYPES.MemorySaved) { + return ; + } const offer = capabilityOffersByMessageId.get(message.id); if (!offer) return undefined; const introMessage: AcpUiMessage = { diff --git a/apps/web/src/app/(sandbox)/task/[taskId]/Messages.tsx b/apps/web/src/app/(sandbox)/task/[taskId]/Messages.tsx index 7e194d494d..c5332eda9b 100644 --- a/apps/web/src/app/(sandbox)/task/[taskId]/Messages.tsx +++ b/apps/web/src/app/(sandbox)/task/[taskId]/Messages.tsx @@ -14,12 +14,14 @@ import { useStickToBottomContext, type ScrollToBottom, } from 'use-stick-to-bottom'; +import { ACP_ENVELOPE_EVENT_TYPES } from '@roomote/types'; import { Conversation, ConversationContent, ConversationScrollButton, } from '@/components/ai-elements'; +import { MemorySavedMessage } from '@/components/ai-elements/MemorySavedMessage'; import { MessageUiOptionsProvider, type MessageUiOptions, @@ -40,6 +42,7 @@ import { useSandboxTaskPhase, type TaskSession, } from './hooks'; +import type { AcpUiMessage } from './types'; import { useInternalTranscriptRowsVisible } from './useInternalTranscriptRowsVisible'; import { SleepWakeMessages } from './messages/index'; @@ -303,6 +306,13 @@ const MessagesBase = ({ resetKey: session.taskId, isWorking: taskPhase === 'running', }); + const renderMemoryMessage = useCallback( + (message: AcpUiMessage) => + message.updateType === ACP_ENVELOPE_EVENT_TYPES.MemorySaved ? ( + + ) : undefined, + [], + ); const shouldShowWorking = taskPhase === 'running' && !hasVisibleAssistantOutput(renderBlocks); @@ -332,6 +342,7 @@ const MessagesBase = ({ blocks={renderBlocks} showInternalMessages={showInternalMessages} onSuppress={suppressMessage} + renderMessage={renderMemoryMessage} /> {session.taskRun && } {shouldShowWorking && } diff --git a/apps/web/src/components/ai-elements/MemorySavedMessage.client.test.tsx b/apps/web/src/components/ai-elements/MemorySavedMessage.client.test.tsx new file mode 100644 index 0000000000..fa47e5a66a --- /dev/null +++ b/apps/web/src/components/ai-elements/MemorySavedMessage.client.test.tsx @@ -0,0 +1,35 @@ +import { fireEvent, render, screen } from '@testing-library/react'; +import { ACP_ENVELOPE_EVENT_TYPES } from '@roomote/types'; + +import { MemorySavedMessage } from './MemorySavedMessage'; + +describe('MemorySavedMessage', () => { + it('renders task and session save facts behind an expandable indication', () => { + render( + , + ); + + const summary = screen.getByText('Saved to memory'); + expect(summary.closest('details')).not.toHaveAttribute('open'); + fireEvent.click(summary); + expect(summary.closest('details')).toHaveAttribute('open'); + expect( + screen.getByText('The task retries are capped at three attempts.'), + ).toBeInTheDocument(); + }); +}); diff --git a/apps/web/src/components/ai-elements/MemorySavedMessage.tsx b/apps/web/src/components/ai-elements/MemorySavedMessage.tsx new file mode 100644 index 0000000000..d72738b68f --- /dev/null +++ b/apps/web/src/components/ai-elements/MemorySavedMessage.tsx @@ -0,0 +1,39 @@ +import { + ACP_ENVELOPE_EVENT_TYPES, + MEMORY_SAVED_EVENT_TEXT, + parseMemorySavedEventPayload, +} from '@roomote/types'; + +import { Brain } from '@/components/system'; +import type { AcpUiMessage } from '@/app/(sandbox)/task/[taskId]/types'; + +export function MemorySavedMessage({ message }: { message: AcpUiMessage }) { + if (message.updateType !== ACP_ENVELOPE_EVENT_TYPES.MemorySaved) { + return null; + } + + const payload = parseMemorySavedEventPayload(message.data); + + return ( +
+ + + {MEMORY_SAVED_EVENT_TEXT} + {payload ? ( + + {payload.memories.length} + + ) : null} + + {payload ? ( +
    + {payload.memories.map((memory) => ( +
  • + {memory} +
  • + ))} +
+ ) : null} +
+ ); +} diff --git a/apps/web/src/lib/server/task-messages.test.ts b/apps/web/src/lib/server/task-messages.test.ts index 26076a2ee2..acf09947e2 100644 --- a/apps/web/src/lib/server/task-messages.test.ts +++ b/apps/web/src/lib/server/task-messages.test.ts @@ -48,6 +48,34 @@ describe('getTaskMessageEnvelopes', () => { }); }); + it('delivers automatic memory-save events through task transcript history', async () => { + const task = await taskFactory.create({ + id: 'task-message-memory-saved', + title: 'Memory save delivery', + }); + const run = await runFactory.create({ taskId: task.id }); + + await db.insert(taskMessages).values({ + runId: run.id, + taskId: task.id, + ts: run.id, + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + role: 'system', + protocol: ROOMOTE_RUNTIME_TASK_MESSAGE_PROTOCOL, + contentBlocks: [{ type: 'text', text: 'Saved to memory' }], + metadata: { visibleInTranscript: true, memorySave: true }, + payload: { memories: ['The task retries are capped at three attempts.'] }, + }); + + const [message] = await getTaskMessageEnvelopes({ taskId: task.id }); + + expect(message).toMatchObject({ + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + role: 'system', + payload: { memories: ['The task retries are capped at three attempts.'] }, + }); + }); + it('pages backward through a large transcript without gaps or duplicates', async () => { const task = await taskFactory.create({ id: 'task-message-paginated-history', diff --git a/apps/web/src/testing/exclusive-automation-settings-database-lock.ts b/apps/web/src/testing/exclusive-automation-settings-database-lock.ts index 1174a74423..b25a05307d 100644 --- a/apps/web/src/testing/exclusive-automation-settings-database-lock.ts +++ b/apps/web/src/testing/exclusive-automation-settings-database-lock.ts @@ -1,6 +1,7 @@ import { db, sql } from '@roomote/db/server'; const AUTOMATION_SETTINGS_TEST_LOCK = [20260910, 2451] as const; +const EXCLUSIVE_LOCK_HOOK_TIMEOUT_MS = 60_000; /** Serializes suites that replace deployment-wide automation settings rows. */ export function registerExclusiveAutomationSettingsDatabaseLock() { @@ -27,10 +28,10 @@ export function registerExclusiveAutomationSettingsDatabaseLock() { }); void lockTransaction.catch((error) => rejectAcquired?.(error)); await acquired; - }); + }, EXCLUSIVE_LOCK_HOOK_TIMEOUT_MS); afterAll(async () => { releaseLock?.(); await lockTransaction; - }); + }, EXCLUSIVE_LOCK_HOOK_TIMEOUT_MS); } diff --git a/packages/cloud-agents/src/server/__tests__/task-run-memory-distillation.test.ts b/packages/cloud-agents/src/server/__tests__/task-run-memory-distillation.test.ts index e25bdc4628..dad5878c19 100644 --- a/packages/cloud-agents/src/server/__tests__/task-run-memory-distillation.test.ts +++ b/packages/cloud-agents/src/server/__tests__/task-run-memory-distillation.test.ts @@ -5,6 +5,8 @@ const { mockIsBrainEnabled, mockIsTaskRunSharedBrainEligible, mockSaveBrainDistilledSummary, + mockInsertTaskMemoryEvent, + mockInsertTaskMemoryValues, mockTurnRows, } = vi.hoisted(() => ({ mockEvaluateDecisionModel: vi.fn(), @@ -13,28 +15,37 @@ const { mockIsBrainEnabled: vi.fn(), mockIsTaskRunSharedBrainEligible: vi.fn(), mockSaveBrainDistilledSummary: vi.fn(), + mockInsertTaskMemoryEvent: vi.fn(), + mockInsertTaskMemoryValues: vi.fn(), mockTurnRows: vi.fn(), })); -vi.mock('@roomote/db/server', () => ({ - and: vi.fn(), - desc: vi.fn(), - eq: vi.fn(), - inArray: vi.fn(), - sql: vi.fn(), - taskMessages: {}, - getBrainMemorySummary: mockGetBrainMemorySummary, - isBrainEnabled: mockIsBrainEnabled, - isTaskRunSharedBrainEligible: mockIsTaskRunSharedBrainEligible, - saveBrainDistilledSummary: mockSaveBrainDistilledSummary, - db: { +vi.mock('@roomote/db/server', () => { + const insert = () => ({ values: mockInsertTaskMemoryValues }); + const database = { + insert, select: () => ({ from: () => ({ where: () => ({ orderBy: () => ({ limit: mockTurnRows }) }), }), }), - }, -})); + transaction: async (callback: (tx: { insert: typeof insert }) => unknown) => + callback({ insert }), + }; + return { + and: vi.fn(), + desc: vi.fn(), + eq: vi.fn(), + inArray: vi.fn(), + sql: vi.fn(), + taskMessages: {}, + getBrainMemorySummary: mockGetBrainMemorySummary, + isBrainEnabled: mockIsBrainEnabled, + isTaskRunSharedBrainEligible: mockIsTaskRunSharedBrainEligible, + saveBrainDistilledSummary: mockSaveBrainDistilledSummary, + db: database, + }; +}); vi.mock('../typesafe-judgment', () => ({ evaluateDecisionModel: mockEvaluateDecisionModel, @@ -51,14 +62,16 @@ import { ACP_ENVELOPE_EVENT_TYPES } from '@roomote/types'; import { distillTaskRunTurnMemory } from '../task-run-memory-distillation'; -const row = (eventType: string, text: string) => ({ +const row = (eventType: string, text: string, ts = 1) => ({ eventType, contentBlocks: [{ type: 'text', text }], payload: null, + ts, }); -const assistant = (text: string) => - row(ACP_ENVELOPE_EVENT_TYPES.AssistantMessage, text); -const user = (text: string) => row(ACP_ENVELOPE_EVENT_TYPES.UserPrompt, text); +const assistant = (text: string, ts?: number) => + row(ACP_ENVELOPE_EVENT_TYPES.AssistantMessage, text, ts); +const user = (text: string, ts?: number) => + row(ACP_ENVELOPE_EVENT_TYPES.UserPrompt, text, ts); function answers( overrides: Partial< @@ -105,6 +118,10 @@ describe('distillTaskRunTurnMemory', () => { mockIsTaskRunSharedBrainEligible.mockResolvedValue(true); mockGetBrainMemorySummary.mockResolvedValue(null); mockSaveBrainDistilledSummary.mockResolvedValue(true); + mockInsertTaskMemoryEvent.mockResolvedValue(undefined); + mockInsertTaskMemoryValues.mockImplementation(() => ({ + onConflictDoNothing: mockInsertTaskMemoryEvent, + })); // Newest first, as the query returns them. Only the latest turn counts. mockTurnRows.mockResolvedValue([ assistant('Retries are capped at 3 and skip 4xx responses.'), @@ -135,10 +152,30 @@ describe('distillTaskRunTurnMemory', () => { request: 'Do not retry 4xx; the partner API bans replayed requests.', report: 'Updating the retry policy.\n\nRetries are capped at 3 and skip 4xx responses.', + turnTs: 1, existing_memory: '', }, }), ); + expect(mockInsertTaskMemoryValues).toHaveBeenCalledWith({ + runId: 102, + taskId: 'task-1', + userId: 'user-1', + ts: 1, + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + role: 'system', + protocol: 'roomote_runtime', + contentBlocks: [{ type: 'text', text: 'Saved to memory' }], + metadata: { + visibleInTranscript: true, + memorySave: true, + automatic: true, + }, + payload: { + memories: [expect.stringContaining('Webhook retries are capped')], + }, + source: 'roomote', + }); expect(summary).toContain('## Outcome\n\nWebhook retries are capped at 3'); expect(summary?.endsWith(DISTILLED_NOTE)).toBe(true); expect(mockSaveBrainDistilledSummary).toHaveBeenCalledWith( @@ -150,6 +187,43 @@ describe('distillTaskRunTurnMemory', () => { ); }); + it('uses the turn key across fallback retries and follow-up turns', async () => { + await distillTaskRunTurnMemory(run); + mockTurnRows.mockResolvedValue([ + assistant('Follow-up found a second durable fact.', 2), + user('Also remember the follow-up decision.', 2), + ]); + mockGetBrainMemorySummary.mockResolvedValueOnce( + `## Outcome\n\nWebhook retries are capped.\n\n${DISTILLED_NOTE}`, + ); + await distillTaskRunTurnMemory(run); + + expect(mockInsertTaskMemoryValues).toHaveBeenCalledTimes(2); + expect(mockInsertTaskMemoryValues.mock.calls[0]?.[0]).toMatchObject({ + taskId: 'task-1', + ts: 1, + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + }); + expect(mockInsertTaskMemoryValues.mock.calls[1]?.[0]).toMatchObject({ + taskId: 'task-1', + ts: 2, + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + }); + }); + + it('rolls back the summary when event publication fails so retry can publish it', async () => { + mockInsertTaskMemoryEvent.mockRejectedValueOnce(new Error('database busy')); + + await expect(distillTaskRunTurnMemory(run)).resolves.toBeNull(); + expect(mockSaveBrainDistilledSummary).toHaveBeenCalledOnce(); + + mockInsertTaskMemoryEvent.mockResolvedValueOnce(undefined); + await expect(distillTaskRunTurnMemory(run)).resolves.toContain( + 'Webhook retries are capped', + ); + expect(mockInsertTaskMemoryEvent).toHaveBeenCalledTimes(2); + }); + it('builds on its own earlier memory and saves against that exact text', async () => { const existing = `## Outcome\n\nAdded webhook retries.\n\n${DISTILLED_NOTE}`; mockGetBrainMemorySummary.mockResolvedValue(existing); diff --git a/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-post-turn-memory.test.ts b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-post-turn-memory.test.ts index fcf91c749f..5594da2fe5 100644 --- a/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-post-turn-memory.test.ts +++ b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-post-turn-memory.test.ts @@ -1,11 +1,13 @@ const { mockAppendFastAgentMemory, + mockAppendFastAgentMemorySavedEvent, mockEvaluateDecisionModel, mockGenerateTrackedNonTaskObject, mockGetFastAgentConversationMemory, mockIsBrainEnabled, } = vi.hoisted(() => ({ mockAppendFastAgentMemory: vi.fn(), + mockAppendFastAgentMemorySavedEvent: vi.fn(), mockEvaluateDecisionModel: vi.fn(), mockGenerateTrackedNonTaskObject: vi.fn(), mockGetFastAgentConversationMemory: vi.fn(), @@ -30,6 +32,10 @@ vi.mock('../../non-task-provider-usage', () => ({ }, })); +vi.mock('../fast-agent-session', () => ({ + appendFastAgentMemorySavedEvent: mockAppendFastAgentMemorySavedEvent, +})); + import { saveFastAgentPostTurnMemory } from '../fast-agent-post-turn-memory'; function answers( @@ -61,6 +67,7 @@ function answers( const turn = { conversationId: 'conversation-1', + turnId: 'turn-1', userId: 'user-1', request: 'From now on we deploy staging from the release branch, not main.', reply: 'Understood. I will deploy staging from the release branch.', @@ -82,6 +89,7 @@ describe('saveFastAgentPostTurnMemory', () => { }, }); mockAppendFastAgentMemory.mockResolvedValue({ saved: true }); + mockAppendFastAgentMemorySavedEvent.mockResolvedValue(undefined); }); it('distills and saves a turn the decision model is confident about', async () => { @@ -108,6 +116,11 @@ describe('saveFastAgentPostTurnMemory', () => { 'conversation-1', 'Staging deploys come from the release branch, not main.', ); + expect(mockAppendFastAgentMemorySavedEvent).toHaveBeenCalledWith({ + sessionId: 'conversation-1', + turnId: 'turn-1', + memories: ['Staging deploys come from the release branch, not main.'], + }); }); it('saves on an explicit remember request or an established finding alone', async () => { @@ -305,6 +318,28 @@ describe('saveFastAgentPostTurnMemory', () => { }); }); + it('only lists facts whose outbox appends succeeded', async () => { + mockEvaluateDecisionModel.mockResolvedValueOnce( + answers({ statedDurable: 0.95 }), + ); + mockGenerateTrackedNonTaskObject.mockResolvedValueOnce({ + object: { memories: ['First fact.', 'Second fact.'] }, + }); + mockAppendFastAgentMemory + .mockResolvedValueOnce({ saved: true }) + .mockResolvedValueOnce({ saved: false, reason: 'memory_full' }); + + await expect(saveFastAgentPostTurnMemory(turn)).resolves.toEqual({ + status: 'saved', + saved: 1, + }); + expect(mockAppendFastAgentMemorySavedEvent).toHaveBeenCalledWith({ + sessionId: 'conversation-1', + turnId: 'turn-1', + memories: ['First fact.'], + }); + }); + it('never throws when the decision model or the helper model fails', async () => { mockEvaluateDecisionModel.mockRejectedValueOnce(new Error('jev timeout')); await expect(saveFastAgentPostTurnMemory(turn)).resolves.toEqual({ diff --git a/packages/cloud-agents/src/server/fast-agent/fast-agent-post-turn-memory.ts b/packages/cloud-agents/src/server/fast-agent/fast-agent-post-turn-memory.ts index 68a5bd66e2..e9bbd70aca 100644 --- a/packages/cloud-agents/src/server/fast-agent/fast-agent-post-turn-memory.ts +++ b/packages/cloud-agents/src/server/fast-agent/fast-agent-post-turn-memory.ts @@ -9,6 +9,7 @@ import { import { FAST_AGENT_MEMORY_FACT_MAX_CHARS } from '@roomote/types'; import { scrubForMemoryCheck } from '../memory-check-scrub'; +import { appendFastAgentMemorySavedEvent } from './fast-agent-session'; import { generateTrackedNonTaskObject, NON_TASK_INFERENCE_SURFACES, @@ -167,6 +168,7 @@ type FastAgentPostTurnMemoryResult = */ export async function saveFastAgentPostTurnMemory(input: { conversationId: string; + turnId: string; userId: string; request: string; /** Messages the same person steered into the turn while it ran. */ @@ -252,6 +254,7 @@ export async function saveFastAgentPostTurnMemory(input: { }); let saved = 0; + const savedFacts: string[] = []; for (const memory of object.memories) { const result = await appendFastAgentMemory( @@ -264,6 +267,7 @@ export async function saveFastAgentPostTurnMemory(input: { if (saved > 0) break; return { status: 'skipped', reason: result.reason }; } + savedFacts.push(memory); saved += 1; } @@ -271,6 +275,20 @@ export async function saveFastAgentPostTurnMemory(input: { return { status: 'skipped', reason: 'nothing_distilled' }; } + try { + await appendFastAgentMemorySavedEvent({ + sessionId: input.conversationId, + turnId: input.turnId, + memories: savedFacts.map((memory) => scrubForMemoryCheck(memory)), + }); + } catch (error) { + // Memory persistence already succeeded; a transcript event is best + // effort and must not turn a successful save into a reported failure. + console.warn( + `[FastPostTurnMemory] Failed to publish save event. conversationId="${input.conversationId}" error="${error instanceof Error ? error.message : String(error)}"`, + ); + } + console.info( `[FastPostTurnMemory] Saved ${saved} memory fact(s). conversationId="${input.conversationId}" worthSaving=${worthSaving.toFixed(2)}`, ); diff --git a/packages/cloud-agents/src/server/fast-agent/fast-agent-service.ts b/packages/cloud-agents/src/server/fast-agent/fast-agent-service.ts index 8dd5bc7b3e..4f1f8292eb 100644 --- a/packages/cloud-agents/src/server/fast-agent/fast-agent-service.ts +++ b/packages/cloud-agents/src/server/fast-agent/fast-agent-service.ts @@ -6841,8 +6841,11 @@ export async function answerFastAgentQuestion({ !setupSession && currentSessionPrivacy === 'shared' ) { + // Reply delivery has already settled before this detached best-effort + // pass starts; Jev and distillation never gate the visible response. void saveFastAgentPostTurnMemory({ conversationId: session.id, + turnId, userId, // A platform event's own text is not something a person said. request: substantiveHumanInput ? question : '', diff --git a/packages/cloud-agents/src/server/fast-agent/fast-agent-session.ts b/packages/cloud-agents/src/server/fast-agent/fast-agent-session.ts index efb16171c2..424e4ffea6 100644 --- a/packages/cloud-agents/src/server/fast-agent/fast-agent-session.ts +++ b/packages/cloud-agents/src/server/fast-agent/fast-agent-session.ts @@ -13,12 +13,14 @@ import { taskRuns, tasks, } from '@roomote/db/server'; -import type { - FastAgentConversationOwner, - ReasoningEffort, - RunStatus, - TaskSurface, - TaskTrigger, +import { + ACP_ENVELOPE_EVENT_TYPES, + MEMORY_SAVED_EVENT_TEXT, + type FastAgentConversationOwner, + type ReasoningEffort, + type RunStatus, + type TaskSurface, + type TaskTrigger, } from '@roomote/types'; import { captureUserStartedSessionCreated } from '../session-telemetry'; import type { FastAgentConversation } from './fast-agent-conversation'; @@ -227,6 +229,33 @@ export async function upsertFastAgentMessage({ throw lastError; } +export async function appendFastAgentMemorySavedEvent({ + sessionId, + turnId, + memories, +}: { + sessionId: string; + turnId: string; + memories: string[]; +}): Promise { + await upsertFastAgentMessage({ + sessionId, + insertOnly: true, + message: { + eventId: `${turnId}:memory-saved`, + turnId, + turnSeq: 2_000_000_001, + ts: Date.now(), + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + role: 'system', + contentBlocks: [{ type: 'text', text: MEMORY_SAVED_EVENT_TEXT }], + metadata: { visibleInTranscript: true, memorySave: true }, + payload: { memories }, + source: 'roomote', + }, + }); +} + export async function setFastAgentOpenCodeSession({ sessionId, openCodeSessionId, diff --git a/packages/cloud-agents/src/server/task-run-memory-distillation.ts b/packages/cloud-agents/src/server/task-run-memory-distillation.ts index 57ceaeb850..ac309aab3b 100644 --- a/packages/cloud-agents/src/server/task-run-memory-distillation.ts +++ b/packages/cloud-agents/src/server/task-run-memory-distillation.ts @@ -11,8 +11,11 @@ import { sql, taskMessages, } from '@roomote/db/server'; +import type { DatabaseOrTransaction } from '@roomote/db/server'; import { ACP_ENVELOPE_EVENT_TYPES, + MEMORY_SAVED_EVENT_TEXT, + ROOMOTE_RUNTIME_TASK_MESSAGE_PROTOCOL, extractAcpMessageText, extractVisibleAcpPromptText, isSystemInjectedAcpPromptText, @@ -125,12 +128,13 @@ function clip(text: string, maxChars: number): string { */ async function loadLatestTurn( runId: number, -): Promise<{ request: string; report: string } | null> { +): Promise<{ request: string; report: string; turnTs: number } | null> { const rows = await db .select({ eventType: taskMessages.eventType, contentBlocks: taskMessages.contentBlocks, payload: taskMessages.payload, + ts: taskMessages.ts, }) .from(taskMessages) .where( @@ -168,7 +172,12 @@ async function loadLatestTurn( : raw, ACP_ENVELOPE_EVENT_TYPES.UserPrompt, )?.trim() ?? ''; - break; + if (reportMessages.length === 0) return null; + return { + request: clip(request, REQUEST_MAX_CHARS), + report: reportMessages.join('\n\n'), + turnTs: row.ts, + }; } if (reportMessages.length < REPORT_MESSAGE_LIMIT && remaining > 0) { @@ -182,10 +191,88 @@ async function loadLatestTurn( ? { request: clip(request, REQUEST_MAX_CHARS), report: reportMessages.join('\n\n'), + turnTs: rows[0]?.ts ?? runId, } : null; } +async function publishTaskMemorySavedEvent(input: { + database: DatabaseOrTransaction; + runId: number; + taskId: string; + userId?: string | null; + turnTs: number; + summary: string; +}): Promise { + const distilled = input.summary.endsWith(DISTILLED_SUMMARY_NOTE) + ? input.summary.slice(0, -DISTILLED_SUMMARY_NOTE.length).trim() + : input.summary; + const memories = [scrubForMemoryCheck(distilled)].filter(Boolean); + + if (memories.length === 0) return; + + await input.database + .insert(taskMessages) + .values({ + runId: input.runId, + taskId: input.taskId, + userId: input.userId ?? null, + // The settled turn timestamp makes retries idempotent while allowing + // later follow-up turns in the same run to publish their own row. + ts: input.turnTs, + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + role: 'system', + protocol: ROOMOTE_RUNTIME_TASK_MESSAGE_PROTOCOL, + contentBlocks: [{ type: 'text', text: MEMORY_SAVED_EVENT_TEXT }], + metadata: { + visibleInTranscript: true, + memorySave: true, + automatic: true, + }, + payload: { memories }, + source: 'roomote', + }) + .onConflictDoNothing({ + target: [ + taskMessages.taskId, + taskMessages.protocol, + taskMessages.ts, + taskMessages.eventType, + ], + }); +} + +async function saveSummaryAndPublishTaskMemoryEvent(input: { + runId: number; + taskId: string; + userId?: string | null; + turnTs: number; + summary: string; + existing: string | null; + requeue: boolean; +}): Promise { + return db.transaction(async (tx) => { + const saved = await saveBrainDistilledSummary( + tx, + input.runId, + input.summary, + input.existing, + { requeue: input.requeue }, + ); + if (!saved) return false; + + await publishTaskMemorySavedEvent({ + database: tx, + runId: input.runId, + taskId: input.taskId, + userId: input.userId, + turnTs: input.turnTs, + summary: input.summary, + }); + return true; + }); +} + /** * After a task turn settles, ask the decision model whether the turn holds * anything a later task could reuse, and only then pay for a helper-model @@ -284,7 +371,13 @@ export async function distillTaskRunTurnMemory(input: { const summary = `${renderTaskMemorySummary(object)}\n\n${DISTILLED_SUMMARY_NOTE}`; if ( - !(await saveBrainDistilledSummary(db, input.runId, summary, existing, { + !(await saveSummaryAndPublishTaskMemoryEvent({ + runId: input.runId, + taskId: input.taskId, + userId: input.userId, + turnTs: turn.turnTs, + summary, + existing, requeue: input.requeue, })) ) { diff --git a/packages/types/src/acp.ts b/packages/types/src/acp.ts index b3ea7e8da5..e7b1087495 100644 --- a/packages/types/src/acp.ts +++ b/packages/types/src/acp.ts @@ -36,6 +36,7 @@ export const ACP_ENVELOPE_EVENT_TYPES = { TaskCancelled: 'roomote_runtime.task_cancelled', /** Voice call lifecycle marker persisted in a Fast Session transcript. */ VoiceCall: 'roomote_runtime.voice_call', + MemorySaved: 'roomote_runtime.memory_saved', } as const; export type AcpEnvelopeEventType = diff --git a/packages/types/src/index.ts b/packages/types/src/index.ts index 23d69b6f9c..d38215e7c1 100644 --- a/packages/types/src/index.ts +++ b/packages/types/src/index.ts @@ -35,6 +35,7 @@ export * from './deploy-marker'; export * from './deployment-access-policy'; export * from './brain'; export * from './memory-mcp'; +export * from './memory-events'; export * from './custom-mcp-servers'; export * from './environment-config'; export * from './environment-recipe'; diff --git a/packages/types/src/memory-events.ts b/packages/types/src/memory-events.ts new file mode 100644 index 0000000000..6e9b26f35e --- /dev/null +++ b/packages/types/src/memory-events.ts @@ -0,0 +1,18 @@ +import { z } from 'zod'; + +export const MEMORY_SAVED_EVENT_TEXT = 'Saved to memory' as const; + +export const memorySavedEventPayloadSchema = z.object({ + memories: z.array(z.string().trim().min(1)).min(1), +}); + +export type MemorySavedEventPayload = z.infer< + typeof memorySavedEventPayloadSchema +>; + +export function parseMemorySavedEventPayload( + payload: unknown, +): MemorySavedEventPayload | null { + const parsed = memorySavedEventPayloadSchema.safeParse(payload); + return parsed.success ? parsed.data : null; +}