import { Context, Effect, Layer, Ref } from 'effect'; import type { Prompt } from '@effect/ai'; import type { ConversationContext, ConversationMetadata, ConversationPolicy, ContextStorage, SessionSummary } from './context'; import { ContextNotFoundError, ContextStorageError } from './errors'; import { normalizeMessage, normalizeMessages } from '../messages'; /** * ContextStorageService interface for Effect-based conversation management */ export interface ContextStorageService { /** * Generate a unique conversation ID */ generateConversationId(): Effect.Effect; /** * Get or create a conversation context * In strict mode, throws if conversationId provided but not found */ getContext( conversationId?: string, options?: { strict?: boolean } ): Effect.Effect; /** * Get conversation context by ID (returns null if not found) */ getContextById(conversationId: string): Effect.Effect; /** * Add a message to the conversation context */ addMessage( conversationId: string, message: Prompt.MessageEncoded ): Effect.Effect; /** * Add multiple messages to the conversation context */ addMessages( conversationId: string, messages: Prompt.MessageEncoded[] ): Effect.Effect; /** * Get conversation history */ getHistory(conversationId: string): Effect.Effect; /** * Update conversation metadata */ updateMetadata( conversationId: string, metadata: Partial ): Effect.Effect; /** * Clear conversation context */ clearContext(conversationId: string): Effect.Effect; /** * Clear and return whether context existed */ resetContext(conversationId: string): Effect.Effect; /** * Clear all conversation contexts */ clearAll(): Effect.Effect; /** * Set default policy for new contexts */ setDefaultPolicy(policy: ConversationPolicy): Effect.Effect; /** * Set a policy for an existing context */ setContextPolicy( conversationId: string, policy: ConversationPolicy ): Effect.Effect; /** * Replace the storage adapter at runtime (e.g. swap in-memory for SQLite/Postgres). */ replaceStorage(adapter: ContextStorage): Effect.Effect; /** * List session summaries from the currently active storage adapter. * Returns an empty array when no persistent adapter is configured * (the default in-memory storage is not queryable this way). */ listSessions(): Effect.Effect; } export const ContextStorageService = Context.GenericTag( 'ContextStorageService' ); /** * In-memory storage for ContextStorageService */ class InMemoryStorage { constructor(public contexts: Ref.Ref>) {} get(id: string): Effect.Effect { const self = this; return Effect.gen(function* () { const contexts = yield* Ref.get(self.contexts); return contexts.get(id) || null; }); } set(id: string, context: ConversationContext): Effect.Effect { const self = this; return Effect.gen(function* () { const contexts = yield* Ref.get(self.contexts); const newContexts = new Map(contexts); newContexts.set(id, context); yield* Ref.set(self.contexts, newContexts); }); } delete(id: string): Effect.Effect { const self = this; return Effect.gen(function* () { const contexts = yield* Ref.get(self.contexts); const newContexts = new Map(contexts); newContexts.delete(id); yield* Ref.set(self.contexts, newContexts); }); } clear(): Effect.Effect { return Ref.set(this.contexts, new Map()); } listSessions(): Effect.Effect { return Effect.succeed([]); } } /** * Wraps a Promise-based ContextStorage into the internal Effect-based storage interface. */ class ExternalStorageAdapter { constructor(private adapter: ContextStorage) {} get(id: string): Effect.Effect { return Effect.tryPromise({ try: () => this.adapter.get(id), catch: (error) => error, }).pipe(Effect.orElseSucceed(() => null)); } set(id: string, context: ConversationContext): Effect.Effect { return Effect.tryPromise({ try: () => this.adapter.set(id, context), catch: (error) => error, }).pipe(Effect.catchAll(() => Effect.void)); } delete(id: string): Effect.Effect { return Effect.tryPromise({ try: () => this.adapter.delete(id), catch: (error) => error, }).pipe(Effect.catchAll(() => Effect.void)); } clear(): Effect.Effect { return Effect.tryPromise({ try: () => this.adapter.clear(), catch: (error) => error, }).pipe(Effect.catchAll(() => Effect.void)); } listSessions(): Effect.Effect { return Effect.tryPromise({ try: () => this.adapter.listSessions(), catch: (error) => error, }).pipe(Effect.catchAll(() => Effect.succeed([]))); } } /** * Implementation of ContextStorageService */ class ContextStorageServiceImpl implements ContextStorageService { constructor( private storage: InMemoryStorage | ExternalStorageAdapter, private defaultMetadata: Ref.Ref>, private defaultPolicy: Ref.Ref ) {} generateConversationId(): Effect.Effect { return Effect.sync(() => `conv_${crypto.randomUUID()}`); } getContext( conversationId?: string, options?: { strict?: boolean } ): Effect.Effect { const self = this; return Effect.gen(function* () { const id = conversationId ?? (yield* self.generateConversationId()); const existing = yield* self.storage.get(id); const policy = yield* Ref.get(self.defaultPolicy); const strict = options?.strict ?? policy.strict; if (existing) { return existing; } if (strict && conversationId) { return yield* Effect.fail(new ContextNotFoundError({ conversationId, message: `Context not found for conversation: ${conversationId}`, })); } const metadata = yield* Ref.get(self.defaultMetadata); const newContext: ConversationContext = { id, messages: [], metadata: { createdAt: new Date(), updatedAt: new Date(), policy, ...metadata, }, }; yield* self.storage.set(id, newContext); return newContext; }); } getContextById(conversationId: string): Effect.Effect { return this.storage.get(conversationId); } addMessage( conversationId: string, message: Prompt.MessageEncoded ): Effect.Effect { const self = this; return Effect.gen(function* () { const normalized = normalizeMessage(message); if (normalized.role === 'system') { return; } const context = yield* self.getContext(conversationId); context.messages.push(normalized); self.applyCaps(context); context.metadata.updatedAt = new Date(); yield* self.storage.set(conversationId, context); }).pipe( Effect.mapError((cause) => self.toStorageError('addMessage', cause)) ); } addMessages( conversationId: string, messages: Prompt.MessageEncoded[] ): Effect.Effect { const self = this; return Effect.gen(function* () { const context = yield* self.getContext(conversationId); const filtered = normalizeMessages(messages).filter(m => m.role !== 'system'); if (filtered.length === 0) return; context.messages.push(...filtered); self.applyCaps(context); context.metadata.updatedAt = new Date(); yield* self.storage.set(conversationId, context); }).pipe( Effect.mapError((cause) => self.toStorageError('addMessages', cause)) ); } getHistory(conversationId: string): Effect.Effect { const self = this; return Effect.gen(function* () { const context = yield* self.getContext(conversationId); return normalizeMessages(context.messages); }).pipe(Effect.orDie); } updateMetadata( conversationId: string, metadata: Partial ): Effect.Effect { const self = this; return Effect.gen(function* () { const context = yield* self.getContext(conversationId); context.metadata = { ...context.metadata, ...metadata, updatedAt: new Date(), }; yield* self.storage.set(conversationId, context); }).pipe( Effect.mapError((cause) => self.toStorageError('updateMetadata', cause)) ); } clearContext(conversationId: string): Effect.Effect { return this.storage.delete(conversationId); } resetContext(conversationId: string): Effect.Effect { const self = this; return Effect.gen(function* () { const existing = yield* self.storage.get(conversationId); if (existing) { yield* self.storage.delete(conversationId); return true; } return false; }); } clearAll(): Effect.Effect { return this.storage.clear(); } setDefaultPolicy(policy: ConversationPolicy): Effect.Effect { return Ref.set(this.defaultPolicy, policy); } setContextPolicy( conversationId: string, policy: ConversationPolicy ): Effect.Effect { const self = this; return Effect.gen(function* () { const context = yield* self.getContext(conversationId); context.metadata.policy = { ...(context.metadata.policy ?? {}), ...policy, }; context.metadata.updatedAt = new Date(); self.applyCaps(context); yield* self.storage.set(conversationId, context); }).pipe( Effect.mapError((cause) => self.toStorageError('setContextPolicy', cause)) ); } replaceStorage(adapter: ContextStorage): Effect.Effect { return Effect.sync(() => { this.storage = new ExternalStorageAdapter(adapter); }); } listSessions(): Effect.Effect { return this.storage.listSessions(); } private toStorageError(operation: string, cause: unknown): ContextStorageError { return new ContextStorageError({ operation, message: this.storageErrorMessage(operation), cause, }); } private storageErrorMessage(operation: string): string { switch (operation) { case 'addMessage': return 'Failed to add message to context'; case 'addMessages': return 'Failed to add messages to context'; case 'updateMetadata': return 'Failed to update context metadata'; case 'setContextPolicy': return 'Failed to update context policy'; default: return 'Context storage operation failed'; } } private applyCaps(context: ConversationContext): void { const policy = context.metadata.policy; if (!policy) return; if (policy.maxMessages !== undefined && policy.maxMessages >= 0) { if (context.messages.length > policy.maxMessages) { context.messages.splice(0, context.messages.length - policy.maxMessages); } } if (policy.maxChars !== undefined && policy.maxChars >= 0) { let totalChars = context.messages.reduce((sum, msg) => sum + this.countMessageChars(msg), 0); while (context.messages.length > 0 && totalChars > policy.maxChars) { const removed = context.messages.shift(); if (removed) { totalChars -= this.countMessageChars(removed); } } } } private countMessageChars(message: Prompt.MessageEncoded): number { const content = message.content; if (typeof content === 'string') return content.length; if (content == null) return 0; if (Array.isArray(content)) { return content.reduce((sum, part) => { if (part && typeof part === 'object' && 'type' in part && part.type === 'text') { const text = (part as { text?: string }).text; return sum + (typeof text === 'string' ? text.length : 0); } return sum + JSON.stringify(part).length; }, 0); } return JSON.stringify(content).length; } } const makeContextStorageService = (adapter?: ContextStorage) => Effect.gen(function* () { const contexts = yield* Ref.make(new Map()); const storage = adapter ? new ExternalStorageAdapter(adapter) : new InMemoryStorage(contexts); const defaultMetadata = yield* Ref.make>({}); const defaultPolicy = yield* Ref.make({}); return new ContextStorageServiceImpl(storage, defaultMetadata, defaultPolicy); }); /** Live layer providing ContextStorageService (in-memory). */ export const ContextStorageServiceLive = Layer.effect( ContextStorageService, makeContextStorageService(), ); /** Live layer backed by an explicitly owned external adapter. */ export const ContextStorageServiceLiveWithAdapter = (adapter: ContextStorage) => Layer.effect( ContextStorageService, makeContextStorageService(adapter), );