From 3e9fbfa5effcd80d255a6f2a618ff8f248d60a76 Mon Sep 17 00:00:00 2001 From: Artur Do Lago Date: Sun, 11 Jan 2026 19:54:03 +0100 Subject: [PATCH] refactor(memory): unify memory abstraction into single Memory class Consolidate 4-layer memory abstraction into unified Memory class: - MemoryStore, QdrantMemoryStore, QdrantMemoryBridge, ContinuityManager are now consolidated into a single src/memory/unified.ts Single Qdrant collection with type field for discrimination: - type: "memory" - facts, preferences, decisions - type: "state" - personas orchestration state - type: "conversation" - conversation continuity - type: "session_chain" - session chain index Key changes: - New Memory class with save/search/get/delete for memories - State persistence (saveState/loadState) - Conversation continuity (startSession, processMessages, etc.) - Cross-session memory injection for personas bootstrap - Updated tiara.ts to use Memory instead of QdrantMemoryBridge - Added bootstrap/personas.ts for lifecycle hooks Legacy APIs preserved for backwards compatibility. Co-Authored-By: Claude Opus 4.5 --- packages/agent-core/src/bootstrap/personas.ts | 222 +++++ src/memory/index.ts | 30 +- src/memory/unified.ts | 928 ++++++++++++++++++ src/personas/tiara.ts | 39 +- 4 files changed, 1204 insertions(+), 15 deletions(-) create mode 100644 packages/agent-core/src/bootstrap/personas.ts create mode 100644 src/memory/unified.ts diff --git a/packages/agent-core/src/bootstrap/personas.ts b/packages/agent-core/src/bootstrap/personas.ts new file mode 100644 index 00000000000..28fa5fbd23d --- /dev/null +++ b/packages/agent-core/src/bootstrap/personas.ts @@ -0,0 +1,222 @@ +/** + * Personas Bootstrap + * + * Initializes persona-specific hooks and services. + * Called by daemon on startup to enable cross-session memory and fact extraction. + */ + +import { LifecycleHooks } from "../hooks/lifecycle" +import { Persistence } from "../session/persistence" +import { Session } from "../session" +import { Log } from "../util/log" + +const log = Log.create({ service: "personas-bootstrap" }) + +// Track initialized state +let initialized = false + +// Memory instance (lazy loaded to avoid import issues) +let memoryPromise: Promise | null = null + +/** + * Get the unified Memory instance (lazy load) + */ +async function getMemoryInstance(): Promise { + if (!memoryPromise) { + memoryPromise = (async () => { + try { + // Dynamic import to avoid build-time dependency on src/memory + const memoryPath = "../../../../src/memory/unified.js" + const memoryModule = await import(memoryPath) + return memoryModule.getMemory() + } catch (e) { + log.debug("Unified Memory not available", { + error: e instanceof Error ? e.message : String(e), + }) + return null + } + })() + } + return memoryPromise +} + +/** + * Initialize persona hooks and services + */ +export async function initPersonas(): Promise { + if (initialized) { + log.debug("Personas already initialized") + return + } + + log.info("Initializing persona hooks") + + // Pre-initialize unified Memory + const memory = await getMemoryInstance() + if (memory) { + log.info("Unified Memory connected for cross-session context") + } + + // Register session start hook for memory injection + LifecycleHooks.on( + LifecycleHooks.SessionLifecycle.Start, + async (payload) => { + if (payload.source === "whatsapp" || payload.source === "telegram") { + await injectCrossSessionMemory(payload.sessionId, payload.persona) + } + } + ) + + // Register session restore hook for memory injection + LifecycleHooks.on( + LifecycleHooks.SessionLifecycle.Restore, + async (payload) => { + if (payload.source === "whatsapp" || payload.source === "telegram") { + await injectCrossSessionMemory(payload.sessionId, payload.persona) + } + } + ) + + // Try to initialize external persona hooks (fact extraction, etc.) + // These are in src/personas/hooks/ which may not be available in all builds + try { + // Dynamic import to avoid build-time dependency + const hooksPath = "../../../../src/personas/hooks/index.js" + const hooks = await import(hooksPath) + if (hooks.initFactExtractionHook) { + hooks.initFactExtractionHook() + log.info("Fact extraction hook initialized") + } + } catch (e) { + log.debug("External persona hooks not available", { + error: e instanceof Error ? e.message : String(e), + }) + } + + initialized = true + log.info("Personas initialized") +} + +/** + * Inject relevant memories from previous sessions + */ +async function injectCrossSessionMemory( + sessionId: string, + persona: "zee" | "stanley" | "johny" +): Promise { + try { + const memory = await getMemoryInstance() + const memories: string[] = [] + + // 1. Get yesterday's session for recent context + const yesterday = new Date() + yesterday.setDate(yesterday.getDate() - 1) + const yesterdaySession = await Persistence.getDailySession(persona, yesterday) + + if (yesterdaySession) { + const previousSession = await Session.get(yesterdaySession.sessionId) + if (previousSession) { + log.info("Found previous session for context", { + persona, + previousSessionId: yesterdaySession.sessionId, + previousDate: yesterday.toISOString().split("T")[0], + }) + memories.push(`[Previous session from ${yesterday.toISOString().split("T")[0]}]`) + } + } + + // 2. Search for relevant memories from Qdrant (semantic search) + if (memory) { + try { + // Search for recent facts and important information using unified Memory API + const personaContext = getPersonaSearchContext(persona) + const searchResults = await memory.searchPersonaMemories( + personaContext, + persona, + { limit: 5, categories: ["fact", "preference", "decision"] } + ) + + if (searchResults.length > 0) { + log.info("Found relevant memories", { + persona, + count: searchResults.length, + }) + + for (const result of searchResults) { + const entry = result.entry + const score = result.score.toFixed(2) + memories.push(`[Memory (${entry.category}, relevance: ${score})]: ${entry.content}`) + } + } + } catch (memoryError) { + log.debug("Memory search failed", { + error: memoryError instanceof Error ? memoryError.message : String(memoryError), + }) + } + } + + // 3. Store context for session (to be used by prompt injection) + if (memories.length > 0) { + // Store the memory context in session metadata for prompt injection + // This will be picked up by the prompt builder + await storeSessionContext(sessionId, memories) + log.info("Injected cross-session context", { + sessionId: sessionId.slice(0, 8), + memoriesCount: memories.length, + }) + } + + } catch (error) { + log.debug("Could not inject cross-session memory", { + sessionId: sessionId.slice(0, 8), + error: error instanceof Error ? error.message : String(error), + }) + } +} + +/** + * Get persona-specific search context for memory retrieval + */ +function getPersonaSearchContext(persona: "zee" | "stanley" | "johny"): string { + switch (persona) { + case "zee": + return "personal assistant tasks reminders calendar schedule preferences contacts" + case "stanley": + return "portfolio investments stocks markets trading analysis positions" + case "johny": + return "learning study knowledge practice concepts understanding progress" + } +} + +/** + * Store cross-session context for prompt injection + * Saved as a special system message in the session + */ +async function storeSessionContext(sessionId: string, memories: string[]): Promise { + // Store as session metadata for prompt injection + // The prompt builder will look for this and inject it + try { + const contextKey = `persona_context_${sessionId}` + const context = { + timestamp: Date.now(), + memories, + } + + // Store in persistence for retrieval by prompt builder + await Persistence.setSessionContext(sessionId, context) + } catch (error) { + log.debug("Could not store session context", { + sessionId: sessionId.slice(0, 8), + error: error instanceof Error ? error.message : String(error), + }) + } +} + +/** + * Cleanup persona hooks + */ +export function cleanupPersonas(): void { + // Currently hooks clean up automatically via garbage collection + initialized = false + log.info("Personas cleanup complete") +} diff --git a/src/memory/index.ts b/src/memory/index.ts index a25ee0f3e58..d7bb59a104b 100644 --- a/src/memory/index.ts +++ b/src/memory/index.ts @@ -3,9 +3,35 @@ * * Unified memory system with Qdrant vector storage, * supporting semantic search, pattern storage, and cross-session context. + * + * Primary API: + * - Memory class (unified.ts) - Single class for all memory operations + * - getMemory() - Get shared Memory instance + * + * Legacy exports (for backwards compatibility): + * - QdrantVectorStorage - Low-level Qdrant client + * - MemoryStore - Legacy store API + * - createEmbeddingProvider - Embedding factory */ +// Primary API - unified Memory class +export { Memory, getMemory, resetMemory, extractKeyFacts } from "./unified"; +export type { + MemoryConfig, + PersonaId, + ConversationState, + PersonasState, + EntryType, +} from "./unified"; + +// Types export * from "./types"; + +// Embedding providers export * from "./embedding"; -export * from "./qdrant"; -export * from "./store"; + +// Low-level storage (for advanced use cases) +export { QdrantVectorStorage, QdrantMemoryStore } from "./qdrant"; + +// Legacy store API (deprecated - use Memory class instead) +export { MemoryStore, getMemoryStore, resetMemoryStore } from "./store"; diff --git a/src/memory/unified.ts b/src/memory/unified.ts new file mode 100644 index 00000000000..6a94b84b6bf --- /dev/null +++ b/src/memory/unified.ts @@ -0,0 +1,928 @@ +/** + * Unified Memory System + * + * Single class that handles all memory operations: + * - Semantic memory storage and search + * - Persona state persistence + * - Conversation continuity (fact extraction, session chaining) + * - Cross-session context injection + * + * Uses a single Qdrant collection with `type` field for discrimination. + */ + +import { randomUUID } from "node:crypto"; +import { QdrantVectorStorage } from "./qdrant"; +import { createEmbeddingProvider, type EmbeddingConfig } from "./embedding"; +import type { + MemoryEntry, + MemoryInput, + MemorySearchParams, + MemorySearchResult, + MemoryCategory, + EmbeddingProvider, +} from "./types"; +import { + QDRANT_URL, + QDRANT_COLLECTION_MEMORY, + CONTINUITY_MAX_KEY_FACTS, +} from "../config/constants"; +import { Log } from "../../packages/agent-core/src/util/log"; + +const log = Log.create({ service: "memory" }); + +// ============================================================================= +// Types +// ============================================================================= + +/** Entry types stored in unified collection */ +export type EntryType = + | "memory" // Regular memories (facts, preferences, etc.) + | "state" // Personas orchestration state + | "conversation" // Conversation continuity state + | "session_chain"; // Session chain index + +/** Persona identifiers */ +export type PersonaId = "zee" | "stanley" | "johny"; + +/** Conversation state for continuity */ +export interface ConversationState { + sessionId: string; + leadPersona: PersonaId; + summary: string; + plan: string; + objectives: string[]; + keyFacts: string[]; + sessionChain: string[]; + updatedAt: number; +} + +/** Personas orchestration state */ +export interface PersonasState { + version: string; + tiaraSwarmId?: string; + workers: Array<{ + id: string; + persona: PersonaId; + role: "queen" | "drone"; + status: string; + paneId?: string; + pid?: number; + currentTask?: string; + createdAt: number; + lastActivityAt: number; + }>; + tasks: Array<{ + id: string; + persona: PersonaId; + description: string; + prompt: string; + status: "pending" | "assigned" | "running" | "completed" | "failed"; + priority?: "low" | "normal" | "high" | "critical"; + workerId?: string; + createdAt: number; + completedAt?: number; + result?: string; + error?: string; + }>; + conversation?: ConversationState; + lastSyncAt: number; + stats: { + totalTasksCompleted: number; + totalDronesSpawned: number; + totalTokensUsed: number; + }; +} + +/** Memory configuration */ +export interface MemoryConfig { + qdrant: { + url?: string; + apiKey?: string; + collection?: string; + }; + embedding: EmbeddingConfig; + namespace?: string; + maxKeyFacts?: number; +} + +// ============================================================================= +// Mock Embedding Provider +// ============================================================================= + +/** + * Mock embedding provider for testing when no API key is available. + */ +class MockEmbeddingProvider implements EmbeddingProvider { + readonly id = "mock"; + readonly model = "mock-embedding"; + readonly dimension = 384; + + async embed(text: string): Promise { + const vector: number[] = new Array(this.dimension).fill(0); + for (let i = 0; i < text.length && i < this.dimension; i++) { + vector[i] = (text.charCodeAt(i) % 100) / 100; + } + const mag = Math.sqrt(vector.reduce((s, v) => s + v * v, 0)); + return vector.map((v) => v / (mag || 1)); + } + + async embedBatch(texts: string[]): Promise { + return Promise.all(texts.map((t) => this.embed(t))); + } +} + +// ============================================================================= +// Utilities +// ============================================================================= + +/** Generate deterministic UUID from string */ +function stringToUUID(str: string): string { + let hash = 0; + for (let i = 0; i < str.length; i++) { + const char = str.charCodeAt(i); + hash = ((hash << 5) - hash) + char; + hash = hash & hash; + } + const hex = Math.abs(hash).toString(16).padStart(8, "0"); + return `${hex.slice(0, 8)}-${hex.slice(0, 4)}-4${hex.slice(1, 4)}-8${hex.slice(0, 3)}-${hex.padEnd(12, "0").slice(0, 12)}`; +} + +/** Generate stable instance ID */ +function generateInstanceId(): string { + // eslint-disable-next-line @typescript-eslint/no-require-imports + const os = require("os"); + const hostname = os.hostname() || "unknown"; + const username = os.userInfo().username || "user"; + return stringToUUID(`memory-${hostname}-${username}`); +} + +/** Extract key facts from text (simple heuristics, use LLM in production) */ +export function extractKeyFacts(message: string): string[] { + const facts: string[] = []; + const sentences = message.split(/[.!?]+/).filter((s) => s.trim().length > 20); + + for (const sentence of sentences) { + const s = sentence.trim().toLowerCase(); + + // Fact-like patterns + if ( + s.includes("is ") || s.includes("are ") || s.includes("was ") || + s.includes("were ") || s.includes("has ") || s.includes("have ") || + s.includes("prefers ") || s.includes("wants ") || s.includes("needs ") || + s.includes("decided ") || s.includes("agreed ") + ) { + facts.push(sentence.trim()); + } + + // Preferences + if ( + s.includes("i like ") || s.includes("i prefer ") || + s.includes("i want ") || s.includes("i need ") + ) { + facts.push(sentence.trim()); + } + + // Decisions + if ( + s.includes("we should ") || s.includes("we will ") || + s.includes("let's ") || s.includes("the plan is ") + ) { + facts.push(sentence.trim()); + } + } + + return Array.from(new Set(facts)).slice(0, 20); +} + +/** Generate summary from messages */ +function generateSummary(messages: string[]): string { + if (messages.length === 0) return ""; + + const recentMessages = messages.slice(-10); + const parts = [ + "## Conversation Summary", + "", + `**Messages:** ${messages.length} total`, + "", + "### Recent Exchange:", + "", + ]; + + for (const msg of recentMessages) { + const truncated = msg.length > 200 ? msg.slice(0, 200) + "..." : msg; + parts.push(`- ${truncated}`); + } + + return parts.join("\n"); +} + +/** Merge facts with deduplication */ +function mergeFacts(existing: string[], newFacts: string[], max: number): string[] { + const allFacts = [...existing, ...newFacts]; + const seen: Record = {}; + const unique: string[] = []; + + for (const fact of allFacts) { + const normalized = fact.toLowerCase().trim(); + if (!seen[normalized]) { + seen[normalized] = true; + unique.push(fact); + } + } + + return unique.slice(-max); +} + +// ============================================================================= +// Unified Memory Class +// ============================================================================= + +/** + * Unified Memory - single class for all memory operations. + * + * Replaces: + * - MemoryStore (store.ts) + * - QdrantMemoryStore (qdrant.ts) + * - QdrantMemoryBridge (memory-bridge.ts) + * - ContinuityManager (continuity.ts) + */ +export class Memory { + private readonly storage: QdrantVectorStorage; + private readonly embedding: EmbeddingProvider; + private readonly namespace: string; + private readonly collection: string; + private readonly instanceId: string; + private readonly maxKeyFacts: number; + private initialized = false; + + // Current conversation state (for continuity) + private currentConversation?: ConversationState; + + constructor(config: Partial = {}) { + const qdrantConfig = { + url: config.qdrant?.url ?? process.env.QDRANT_URL ?? QDRANT_URL, + apiKey: config.qdrant?.apiKey, + collection: config.qdrant?.collection ?? QDRANT_COLLECTION_MEMORY, + }; + + this.collection = qdrantConfig.collection; + this.storage = new QdrantVectorStorage(qdrantConfig); + this.namespace = config.namespace ?? "default"; + this.instanceId = generateInstanceId(); + this.maxKeyFacts = config.maxKeyFacts ?? CONTINUITY_MAX_KEY_FACTS; + + // Use mock embeddings if no API key available + const usesMock = !process.env.OPENAI_API_KEY && !config.embedding?.apiKey; + if (usesMock) { + this.embedding = new MockEmbeddingProvider(); + log.debug("Using mock embeddings (no API key)"); + } else { + this.embedding = createEmbeddingProvider({ + provider: config.embedding?.provider ?? "openai", + model: config.embedding?.model, + dimensions: config.embedding?.dimensions, + apiKey: config.embedding?.apiKey, + baseUrl: config.embedding?.baseUrl, + }); + } + } + + // =========================================================================== + // Initialization + // =========================================================================== + + /** Initialize the memory store */ + async init(): Promise { + if (this.initialized) return; + + await this.storage.createCollection(this.collection, this.embedding.dimension); + this.storage.setCollection(this.collection); + this.initialized = true; + + log.info("Memory initialized", { + collection: this.collection, + namespace: this.namespace, + dimension: this.embedding.dimension, + }); + } + + // =========================================================================== + // Memory Operations (facts, preferences, etc.) + // =========================================================================== + + /** Save a memory entry */ + async save(input: MemoryInput): Promise { + await this.init(); + + const id = randomUUID(); + const now = Date.now(); + const vector = await this.embedding.embed(input.content); + + const entry: MemoryEntry = { + id, + category: input.category, + content: input.content, + summary: input.summary, + embedding: vector, + metadata: input.metadata ?? {}, + createdAt: now, + accessedAt: now, + ttl: input.ttl, + namespace: input.namespace ?? this.namespace, + }; + + await this.storage.insert([{ + id, + vector, + payload: { + type: "memory" as EntryType, + category: entry.category, + content: entry.content, + summary: entry.summary, + metadata: entry.metadata, + createdAt: entry.createdAt, + accessedAt: entry.accessedAt, + ttl: entry.ttl, + namespace: entry.namespace, + }, + }]); + + return entry; + } + + /** Search memories semantically */ + async search(params: MemorySearchParams): Promise { + await this.init(); + + const queryVector = await this.embedding.embed(params.query); + + // Build filter + const filter: Record = { + type: "memory", + namespace: params.namespace ?? this.namespace, + }; + + if (params.category) { + if (Array.isArray(params.category)) { + filter.category = { $in: params.category }; + } else { + filter.category = params.category; + } + } + + if (params.tags?.length) { + filter["metadata.tags"] = { $in: params.tags }; + } + + const results = await this.storage.search(queryVector, { + limit: params.limit ?? 10, + threshold: params.threshold ?? 0.5, + filter, + }); + + return results.map((r) => ({ + entry: { + id: r.id, + category: r.payload.category as MemoryCategory, + content: r.payload.content as string, + summary: r.payload.summary as string | undefined, + metadata: r.payload.metadata as MemoryEntry["metadata"], + createdAt: r.payload.createdAt as number, + accessedAt: r.payload.accessedAt as number, + ttl: r.payload.ttl as number | undefined, + namespace: r.payload.namespace as string | undefined, + }, + score: r.score, + })); + } + + /** Get a memory by ID */ + async get(id: string): Promise { + await this.init(); + + const results = await this.storage.get([id]); + const point = results[0]; + if (!point || point.payload.type !== "memory") return null; + + return { + id: point.id, + category: point.payload.category as MemoryCategory, + content: point.payload.content as string, + summary: point.payload.summary as string | undefined, + embedding: point.vector, + metadata: point.payload.metadata as MemoryEntry["metadata"], + createdAt: point.payload.createdAt as number, + accessedAt: point.payload.accessedAt as number, + ttl: point.payload.ttl as number | undefined, + namespace: point.payload.namespace as string | undefined, + }; + } + + /** Delete a memory by ID */ + async delete(id: string): Promise { + await this.init(); + await this.storage.delete([id]); + } + + /** Delete memories matching filter */ + async deleteWhere(filter: { + category?: MemoryCategory; + namespace?: string; + olderThan?: number; + }): Promise { + await this.init(); + + const qdrantFilter: Record = { type: "memory" }; + if (filter.category) qdrantFilter.category = filter.category; + if (filter.namespace) qdrantFilter.namespace = filter.namespace; + if (filter.olderThan) qdrantFilter.createdAt = { $lt: filter.olderThan }; + + return this.storage.deleteWhere(qdrantFilter); + } + + /** Delete expired memories */ + async deleteExpired(): Promise { + await this.init(); + const now = Date.now(); + return this.storage.deleteWhere({ + type: "memory", + expiresAt: { lt: now, gt: 0 }, + }); + } + + // =========================================================================== + // State Persistence (Personas orchestration state) + // =========================================================================== + + private getStateId(): string { + return stringToUUID(`state-${this.instanceId}`); + } + + /** Save personas state */ + async saveState(state: PersonasState): Promise { + await this.init(); + + const stateJson = JSON.stringify(state); + const stateId = this.getStateId(); + + // Generate summary for embedding + const summary = this.generateStateSummary(state); + const embedding = await this.embedding.embed(summary); + + await this.storage.insert([{ + id: stateId, + vector: embedding, + payload: { + type: "state" as EntryType, + state: stateJson, + summary, + version: state.version, + updatedAt: Date.now(), + }, + }]); + + // Also save conversation state if present + if (state.conversation) { + await this.saveConversation(state.conversation); + } + } + + /** Load personas state */ + async loadState(): Promise { + await this.init(); + + const stateId = this.getStateId(); + const results = await this.storage.get([stateId]); + const result = results[0]; + + if (!result?.payload?.state) return null; + + try { + return JSON.parse(result.payload.state as string) as PersonasState; + } catch { + return null; + } + } + + private generateStateSummary(state: PersonasState): string { + const parts: string[] = [`Personas state v${state.version}`]; + + if (state.workers.length > 0) { + const workersByPersona = state.workers.reduce( + (acc, w) => { + acc[w.persona] = (acc[w.persona] ?? 0) + 1; + return acc; + }, + {} as Record + ); + parts.push(`Workers: ${JSON.stringify(workersByPersona)}`); + } + + if (state.tasks.length > 0) { + const pending = state.tasks.filter((t) => t.status === "pending").length; + const running = state.tasks.filter((t) => t.status === "running").length; + parts.push(`Tasks: ${pending} pending, ${running} running`); + } + + if (state.conversation) { + parts.push(`Lead: ${state.conversation.leadPersona}`); + parts.push(`Summary: ${state.conversation.summary.slice(0, 200)}`); + } + + parts.push(`Stats: ${state.stats.totalTasksCompleted} tasks completed`); + + return parts.join("\n"); + } + + // =========================================================================== + // Conversation Continuity + // =========================================================================== + + /** Save conversation state */ + async saveConversation(state: ConversationState): Promise { + await this.init(); + + const conversationId = stringToUUID(`conversation-${state.sessionId}`); + + // Create rich summary for embedding + const summaryParts = [ + state.summary, + `Plan: ${state.plan}`, + `Objectives: ${state.objectives.join(", ")}`, + `Key facts: ${state.keyFacts.join("; ")}`, + ]; + const fullSummary = summaryParts.filter(Boolean).join("\n"); + const embedding = await this.embedding.embed(fullSummary); + + await this.storage.insert([{ + id: conversationId, + vector: embedding, + payload: { + type: "conversation" as EntryType, + sessionId: state.sessionId, + leadPersona: state.leadPersona, + summary: state.summary, + plan: state.plan, + objectives: state.objectives, + keyFacts: state.keyFacts, + sessionChain: state.sessionChain, + updatedAt: state.updatedAt, + }, + }]); + + // Also store session chain index + const chainId = stringToUUID(`session-chain-${state.sessionId}`); + await this.storage.insert([{ + id: chainId, + vector: embedding, + payload: { + type: "session_chain" as EntryType, + sessionId: state.sessionId, + previousSessions: state.sessionChain, + updatedAt: state.updatedAt, + }, + }]); + } + + /** Load conversation state by session ID */ + async loadConversation(sessionId: string): Promise { + await this.init(); + + const conversationId = stringToUUID(`conversation-${sessionId}`); + const results = await this.storage.get([conversationId]); + const result = results[0]; + + if (!result?.payload || result.payload.type !== "conversation") return null; + + const p = result.payload as Record; + return { + sessionId: p.sessionId as string, + leadPersona: p.leadPersona as PersonaId, + summary: p.summary as string, + plan: (p.plan as string) ?? "", + objectives: (p.objectives as string[]) ?? [], + keyFacts: (p.keyFacts as string[]) ?? [], + sessionChain: (p.sessionChain as string[]) ?? [], + updatedAt: p.updatedAt as number, + }; + } + + /** Find most recent conversation (optionally for specific persona) */ + async findRecentConversation(persona?: PersonaId): Promise { + await this.init(); + + const query = persona + ? `Recent conversation with ${persona}` + : "Recent conversation state"; + + const embedding = await this.embedding.embed(query); + const filter: Record = { type: "conversation" }; + if (persona) filter.leadPersona = persona; + + const results = await this.storage.search(embedding, { + limit: 1, + filter, + }); + + if (results.length === 0) return null; + + const p = results[0].payload as Record; + return { + sessionId: p.sessionId as string, + leadPersona: p.leadPersona as PersonaId, + summary: p.summary as string, + plan: (p.plan as string) ?? "", + objectives: (p.objectives as string[]) ?? [], + keyFacts: (p.keyFacts as string[]) ?? [], + sessionChain: (p.sessionChain as string[]) ?? [], + updatedAt: p.updatedAt as number, + }; + } + + /** Start a new conversation session (with continuity from previous) */ + async startSession( + sessionId: string, + leadPersona: PersonaId, + previousSessionId?: string + ): Promise { + // Try to load previous session + let previousState: ConversationState | null = null; + if (previousSessionId) { + previousState = await this.loadConversation(previousSessionId); + } else { + previousState = await this.findRecentConversation(leadPersona); + } + + // Create new state + this.currentConversation = { + sessionId, + leadPersona, + summary: "", + plan: previousState?.plan ?? "", + objectives: previousState?.objectives ?? [], + keyFacts: previousState?.keyFacts.slice(-this.maxKeyFacts) ?? [], + sessionChain: previousState + ? [...previousState.sessionChain, previousState.sessionId] + : [], + updatedAt: Date.now(), + }; + + await this.saveConversation(this.currentConversation); + return this.currentConversation; + } + + /** Get current conversation state */ + getCurrentConversation(): ConversationState | undefined { + return this.currentConversation; + } + + /** Process messages and extract facts */ + async processMessages(messages: string[]): Promise { + if (!this.currentConversation) { + throw new Error("No active session. Call startSession first."); + } + + // Extract facts from new messages + const newFacts: string[] = []; + for (const msg of messages) { + newFacts.push(...extractKeyFacts(msg)); + } + + // Update state + this.currentConversation = { + ...this.currentConversation, + summary: generateSummary(messages), + keyFacts: mergeFacts( + this.currentConversation.keyFacts, + newFacts, + this.maxKeyFacts + ), + updatedAt: Date.now(), + }; + + // Save to Qdrant + await this.saveConversation(this.currentConversation); + + // Store individual facts as memories (persona-isolated) + if (newFacts.length > 0) { + await this.storeKeyFacts( + newFacts, + this.currentConversation.sessionId, + this.currentConversation.leadPersona + ); + } + + return this.currentConversation; + } + + /** Store key facts as searchable memories */ + async storeKeyFacts(facts: string[], sessionId: string, persona: PersonaId): Promise { + await this.init(); + + for (const fact of facts) { + await this.save({ + category: "fact", + content: fact, + metadata: { + sessionId, + agent: persona, + extra: { extractedAt: Date.now() }, + }, + namespace: `personas:${persona}`, + }); + } + } + + /** Update plan */ + async updatePlan(plan: string): Promise { + if (!this.currentConversation) { + throw new Error("No active session"); + } + + this.currentConversation.plan = plan; + this.currentConversation.updatedAt = Date.now(); + await this.saveConversation(this.currentConversation); + } + + /** Add objective */ + async addObjective(objective: string): Promise { + if (!this.currentConversation) { + throw new Error("No active session"); + } + + this.currentConversation.objectives.push(objective); + this.currentConversation.updatedAt = Date.now(); + await this.saveConversation(this.currentConversation); + } + + /** End session */ + async endSession(): Promise { + if (!this.currentConversation) return; + + await this.saveConversation(this.currentConversation); + this.currentConversation = undefined; + } + + /** Format conversation state for prompt injection */ + formatContextForPrompt(): string { + if (!this.currentConversation) return ""; + + const state = this.currentConversation; + const parts: string[] = ["# Conversation Context (Restored)", ""]; + + if (state.summary) { + parts.push("## Previous Conversation Summary"); + parts.push(state.summary); + parts.push(""); + } + + if (state.plan) { + parts.push("## Current Plan"); + parts.push(state.plan); + parts.push(""); + } + + if (state.objectives.length > 0) { + parts.push("## Active Objectives"); + state.objectives.forEach((obj, i) => { + parts.push(`${i + 1}. ${obj}`); + }); + parts.push(""); + } + + if (state.keyFacts.length > 0) { + parts.push("## Key Facts"); + state.keyFacts.forEach((fact) => { + parts.push(`- ${fact}`); + }); + parts.push(""); + } + + if (state.sessionChain.length > 0) { + parts.push(`_This is session ${state.sessionChain.length + 1} in a continuing conversation._`); + } + + return parts.join("\n"); + } + + // =========================================================================== + // Cross-Session Memory Injection (for bootstrap/personas.ts) + // =========================================================================== + + /** Search memories for a specific persona */ + async searchPersonaMemories( + query: string, + persona: PersonaId, + options?: { limit?: number; categories?: MemoryCategory[] } + ): Promise { + return this.search({ + query, + namespace: `personas:${persona}`, + category: options?.categories, + limit: options?.limit ?? 5, + threshold: 0.6, + }); + } + + /** Get relevant context for a task */ + async getTaskContext( + taskDescription: string, + options?: { + limit?: number; + sessionId?: string; + persona?: PersonaId; + } + ): Promise<{ + relevantMemories: Array<{ content: string; score: number }>; + conversationState?: ConversationState; + }> { + await this.init(); + + const limit = options?.limit ?? 5; + + // Search for relevant memories + const results = await this.searchPersonaMemories( + taskDescription, + options?.persona ?? "zee", + { limit } + ); + + // Load conversation state + let conversationState: ConversationState | undefined; + if (options?.sessionId) { + const state = await this.loadConversation(options.sessionId); + if (state) conversationState = state; + } else { + const recent = await this.findRecentConversation(options?.persona); + if (recent) conversationState = recent; + } + + return { + relevantMemories: results.map((r) => ({ + content: r.entry.content, + score: r.score, + })), + conversationState, + }; + } + + // =========================================================================== + // Statistics + // =========================================================================== + + /** Get memory statistics */ + async stats(): Promise<{ + total: number; + byType: Record; + byCategory: Record; + }> { + await this.init(); + + const total = await this.storage.count(); + + const types: EntryType[] = ["memory", "state", "conversation", "session_chain"]; + const byType: Record = {}; + for (const type of types) { + byType[type] = await this.storage.count({ type }); + } + + const categories: MemoryCategory[] = [ + "conversation", "fact", "preference", "task", + "decision", "relationship", "note", "pattern", + ]; + const byCategory: Record = {}; + for (const cat of categories) { + byCategory[cat] = await this.storage.count({ type: "memory", category: cat }); + } + + return { + total, + byType: byType as Record, + byCategory, + }; + } + + /** Cleanup old entries */ + async cleanup(): Promise { + return this.deleteExpired(); + } +} + +// ============================================================================= +// Singleton +// ============================================================================= + +let _instance: Memory | null = null; + +/** Get the shared Memory instance */ +export function getMemory(config?: Partial): Memory { + if (!_instance) { + _instance = new Memory(config); + } + return _instance; +} + +/** Reset the shared instance (for testing) */ +export function resetMemory(): void { + _instance = null; +} diff --git a/src/personas/tiara.ts b/src/personas/tiara.ts index 0e00dd113cd..1fb0efc31e7 100644 --- a/src/personas/tiara.ts +++ b/src/personas/tiara.ts @@ -13,7 +13,7 @@ import { spawn, type ChildProcess } from "node:child_process"; import { EventEmitter } from "node:events"; import type { PersonasOrchestrator, - PersonasState, + PersonasState as TypesPersonasState, PersonasConfig, PersonasTask, PersonasEvent, @@ -29,11 +29,10 @@ import type { TaskStatus, ConversationState, } from "./types"; -import { QdrantMemoryBridge, createMemoryBridge } from "./memory-bridge"; +import { Memory, getMemory, type PersonasState } from "../memory/unified"; import { WeztermPaneBridge, createWeztermBridge } from "./wezterm"; import { generateDronePrompt } from "./persona"; import { formatAnnouncement, getDroneWaiter, shouldAnnounce, shutdownDroneWaiter } from "./drone-wait"; -import { createConversationState } from "./continuity"; import { Log } from "../../packages/agent-core/src/util/log"; import { QDRANT_URL, @@ -92,7 +91,7 @@ function generateTaskId(): string { */ export class Orchestrator extends EventEmitter implements PersonasOrchestrator { private config: PersonasConfig; - private memoryBridge: QdrantMemoryBridge; + private memory: Memory; private weztermBridge: WeztermPaneBridge; private currentState: PersonasState; private processes = new Map(); @@ -105,7 +104,15 @@ export class Orchestrator extends EventEmitter implements PersonasOrchestrator { constructor(config?: Partial) { super(); this.config = { ...DEFAULT_CONFIG, ...config }; - this.memoryBridge = createMemoryBridge(this.config.qdrant); + // Use unified Memory class instead of separate memory bridge + this.memory = getMemory({ + qdrant: { + url: this.config.qdrant.url, + apiKey: this.config.qdrant.apiKey, + collection: this.config.qdrant.stateCollection, + }, + maxKeyFacts: this.config.continuity.maxKeyFacts, + }); this.weztermBridge = createWeztermBridge(this.config.wezterm); this.currentState = this.createInitialState(); } @@ -136,10 +143,10 @@ export class Orchestrator extends EventEmitter implements PersonasOrchestrator { const leadPersona = this.resolveLeadPersona(); // Initialize memory bridge - await this.memoryBridge.init(); + await this.memory.init(); // Try to load existing state - const existingState = await this.memoryBridge.loadState(); + const existingState = await this.memory.loadState(); if (existingState) { this.currentState = existingState; // Clean up any stale workers @@ -150,10 +157,16 @@ export class Orchestrator extends EventEmitter implements PersonasOrchestrator { // Ensure a lead persona is always present if (!this.currentState.conversation) { - this.currentState.conversation = createConversationState( - `session-${Date.now()}`, - leadPersona - ); + this.currentState.conversation = { + sessionId: `session-${Date.now()}`, + leadPersona, + summary: "", + plan: "", + objectives: [], + keyFacts: [], + sessionChain: [], + updatedAt: Date.now(), + }; } this.ensureQueenPresence(leadPersona); @@ -650,7 +663,7 @@ export class Orchestrator extends EventEmitter implements PersonasOrchestrator { * Restore context from a previous session */ async restoreContext(sessionId: string): Promise { - const state = await this.memoryBridge.loadConversationState(sessionId); + const state = await this.memory.loadConversation(sessionId); if (state) { this.currentState.conversation = state; this.emitEvent("continuity:restored", { sessionId }); @@ -663,7 +676,7 @@ export class Orchestrator extends EventEmitter implements PersonasOrchestrator { */ async saveState(): Promise { this.currentState.lastSyncAt = Date.now(); - await this.memoryBridge.saveState(this.currentState); + await this.memory.saveState(this.currentState); this.emitEvent("state:synced", {}); }