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;
+}