import type { EmbeddingModel } from "ai"; import type { ObservationCompressor } from "../ai/compressor"; import type { ConflictEvaluator } from "../ai/conflict-evaluator"; import type { EntityExtractor } from "../ai/entity-extractor"; import type { SessionSummarizer } from "../ai/summarizer"; import type { EntityRepository } from "../db/entities"; import type { ObservationRepository } from "../db/observations"; import type { PendingMessageRepository } from "../db/pending"; import type { SessionRepository } from "../db/sessions"; import type { SummaryRepository } from "../db/summaries"; /** Whether observations are processed in-process or enqueued for a daemon. */ export type ProcessingMode = "in-process" | "enqueue-only"; /** Configuration for the observation queue processor. */ export interface QueueProcessorConfig { batchSize: number; batchIntervalMs: number; conflictResolutionEnabled?: boolean; conflictSimilarityBandLow?: number; conflictSimilarityBandHigh?: number; entityExtractionEnabled?: boolean; } /** Optional observer for queue lifecycle and performance instrumentation. */ export interface QueueObserver { onEnqueue?(payload: { sessionId: string; toolName: string; createdAt: string; }): void; onBatchStart?(payload: { pending: number; mode: ProcessingMode; startedAt: string; }): void; onBatchEnd?(payload: { processed: number; failed: number; durationMs: number; finishedAt: string; }): void; onItemFailed?(payload: { pendingId: string; error: string; failedAt: string; }): void; } /** * Orchestrates asynchronous observation processing: * 1. Dequeues pending tool outputs from SQLite * 2. Compresses them via the AI compressor (or falls back) * 3. Stores resulting observations * 4. Optionally summarizes completed sessions * * Processing can be triggered by `session.idle` events or a periodic timer. */ export declare class QueueProcessor { private config; private compressor; private summarizer; private pendingRepo; private observationRepo; private sessionRepo; private summaryRepo; private embeddingModel; private conflictEvaluator; private entityExtractor; private entityRepo; private observer; private processing; private timer; private mode; private onEnqueue; constructor(config: QueueProcessorConfig, compressor: ObservationCompressor, summarizer: SessionSummarizer, pendingRepo: PendingMessageRepository, observationRepo: ObservationRepository, sessionRepo: SessionRepository, summaryRepo: SummaryRepository, embeddingModel?: EmbeddingModel | null, conflictEvaluator?: ConflictEvaluator | null, entityExtractor?: EntityExtractor | null, entityRepo?: EntityRepository | null, observer?: QueueObserver | null); setMode(mode: ProcessingMode): void; getMode(): ProcessingMode; setOnEnqueue(callback: (() => void) | null): void; /** Add a new pending message to the queue */ enqueue(sessionId: string, toolName: string, toolOutput: string, callId: string): void; /** * Process up to `batchSize` pending messages. Returns the number * of items successfully processed. Concurrent calls are serialized * via a simple `processing` flag. */ processBatch(): Promise; /** Generate and store a summary for the given session */ summarizeSession(sessionId: string): Promise; /** Start periodic batch processing */ start(): void; /** Stop the periodic timer */ stop(): void; get isRunning(): boolean; get isProcessing(): boolean; getStats(): { pending: number; processing: boolean; }; } //# sourceMappingURL=processor.d.ts.map