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',
+ };
+}