From 3db32e0e5d97c97bf5fd07c894f468557fc1a6ad Mon Sep 17 00:00:00 2001 From: "@mrubens" <2600+mrubens@users.noreply.github.com> Date: Mon, 21 Sep 2026 19:50:00 +0000 Subject: [PATCH 1/6] [Improve] Show automatic memory saves and cache repository skills --- .../FastSessionTranscript.client.test.tsx | 66 +++++++++++++++++++ .../[sessionId]/FastSessionTranscript.tsx | 36 ++++++++++ .../fast-agent-post-turn-memory.test.ts | 13 ++++ .../fast-agent-prompt-skill-catalog.test.ts | 32 +++++++++ .../fast-agent/fast-agent-post-turn-memory.ts | 16 +++++ .../fast-agent-prompt-skill-catalog.ts | 63 ++++++++++++++++-- .../server/fast-agent/fast-agent-service.ts | 9 +++ .../server/fast-agent/fast-agent-session.ts | 40 +++++++++-- packages/types/src/acp.ts | 1 + 9 files changed, 264 insertions(+), 12 deletions(-) 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..a521322e74 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 { Brain } from '@/components/system'; import { AcpMessageItem, @@ -127,6 +128,39 @@ type TranscriptOwner = { const ROOMOTE_KICKOFF_LINK = /\r?\n\r?\n\[Open in Roomote\]\([^\r\n]+\)\s*$/; +function renderMemorySavedMessage(message: AcpUiMessage) { + if (message.updateType !== ACP_ENVELOPE_EVENT_TYPES.MemorySaved) { + return undefined; + } + + const payload = message.data as Record; + const rawMemories = payload.memories; + const memories = Array.isArray(rawMemories) + ? rawMemories.filter( + (memory): memory is string => typeof memory === 'string', + ) + : []; + + return ( +
+ + + Saved to memory + {memories.length > 0 ? ( + {memories.length} + ) : null} + + {memories.length > 0 ? ( +
    + {memories.map((memory) => ( +
  • {memory}
  • + ))} +
+ ) : null} +
+ ); +} + function getTranscriptMessageText(message: TranscriptMessage) { const text = getTextFromContentBlocks(message.contentBlocks) ?? undefined; const payload = message.payload as { kickoff?: unknown } | null; @@ -941,6 +975,8 @@ export function FastSessionTranscript({ ); const renderCapabilityOfferMessage = useCallback( (message: AcpUiMessage) => { + const memorySaved = renderMemorySavedMessage(message); + if (memorySaved !== undefined) return memorySaved; const offer = capabilityOffersByMessageId.get(message.id); if (!offer) return undefined; const introMessage: AcpUiMessage = { 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..5c22dfd015 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 () => { diff --git a/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts index b7bdf992f4..d7bfb55402 100644 --- a/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts +++ b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts @@ -127,6 +127,38 @@ describe('loadFastAgentPromptSkillCatalog', () => { ); }); + it('reuses a scoped catalog during the short per-turn cache window', async () => { + const instanceSkills = { + list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), + }; + const settingsSkills = { + listPromptCatalog: vi.fn().mockResolvedValue({ + marketplaceSources: [], + skills: [], + warnings: [], + }), + dispose: vi.fn().mockResolvedValue(undefined), + }; + const repositorySkills = { + list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), + dispose: vi.fn().mockResolvedValue(undefined), + }; + const sources = { instanceSkills, settingsSkills, repositorySkills }; + + await loadFastAgentPromptSkillCatalog(sources, { + repositoryCacheKey: 'cache-test-user:env-1', + }); + await loadFastAgentPromptSkillCatalog(sources, { + repositoryCacheKey: 'cache-test-user:env-1', + }); + + expect(instanceSkills.list).toHaveBeenCalledTimes(2); + expect(settingsSkills.listPromptCatalog).toHaveBeenCalledTimes(2); + expect(repositorySkills.list).toHaveBeenCalledOnce(); + expect(settingsSkills.dispose).toHaveBeenCalledTimes(2); + expect(repositorySkills.dispose).toHaveBeenCalledTimes(2); + }); + it('drops custom skills that collide with a packaged skill name', async () => { const catalog = await loadFastAgentPromptSkillCatalog({ instanceSkills: { 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..d2d47bb186 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. */ @@ -271,6 +273,20 @@ export async function saveFastAgentPostTurnMemory(input: { return { status: 'skipped', reason: 'nothing_distilled' }; } + try { + await appendFastAgentMemorySavedEvent({ + sessionId: input.conversationId, + turnId: input.turnId, + memories: object.memories.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-prompt-skill-catalog.ts b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts index 1496598259..f2a8537c1d 100644 --- a/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts +++ b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts @@ -27,6 +27,7 @@ export type FastAgentPromptSkillCatalog = { }; export const FAST_AGENT_PROMPT_SKILL_LIMIT = 64; +const PROMPT_SKILL_CATALOG_CACHE_TTL_MS = 30_000; type PromptSkillCatalogSources = { instanceSkills: Pick; @@ -74,17 +75,67 @@ async function settlePromptSkillSource( export async function loadFastAgentPromptSkillCatalog( sources: PromptSkillCatalogSources, + options: { repositoryCacheKey?: string } = {}, ): Promise { - const [instance, settings, repository] = await Promise.all([ + const repository = sources.repositorySkills + ? loadRepositorySkillSource( + sources.repositorySkills, + options.repositoryCacheKey, + ) + : Promise.resolve< + PromiseSettledResult | undefined + >(undefined); + const [instance, settings, repositoryResult] = await Promise.all([ settlePromptSkillSource(() => sources.instanceSkills.list()), settlePromptSkillSource(() => sources.settingsSkills.listPromptCatalog()), - sources.repositorySkills - ? settlePromptSkillSource(() => sources.repositorySkills!.list()) - : Promise.resolve< - PromiseSettledResult | undefined - >(undefined), + repository, ]); + return mergeFastAgentPromptSkillCatalog( + instance, + settings, + repositoryResult, + sources, + ); +} + +const repositorySkillSourceCache = new Map< + string, + { + expiresAt: number; + promise: Promise>; + } +>(); + +async function loadRepositorySkillSource( + source: NonNullable, + cacheKey?: string, +): Promise> { + if (!cacheKey) return settlePromptSkillSource(() => source.list()); + const cached = repositorySkillSourceCache.get(cacheKey); + if (cached && cached.expiresAt > Date.now()) return cached.promise; + + const promise = settlePromptSkillSource(() => source.list()); + repositorySkillSourceCache.set(cacheKey, { + expiresAt: Date.now() + PROMPT_SKILL_CATALOG_CACHE_TTL_MS, + promise, + }); + promise.then((result) => { + if (result.status === 'rejected') { + const current = repositorySkillSourceCache.get(cacheKey); + if (current?.promise === promise) + repositorySkillSourceCache.delete(cacheKey); + } + }); + return promise; +} + +async function mergeFastAgentPromptSkillCatalog( + instance: PromiseSettledResult, + settings: PromiseSettledResult, + repository: PromiseSettledResult | undefined, + sources: PromptSkillCatalogSources, +): Promise { try { // A partial failure degrades to a warning so the surviving source still // reaches the prompt. When nothing loaded, throw instead of rendering an 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..57dde7490b 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 @@ -3518,6 +3518,12 @@ export async function answerFastAgentQuestion({ ), userId, }), + { + repositoryCacheKey: `${userId}:${availableEnvironments + .map((environment) => environment.id) + .sort() + .join(',')}`, + }, ) .then((catalog) => { for (const warning of catalog.warnings) { @@ -6841,8 +6847,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..186f17b5b4 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,13 @@ import { taskRuns, tasks, } from '@roomote/db/server'; -import type { - FastAgentConversationOwner, - ReasoningEffort, - RunStatus, - TaskSurface, - TaskTrigger, +import { + ACP_ENVELOPE_EVENT_TYPES, + 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 +228,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: 'Saved to memory' }], + metadata: { visibleInTranscript: true, memorySave: true }, + payload: { memories }, + source: 'roomote', + }, + }); +} + export async function setFastAgentOpenCodeSession({ sessionId, openCodeSessionId, 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 = From fd4864bf426d3d841fb438baed9d43c2c36f3ad6 Mon Sep 17 00:00:00 2001 From: "@mrubens" <2600+mrubens@users.noreply.github.com> Date: Mon, 21 Sep 2026 20:12:07 +0000 Subject: [PATCH 2/6] [Fix] Address memory save review and test contention --- ...usive-automation-settings-database-lock.ts | 5 ++- .../fast-agent-post-turn-memory.test.ts | 22 ++++++++++ .../fast-agent-prompt-skill-catalog.test.ts | 44 +++++++++++++++++++ .../fast-agent/fast-agent-post-turn-memory.ts | 4 +- .../fast-agent-prompt-skill-catalog.ts | 6 ++- 5 files changed, 77 insertions(+), 4 deletions(-) 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/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 5c22dfd015..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 @@ -318,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/__tests__/fast-agent-prompt-skill-catalog.test.ts b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts index d7bfb55402..4fcf57ec06 100644 --- a/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts +++ b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts @@ -159,6 +159,50 @@ describe('loadFastAgentPromptSkillCatalog', () => { expect(repositorySkills.dispose).toHaveBeenCalledTimes(2); }); + it('prunes expired repository scopes when a later scope is loaded', async () => { + vi.useFakeTimers(); + try { + const makeSources = () => ({ + instanceSkills: { + list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), + }, + settingsSkills: { + listPromptCatalog: vi.fn().mockResolvedValue({ + marketplaceSources: [], + skills: [], + warnings: [], + }), + }, + repositorySkills: { + list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), + }, + }); + const first = makeSources(); + const second = makeSources(); + const third = makeSources(); + + await loadFastAgentPromptSkillCatalog(first, { + repositoryCacheKey: 'cache-expiry-first', + }); + await loadFastAgentPromptSkillCatalog(second, { + repositoryCacheKey: 'cache-expiry-second', + }); + vi.advanceTimersByTime(30_001); + await loadFastAgentPromptSkillCatalog(third, { + repositoryCacheKey: 'cache-expiry-third', + }); + await loadFastAgentPromptSkillCatalog(first, { + repositoryCacheKey: 'cache-expiry-first', + }); + + expect(first.repositorySkills.list).toHaveBeenCalledTimes(2); + expect(second.repositorySkills.list).toHaveBeenCalledOnce(); + expect(third.repositorySkills.list).toHaveBeenCalledOnce(); + } finally { + vi.useRealTimers(); + } + }); + it('drops custom skills that collide with a packaged skill name', async () => { const catalog = await loadFastAgentPromptSkillCatalog({ instanceSkills: { 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 d2d47bb186..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 @@ -254,6 +254,7 @@ export async function saveFastAgentPostTurnMemory(input: { }); let saved = 0; + const savedFacts: string[] = []; for (const memory of object.memories) { const result = await appendFastAgentMemory( @@ -266,6 +267,7 @@ export async function saveFastAgentPostTurnMemory(input: { if (saved > 0) break; return { status: 'skipped', reason: result.reason }; } + savedFacts.push(memory); saved += 1; } @@ -277,7 +279,7 @@ export async function saveFastAgentPostTurnMemory(input: { await appendFastAgentMemorySavedEvent({ sessionId: input.conversationId, turnId: input.turnId, - memories: object.memories.map((memory) => scrubForMemoryCheck(memory)), + memories: savedFacts.map((memory) => scrubForMemoryCheck(memory)), }); } catch (error) { // Memory persistence already succeeded; a transcript event is best diff --git a/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts index f2a8537c1d..01f36202d6 100644 --- a/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts +++ b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts @@ -111,9 +111,13 @@ async function loadRepositorySkillSource( source: NonNullable, cacheKey?: string, ): Promise> { + const now = Date.now(); + for (const [key, entry] of repositorySkillSourceCache) { + if (entry.expiresAt <= now) repositorySkillSourceCache.delete(key); + } if (!cacheKey) return settlePromptSkillSource(() => source.list()); const cached = repositorySkillSourceCache.get(cacheKey); - if (cached && cached.expiresAt > Date.now()) return cached.promise; + if (cached) return cached.promise; const promise = settlePromptSkillSource(() => source.list()); repositorySkillSourceCache.set(cacheKey, { From e51101251b6efc2991305f4f0a82b8ded34b5a77 Mon Sep 17 00:00:00 2001 From: "@mrubens" <2600+mrubens@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:41:53 +0000 Subject: [PATCH 3/6] [Feat] Show automatic task memory saves --- .../[sessionId]/FastSessionTranscript.tsx | 40 +--------- .../app/(sandbox)/task/[taskId]/Messages.tsx | 11 +++ .../MemorySavedMessage.client.test.tsx | 35 +++++++++ .../ai-elements/MemorySavedMessage.tsx | 39 ++++++++++ apps/web/src/lib/server/task-messages.test.ts | 28 +++++++ .../task-run-memory-distillation.test.ts | 57 ++++++++++++++ .../fast-agent-prompt-skill-catalog.test.ts | 76 ------------------- .../fast-agent-prompt-skill-catalog.ts | 55 ++------------ .../server/fast-agent/fast-agent-service.ts | 6 -- .../server/fast-agent/fast-agent-session.ts | 3 +- .../server/task-run-memory-distillation.ts | 61 +++++++++++++++ packages/types/src/index.ts | 1 + packages/types/src/memory-events.ts | 18 +++++ 13 files changed, 263 insertions(+), 167 deletions(-) create mode 100644 apps/web/src/components/ai-elements/MemorySavedMessage.client.test.tsx create mode 100644 apps/web/src/components/ai-elements/MemorySavedMessage.tsx create mode 100644 packages/types/src/memory-events.ts diff --git a/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx b/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx index a521322e74..a2a592f64e 100644 --- a/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx +++ b/apps/web/src/app/(sandbox)/sessions/[sessionId]/FastSessionTranscript.tsx @@ -95,7 +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 { Brain } from '@/components/system'; +import { MemorySavedMessage } from '@/components/ai-elements/MemorySavedMessage'; import { AcpMessageItem, @@ -128,39 +128,6 @@ type TranscriptOwner = { const ROOMOTE_KICKOFF_LINK = /\r?\n\r?\n\[Open in Roomote\]\([^\r\n]+\)\s*$/; -function renderMemorySavedMessage(message: AcpUiMessage) { - if (message.updateType !== ACP_ENVELOPE_EVENT_TYPES.MemorySaved) { - return undefined; - } - - const payload = message.data as Record; - const rawMemories = payload.memories; - const memories = Array.isArray(rawMemories) - ? rawMemories.filter( - (memory): memory is string => typeof memory === 'string', - ) - : []; - - return ( -
- - - Saved to memory - {memories.length > 0 ? ( - {memories.length} - ) : null} - - {memories.length > 0 ? ( -
    - {memories.map((memory) => ( -
  • {memory}
  • - ))} -
- ) : null} -
- ); -} - function getTranscriptMessageText(message: TranscriptMessage) { const text = getTextFromContentBlocks(message.contentBlocks) ?? undefined; const payload = message.payload as { kickoff?: unknown } | null; @@ -975,8 +942,9 @@ export function FastSessionTranscript({ ); const renderCapabilityOfferMessage = useCallback( (message: AcpUiMessage) => { - const memorySaved = renderMemorySavedMessage(message); - if (memorySaved !== undefined) return memorySaved; + 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/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..992c80541a 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,6 +15,8 @@ const { mockIsBrainEnabled: vi.fn(), mockIsTaskRunSharedBrainEligible: vi.fn(), mockSaveBrainDistilledSummary: vi.fn(), + mockInsertTaskMemoryEvent: vi.fn(), + mockInsertTaskMemoryValues: vi.fn(), mockTurnRows: vi.fn(), })); @@ -28,6 +32,7 @@ vi.mock('@roomote/db/server', () => ({ isTaskRunSharedBrainEligible: mockIsTaskRunSharedBrainEligible, saveBrainDistilledSummary: mockSaveBrainDistilledSummary, db: { + insert: () => ({ values: mockInsertTaskMemoryValues }), select: () => ({ from: () => ({ where: () => ({ orderBy: () => ({ limit: mockTurnRows }) }), @@ -105,6 +110,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.'), @@ -139,6 +148,25 @@ describe('distillTaskRunTurnMemory', () => { }, }), ); + expect(mockInsertTaskMemoryValues).toHaveBeenCalledWith({ + runId: 102, + taskId: 'task-1', + userId: 'user-1', + ts: 102, + 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 +178,35 @@ describe('distillTaskRunTurnMemory', () => { ); }); + it('uses a deterministic task event key across fallback retries', async () => { + await distillTaskRunTurnMemory(run); + 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: 102, + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + }); + expect(mockInsertTaskMemoryValues.mock.calls[1]?.[0]).toMatchObject({ + taskId: 'task-1', + ts: 102, + eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, + }); + }); + + it('keeps a successful summary save when event publication is temporarily unavailable', async () => { + mockInsertTaskMemoryEvent.mockRejectedValueOnce(new Error('database busy')); + + await expect(distillTaskRunTurnMemory(run)).resolves.toContain( + 'Webhook retries are capped', + ); + expect(mockSaveBrainDistilledSummary).toHaveBeenCalledOnce(); + }); + 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-prompt-skill-catalog.test.ts b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts index 4fcf57ec06..b7bdf992f4 100644 --- a/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts +++ b/packages/cloud-agents/src/server/fast-agent/__tests__/fast-agent-prompt-skill-catalog.test.ts @@ -127,82 +127,6 @@ describe('loadFastAgentPromptSkillCatalog', () => { ); }); - it('reuses a scoped catalog during the short per-turn cache window', async () => { - const instanceSkills = { - list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), - }; - const settingsSkills = { - listPromptCatalog: vi.fn().mockResolvedValue({ - marketplaceSources: [], - skills: [], - warnings: [], - }), - dispose: vi.fn().mockResolvedValue(undefined), - }; - const repositorySkills = { - list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), - dispose: vi.fn().mockResolvedValue(undefined), - }; - const sources = { instanceSkills, settingsSkills, repositorySkills }; - - await loadFastAgentPromptSkillCatalog(sources, { - repositoryCacheKey: 'cache-test-user:env-1', - }); - await loadFastAgentPromptSkillCatalog(sources, { - repositoryCacheKey: 'cache-test-user:env-1', - }); - - expect(instanceSkills.list).toHaveBeenCalledTimes(2); - expect(settingsSkills.listPromptCatalog).toHaveBeenCalledTimes(2); - expect(repositorySkills.list).toHaveBeenCalledOnce(); - expect(settingsSkills.dispose).toHaveBeenCalledTimes(2); - expect(repositorySkills.dispose).toHaveBeenCalledTimes(2); - }); - - it('prunes expired repository scopes when a later scope is loaded', async () => { - vi.useFakeTimers(); - try { - const makeSources = () => ({ - instanceSkills: { - list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), - }, - settingsSkills: { - listPromptCatalog: vi.fn().mockResolvedValue({ - marketplaceSources: [], - skills: [], - warnings: [], - }), - }, - repositorySkills: { - list: vi.fn().mockResolvedValue({ skills: [], warnings: [] }), - }, - }); - const first = makeSources(); - const second = makeSources(); - const third = makeSources(); - - await loadFastAgentPromptSkillCatalog(first, { - repositoryCacheKey: 'cache-expiry-first', - }); - await loadFastAgentPromptSkillCatalog(second, { - repositoryCacheKey: 'cache-expiry-second', - }); - vi.advanceTimersByTime(30_001); - await loadFastAgentPromptSkillCatalog(third, { - repositoryCacheKey: 'cache-expiry-third', - }); - await loadFastAgentPromptSkillCatalog(first, { - repositoryCacheKey: 'cache-expiry-first', - }); - - expect(first.repositorySkills.list).toHaveBeenCalledTimes(2); - expect(second.repositorySkills.list).toHaveBeenCalledOnce(); - expect(third.repositorySkills.list).toHaveBeenCalledOnce(); - } finally { - vi.useRealTimers(); - } - }); - it('drops custom skills that collide with a packaged skill name', async () => { const catalog = await loadFastAgentPromptSkillCatalog({ instanceSkills: { diff --git a/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts index 01f36202d6..7a0c18117a 100644 --- a/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts +++ b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts @@ -27,7 +27,6 @@ export type FastAgentPromptSkillCatalog = { }; export const FAST_AGENT_PROMPT_SKILL_LIMIT = 64; -const PROMPT_SKILL_CATALOG_CACHE_TTL_MS = 30_000; type PromptSkillCatalogSources = { instanceSkills: Pick; @@ -75,65 +74,25 @@ async function settlePromptSkillSource( export async function loadFastAgentPromptSkillCatalog( sources: PromptSkillCatalogSources, - options: { repositoryCacheKey?: string } = {}, ): Promise { - const repository = sources.repositorySkills - ? loadRepositorySkillSource( - sources.repositorySkills, - options.repositoryCacheKey, - ) - : Promise.resolve< - PromiseSettledResult | undefined - >(undefined); - const [instance, settings, repositoryResult] = await Promise.all([ + const [instance, settings, repository] = await Promise.all([ settlePromptSkillSource(() => sources.instanceSkills.list()), settlePromptSkillSource(() => sources.settingsSkills.listPromptCatalog()), - repository, + sources.repositorySkills + ? settlePromptSkillSource(() => sources.repositorySkills!.list()) + : Promise.resolve< + PromiseSettledResult | undefined + >(undefined), ]); return mergeFastAgentPromptSkillCatalog( instance, settings, - repositoryResult, + repository, sources, ); } -const repositorySkillSourceCache = new Map< - string, - { - expiresAt: number; - promise: Promise>; - } ->(); - -async function loadRepositorySkillSource( - source: NonNullable, - cacheKey?: string, -): Promise> { - const now = Date.now(); - for (const [key, entry] of repositorySkillSourceCache) { - if (entry.expiresAt <= now) repositorySkillSourceCache.delete(key); - } - if (!cacheKey) return settlePromptSkillSource(() => source.list()); - const cached = repositorySkillSourceCache.get(cacheKey); - if (cached) return cached.promise; - - const promise = settlePromptSkillSource(() => source.list()); - repositorySkillSourceCache.set(cacheKey, { - expiresAt: Date.now() + PROMPT_SKILL_CATALOG_CACHE_TTL_MS, - promise, - }); - promise.then((result) => { - if (result.status === 'rejected') { - const current = repositorySkillSourceCache.get(cacheKey); - if (current?.promise === promise) - repositorySkillSourceCache.delete(cacheKey); - } - }); - return promise; -} - async function mergeFastAgentPromptSkillCatalog( instance: PromiseSettledResult, settings: PromiseSettledResult, 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 57dde7490b..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 @@ -3518,12 +3518,6 @@ export async function answerFastAgentQuestion({ ), userId, }), - { - repositoryCacheKey: `${userId}:${availableEnvironments - .map((environment) => environment.id) - .sort() - .join(',')}`, - }, ) .then((catalog) => { for (const warning of catalog.warnings) { 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 186f17b5b4..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 @@ -15,6 +15,7 @@ import { } from '@roomote/db/server'; import { ACP_ENVELOPE_EVENT_TYPES, + MEMORY_SAVED_EVENT_TEXT, type FastAgentConversationOwner, type ReasoningEffort, type RunStatus, @@ -247,7 +248,7 @@ export async function appendFastAgentMemorySavedEvent({ ts: Date.now(), eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, role: 'system', - contentBlocks: [{ type: 'text', text: 'Saved to memory' }], + contentBlocks: [{ type: 'text', text: MEMORY_SAVED_EVENT_TEXT }], metadata: { visibleInTranscript: true, memorySave: true }, payload: { memories }, source: 'roomote', 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..9aa72f41a3 100644 --- a/packages/cloud-agents/src/server/task-run-memory-distillation.ts +++ b/packages/cloud-agents/src/server/task-run-memory-distillation.ts @@ -13,6 +13,8 @@ import { } from '@roomote/db/server'; import { ACP_ENVELOPE_EVENT_TYPES, + MEMORY_SAVED_EVENT_TEXT, + ROOMOTE_RUNTIME_TASK_MESSAGE_PROTOCOL, extractAcpMessageText, extractVisibleAcpPromptText, isSystemInjectedAcpPromptText, @@ -186,6 +188,50 @@ async function loadLatestTurn( : null; } +async function publishTaskMemorySavedEvent(input: { + runId: number; + taskId: string; + userId?: string | null; + 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 db + .insert(taskMessages) + .values({ + runId: input.runId, + taskId: input.taskId, + userId: input.userId ?? null, + // The run id makes retries idempotent even when the completion timestamp + // is unavailable to this background projection. + ts: input.runId, + 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, + ], + }); +} + /** * 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 @@ -291,6 +337,21 @@ export async function distillTaskRunTurnMemory(input: { return null; } + try { + await publishTaskMemorySavedEvent({ + runId: input.runId, + taskId: input.taskId, + userId: input.userId, + summary, + }); + } catch (error) { + // The Brain summary is already persisted; a transcript event can be + // retried by the durable drainer without turning the save into failure. + console.warn( + `[TaskRunMemoryDistillation] Failed to publish save event. runId=${input.runId} error="${error instanceof Error ? error.message : String(error)}"`, + ); + } + console.info( `[TaskRunMemoryDistillation] Updated the run memory. runId=${input.runId} worthSaving=${worthSaving.toFixed(2)}`, ); 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; +} From 1d51787835bea98bf3821d59082087152d13c1f1 Mon Sep 17 00:00:00 2001 From: "@mrubens" <2600+mrubens@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:43:57 +0000 Subject: [PATCH 4/6] [Feat] Show automatic task memory saves --- .../fast-agent/fast-agent-prompt-skill-catalog.ts | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts index 7a0c18117a..1496598259 100644 --- a/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts +++ b/packages/cloud-agents/src/server/fast-agent/fast-agent-prompt-skill-catalog.ts @@ -85,20 +85,6 @@ export async function loadFastAgentPromptSkillCatalog( >(undefined), ]); - return mergeFastAgentPromptSkillCatalog( - instance, - settings, - repository, - sources, - ); -} - -async function mergeFastAgentPromptSkillCatalog( - instance: PromiseSettledResult, - settings: PromiseSettledResult, - repository: PromiseSettledResult | undefined, - sources: PromptSkillCatalogSources, -): Promise { try { // A partial failure degrades to a warning so the surviving source still // reaches the prompt. When nothing loaded, throw instead of rendering an From 68039caf7c77accbd603166ae919b66068b23da6 Mon Sep 17 00:00:00 2001 From: "@mrubens" <2600+mrubens@users.noreply.github.com> Date: Mon, 21 Sep 2026 21:52:56 +0000 Subject: [PATCH 5/6] [Fix] Use settled task turn keys for memory events --- .../task-run-memory-distillation.test.ts | 23 ++++++++++++------- .../server/task-run-memory-distillation.ts | 19 +++++++++++---- 2 files changed, 29 insertions(+), 13 deletions(-) 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 992c80541a..d8117e0834 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 @@ -56,14 +56,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< @@ -144,6 +146,7 @@ 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: '', }, }), @@ -152,7 +155,7 @@ describe('distillTaskRunTurnMemory', () => { runId: 102, taskId: 'task-1', userId: 'user-1', - ts: 102, + ts: 1, eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, role: 'system', protocol: 'roomote_runtime', @@ -178,8 +181,12 @@ describe('distillTaskRunTurnMemory', () => { ); }); - it('uses a deterministic task event key across fallback retries', async () => { + 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}`, ); @@ -188,12 +195,12 @@ describe('distillTaskRunTurnMemory', () => { expect(mockInsertTaskMemoryValues).toHaveBeenCalledTimes(2); expect(mockInsertTaskMemoryValues.mock.calls[0]?.[0]).toMatchObject({ taskId: 'task-1', - ts: 102, + ts: 1, eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, }); expect(mockInsertTaskMemoryValues.mock.calls[1]?.[0]).toMatchObject({ taskId: 'task-1', - ts: 102, + ts: 2, eventType: ACP_ENVELOPE_EVENT_TYPES.MemorySaved, }); }); 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 9aa72f41a3..76261a3de8 100644 --- a/packages/cloud-agents/src/server/task-run-memory-distillation.ts +++ b/packages/cloud-agents/src/server/task-run-memory-distillation.ts @@ -127,12 +127,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( @@ -170,7 +171,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) { @@ -184,6 +190,7 @@ async function loadLatestTurn( ? { request: clip(request, REQUEST_MAX_CHARS), report: reportMessages.join('\n\n'), + turnTs: rows[0]?.ts ?? runId, } : null; } @@ -192,6 +199,7 @@ async function publishTaskMemorySavedEvent(input: { runId: number; taskId: string; userId?: string | null; + turnTs: number; summary: string; }): Promise { const distilled = input.summary.endsWith(DISTILLED_SUMMARY_NOTE) @@ -207,9 +215,9 @@ async function publishTaskMemorySavedEvent(input: { runId: input.runId, taskId: input.taskId, userId: input.userId ?? null, - // The run id makes retries idempotent even when the completion timestamp - // is unavailable to this background projection. - ts: input.runId, + // 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, @@ -342,6 +350,7 @@ export async function distillTaskRunTurnMemory(input: { runId: input.runId, taskId: input.taskId, userId: input.userId, + turnTs: turn.turnTs, summary, }); } catch (error) { From e2bffad250a54dded9c39c3a19a9502c1a1193cb Mon Sep 17 00:00:00 2001 From: "@mrubens" <2600+mrubens@users.noreply.github.com> Date: Mon, 21 Sep 2026 22:06:45 +0000 Subject: [PATCH 6/6] [Fix] Retry task memory transcript events atomically --- .../task-run-memory-distillation.test.ts | 44 ++++++++------ .../server/task-run-memory-distillation.ts | 57 +++++++++++++------ 2 files changed, 67 insertions(+), 34 deletions(-) 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 d8117e0834..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 @@ -20,26 +20,32 @@ const { 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: { - insert: () => ({ values: mockInsertTaskMemoryValues }), +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, @@ -205,13 +211,17 @@ describe('distillTaskRunTurnMemory', () => { }); }); - it('keeps a successful summary save when event publication is temporarily unavailable', async () => { + 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(mockSaveBrainDistilledSummary).toHaveBeenCalledOnce(); + expect(mockInsertTaskMemoryEvent).toHaveBeenCalledTimes(2); }); it('builds on its own earlier memory and saves against that exact text', async () => { 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 76261a3de8..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,6 +11,7 @@ import { sql, taskMessages, } from '@roomote/db/server'; +import type { DatabaseOrTransaction } from '@roomote/db/server'; import { ACP_ENVELOPE_EVENT_TYPES, MEMORY_SAVED_EVENT_TEXT, @@ -196,6 +197,7 @@ async function loadLatestTurn( } async function publishTaskMemorySavedEvent(input: { + database: DatabaseOrTransaction; runId: number; taskId: string; userId?: string | null; @@ -209,7 +211,7 @@ async function publishTaskMemorySavedEvent(input: { if (memories.length === 0) return; - await db + await input.database .insert(taskMessages) .values({ runId: input.runId, @@ -240,6 +242,37 @@ async function publishTaskMemorySavedEvent(input: { }); } +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 @@ -338,27 +371,17 @@ export async function distillTaskRunTurnMemory(input: { const summary = `${renderTaskMemorySummary(object)}\n\n${DISTILLED_SUMMARY_NOTE}`; if ( - !(await saveBrainDistilledSummary(db, input.runId, summary, existing, { - requeue: input.requeue, - })) - ) { - return null; - } - - try { - await publishTaskMemorySavedEvent({ + !(await saveSummaryAndPublishTaskMemoryEvent({ runId: input.runId, taskId: input.taskId, userId: input.userId, turnTs: turn.turnTs, summary, - }); - } catch (error) { - // The Brain summary is already persisted; a transcript event can be - // retried by the durable drainer without turning the save into failure. - console.warn( - `[TaskRunMemoryDistillation] Failed to publish save event. runId=${input.runId} error="${error instanceof Error ? error.message : String(error)}"`, - ); + existing, + requeue: input.requeue, + })) + ) { + return null; } console.info(