From 99f9ac84fb1f0f199046fc70371434ea3aac5f11 Mon Sep 17 00:00:00 2001 From: maaz519 Date: Sat, 18 Jul 2026 01:02:00 +0530 Subject: [PATCH] feat: add useMessages hook with explicit ownership and optimistic send Co-Authored-By: Claude Opus 4.8 (1M context) --- .../src/hooks/use-messages.test.tsx | 160 ++++++++++++++ .../src/hooks/use-messages.ts | 205 ++++++++++++++++++ 2 files changed, 365 insertions(+) create mode 100644 packages/iios-messaging-ui/src/hooks/use-messages.test.tsx create mode 100644 packages/iios-messaging-ui/src/hooks/use-messages.ts diff --git a/packages/iios-messaging-ui/src/hooks/use-messages.test.tsx b/packages/iios-messaging-ui/src/hooks/use-messages.test.tsx new file mode 100644 index 0000000..f6d0771 --- /dev/null +++ b/packages/iios-messaging-ui/src/hooks/use-messages.test.tsx @@ -0,0 +1,160 @@ +import { describe, it, expect, vi } from 'vitest'; +import { renderHook, waitFor, act } from '@testing-library/react'; +import type { ReactNode } from 'react'; +import { MessagingProvider } from '../provider'; +import { MockAdapter } from '../adapters/mock'; +import { useMessages } from './use-messages'; +import type { MessagingAdapter } from '../adapter'; +import type { Message } from '../types'; + +const wrap = (adapter: MessagingAdapter) => + function Wrapper({ children }: { children: ReactNode }) { + return {children}; + }; + +describe('useMessages', () => { + it('loads history for the thread', async () => { + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(new MockAdapter()) }); + await waitFor(() => expect(result.current.loading).toBe(false)); + expect(result.current.messages).toHaveLength(1); + expect(result.current.messages[0]!.text).toBe('Can you review the Henderson estimate?'); + }); + + // REGRESSION: the CRM inferred actor identity by scanning for a sent message, so + // before you had spoken in a thread EVERY message rendered as not-yours. + it('marks ownership correctly before the user has sent anything', async () => { + const adapter = new MockAdapter(); + await adapter.send('th_mock_1', 'an earlier message of mine'); + + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + await waitFor(() => expect(result.current.messages).toHaveLength(2)); + + // Never sent anything via the hook — ownership still resolves from currentActorId(). + expect(result.current.messages[0]!.mine).toBe(false); // from pp_sofia + expect(result.current.messages[1]!.mine).toBe(true); // from me + }); + + it('appends an optimistic message immediately on send', async () => { + const adapter = new MockAdapter(); + let release!: () => void; + vi.spyOn(adapter, 'send').mockImplementation( + () => new Promise((res) => { release = () => res({ id: 'srv_1', actorId: 'me', text: 'hi', at: '2026-07-17T10:00:00.000Z' }); }), + ); + + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + await waitFor(() => expect(result.current.loading).toBe(false)); + + act(() => { void result.current.send('hi'); }); + + await waitFor(() => expect(result.current.messages).toHaveLength(2)); + expect(result.current.messages[1]!.pending).toBe(true); + expect(result.current.messages[1]!.mine).toBe(true); + + await act(async () => { release(); }); + await waitFor(() => expect(result.current.messages[1]!.pending).toBeFalsy()); + }); + + it('rolls back the optimistic message and reports error when send fails', async () => { + const adapter = new MockAdapter(); + vi.spyOn(adapter, 'send').mockRejectedValue(new Error('offline')); + + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + await waitFor(() => expect(result.current.loading).toBe(false)); + + await act(async () => { + await expect(result.current.send('doomed')).rejects.toThrow('offline'); + }); + + expect(result.current.messages).toHaveLength(1); + expect(result.current.messages.some((m) => m.text === 'doomed')).toBe(false); + expect(result.current.error).toBe('offline'); + }); + + it('does not duplicate a message when the transport echoes it back', async () => { + const adapter = new MockAdapter(); + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + await waitFor(() => expect(result.current.loading).toBe(false)); + + // MockAdapter.send emits a 'message' event AND resolves with the same message. + await act(async () => { await result.current.send('echo once'); }); + + expect(result.current.messages.filter((m) => m.text === 'echo once')).toHaveLength(1); + }); + + it('collects typing user ids from subscribe events', async () => { + const adapter = new MockAdapter(); + let emit!: (userId: string) => void; + vi.spyOn(adapter, 'subscribe').mockImplementation((_t, cb) => { + emit = (userId) => cb({ kind: 'typing', userId }); + return () => {}; + }); + + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + await waitFor(() => expect(result.current.loading).toBe(false)); + + act(() => emit('pp_sofia')); + expect(result.current.typingUserIds).toEqual(['pp_sofia']); + }); + + it('unsubscribes on unmount', async () => { + const adapter = new MockAdapter(); + const off = vi.fn(); + vi.spyOn(adapter, 'subscribe').mockReturnValue(off); + + const { unmount, result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + await waitFor(() => expect(result.current.loading).toBe(false)); + unmount(); + + expect(off).toHaveBeenCalled(); + }); + + it('keeps a live message that arrives before history resolves', async () => { + const adapter = new MockAdapter(); + let resolveHistory!: (msgs: Message[]) => void; + vi.spyOn(adapter, 'history').mockImplementation( + () => new Promise((res) => { resolveHistory = res; }), + ); + let emit!: (m: Message) => void; + vi.spyOn(adapter, 'subscribe').mockImplementation((_t, cb) => { + emit = (m) => cb({ kind: 'message', message: m }); + return () => {}; + }); + + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + + // A live message arrives while history() is still pending. + act(() => emit({ id: 'live_1', actorId: 'pp_sofia', text: 'ping before history', at: '2026-07-17T10:00:00.000Z' })); + + // History resolves afterwards with an older message. + await act(async () => { + resolveHistory([{ id: 'hist_1', actorId: 'pp_sofia', text: 'older', at: '2026-07-17T09:00:00.000Z' }]); + }); + + const texts = result.current.messages.map((m) => m.text); + expect(texts).toContain('older'); + expect(texts).toContain('ping before history'); // must NOT be clobbered by history load + }); + + it('clears a typing indicator after its TTL elapses', async () => { + vi.useFakeTimers(); + try { + const adapter = new MockAdapter(); + let emit!: (userId: string) => void; + vi.spyOn(adapter, 'subscribe').mockImplementation((_t, cb) => { + emit = (userId) => cb({ kind: 'typing', userId }); + return () => {}; + }); + + const { result } = renderHook(() => useMessages('th_mock_1'), { wrapper: wrap(adapter) }); + await act(async () => { await vi.advanceTimersByTimeAsync(0); }); // flush history microtask + + act(() => emit('pp_sofia')); + expect(result.current.typingUserIds).toEqual(['pp_sofia']); + + await act(async () => { await vi.advanceTimersByTimeAsync(3600); }); + expect(result.current.typingUserIds).toEqual([]); + } finally { + vi.useRealTimers(); + } + }); +}); diff --git a/packages/iios-messaging-ui/src/hooks/use-messages.ts b/packages/iios-messaging-ui/src/hooks/use-messages.ts new file mode 100644 index 0000000..d50954d --- /dev/null +++ b/packages/iios-messaging-ui/src/hooks/use-messages.ts @@ -0,0 +1,205 @@ +import { useCallback, useEffect, useMemo, useRef, useState } from 'react'; +import { useAdapter } from '../provider'; +import { isOwnMessage } from '../types'; +import type { Message, SendOpts } from '../types'; + +const TYPING_TTL_MS = 3500; + +export interface UiMessage extends Message { + mine: boolean; +} + +export interface MessagesState { + messages: UiMessage[]; + loading: boolean; + error: string | null; + send: (content: string, opts?: SendOpts) => Promise; + react: (messageId: string, emoji: string) => Promise; + typingUserIds: string[]; + seenIds: Set; + sendTyping: () => void; + canReact: boolean; + canUpload: boolean; +} + +let optimisticSeq = 0; + +export function useMessages(threadId: string | null): MessagesState { + const adapter = useAdapter(); + const [raw, setRaw] = useState([]); + const [loading, setLoading] = useState(true); + const [error, setError] = useState(null); + const [typing, setTyping] = useState>({}); + const [seenIds, setSeenIds] = useState>(new Set()); + + const currentActorId = adapter.currentActorId(); + const actorRef = useRef(currentActorId); + actorRef.current = currentActorId; + + // Load history, then subscribe. Reconciliation is by message id, so an echoed + // send never duplicates the optimistic row. + useEffect(() => { + if (!threadId) { + setRaw([]); + setLoading(false); + return; + } + let alive = true; + setLoading(true); + setRaw([]); + setError(null); + setSeenIds(new Set()); + setTyping({}); + + adapter + .history(threadId) + .then((h) => { + if (!alive) return; + // Merge, don't clobber: a live message can arrive via subscribe while this + // history fetch is still in flight. Blindly setting raw = h would drop it. + setRaw((live) => { + const histIds = new Set(h.map((m) => m.id)); + const extras = live.filter((m) => !histIds.has(m.id)); + return extras.length ? [...h, ...extras] : h; + }); + setError(null); + }) + .catch((e: unknown) => { + if (alive) setError(e instanceof Error ? e.message : String(e)); + }) + .finally(() => { + if (alive) setLoading(false); + }); + + const off = adapter.subscribe(threadId, (e) => { + if (!alive) return; + switch (e.kind) { + case 'message': + setRaw((l) => (l.some((m) => m.id === e.message.id) ? l : [...l, e.message])); + break; + case 'typing': + if (e.userId !== actorRef.current) { + setTyping((t) => ({ ...t, [e.userId]: Date.now() + TYPING_TTL_MS })); + } + break; + case 'receipt': + // Only the OTHER side reading my message counts as "seen". + if (e.actorId !== actorRef.current) { + setSeenIds((s) => (s.has(e.messageId) ? s : new Set(s).add(e.messageId))); + } + break; + case 'reaction': + setRaw((l) => l.map((m) => (m.id === e.messageId ? { ...m, reactions: e.reactions } : m))); + break; + } + }); + + return () => { + alive = false; + off(); + }; + }, [adapter, threadId]); + + const messages: UiMessage[] = useMemo( + () => raw.map((m) => ({ ...m, mine: isOwnMessage(m, currentActorId) })), + [raw, currentActorId], + ); + + const send = useCallback( + async (content: string, opts?: SendOpts) => { + if (!threadId) return; + const tempId = `optimistic_${optimisticSeq++}`; + const optimistic: Message = { + id: tempId, + actorId: actorRef.current, + text: content, + at: new Date().toISOString(), + pending: true, + reactions: [], + ...(opts?.parentInteractionId ? { parentInteractionId: opts.parentInteractionId } : {}), + ...(opts?.attachment ? { attachment: opts.attachment } : {}), + }; + setRaw((l) => [...l, optimistic]); + + try { + const saved = await adapter.send(threadId, content, opts); + setError(null); + // Replace the optimistic row with the server's. If the subscribe echo already + // added the real message, just drop the optimistic one. + setRaw((l) => { + const withoutTemp = l.filter((m) => m.id !== tempId); + return withoutTemp.some((m) => m.id === saved.id) ? withoutTemp : [...withoutTemp, saved]; + }); + } catch (e: unknown) { + setRaw((l) => l.filter((m) => m.id !== tempId)); + setError(e instanceof Error ? e.message : String(e)); + throw e; + } + }, + [adapter, threadId], + ); + + const react = useCallback( + async (messageId: string, emoji: string) => { + if (!threadId || !adapter.react) return; + await adapter.react(threadId, messageId, emoji); + }, + [adapter, threadId], + ); + + const sendTyping = useCallback(() => { + if (threadId) adapter.sendTyping(threadId); + }, [adapter, threadId]); + + // The newest acknowledged (non-pending) message id — what we report as read. + const lastReadableId = useMemo(() => { + for (let i = raw.length - 1; i >= 0; i--) { + if (!raw[i]!.pending) return raw[i]!.id; + } + return null; + }, [raw]); + + // Report my read of the newest message (drives the other side's "seen" tick). + // Keyed on the id, not the whole array, so reaction/optimistic churn doesn't re-fire it. + useEffect(() => { + if (!threadId || !lastReadableId) return; + void adapter.markRead(threadId, lastReadableId).catch(() => {}); + }, [adapter, threadId, lastReadableId]); + + const typingUserIds = useMemo(() => { + const now = Date.now(); + return Object.entries(typing) + .filter(([, exp]) => exp > now) + .map(([u]) => u); + }, [typing]); + + // Expire stale typing entries. Bumping `typing` to a new reference forces the + // memo above to recompute with a fresh `now`, dropping entries past their TTL. + // (A bump of unrelated state can't do this — the memo is keyed on `typing`, so it + // would return its cached array and the indicator would stick forever.) + useEffect(() => { + if (typingUserIds.length === 0) return; + const t = setTimeout(() => setTyping((p) => ({ ...p })), TYPING_TTL_MS); + return () => clearTimeout(t); + }, [typingUserIds.length, typing]); + + // Only my messages that the other side has read. + const seenMine = useMemo(() => { + const out = new Set(); + for (const id of seenIds) if (messages.some((m) => m.id === id && m.mine)) out.add(id); + return out; + }, [seenIds, messages]); + + return { + messages, + loading, + error, + send, + react, + typingUserIds, + seenIds: seenMine, + sendTyping, + canReact: typeof adapter.react === 'function', + canUpload: typeof adapter.upload === 'function', + }; +}