diff --git a/app/src/hooks/__tests__/useMemoryIngestionStatus.test.ts b/app/src/hooks/__tests__/useMemoryIngestionStatus.test.ts new file mode 100644 index 000000000..61dc2f3ca --- /dev/null +++ b/app/src/hooks/__tests__/useMemoryIngestionStatus.test.ts @@ -0,0 +1,80 @@ +import { act, renderHook, waitFor } from '@testing-library/react'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +import { useMemoryIngestionStatus } from '../useMemoryIngestionStatus'; + +const mockCallCoreRpc = vi.fn(); + +vi.mock('../../services/coreRpcClient', () => ({ + callCoreRpc: (args: unknown) => mockCallCoreRpc(args), +})); + +describe('useMemoryIngestionStatus', () => { + beforeEach(() => { + mockCallCoreRpc.mockReset(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('maps the snake_case RPC envelope into camelCase status', async () => { + mockCallCoreRpc.mockResolvedValue({ + running: true, + current_document_id: 'doc-1', + current_title: 'Notes', + current_namespace: 'global', + queue_depth: 2, + last_completed_at: 1700000000000, + last_document_id: 'doc-0', + last_success: true, + }); + + const { result } = renderHook(() => useMemoryIngestionStatus()); + + await waitFor(() => expect(result.current.loading).toBe(false)); + + expect(result.current.status).toEqual({ + running: true, + currentDocumentId: 'doc-1', + currentTitle: 'Notes', + currentNamespace: 'global', + queueDepth: 2, + lastCompletedAt: 1700000000000, + lastDocumentId: 'doc-0', + lastSuccess: true, + }); + expect(result.current.error).toBeNull(); + expect(mockCallCoreRpc).toHaveBeenCalledWith({ method: 'openhuman.memory_ingestion_status' }); + }); + + it('reports an error when the RPC fails and keeps idle defaults', async () => { + mockCallCoreRpc.mockRejectedValue(new Error('boom')); + + const { result } = renderHook(() => useMemoryIngestionStatus()); + + await waitFor(() => expect(result.current.loading).toBe(false)); + + expect(result.current.status.running).toBe(false); + expect(result.current.status.queueDepth).toBe(0); + expect(result.current.error).toBe('boom'); + }); + + it('refresh() re-issues the RPC and updates status', async () => { + mockCallCoreRpc + .mockResolvedValueOnce({ running: false, queue_depth: 0 }) + .mockResolvedValueOnce({ running: true, queue_depth: 1, current_document_id: 'doc-2' }); + + const { result } = renderHook(() => useMemoryIngestionStatus()); + await waitFor(() => expect(result.current.loading).toBe(false)); + expect(result.current.status.running).toBe(false); + + await act(async () => { + await result.current.refresh(); + }); + + expect(result.current.status.running).toBe(true); + expect(result.current.status.queueDepth).toBe(1); + expect(result.current.status.currentDocumentId).toBe('doc-2'); + }); +}); diff --git a/app/src/hooks/useMemoryIngestionStatus.ts b/app/src/hooks/useMemoryIngestionStatus.ts new file mode 100644 index 000000000..159da2834 --- /dev/null +++ b/app/src/hooks/useMemoryIngestionStatus.ts @@ -0,0 +1,97 @@ +import { useCallback, useEffect, useRef, useState } from 'react'; + +import { callCoreRpc } from '../services/coreRpcClient'; + +export interface MemoryIngestionStatus { + running: boolean; + currentDocumentId?: string; + currentTitle?: string; + currentNamespace?: string; + queueDepth: number; + lastCompletedAt?: number; + lastDocumentId?: string; + lastSuccess?: boolean; +} + +interface IngestionStatusEnvelope { + running: boolean; + current_document_id?: string; + current_title?: string; + current_namespace?: string; + queue_depth: number; + last_completed_at?: number; + last_document_id?: string; + last_success?: boolean; +} + +const DEFAULT_POLL_MS = 4000; +const FAST_POLL_MS = 1500; + +const EMPTY_STATUS: MemoryIngestionStatus = { running: false, queueDepth: 0 }; + +/** + * Polls `openhuman.memory_ingestion_status`. Polls faster while a job is + * running or queued so the UI reacts quickly when ingestion finishes; + * relaxes to a slower cadence at idle. + */ +export function useMemoryIngestionStatus(): { + status: MemoryIngestionStatus; + loading: boolean; + error: string | null; + refresh: () => void; +} { + const [status, setStatus] = useState(EMPTY_STATUS); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + const cancelledRef = useRef(false); + + const fetchOnce = useCallback(async () => { + try { + const env = await callCoreRpc({ + method: 'openhuman.memory_ingestion_status', + }); + if (cancelledRef.current) return; + setStatus({ + running: env.running, + currentDocumentId: env.current_document_id, + currentTitle: env.current_title, + currentNamespace: env.current_namespace, + queueDepth: env.queue_depth ?? 0, + lastCompletedAt: env.last_completed_at, + lastDocumentId: env.last_document_id, + lastSuccess: env.last_success, + }); + setError(null); + } catch (err) { + if (cancelledRef.current) return; + setError(err instanceof Error ? err.message : String(err)); + } finally { + if (!cancelledRef.current) setLoading(false); + } + }, []); + + const statusRef = useRef(status); + statusRef.current = status; + + useEffect(() => { + cancelledRef.current = false; + let timer: ReturnType | null = null; + + const tick = async () => { + await fetchOnce(); + if (cancelledRef.current) return; + const live = statusRef.current; + const delay = live.running || live.queueDepth > 0 ? FAST_POLL_MS : DEFAULT_POLL_MS; + timer = setTimeout(tick, delay); + }; + + void tick(); + + return () => { + cancelledRef.current = true; + if (timer) clearTimeout(timer); + }; + }, [fetchOnce]); + + return { status, loading, error, refresh: fetchOnce }; +} diff --git a/app/src/pages/Intelligence.tsx b/app/src/pages/Intelligence.tsx index b9f38024b..4e8d70845 100644 --- a/app/src/pages/Intelligence.tsx +++ b/app/src/pages/Intelligence.tsx @@ -17,6 +17,7 @@ import { useIntelligenceSocketManager, } from '../hooks/useIntelligenceSocket'; import { useIntelligenceStats } from '../hooks/useIntelligenceStats'; +import { useMemoryIngestionStatus } from '../hooks/useMemoryIngestionStatus'; import { useScreenIntelligenceItems } from '../hooks/useScreenIntelligenceItems'; import { useSubconscious } from '../hooks/useSubconscious'; import type { @@ -31,6 +32,7 @@ type IntelligenceTab = 'memory' | 'subconscious' | 'dreams'; export default function Intelligence() { const { aiStatus } = useIntelligenceStats(); + const { status: ingestionStatus } = useMemoryIngestionStatus(); const [activeTab, setActiveTab] = useState('memory'); const [sourceFilter, setSourceFilter] = useState('all'); @@ -312,6 +314,24 @@ export default function Intelligence() { {systemStatusLabel} )} + {activeTab === 'memory' && + (ingestionStatus.running || ingestionStatus.queueDepth > 0) && ( +
+
+ + {ingestionStatus.running ? 'Ingesting' : 'Queued'} + {ingestionStatus.queueDepth > 0 && ` · ${ingestionStatus.queueDepth}`} + +
+ )} {activeTab === 'memory' && (