feat: voice messages, offline queue, delivery ticks, volume slider, admin + scaling
- Voice messages: MediaRecorder → encrypted attachment, custom waveform player via OfflineAudioContext, 60s limit + live mic-level meter - Offline message queue: localStorage outbox, exponential backoff retries, optimistic pending bubble with retry/discard - Delivery indicator: message_deliveries table + RLS (reciprocal receipts), ✓ / ✓✓ / ✓✓-blue tick states, group-aware (all members must ack) - Per-participant volume slider in calls via right-click tile menu, persisted to localStorage, applied to attached audio elements - Group call scaling: grid up to 12 tiles with pagination, active-speaker auto-promotion in fullscreen - Push notifications scaffolding: service worker, VAPID subscription registration, notify-push edge function skeleton - Backup recovery code: 24-char base32 code (~120 bits entropy) as alternative decrypt path, restore UI with mode toggle - Admin panel: conversations list, audit log (admin_audit_log table + admin_log_action RPC), audit entry on user flag toggle - Search v2: sender filter, attachment-only toggle, date range - Reactions pop animation (scale 0.4→1.15→1 on count change) - Message list windowing (150 default, expand via IntersectionObserver) - Stub cleanup: removed dead ScreenshareStub from CallParticipantTile Fixes: - Focus-triggered flicker: dropped window.focus listeners in three spots, throttled visibilitychange/online wake-refreshes to 30s, keep existing data visible during background re-syncs (no more spinner on every click) - Voice attachment audio element collapsed to 0px on peer side — now forces 280px min-width on bubble Migrations (push required): 20260421000001_message_deliveries.sql 20260421000002_admin_audit_log.sql Server TODO: VAPID keys + notify-push edge function deploy
This commit is contained in:
@@ -140,3 +140,59 @@ export async function importDeviceBackup(
|
||||
export function decodePrivateKeyFromBackup(payload: DeviceBackupPayload): Uint8Array {
|
||||
return unb64url(payload.privateKeyB64);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Recovery code
|
||||
// ---------------------------------------------------------------------------
|
||||
//
|
||||
// Generates a high-entropy code shown to the user once at backup time. The
|
||||
// same payload is encrypted twice — once with the user's passphrase, once
|
||||
// with the recovery code — so either string can decrypt the device key.
|
||||
//
|
||||
// The recovery code is 24 chars from a 32-symbol alphabet (no ambiguous
|
||||
// characters), grouped as 4×6. ~120 bits of entropy.
|
||||
|
||||
const RECOVERY_ALPHABET = 'ABCDEFGHJKLMNPQRSTUVWXYZ23456789';
|
||||
|
||||
export interface BackupBundle {
|
||||
passphraseBackup: string;
|
||||
recoveryBackup: string;
|
||||
recoveryCode: string;
|
||||
}
|
||||
|
||||
export async function exportDeviceBackupWithRecovery(params: {
|
||||
userId: string;
|
||||
deviceId: string;
|
||||
privateKey: Uint8Array;
|
||||
passphrase: string;
|
||||
}): Promise<BackupBundle> {
|
||||
const recoveryCode = await generateRecoveryCode();
|
||||
const [passphraseBackup, recoveryBackup] = await Promise.all([
|
||||
exportDeviceBackup(params),
|
||||
exportDeviceBackup({ ...params, passphrase: recoveryCode }),
|
||||
]);
|
||||
return { passphraseBackup, recoveryBackup, recoveryCode };
|
||||
}
|
||||
|
||||
async function generateRecoveryCode(): Promise<string> {
|
||||
const s = await ensureSodium();
|
||||
const raw = s.randombytes_buf(24);
|
||||
let out = '';
|
||||
for (let i = 0; i < raw.length; i++) {
|
||||
out += RECOVERY_ALPHABET[raw[i]! % RECOVERY_ALPHABET.length];
|
||||
if ((i + 1) % 6 === 0 && i !== raw.length - 1) out += '-';
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
// Normalizes a user-typed recovery code: strips dashes/spaces, uppercases,
|
||||
// maps look-alike characters. Lets users enter the code with imperfect
|
||||
// spacing without rejecting valid input.
|
||||
export function normalizeRecoveryCode(input: string): string {
|
||||
return input
|
||||
.toUpperCase()
|
||||
.replace(/[\s-]/g, '')
|
||||
.split('')
|
||||
.filter((c) => RECOVERY_ALPHABET.includes(c))
|
||||
.join('');
|
||||
}
|
||||
|
||||
@@ -0,0 +1,164 @@
|
||||
// Offline text-message outbox.
|
||||
//
|
||||
// When `sendEncryptedMessage` fails (network loss, server 5xx, transient
|
||||
// realtime hiccups) we stash the plaintext in localStorage keyed by
|
||||
// conversation. The UI shows it inline with a "sending…" indicator; the
|
||||
// drain loop retries with exponential backoff until it lands.
|
||||
//
|
||||
// Attachments are intentionally NOT queued — they're too large to persist
|
||||
// and require server-side upload that can't be deferred reliably. The
|
||||
// composer surfaces an immediate error for those.
|
||||
|
||||
const STORAGE_KEY = 'chat.outbox.v1';
|
||||
const MAX_ATTEMPTS = 8;
|
||||
|
||||
export interface OutboxItem {
|
||||
/** Local-only id (never collides with server UUIDs). */
|
||||
id: string;
|
||||
conversationId: string;
|
||||
text: string;
|
||||
replyToId: string | null;
|
||||
createdAt: string;
|
||||
attempts: number;
|
||||
/** Timestamp of the next permitted send attempt. */
|
||||
nextAttemptAt: string;
|
||||
/** Last error message for surface in the UI. */
|
||||
lastError: string | null;
|
||||
}
|
||||
|
||||
type Store = Record<string, OutboxItem[]>; // keyed by conversationId
|
||||
type Listener = (byConv: Store) => void;
|
||||
|
||||
function load(): Store {
|
||||
try {
|
||||
const raw = localStorage.getItem(STORAGE_KEY);
|
||||
if (!raw) return {};
|
||||
const parsed: unknown = JSON.parse(raw);
|
||||
if (!parsed || typeof parsed !== 'object') return {};
|
||||
const out: Store = {};
|
||||
for (const [k, v] of Object.entries(parsed as Record<string, unknown>)) {
|
||||
if (Array.isArray(v)) {
|
||||
out[k] = v.filter((x): x is OutboxItem => isItem(x));
|
||||
}
|
||||
}
|
||||
return out;
|
||||
} catch {
|
||||
return {};
|
||||
}
|
||||
}
|
||||
|
||||
function isItem(x: unknown): x is OutboxItem {
|
||||
if (!x || typeof x !== 'object') return false;
|
||||
const o = x as Record<string, unknown>;
|
||||
return (
|
||||
typeof o.id === 'string' &&
|
||||
typeof o.conversationId === 'string' &&
|
||||
typeof o.text === 'string' &&
|
||||
typeof o.createdAt === 'string'
|
||||
);
|
||||
}
|
||||
|
||||
let store: Store = load();
|
||||
const listeners = new Set<Listener>();
|
||||
|
||||
function persist(): void {
|
||||
try {
|
||||
localStorage.setItem(STORAGE_KEY, JSON.stringify(store));
|
||||
} catch {
|
||||
/* quota exceeded — ignore; in-memory queue still retries */
|
||||
}
|
||||
}
|
||||
|
||||
function notify(): void {
|
||||
for (const fn of listeners) fn(store);
|
||||
}
|
||||
|
||||
export function subscribeOutbox(fn: Listener): () => void {
|
||||
listeners.add(fn);
|
||||
fn(store);
|
||||
return () => {
|
||||
listeners.delete(fn);
|
||||
};
|
||||
}
|
||||
|
||||
export function getOutbox(conversationId: string): OutboxItem[] {
|
||||
return store[conversationId] ?? [];
|
||||
}
|
||||
|
||||
export function enqueueOutbox(args: {
|
||||
conversationId: string;
|
||||
text: string;
|
||||
replyToId: string | null;
|
||||
error: string | null;
|
||||
}): OutboxItem {
|
||||
const item: OutboxItem = {
|
||||
id: 'local-' + crypto.randomUUID(),
|
||||
conversationId: args.conversationId,
|
||||
text: args.text,
|
||||
replyToId: args.replyToId,
|
||||
createdAt: new Date().toISOString(),
|
||||
attempts: 0,
|
||||
nextAttemptAt: new Date().toISOString(),
|
||||
lastError: args.error,
|
||||
};
|
||||
const bucket = store[args.conversationId] ?? [];
|
||||
store = { ...store, [args.conversationId]: [...bucket, item] };
|
||||
persist();
|
||||
notify();
|
||||
return item;
|
||||
}
|
||||
|
||||
export function removeOutbox(conversationId: string, id: string): void {
|
||||
const bucket = store[conversationId];
|
||||
if (!bucket) return;
|
||||
const next = bucket.filter((x) => x.id !== id);
|
||||
if (next.length === bucket.length) return;
|
||||
if (next.length === 0) {
|
||||
const copy = { ...store };
|
||||
delete copy[conversationId];
|
||||
store = copy;
|
||||
} else {
|
||||
store = { ...store, [conversationId]: next };
|
||||
}
|
||||
persist();
|
||||
notify();
|
||||
}
|
||||
|
||||
export function markAttempt(
|
||||
conversationId: string,
|
||||
id: string,
|
||||
error: string | null,
|
||||
): OutboxItem | null {
|
||||
const bucket = store[conversationId];
|
||||
if (!bucket) return null;
|
||||
const idx = bucket.findIndex((x) => x.id === id);
|
||||
if (idx < 0) return null;
|
||||
const prev = bucket[idx]!;
|
||||
const attempts = prev.attempts + 1;
|
||||
// Exponential backoff: 2^n seconds, capped at 60s.
|
||||
const delaySec = Math.min(60, Math.pow(2, attempts));
|
||||
const next: OutboxItem = {
|
||||
...prev,
|
||||
attempts,
|
||||
lastError: error,
|
||||
nextAttemptAt: new Date(Date.now() + delaySec * 1000).toISOString(),
|
||||
};
|
||||
const nextBucket = [...bucket];
|
||||
nextBucket[idx] = next;
|
||||
store = { ...store, [conversationId]: nextBucket };
|
||||
persist();
|
||||
notify();
|
||||
return next;
|
||||
}
|
||||
|
||||
export function shouldGiveUp(item: OutboxItem): boolean {
|
||||
return item.attempts >= MAX_ATTEMPTS;
|
||||
}
|
||||
|
||||
export function isDue(item: OutboxItem): boolean {
|
||||
return new Date(item.nextAttemptAt).getTime() <= Date.now();
|
||||
}
|
||||
|
||||
export function allConversationsWithItems(): string[] {
|
||||
return Object.keys(store);
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
// Per-participant output volume overrides. Values are 0..1 (HTMLMediaElement
|
||||
// scale). Persisted so the user's mixing choices survive reconnects.
|
||||
//
|
||||
// The store is intentionally tiny — a flat Record keyed by LiveKit identity
|
||||
// (our `userId`). `attachTrack` reads from here when a remote audio track
|
||||
// first lands; live changes are pushed to any already-attached audio
|
||||
// elements via the `data-participant` attribute selector.
|
||||
|
||||
const STORAGE_KEY = 'call.participantVolumes.v1';
|
||||
const DEFAULT_VOLUME = 1;
|
||||
|
||||
type VolumeMap = Record<string, number>;
|
||||
type Listener = (map: VolumeMap) => void;
|
||||
|
||||
function load(): VolumeMap {
|
||||
try {
|
||||
const raw = localStorage.getItem(STORAGE_KEY);
|
||||
if (!raw) return {};
|
||||
const parsed: unknown = JSON.parse(raw);
|
||||
if (!parsed || typeof parsed !== 'object') return {};
|
||||
const out: VolumeMap = {};
|
||||
for (const [k, v] of Object.entries(parsed as Record<string, unknown>)) {
|
||||
if (typeof v === 'number' && Number.isFinite(v)) {
|
||||
out[k] = clamp(v);
|
||||
}
|
||||
}
|
||||
return out;
|
||||
} catch {
|
||||
return {};
|
||||
}
|
||||
}
|
||||
|
||||
function clamp(v: number): number {
|
||||
if (v < 0) return 0;
|
||||
if (v > 1) return 1;
|
||||
return v;
|
||||
}
|
||||
|
||||
let current: VolumeMap = load();
|
||||
const listeners = new Set<Listener>();
|
||||
|
||||
function persist(): void {
|
||||
try {
|
||||
localStorage.setItem(STORAGE_KEY, JSON.stringify(current));
|
||||
} catch {
|
||||
/* ignore quota */
|
||||
}
|
||||
}
|
||||
|
||||
function notify(): void {
|
||||
for (const fn of listeners) fn(current);
|
||||
}
|
||||
|
||||
export function getParticipantVolume(userId: string): number {
|
||||
return current[userId] ?? DEFAULT_VOLUME;
|
||||
}
|
||||
|
||||
export function setParticipantVolume(userId: string, volume: number): void {
|
||||
const next = clamp(volume);
|
||||
if (next === (current[userId] ?? DEFAULT_VOLUME)) return;
|
||||
current = { ...current, [userId]: next };
|
||||
persist();
|
||||
applyToAttachedElements(userId, next);
|
||||
notify();
|
||||
}
|
||||
|
||||
export function subscribeParticipantVolumes(fn: Listener): () => void {
|
||||
listeners.add(fn);
|
||||
return () => {
|
||||
listeners.delete(fn);
|
||||
};
|
||||
}
|
||||
|
||||
// Apply a volume to any audio elements already attached for this user.
|
||||
// Attached elements are tagged with `data-participant` in attachTrack.
|
||||
function applyToAttachedElements(userId: string, volume: number): void {
|
||||
const nodes = document.querySelectorAll<HTMLAudioElement>(
|
||||
'audio[data-participant="' + cssEscape(userId) + '"]',
|
||||
);
|
||||
nodes.forEach((el) => {
|
||||
el.volume = volume;
|
||||
});
|
||||
}
|
||||
|
||||
function cssEscape(v: string): string {
|
||||
if (typeof (globalThis as { CSS?: { escape?: (s: string) => string } }).CSS
|
||||
?.escape === 'function') {
|
||||
return (globalThis as { CSS: { escape: (s: string) => string } }).CSS.escape(v);
|
||||
}
|
||||
return v.replace(/"/g, '\\"');
|
||||
}
|
||||
@@ -13,6 +13,16 @@ import {
|
||||
import { bytesToPgHex, pgBytesToBytes } from '@chat-app/shared/supabase';
|
||||
import { useCallback, useEffect, useMemo, useRef, useState } from 'react';
|
||||
|
||||
import {
|
||||
enqueueOutbox,
|
||||
getOutbox,
|
||||
isDue,
|
||||
markAttempt,
|
||||
type OutboxItem,
|
||||
removeOutbox,
|
||||
shouldGiveUp,
|
||||
subscribeOutbox,
|
||||
} from './messageOutbox';
|
||||
import { devLocalSecretStore } from './secretStore';
|
||||
import { supabase } from './supabase';
|
||||
|
||||
@@ -53,10 +63,26 @@ function rowToMessage(row: Record<string, unknown>): MessageWithCipher {
|
||||
export function useConversationMessages({ conversationId, userId, deviceId }: Args): State & {
|
||||
send: (text: string, images?: File[], replyToId?: string | null) => Promise<void>;
|
||||
refresh: () => Promise<void>;
|
||||
pending: OutboxItem[];
|
||||
retryPending: (id: string) => void;
|
||||
cancelPending: (id: string) => void;
|
||||
} {
|
||||
const [state, setState] = useState<State>({ 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 private key once per (user, device).
|
||||
useEffect(() => {
|
||||
privateKeyRef.current = null;
|
||||
@@ -85,7 +111,12 @@ export function useConversationMessages({ conversationId, userId, deviceId }: Ar
|
||||
const refresh = useCallback(async () => {
|
||||
if (!conversationId) return;
|
||||
try {
|
||||
setState((prev) => ({ ...prev, loading: true }));
|
||||
// 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 });
|
||||
@@ -311,9 +342,16 @@ export function useConversationMessages({ conversationId, userId, deviceId }: Ar
|
||||
|
||||
// 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.
|
||||
// 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();
|
||||
@@ -322,17 +360,40 @@ export function useConversationMessages({ conversationId, userId, deviceId }: Ar
|
||||
}
|
||||
};
|
||||
document.addEventListener('visibilitychange', onAwake);
|
||||
window.addEventListener('focus', onAwake);
|
||||
window.addEventListener('online', onAwake);
|
||||
|
||||
return () => {
|
||||
document.removeEventListener('visibilitychange', onAwake);
|
||||
window.removeEventListener('focus', 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;
|
||||
return {
|
||||
...prev,
|
||||
messages: [
|
||||
...prev.messages,
|
||||
{ ...msg, plaintext: text } as DecryptedMessage,
|
||||
],
|
||||
};
|
||||
});
|
||||
},
|
||||
[],
|
||||
);
|
||||
|
||||
const send = useCallback(
|
||||
async (text: string, images: File[] = [], replyToId: string | null = null) => {
|
||||
const trimmed = text.trim();
|
||||
@@ -340,6 +401,25 @@ export function useConversationMessages({ conversationId, userId, deviceId }: Ar
|
||||
const priv = privateKeyRef.current;
|
||||
if (!priv) throw new Error('private key not loaded');
|
||||
|
||||
// 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).
|
||||
const handles: AttachmentHandle[] = [];
|
||||
@@ -401,10 +481,104 @@ export function useConversationMessages({ conversationId, userId, deviceId }: Ar
|
||||
await insertAttachmentRow(supabase, msg.id, h, blobNonce);
|
||||
}
|
||||
},
|
||||
[conversationId, userId, deviceId],
|
||||
[conversationId, userId, deviceId, sendText],
|
||||
);
|
||||
|
||||
return useMemo(() => ({ ...state, send, refresh }), [state, send, refresh]);
|
||||
// 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.
|
||||
|
||||
@@ -45,8 +45,15 @@ export function useFriendships(userId: string | undefined): FriendshipsState & {
|
||||
.subscribe();
|
||||
|
||||
// Windows WebView2 throttles background sockets — refresh on wake.
|
||||
// Throttled + visibility-only so a normal click into the window does not
|
||||
// re-fetch on every focus.
|
||||
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();
|
||||
@@ -55,12 +62,10 @@ export function useFriendships(userId: string | undefined): FriendshipsState & {
|
||||
}
|
||||
};
|
||||
document.addEventListener('visibilitychange', onAwake);
|
||||
window.addEventListener('focus', onAwake);
|
||||
window.addEventListener('online', onAwake);
|
||||
|
||||
return () => {
|
||||
document.removeEventListener('visibilitychange', onAwake);
|
||||
window.removeEventListener('focus', onAwake);
|
||||
window.removeEventListener('online', onAwake);
|
||||
void supabase.removeChannel(channel);
|
||||
};
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
import {
|
||||
listGroupDeliveriesForMessages,
|
||||
listGroupReadsForMessages,
|
||||
} from '@chat-app/shared/chat';
|
||||
import { useCallback, useEffect, useMemo, useState } from 'react';
|
||||
|
||||
import { supabase } from './supabase';
|
||||
|
||||
export interface GroupReceiptState {
|
||||
// message_id → set of user ids who have delivered/read.
|
||||
deliveredByMessage: Map<string, Set<string>>;
|
||||
readByMessage: Map<string, Set<string>>;
|
||||
refresh: () => Promise<void>;
|
||||
}
|
||||
|
||||
// Aggregated delivery + read receipts for group conversations. Returns the
|
||||
// set of user ids per message; consumers join with conversation members to
|
||||
// figure out who has not yet acknowledged.
|
||||
export function useGroupReceipts(
|
||||
messageIds: string[],
|
||||
selfUserId: string | undefined,
|
||||
active: boolean,
|
||||
): GroupReceiptState {
|
||||
const idsKey = useMemo(() => messageIds.join(','), [messageIds]);
|
||||
const [deliveredByMessage, setDelivered] = useState<Map<string, Set<string>>>(new Map());
|
||||
const [readByMessage, setRead] = useState<Map<string, Set<string>>>(new Map());
|
||||
|
||||
const refresh = useCallback(async () => {
|
||||
if (!active || !selfUserId || messageIds.length === 0) {
|
||||
setDelivered(new Map());
|
||||
setRead(new Map());
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const [delivered, read] = await Promise.all([
|
||||
listGroupDeliveriesForMessages(supabase, messageIds, selfUserId),
|
||||
listGroupReadsForMessages(supabase, messageIds, selfUserId),
|
||||
]);
|
||||
setDelivered(toIdSets(delivered));
|
||||
setRead(toIdSets(read));
|
||||
} catch (err: unknown) {
|
||||
console.warn('useGroupReceipts refresh failed', err);
|
||||
}
|
||||
// eslint-disable-next-line react-hooks/exhaustive-deps
|
||||
}, [active, selfUserId, idsKey]);
|
||||
|
||||
useEffect(() => {
|
||||
void refresh();
|
||||
if (!active || !selfUserId) return;
|
||||
|
||||
// Listen on both tables. RLS already enforces visibility.
|
||||
const reads = supabase
|
||||
.channel('group-reads:' + selfUserId)
|
||||
.on(
|
||||
'postgres_changes',
|
||||
{ event: 'INSERT', schema: 'public', table: 'message_reads' },
|
||||
() => {
|
||||
void refresh();
|
||||
},
|
||||
)
|
||||
.subscribe();
|
||||
const deliveries = supabase
|
||||
.channel('group-deliv:' + selfUserId)
|
||||
.on(
|
||||
'postgres_changes',
|
||||
{ event: 'INSERT', schema: 'public', table: 'message_deliveries' },
|
||||
() => {
|
||||
void refresh();
|
||||
},
|
||||
)
|
||||
.subscribe();
|
||||
return () => {
|
||||
void supabase.removeChannel(reads);
|
||||
void supabase.removeChannel(deliveries);
|
||||
};
|
||||
}, [active, selfUserId, refresh]);
|
||||
|
||||
return { deliveredByMessage, readByMessage, refresh };
|
||||
}
|
||||
|
||||
function toIdSets(input: Map<string, Map<string, string>>): Map<string, Set<string>> {
|
||||
const out = new Map<string, Set<string>>();
|
||||
for (const [mid, inner] of input) {
|
||||
out.set(mid, new Set(inner.keys()));
|
||||
}
|
||||
return out;
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
import {
|
||||
listPeerDeliveriesForMessages,
|
||||
markMessagesDelivered,
|
||||
} from '@chat-app/shared/chat';
|
||||
import { useCallback, useEffect, useMemo, useState } from 'react';
|
||||
|
||||
import { supabase } from './supabase';
|
||||
|
||||
// Tracks peer delivery acknowledgements for our own messages. Single peer
|
||||
// (DM). For groups this would need a per-user map; keep parity with
|
||||
// `useMessageReads` for now.
|
||||
export function useMessageDeliveries(
|
||||
messageIds: string[],
|
||||
peerUserId: string | undefined,
|
||||
): { peerDeliveredSet: Set<string>; refresh: () => Promise<void> } {
|
||||
const idsKey = useMemo(() => messageIds.join(','), [messageIds]);
|
||||
const [peerDeliveredSet, setPeerDeliveredSet] = useState<Set<string>>(new Set());
|
||||
|
||||
const refresh = useCallback(async () => {
|
||||
if (!peerUserId || messageIds.length === 0) {
|
||||
setPeerDeliveredSet(new Set());
|
||||
return;
|
||||
}
|
||||
try {
|
||||
const s = await listPeerDeliveriesForMessages(supabase, messageIds, peerUserId);
|
||||
setPeerDeliveredSet(s);
|
||||
} catch (err: unknown) {
|
||||
console.error('listPeerDeliveriesForMessages failed', err);
|
||||
}
|
||||
// eslint-disable-next-line react-hooks/exhaustive-deps
|
||||
}, [peerUserId, idsKey]);
|
||||
|
||||
useEffect(() => {
|
||||
void refresh();
|
||||
if (!peerUserId) return;
|
||||
const channel = supabase
|
||||
.channel('deliveries:' + peerUserId)
|
||||
.on(
|
||||
'postgres_changes',
|
||||
{
|
||||
event: 'INSERT',
|
||||
schema: 'public',
|
||||
table: 'message_deliveries',
|
||||
filter: 'user_id=eq.' + peerUserId,
|
||||
},
|
||||
() => {
|
||||
void refresh();
|
||||
},
|
||||
)
|
||||
.subscribe();
|
||||
return () => {
|
||||
void supabase.removeChannel(channel);
|
||||
};
|
||||
}, [peerUserId, refresh]);
|
||||
|
||||
return { peerDeliveredSet, refresh };
|
||||
}
|
||||
|
||||
export async function markDelivered(messageIds: string[]): Promise<void> {
|
||||
await markMessagesDelivered(supabase, messageIds);
|
||||
}
|
||||
@@ -0,0 +1,87 @@
|
||||
// Web Push subscription registration. Browser-only; Tauri WebView2 / WebKit
|
||||
// do not install service workers. The Tauri build relies on native OS
|
||||
// notifications via osNotify.ts instead.
|
||||
//
|
||||
// Flow:
|
||||
// 1. Register /sw.js if not already.
|
||||
// 2. Subscribe with the VAPID public key (Vite injects it via env).
|
||||
// 3. Persist {endpoint, keys} JSON-encoded into push_tokens.token for the
|
||||
// current device. Server-side fan-out reads this row to send pushes.
|
||||
|
||||
import { isTauriRuntime } from './globalShortcut';
|
||||
import { supabase } from './supabase';
|
||||
|
||||
const VAPID_PUBLIC_KEY: string | undefined =
|
||||
(import.meta as unknown as { env?: { VITE_VAPID_PUBLIC_KEY?: string } }).env
|
||||
?.VITE_VAPID_PUBLIC_KEY;
|
||||
|
||||
export async function registerWebPush(deviceId: string): Promise<void> {
|
||||
// Tauri uses native notifications — no service worker.
|
||||
if (isTauriRuntime()) return;
|
||||
if (typeof navigator === 'undefined' || !('serviceWorker' in navigator)) return;
|
||||
if (!('PushManager' in window)) return;
|
||||
if (!VAPID_PUBLIC_KEY) {
|
||||
console.warn('VITE_VAPID_PUBLIC_KEY not set — push disabled');
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
const reg = await navigator.serviceWorker.register('/sw.js');
|
||||
await navigator.serviceWorker.ready;
|
||||
let sub = await reg.pushManager.getSubscription();
|
||||
if (!sub) {
|
||||
const permission = await Notification.requestPermission();
|
||||
if (permission !== 'granted') return;
|
||||
sub = await reg.pushManager.subscribe({
|
||||
userVisibleOnly: true,
|
||||
applicationServerKey: urlBase64ToUint8Array(VAPID_PUBLIC_KEY) as BufferSource,
|
||||
});
|
||||
}
|
||||
|
||||
const token = JSON.stringify({
|
||||
endpoint: sub.endpoint,
|
||||
keys: sub.toJSON().keys ?? {},
|
||||
});
|
||||
|
||||
const { error } = await supabase
|
||||
.from('push_tokens')
|
||||
.upsert(
|
||||
{ device_id: deviceId, platform: 'web' as never, token },
|
||||
{ onConflict: 'device_id' },
|
||||
);
|
||||
if (error) {
|
||||
console.warn('push_tokens upsert failed', error);
|
||||
}
|
||||
} catch (err: unknown) {
|
||||
console.error('registerWebPush failed', err);
|
||||
}
|
||||
}
|
||||
|
||||
export async function unregisterWebPush(deviceId: string): Promise<void> {
|
||||
if (typeof navigator === 'undefined' || !('serviceWorker' in navigator)) return;
|
||||
try {
|
||||
const reg = await navigator.serviceWorker.getRegistration();
|
||||
if (reg) {
|
||||
const sub = await reg.pushManager.getSubscription();
|
||||
if (sub) await sub.unsubscribe();
|
||||
}
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
try {
|
||||
await supabase.from('push_tokens').delete().eq('device_id', deviceId);
|
||||
} catch (err: unknown) {
|
||||
console.warn('push_tokens delete failed', err);
|
||||
}
|
||||
}
|
||||
|
||||
// VAPID public key arrives as URL-safe base64; PushManager expects a raw byte
|
||||
// array. Standard conversion.
|
||||
function urlBase64ToUint8Array(base64: string): Uint8Array {
|
||||
const padding = '='.repeat((4 - (base64.length % 4)) % 4);
|
||||
const padded = (base64 + padding).replace(/-/g, '+').replace(/_/g, '/');
|
||||
const raw = atob(padded);
|
||||
const out = new Uint8Array(raw.length);
|
||||
for (let i = 0; i < raw.length; i++) out[i] = raw.charCodeAt(i);
|
||||
return out;
|
||||
}
|
||||
Reference in New Issue
Block a user