dbf8030e93
Friend-DM "Nachricht nicht lesbar" recurred even after v0.21.3 because the proactive sweep called shareConvKeyToUser, which reads the module-level conv-key cache first. After a server-side cleanup the cache still held the stale locally-bootstrapped key, so each side wrapped its own different key for the peer and the bundles diverged anew. Switch the sweep to rotate_conv_key when any peer's user-id bundle is missing at the active version: a fresh symmetric key is generated, wrapped for every member at their CURRENT pubkey, and the active version is bumped under a row-level FOR UPDATE lock. Concurrent rotations are race-safe — the loser sees "new version must be greater" and bails; the winner's bundles propagate via realtime. Realtime conversation_keys subscription now invalidates the cache for the affected (conversationId, key_version) on any INSERT/UPDATE/DELETE — so admin cleanups, peer rotations, or device wraps can no longer leave a stale entry in this client's session cache. queueInsert now fires the first event of a quiet period immediately and only collapses follow-up bursts. BATCH_WINDOW_MS dropped 250 → 80 ms. This closes the ~250 ms gap between the notification sound and the message body appearing. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
925 lines
34 KiB
TypeScript
925 lines
34 KiB
TypeScript
import {
|
||
type AttachmentHandle,
|
||
clearConvKeyCache,
|
||
type DecryptedMessage,
|
||
decryptMessages,
|
||
encryptAndUploadAttachment,
|
||
fetchConversationMessages,
|
||
getOrCreateConvKey,
|
||
insertAttachmentRow,
|
||
MAX_ATTACHMENT_BYTES,
|
||
type MessageWithCipher,
|
||
rotateConvKey,
|
||
sendEncryptedMessage,
|
||
} from '@chat-app/shared/chat';
|
||
import { bytesToPgHex, pgBytesToBytes } from '@chat-app/shared/supabase';
|
||
import { useCallback, useEffect, useMemo, useRef, useState } from 'react';
|
||
|
||
import { decryptBatch as decryptBatchWorker } from './decryptWorker';
|
||
import { generateWebPThumb } from './imageCompress';
|
||
import {
|
||
loadCachedMessages,
|
||
persistMessages,
|
||
pruneCache,
|
||
deleteCachedMessage,
|
||
} from './messageCache';
|
||
import {
|
||
enqueueOutbox,
|
||
getOutbox,
|
||
isDue,
|
||
markAttempt,
|
||
type OutboxItem,
|
||
removeOutbox,
|
||
shouldGiveUp,
|
||
subscribeOutbox,
|
||
} from './messageOutbox';
|
||
import {
|
||
getCachedMessages,
|
||
hasCachedMessages,
|
||
setCachedMessages,
|
||
} from './messageMemoryCache';
|
||
import { supabase } from './supabase';
|
||
import { cachedUserKey } from './userIdentity';
|
||
|
||
interface State {
|
||
messages: DecryptedMessage[];
|
||
loading: boolean;
|
||
error: string | null;
|
||
}
|
||
|
||
interface Args {
|
||
conversationId: string | undefined;
|
||
userId: string | undefined;
|
||
deviceId: string | undefined;
|
||
}
|
||
|
||
type MessageChangePayload = {
|
||
eventType: 'INSERT' | 'UPDATE' | 'DELETE';
|
||
new: Record<string, unknown>;
|
||
old: Record<string, unknown>;
|
||
};
|
||
|
||
function rowToMessage(row: Record<string, unknown>): MessageWithCipher {
|
||
return {
|
||
id: String(row.id),
|
||
conversationId: String(row.conversation_id),
|
||
senderId: String(row.sender_id),
|
||
senderDeviceId: row.sender_device_id ? String(row.sender_device_id) : null,
|
||
replyToId: row.reply_to_id ? String(row.reply_to_id) : null,
|
||
editedAt: row.edited_at ? String(row.edited_at) : null,
|
||
deletedAt: row.deleted_at ? String(row.deleted_at) : null,
|
||
createdAt: String(row.created_at),
|
||
ciphertext: pgBytesToBytes(String(row.ciphertext ?? '\\x')),
|
||
nonce: pgBytesToBytes(String(row.nonce ?? '\\x')),
|
||
keyVersion: typeof row.key_version === 'number' ? row.key_version : 1,
|
||
};
|
||
}
|
||
|
||
export function useConversationMessages({ conversationId, userId, deviceId }: Args): State & {
|
||
send: (
|
||
text: string,
|
||
images?: File[],
|
||
replyToId?: string | null,
|
||
opts?: { viewOnceFlags?: boolean[] },
|
||
) => Promise<void>;
|
||
refresh: () => Promise<void>;
|
||
pending: OutboxItem[];
|
||
retryPending: (id: string) => void;
|
||
cancelPending: (id: string) => void;
|
||
} {
|
||
// Initialize from the in-memory cache so a previously-viewed chat shows
|
||
// content on the very first render after the parent remounts on `:id`
|
||
// change. `loading` stays true ONLY for never-seen conversations (cache
|
||
// miss) so the spinner doesn't flash on every chat switch.
|
||
const [state, setState] = useState<State>(() => {
|
||
if (conversationId && hasCachedMessages(conversationId)) {
|
||
return {
|
||
messages: getCachedMessages(conversationId),
|
||
loading: false,
|
||
error: null,
|
||
};
|
||
}
|
||
return { messages: [], loading: true, error: null };
|
||
});
|
||
const [pending, setPending] = useState<OutboxItem[]>(() =>
|
||
conversationId ? getOutbox(conversationId) : [],
|
||
);
|
||
const privateKeyRef = useRef<Uint8Array | null>(null);
|
||
|
||
useEffect(() => {
|
||
if (!conversationId) {
|
||
setPending([]);
|
||
return;
|
||
}
|
||
return subscribeOutbox((byConv) => {
|
||
setPending(byConv[conversationId] ?? []);
|
||
});
|
||
}, [conversationId]);
|
||
|
||
// Load own user private key once per user.
|
||
useEffect(() => {
|
||
privateKeyRef.current = null;
|
||
if (!userId) return;
|
||
void cachedUserKey(userId).then((pk) => {
|
||
privateKeyRef.current = pk;
|
||
});
|
||
}, [userId]);
|
||
|
||
// Proactive rewrap sweep: when a conversation opens, ensure the active
|
||
// conv-key has a `recipient_user_id` bundle for every accepted member.
|
||
//
|
||
// If any peer is missing a bundle at the active version, the previous
|
||
// implementation called `shareConvKeyToUser` for each missing peer —
|
||
// that helper reads from the module-level conv-key cache first, and if
|
||
// the cache held a STALE locally-generated key (from a buggy bootstrap
|
||
// race in an earlier app version), the stale key got propagated to the
|
||
// peer's row. Both sides then encrypt with mutually un-mergeable keys
|
||
// and every message is "Nachricht nicht lesbar" forever (incident:
|
||
// conv aae12d84).
|
||
//
|
||
// The replacement: when any peer is missing, call `rotateConvKey` once.
|
||
// Rotation generates a fresh symmetric key locally, fetches each member's
|
||
// CURRENT pubkey, wraps the fresh key for everyone, and atomically bumps
|
||
// `active_key_version` via the `rotate_conv_key` RPC (FOR UPDATE lock
|
||
// serialises concurrent rotations). This bypasses the cache entirely:
|
||
// the new version's cache entry is the just-rotated key, and the stale
|
||
// entry at the old version is irrelevant because nobody reads it any more.
|
||
useEffect(() => {
|
||
if (!conversationId || !userId) return;
|
||
let cancelled = false;
|
||
const run = async (): Promise<void> => {
|
||
// Wait for the private key ref to be populated. Loop with backoff
|
||
// because it's set asynchronously by another effect; bail if the
|
||
// conversation switches.
|
||
for (let i = 0; i < 20 && !privateKeyRef.current && !cancelled; i++) {
|
||
await new Promise((r) => setTimeout(r, 100));
|
||
}
|
||
const priv = privateKeyRef.current;
|
||
if (!priv || cancelled) return;
|
||
|
||
try {
|
||
const { data: convRow, error: convErr } = await supabase
|
||
.from('conversations')
|
||
.select('active_key_version')
|
||
.eq('id', conversationId)
|
||
.single();
|
||
if (convErr || !convRow) return;
|
||
// db-types snapshot predates the active_key_version column; cast via unknown.
|
||
const version = (convRow as unknown as { active_key_version: number }).active_key_version;
|
||
|
||
// First, make sure we have a usable handle for the active version
|
||
// (this auto-rotates if we're locked out of our own bundle — the
|
||
// recovery path added in v0.21.1/v0.21.2).
|
||
const handle = await getOrCreateConvKey(supabase, conversationId, {
|
||
userId,
|
||
privateKey: priv,
|
||
});
|
||
if (cancelled) return;
|
||
if (handle.keyVersion > version) return; // already rotated by helper
|
||
|
||
// Check membership state on the server.
|
||
const { data: members, error: mErr } = await supabase
|
||
.from('conversation_members')
|
||
.select('user_id, accepted')
|
||
.eq('conversation_id', conversationId);
|
||
if (mErr || !members) return;
|
||
const peerIds = (members as Array<{ user_id: string; accepted: boolean }>)
|
||
.filter((m) => m.accepted && m.user_id !== userId)
|
||
.map((m) => m.user_id);
|
||
if (peerIds.length === 0) return;
|
||
|
||
// Count how many of the peers have a recipient_user_id bundle at
|
||
// the active version. If any are missing, rotate to V+1 — the
|
||
// rotation will wrap a fresh key for every accepted member with a
|
||
// user_keys row.
|
||
const { data: existingRows, error: rowsErr } = await (
|
||
supabase as unknown as {
|
||
from: (t: string) => {
|
||
select: (s: string) => {
|
||
eq: (c: string, v: string) => {
|
||
eq: (c: string, v: number) => {
|
||
in: (c: string, v: string[]) => Promise<{
|
||
data: Array<{ recipient_user_id: string }> | null;
|
||
error: unknown;
|
||
}>;
|
||
};
|
||
};
|
||
};
|
||
};
|
||
}
|
||
)
|
||
.from('conversation_keys')
|
||
.select('recipient_user_id')
|
||
.eq('conversation_id', conversationId)
|
||
.eq('key_version', version)
|
||
.in('recipient_user_id', peerIds);
|
||
if (rowsErr) return;
|
||
const wrappedPeerIds = new Set(
|
||
(existingRows ?? []).map((r) => r.recipient_user_id),
|
||
);
|
||
const missing = peerIds.filter((id) => !wrappedPeerIds.has(id));
|
||
if (missing.length === 0) return;
|
||
|
||
// At least one peer is missing a bundle — rotate. We deliberately do
|
||
// NOT use the cached conv-key here. The rotation generates a fresh
|
||
// key wrapped to every current member's CURRENT pubkey, so any
|
||
// staleness in the local cache for the OLD version is irrelevant
|
||
// going forward.
|
||
try {
|
||
await rotateConvKey(supabase, conversationId, {
|
||
userId,
|
||
privateKey: priv,
|
||
});
|
||
} catch (err) {
|
||
// Most likely cause: a concurrent peer also called rotate and
|
||
// won the race; their bumped active_key_version makes our
|
||
// `p_new_version <= cur_version` and the RPC raises. That's fine —
|
||
// the next chat-open / send will fetch the new active version and
|
||
// unwrap the bundle that peer wrapped for us.
|
||
console.warn('proactive rotate failed (likely concurrent rotation)', err);
|
||
}
|
||
} catch (err) {
|
||
console.warn('proactive rewrap sweep failed', err);
|
||
}
|
||
};
|
||
void run();
|
||
return () => {
|
||
cancelled = true;
|
||
};
|
||
}, [conversationId, userId]);
|
||
|
||
const decryptBatch = useCallback(
|
||
async (messages: MessageWithCipher[]): Promise<DecryptedMessage[]> => {
|
||
const priv = privateKeyRef.current;
|
||
if (!priv || !userId || messages.length === 0) {
|
||
return messages.map((m) => ({ ...m, plaintext: null }));
|
||
}
|
||
return decryptMessages({
|
||
client: supabase,
|
||
messages,
|
||
ownUserId: userId,
|
||
ownPrivateKey: priv,
|
||
// Offload the symmetric decrypt + utf-8 decode to a Web Worker so
|
||
// the main thread stays responsive during bulk operations (initial
|
||
// fetch, backfill after sleep).
|
||
aeadBatchDelegate: decryptBatchWorker,
|
||
});
|
||
},
|
||
[userId],
|
||
);
|
||
|
||
const refresh = useCallback(async () => {
|
||
if (!conversationId) return;
|
||
try {
|
||
// Only show the loading spinner on the *initial* fetch — subsequent
|
||
// refreshes (focus/visibility) replace messages in-place to avoid
|
||
// flickering an empty state on every wake.
|
||
setState((prev) =>
|
||
prev.messages.length === 0 ? { ...prev, loading: true } : prev,
|
||
);
|
||
const rows = await fetchConversationMessages(supabase, conversationId);
|
||
const decrypted = await decryptBatch(rows);
|
||
setState({ messages: decrypted, loading: false, error: null });
|
||
setCachedMessages(conversationId, decrypted);
|
||
// Persist the fresh batch to the local cache so next conversation
|
||
// switch / app start can hydrate instantly. Fire-and-forget — cache
|
||
// write failure is never user-visible.
|
||
void persistMessages(conversationId, decrypted);
|
||
} catch (err: unknown) {
|
||
setState((prev) => ({
|
||
...prev,
|
||
loading: false,
|
||
error: err instanceof Error ? err.message : 'failed to load messages',
|
||
}));
|
||
}
|
||
}, [conversationId, decryptBatch]);
|
||
|
||
// Hydrate from the local SQLite cache the moment the conversation id
|
||
// changes. Runs in parallel with the network fetch — whichever resolves
|
||
// first populates the UI, and `refresh` will replace stale cache data
|
||
// when the server response lands. On cache-miss this is a ~5ms no-op.
|
||
useEffect(() => {
|
||
if (!conversationId) return;
|
||
// Memory cache already populated state synchronously — skip the disk
|
||
// round-trip entirely. The canonical data lands shortly via refresh();
|
||
// the SQLite cache only matters for cold-start hydration.
|
||
if (hasCachedMessages(conversationId)) return;
|
||
let cancelled = false;
|
||
void loadCachedMessages(conversationId).then((cached) => {
|
||
if (cancelled || cached.length === 0) return;
|
||
setState((prev) => {
|
||
// Don't clobber a fresh server response that already landed.
|
||
if (prev.messages.length > 0) return prev;
|
||
setCachedMessages(conversationId, cached);
|
||
return { messages: cached, loading: false, error: null };
|
||
});
|
||
});
|
||
return () => {
|
||
cancelled = true;
|
||
};
|
||
}, [conversationId]);
|
||
|
||
// Prune cache once per app session.
|
||
useEffect(() => {
|
||
void pruneCache();
|
||
}, []);
|
||
|
||
// Realtime INSERT handler — refetches the row via REST so we get the
|
||
// canonical bytea encoding (postgres_changes payloads serialize bytea
|
||
// differently and decoding them inline is brittle). Then decrypt + append.
|
||
// Skips if the message is already in state (e.g. optimistic insert from our
|
||
// own send), so the sender's cached copy isn't overwritten with a flicker.
|
||
const handleInsert = useCallback(
|
||
async (row: Record<string, unknown>) => {
|
||
if (!conversationId || !deviceId) return;
|
||
const id = String(row.id);
|
||
|
||
let alreadyHave = false;
|
||
setState((prev) => {
|
||
if (prev.messages.some((m) => m.id === id)) alreadyHave = true;
|
||
return prev;
|
||
});
|
||
if (alreadyHave) return;
|
||
|
||
let decrypted: DecryptedMessage | null = null;
|
||
for (let attempt = 0; attempt < 6; attempt++) {
|
||
const { data, error } = await supabase
|
||
.from('messages')
|
||
.select(
|
||
'id, conversation_id, sender_id, sender_device_id, reply_to_id, edited_at, deleted_at, created_at, ciphertext, nonce, key_version',
|
||
)
|
||
.eq('id', id)
|
||
.maybeSingle();
|
||
if (error) {
|
||
console.warn('handleInsert refetch failed', error);
|
||
return;
|
||
}
|
||
if (!data) {
|
||
await new Promise((r) => window.setTimeout(r, 120 * (attempt + 1)));
|
||
continue;
|
||
}
|
||
// db-types snapshot predates the sender-key columns; cast to bypass.
|
||
const r = data as unknown as {
|
||
id: string;
|
||
conversation_id: string;
|
||
sender_id: string;
|
||
sender_device_id: string | null;
|
||
reply_to_id: string | null;
|
||
edited_at: string | null;
|
||
deleted_at: string | null;
|
||
created_at: string;
|
||
ciphertext: string;
|
||
nonce: string;
|
||
key_version: number;
|
||
};
|
||
const msg: MessageWithCipher = {
|
||
id: r.id,
|
||
conversationId: r.conversation_id,
|
||
senderId: r.sender_id,
|
||
senderDeviceId: r.sender_device_id,
|
||
replyToId: r.reply_to_id,
|
||
editedAt: r.edited_at,
|
||
deletedAt: r.deleted_at,
|
||
createdAt: r.created_at,
|
||
ciphertext: pgBytesToBytes(String(r.ciphertext)),
|
||
nonce: pgBytesToBytes(String(r.nonce)),
|
||
keyVersion: r.key_version,
|
||
};
|
||
const [d] = await decryptBatch([msg]);
|
||
if (d) {
|
||
decrypted = d;
|
||
if (d.plaintext !== null) break;
|
||
}
|
||
await new Promise((r) => window.setTimeout(r, 200 * (attempt + 1)));
|
||
}
|
||
if (!decrypted) return;
|
||
setState((prev) => {
|
||
if (prev.messages.some((m) => m.id === decrypted!.id)) return prev;
|
||
const next = [...prev.messages, decrypted!];
|
||
if (conversationId) setCachedMessages(conversationId, next);
|
||
return { ...prev, messages: next };
|
||
});
|
||
},
|
||
[conversationId, deviceId, decryptBatch],
|
||
);
|
||
|
||
const handleUpdate = useCallback(
|
||
async (row: Record<string, unknown>) => {
|
||
const partial = rowToMessage(row);
|
||
setState((prev) => {
|
||
const idx = prev.messages.findIndex((m) => m.id === partial.id);
|
||
if (idx === -1) return prev;
|
||
const existing = prev.messages[idx];
|
||
if (!existing) return prev;
|
||
const next = [...prev.messages];
|
||
next[idx] = {
|
||
...existing,
|
||
editedAt: partial.editedAt,
|
||
deletedAt: partial.deletedAt,
|
||
};
|
||
if (conversationId) setCachedMessages(conversationId, next);
|
||
return { ...prev, messages: next };
|
||
});
|
||
if (partial.editedAt && !partial.deletedAt) {
|
||
// Realtime bytea encoding varies (base64 vs hex, even null for
|
||
// unchanged columns on some configs). Refetch via REST to get the
|
||
// canonical `\x…` hex then decrypt — same pattern as handleInsert.
|
||
if (!conversationId || !deviceId) return;
|
||
let decrypted: DecryptedMessage | null = null;
|
||
for (let attempt = 0; attempt < 6; attempt++) {
|
||
const { data, error } = await supabase
|
||
.from('messages')
|
||
.select(
|
||
'id, conversation_id, sender_id, sender_device_id, reply_to_id, edited_at, deleted_at, created_at, ciphertext, nonce, key_version',
|
||
)
|
||
.eq('id', partial.id)
|
||
.maybeSingle();
|
||
if (error) {
|
||
console.warn('handleUpdate refetch failed', error);
|
||
return;
|
||
}
|
||
if (!data) {
|
||
await new Promise((r) => window.setTimeout(r, 120 * (attempt + 1)));
|
||
continue;
|
||
}
|
||
const r = data as unknown as {
|
||
id: string;
|
||
conversation_id: string;
|
||
sender_id: string;
|
||
sender_device_id: string | null;
|
||
reply_to_id: string | null;
|
||
edited_at: string | null;
|
||
deleted_at: string | null;
|
||
created_at: string;
|
||
ciphertext: string;
|
||
nonce: string;
|
||
key_version: number;
|
||
};
|
||
const msg: MessageWithCipher = {
|
||
id: r.id,
|
||
conversationId: r.conversation_id,
|
||
senderId: r.sender_id,
|
||
senderDeviceId: r.sender_device_id,
|
||
replyToId: r.reply_to_id,
|
||
editedAt: r.edited_at,
|
||
deletedAt: r.deleted_at,
|
||
createdAt: r.created_at,
|
||
ciphertext: pgBytesToBytes(String(r.ciphertext)),
|
||
nonce: pgBytesToBytes(String(r.nonce)),
|
||
keyVersion: r.key_version,
|
||
};
|
||
const [d] = await decryptBatch([msg]);
|
||
if (d) {
|
||
decrypted = d;
|
||
if (d.plaintext !== null) break;
|
||
}
|
||
await new Promise((r) => window.setTimeout(r, 200 * (attempt + 1)));
|
||
}
|
||
if (!decrypted) return;
|
||
setState((prev) => {
|
||
const idx = prev.messages.findIndex((m) => m.id === decrypted!.id);
|
||
if (idx === -1) return prev;
|
||
const next = [...prev.messages];
|
||
next[idx] = decrypted!;
|
||
if (conversationId) setCachedMessages(conversationId, next);
|
||
return { ...prev, messages: next };
|
||
});
|
||
}
|
||
},
|
||
[conversationId, deviceId, decryptBatch],
|
||
);
|
||
|
||
const handleDelete = useCallback(
|
||
(row: Record<string, unknown>) => {
|
||
const id = String(row.id);
|
||
setState((prev) => {
|
||
const next = prev.messages.filter((m) => m.id !== id);
|
||
if (conversationId) setCachedMessages(conversationId, next);
|
||
return { ...prev, messages: next };
|
||
});
|
||
void deleteCachedMessage(id);
|
||
},
|
||
[conversationId],
|
||
);
|
||
|
||
useEffect(() => {
|
||
if (!conversationId || !userId || !deviceId) return;
|
||
void refresh();
|
||
|
||
// Batch INSERT bursts so a paste / backfill doesn't fire N parallel
|
||
// refetches + decrypts. The first event in a quiet period fires
|
||
// `handleInsert` immediately so single incoming messages don't sit
|
||
// behind a debounce timer (previous behaviour: 250 ms blank between
|
||
// notification-sound and message body). Subsequent events arriving
|
||
// within BATCH_WINDOW_MS of the first are buffered; if the burst grows
|
||
// past BATCH_BURST_THRESHOLD the buffered tail collapses into one
|
||
// `refresh()` instead of N individual refetches.
|
||
const BATCH_WINDOW_MS = 80;
|
||
const BATCH_BURST_THRESHOLD = 3;
|
||
let burstBuffer: Array<Record<string, unknown>> = [];
|
||
let burstTimer: number | null = null;
|
||
const flushBurst = () => {
|
||
const buf = burstBuffer;
|
||
burstBuffer = [];
|
||
if (burstTimer !== null) {
|
||
window.clearTimeout(burstTimer);
|
||
burstTimer = null;
|
||
}
|
||
if (buf.length === 0) return;
|
||
if (buf.length > BATCH_BURST_THRESHOLD) {
|
||
void refresh();
|
||
} else {
|
||
for (const row of buf) void handleInsert(row);
|
||
}
|
||
};
|
||
const queueInsert = (row: Record<string, unknown>) => {
|
||
if (burstBuffer.length === 0 && burstTimer === null) {
|
||
// First event in a quiet period — fire immediately so the user sees
|
||
// the message right when they hear the notification sound. Arm a
|
||
// short window in case a burst follows; follow-ups go through the
|
||
// buffer and may collapse into a refresh.
|
||
void handleInsert(row);
|
||
burstTimer = window.setTimeout(flushBurst, BATCH_WINDOW_MS);
|
||
return;
|
||
}
|
||
burstBuffer.push(row);
|
||
if (burstTimer === null) {
|
||
burstTimer = window.setTimeout(flushBurst, BATCH_WINDOW_MS);
|
||
}
|
||
};
|
||
|
||
const channel = supabase
|
||
.channel('conv:' + conversationId)
|
||
.on(
|
||
'postgres_changes',
|
||
{
|
||
event: '*',
|
||
schema: 'public',
|
||
table: 'messages',
|
||
filter: 'conversation_id=eq.' + conversationId,
|
||
},
|
||
(payload: MessageChangePayload) => {
|
||
if (payload.eventType === 'INSERT') {
|
||
queueInsert(payload.new);
|
||
} else if (payload.eventType === 'UPDATE') {
|
||
void handleUpdate(payload.new);
|
||
} else if (payload.eventType === 'DELETE') {
|
||
handleDelete(payload.old);
|
||
}
|
||
},
|
||
)
|
||
// Any conversation_keys change for this conv invalidates the cached
|
||
// conv-key for the affected version. The module-level cache in
|
||
// shared/chat/convKeys.ts otherwise holds the previously-unwrapped key
|
||
// forever within a session — which is exactly what propagated the
|
||
// stale local bootstrap key in conv aae12d84, recreating divergent
|
||
// bundles after a server-side cleanup. Clearing on any INSERT/UPDATE/
|
||
// DELETE for the conv forces the next `getOrCreateConvKey` /
|
||
// `tryGetConvKey` call to re-fetch the canonical bundle from the
|
||
// server. Cheap (a single Map.delete), defensive, and avoids stale-
|
||
// cache propagation across all of {peer rotation, device wrap, admin
|
||
// cleanup}.
|
||
//
|
||
// We also keep the historical "device wrap → refresh" trigger so a
|
||
// freshly-registered device of our own re-decrypts in place.
|
||
.on(
|
||
'postgres_changes',
|
||
{
|
||
event: '*',
|
||
schema: 'public',
|
||
table: 'conversation_keys',
|
||
filter: 'conversation_id=eq.' + conversationId,
|
||
},
|
||
(payload: {
|
||
eventType: 'INSERT' | 'UPDATE' | 'DELETE';
|
||
new: { recipient_device_id?: string; key_version?: number };
|
||
old: { recipient_device_id?: string; key_version?: number };
|
||
}) => {
|
||
const v =
|
||
payload.eventType === 'DELETE'
|
||
? payload.old?.key_version
|
||
: payload.new?.key_version;
|
||
if (typeof v === 'number') {
|
||
clearConvKeyCache(conversationId, v);
|
||
} else {
|
||
clearConvKeyCache(conversationId);
|
||
}
|
||
if (
|
||
payload.eventType === 'INSERT' &&
|
||
payload.new?.recipient_device_id === deviceId
|
||
) {
|
||
void refresh();
|
||
}
|
||
},
|
||
)
|
||
.subscribe();
|
||
|
||
// Refresh + reconnect on wake from background throttle (mostly Windows
|
||
// WebView2). Without this, messages inserted while the window is
|
||
// minimised never arrive until the user explicitly reloads. Throttled to
|
||
// at most once per AWAKE_THROTTLE_MS so the inevitable cluster of
|
||
// visibility/online events on focus does not cause a flicker storm.
|
||
let lastAwakeRefresh = 0;
|
||
const AWAKE_THROTTLE_MS = 30_000;
|
||
const onAwake = () => {
|
||
if (document.visibilityState !== 'visible') return;
|
||
const now = Date.now();
|
||
if (now - lastAwakeRefresh < AWAKE_THROTTLE_MS) return;
|
||
lastAwakeRefresh = now;
|
||
void refresh();
|
||
try {
|
||
channel.subscribe();
|
||
} catch {
|
||
/* already live */
|
||
}
|
||
};
|
||
document.addEventListener('visibilitychange', onAwake);
|
||
window.addEventListener('online', onAwake);
|
||
|
||
return () => {
|
||
if (burstTimer !== null) window.clearTimeout(burstTimer);
|
||
document.removeEventListener('visibilitychange', onAwake);
|
||
window.removeEventListener('online', onAwake);
|
||
void supabase.removeChannel(channel);
|
||
};
|
||
}, [conversationId, userId, deviceId, refresh, handleInsert, handleUpdate, handleDelete]);
|
||
|
||
const sendText = useCallback(
|
||
async (convId: string, uid: string, did: string, priv: Uint8Array, text: string, replyToId: string | null): Promise<void> => {
|
||
const msg = await sendEncryptedMessage({
|
||
client: supabase,
|
||
conversationId: convId,
|
||
plaintext: text,
|
||
senderUserId: uid,
|
||
senderDeviceId: did,
|
||
senderPrivateKey: priv,
|
||
...(replyToId ? { replyToId } : {}),
|
||
});
|
||
setState((prev) => {
|
||
if (prev.messages.some((m) => m.id === msg.id)) return prev;
|
||
const next = [
|
||
...prev.messages,
|
||
{ ...msg, plaintext: text } as DecryptedMessage,
|
||
];
|
||
setCachedMessages(convId, next);
|
||
return { ...prev, messages: next };
|
||
});
|
||
},
|
||
[],
|
||
);
|
||
|
||
const send = useCallback(
|
||
async (
|
||
text: string,
|
||
images: File[] = [],
|
||
replyToId: string | null = null,
|
||
opts: { viewOnceFlags?: boolean[] } = {},
|
||
) => {
|
||
const trimmed = text.trim();
|
||
if ((!trimmed && images.length === 0) || !conversationId || !userId || !deviceId) return;
|
||
const priv = privateKeyRef.current;
|
||
if (!priv) throw new Error('private key not loaded');
|
||
|
||
// Slash-command: /tempmsg <seconds> <text> sends an ephemeral message
|
||
// that the sender auto-deletes after the window elapses. Both peers
|
||
// see the countdown via the expireMs field embedded in the plaintext
|
||
// payload — no server support required.
|
||
const tempMatch = /^\/tempmsg\s+(\d+)\s+([\s\S]+)$/i.exec(trimmed);
|
||
if (tempMatch && images.length === 0) {
|
||
const seconds = Math.min(3600, Math.max(5, parseInt(tempMatch[1]!, 10)));
|
||
const body = tempMatch[2]!.trim();
|
||
const payload = JSON.stringify({
|
||
v: 1,
|
||
type: 'text',
|
||
text: body,
|
||
attachments: [],
|
||
expireMs: seconds * 1000,
|
||
});
|
||
try {
|
||
await sendText(conversationId, userId, deviceId, priv, payload, replyToId);
|
||
} catch (err: unknown) {
|
||
const msg = err instanceof Error ? err.message : 'send failed';
|
||
enqueueOutbox({
|
||
conversationId,
|
||
text: payload,
|
||
replyToId,
|
||
error: msg,
|
||
});
|
||
}
|
||
return;
|
||
}
|
||
|
||
// Text-only path is retryable — if the network is down or the server
|
||
// rejects transiently, stash in the outbox and keep the UI optimistic.
|
||
// Attachments can't be deferred (large payloads, uploaded separately),
|
||
// so those still surface the error immediately.
|
||
if (images.length === 0) {
|
||
try {
|
||
await sendText(conversationId, userId, deviceId, priv, trimmed, replyToId);
|
||
} catch (err: unknown) {
|
||
const msg = err instanceof Error ? err.message : 'send failed';
|
||
enqueueOutbox({
|
||
conversationId,
|
||
text: trimmed,
|
||
replyToId,
|
||
error: msg,
|
||
});
|
||
}
|
||
return;
|
||
}
|
||
|
||
// 1. Upload + encrypt each image. Collect handles + raw blob nonces
|
||
// (so the public attachment row can reference the blob-level nonce).
|
||
// P7.T4: view-once is now a per-attachment flag rather than a
|
||
// composer-wide toggle. `opts.viewOnceFlags` is a parallel array;
|
||
// missing entries (or whole-array absence) default to false.
|
||
const handles: AttachmentHandle[] = [];
|
||
const blobNonceHexByHandleId = new Map<string, string>();
|
||
for (let i = 0; i < images.length; i++) {
|
||
const file = images[i]!;
|
||
if (file.size > MAX_ATTACHMENT_BYTES) {
|
||
throw new Error('attachment exceeds max size (10 MB)');
|
||
}
|
||
const dims = await readImageDimensions(file);
|
||
// Phase 6B: generate a small WebP preview thumb so the receiver's
|
||
// bubble loads fast (typical 320×240 WebP is 10–30KB vs the full
|
||
// image's 1–10MB). `generateWebPThumb` short-circuits to null on
|
||
// non-images, animated formats, and small files — and on failure;
|
||
// the upload helper then just skips the second upload.
|
||
const thumbBlob = await generateWebPThumb(file);
|
||
const res = await encryptAndUploadAttachment({
|
||
client: supabase,
|
||
conversationId,
|
||
file,
|
||
mimeType: file.type || 'application/octet-stream',
|
||
sizeBytes: file.size,
|
||
...(dims.width !== undefined ? { width: dims.width } : {}),
|
||
...(dims.height !== undefined ? { height: dims.height } : {}),
|
||
...(thumbBlob ? { thumbBlob } : {}),
|
||
});
|
||
// Stamp the view-once flag on each handle the caller flagged. The
|
||
// flag rides inside the encrypted payload (so peers can render the
|
||
// locked card without leaking who-sent-what to the server) AND
|
||
// lands on the public message_attachments row via insertAttachmentRow
|
||
// below (where the mark-viewed RPC enforces it).
|
||
if (opts.viewOnceFlags?.[i]) {
|
||
res.handle.viewOnce = true;
|
||
}
|
||
handles.push(res.handle);
|
||
blobNonceHexByHandleId.set(res.handle.id, bytesToPgHex(res.nonce));
|
||
}
|
||
|
||
// 2. Send message (inserts messages + per-conversation key bundles).
|
||
const msg = await sendEncryptedMessage({
|
||
client: supabase,
|
||
conversationId,
|
||
plaintext: trimmed,
|
||
senderUserId: userId,
|
||
senderDeviceId: deviceId,
|
||
senderPrivateKey: priv,
|
||
...(handles.length > 0 ? { attachmentHandles: handles } : {}),
|
||
...(replyToId ? { replyToId } : {}),
|
||
});
|
||
|
||
// 3. Optimistic insert — we already have the plaintext in hand and the
|
||
// server returned the row id, so add the message to local state
|
||
// immediately. Realtime will then no-op (handleInsert dedupes by id).
|
||
const attachmentsPayload =
|
||
handles.length === 0
|
||
? trimmed
|
||
: JSON.stringify({ v: 1, text: trimmed, attachments: handles });
|
||
setState((prev) => {
|
||
if (prev.messages.some((m) => m.id === msg.id)) return prev;
|
||
const next = [
|
||
...prev.messages,
|
||
{
|
||
...msg,
|
||
plaintext: attachmentsPayload,
|
||
} as DecryptedMessage,
|
||
];
|
||
setCachedMessages(conversationId, next);
|
||
return { ...prev, messages: next };
|
||
});
|
||
|
||
// 4. Insert public attachment metadata rows pointing at the new message.
|
||
for (const h of handles) {
|
||
const blobNonce = blobNonceHexByHandleId.get(h.id) ?? '\\x';
|
||
await insertAttachmentRow(supabase, msg.id, h, blobNonce);
|
||
}
|
||
},
|
||
[conversationId, userId, deviceId, sendText],
|
||
);
|
||
|
||
// Drain outbox: retry due items, remove on success, record attempt on fail.
|
||
// Runs on online event, window focus, and a 10s interval.
|
||
useEffect(() => {
|
||
if (!conversationId || !userId || !deviceId) return;
|
||
|
||
let draining = false;
|
||
const drain = async (): Promise<void> => {
|
||
if (draining) return;
|
||
const priv = privateKeyRef.current;
|
||
if (!priv) return;
|
||
if (typeof navigator !== 'undefined' && navigator.onLine === false) return;
|
||
draining = true;
|
||
try {
|
||
const items = getOutbox(conversationId).filter(isDue);
|
||
for (const item of items) {
|
||
if (shouldGiveUp(item)) continue;
|
||
try {
|
||
await sendText(
|
||
conversationId,
|
||
userId,
|
||
deviceId,
|
||
priv,
|
||
item.text,
|
||
item.replyToId,
|
||
);
|
||
removeOutbox(conversationId, item.id);
|
||
} catch (err: unknown) {
|
||
markAttempt(
|
||
conversationId,
|
||
item.id,
|
||
err instanceof Error ? err.message : 'send failed',
|
||
);
|
||
}
|
||
}
|
||
} finally {
|
||
draining = false;
|
||
}
|
||
};
|
||
|
||
const interval = window.setInterval(() => {
|
||
void drain();
|
||
}, 10_000);
|
||
const onOnline = () => void drain();
|
||
window.addEventListener('online', onOnline);
|
||
window.addEventListener('focus', onOnline);
|
||
// Kick once immediately on mount for stale queued items.
|
||
void drain();
|
||
|
||
return () => {
|
||
window.clearInterval(interval);
|
||
window.removeEventListener('online', onOnline);
|
||
window.removeEventListener('focus', onOnline);
|
||
};
|
||
}, [conversationId, userId, deviceId, sendText]);
|
||
|
||
const retryPending = useCallback(
|
||
(id: string) => {
|
||
if (!conversationId) return;
|
||
// Force due now; the drain loop will pick it up on next tick.
|
||
markAttempt(conversationId, id, null);
|
||
// Manual kick: mutate nextAttemptAt by re-enqueuing? Simpler — just
|
||
// trigger a drain-ish by queueing a microtask. The interval picks up
|
||
// due items within 10s, but for UX we also eagerly try here.
|
||
const priv = privateKeyRef.current;
|
||
if (!priv || !userId || !deviceId) return;
|
||
const item = getOutbox(conversationId).find((x) => x.id === id);
|
||
if (!item) return;
|
||
void (async () => {
|
||
try {
|
||
await sendText(conversationId, userId, deviceId, priv, item.text, item.replyToId);
|
||
removeOutbox(conversationId, id);
|
||
} catch (err: unknown) {
|
||
markAttempt(
|
||
conversationId,
|
||
id,
|
||
err instanceof Error ? err.message : 'send failed',
|
||
);
|
||
}
|
||
})();
|
||
},
|
||
[conversationId, userId, deviceId, sendText],
|
||
);
|
||
|
||
const cancelPending = useCallback(
|
||
(id: string) => {
|
||
if (!conversationId) return;
|
||
removeOutbox(conversationId, id);
|
||
},
|
||
[conversationId],
|
||
);
|
||
|
||
return useMemo(
|
||
() => ({ ...state, send, refresh, pending, retryPending, cancelPending }),
|
||
[state, send, refresh, pending, retryPending, cancelPending],
|
||
);
|
||
}
|
||
|
||
// Best-effort image dimension probe. Falls back silently on non-images.
|
||
async function readImageDimensions(file: File): Promise<{ width?: number; height?: number }> {
|
||
if (!file.type.startsWith('image/')) return {};
|
||
const url = URL.createObjectURL(file);
|
||
try {
|
||
return await new Promise<{ width?: number; height?: number }>((resolve) => {
|
||
const img = new Image();
|
||
img.onload = () => resolve({ width: img.naturalWidth, height: img.naturalHeight });
|
||
img.onerror = () => resolve({});
|
||
img.src = url;
|
||
});
|
||
} finally {
|
||
URL.revokeObjectURL(url);
|
||
}
|
||
}
|