diff --git a/src/app/features/outbox/OutboxFeature.tsx b/src/app/features/outbox/OutboxFeature.tsx new file mode 100644 index 000000000..c8c0d3a89 --- /dev/null +++ b/src/app/features/outbox/OutboxFeature.tsx @@ -0,0 +1,268 @@ +import { useSetAtom } from 'jotai'; +import { useEffect } from 'react'; +import { + ClientEvent, + ClientEventHandlerMap, + ConnectionError, + EventStatus, + KnownMembership, + MatrixClient, + MatrixError, + MatrixEvent, + Room, + RoomEvent, + RoomEventHandlerMap, + SyncState, +} from 'matrix-js-sdk'; +import { useMatrixClient } from '../../hooks/useMatrixClient'; +import { sendOfflineAtom } from '../../state/sendOffline'; +import { + OUTBOX_MAX_AUTO_RETRIES, + OutboxEntry, + isRetryableSendError, + loadOutbox, + outboxRetryDelayMs, + planRestore, + saveOutbox, + shouldKeepInOutbox, + withEntry, + withoutEntry, +} from '../../utils/outbox'; + +const PENDING: ReadonlySet = new Set([ + EventStatus.SENDING, + EventStatus.ENCRYPTING, + EventStatus.QUEUED, + EventStatus.NOT_SENT, +]); +const ONLINE: ReadonlySet = new Set([ + SyncState.Prepared, + SyncState.Syncing, + SyncState.Catchup, +]); +const OFFLINE: ReadonlySet = new Set([SyncState.Reconnecting, SyncState.Error]); + +/** The server already has it: its transaction id came back down /sync. */ +const isDelivered = (room: Room, txnId: string): boolean => { + const hasTxn = (events: MatrixEvent[]) => + events.some((e) => e.getUnsigned().transaction_id === txnId); + return ( + hasTxn(room.getLiveTimeline().getEvents()) || + room.getThreads().some((t) => hasTxn(t.liveTimeline.getEvents())) + ); +}; + +type OutboxSession = { + /** Own pending events of this session, by txnId. */ + live: Map; + autoRetries: Map; + restored: boolean; + retrying: boolean; +}; + +// Per client, not per mount: a remount must neither restore twice nor forget +// the events it is tracking. +const sessions = new WeakMap(); +const getSession = (mx: MatrixClient): OutboxSession => { + let session = sessions.get(mx); + if (!session) { + session = { live: new Map(), autoRetries: new Map(), restored: false, retrying: false }; + sessions.set(mx, session); + } + return session; +}; + +/** + * [Gitea #112] Offline outbox (see utils/outbox.ts): mirrors own message sends + * into localStorage until the server confirms them, puts unsent ones back + * after a reload, and re-sends network failures when the connection returns. + */ +export function OutboxFeature() { + const mx = useMatrixClient(); + const setOffline = useSetAtom(sendOfflineAtom); + + useEffect(() => { + const userId = mx.getSafeUserId(); + let entries = loadOutbox(userId); + const session = getSession(mx); + const { live, autoRetries } = session; + let disposed = false; + let offlineNow = false; + const timers = new Set>(); + + const isOffline = () => + OFFLINE.has(mx.getSyncState()) || + (typeof navigator !== 'undefined' && navigator.onLine === false); + + const persist = (next: OutboxEntry[]) => { + if (next === entries) return; + entries = next; + saveOutbox(userId, entries); + }; + + /** + * A network failure while we think we're online (a blip neither sync nor + * the browser noticed): try again shortly, backing off. Real outages are + * handled by the reconnect / back-online triggers below. + */ + const scheduleBlipRetry = (event: MatrixEvent) => { + const attempts = autoRetries.get(event.getTxnId() ?? '') ?? 0; + if (!isRetryableSendError(event.error) || attempts >= OUTBOX_MAX_AUTO_RETRIES || isOffline()) + return; + const timer = setTimeout(() => { + timers.delete(timer); + retryQueued(); + }, outboxRetryDelayMs(attempts)); + timers.add(timer); + }; + + const onLocalEcho: RoomEventHandlerMap[RoomEvent.LocalEchoUpdated] = (event, room) => { + const txnId = event.getTxnId(); + if (!txnId || event.getSender() !== userId) return; + if (PENDING.has(event.status)) { + if (!shouldKeepInOutbox(event.getType(), event.getContent())) return; + live.set(txnId, { event, room }); + if (event.status === EventStatus.NOT_SENT) scheduleBlipRetry(event); + persist( + withEntry(entries, { + txnId, + roomId: room.roomId, + threadId: event.threadRootId ?? null, + type: event.getType(), + content: event.getContent(), + ts: event.getTs(), + }), + ); + return; + } + // SENT, CANCELLED, or the remote echo replaced it (status null). + live.delete(txnId); + autoRetries.delete(txnId); + persist(withoutEntry(entries, txnId)); + }; + + /** Re-send network failures, oldest first; a room stops at its first failure. */ + const retryQueued = async () => { + if (session.retrying || disposed) return; + session.retrying = true; + try { + const byRoom = new Map(); + Array.from(live.values()) + .filter( + ({ event }) => + event.status === EventStatus.NOT_SENT && + isRetryableSendError(event.error) && + (autoRetries.get(event.getTxnId() ?? '') ?? 0) < OUTBOX_MAX_AUTO_RETRIES, + ) + .sort((a, b) => a.event.getTs() - b.event.getTs()) + .forEach(({ event, room }) => { + byRoom.set(room, [...(byRoom.get(room) ?? []), event]); + }); + await Promise.all( + Array.from(byRoom.entries()).map(async ([room, events]) => { + for (const event of events) { + if (disposed || event.status !== EventStatus.NOT_SENT) continue; + const txnId = event.getTxnId() ?? ''; + autoRetries.set(txnId, (autoRetries.get(txnId) ?? 0) + 1); + try { + // eslint-disable-next-line no-await-in-loop + await mx.resendEvent(event, room); + } catch (e) { + // Still offline: keep the rest of this room in order for the next try. + if (isRetryableSendError(e)) break; + } + } + }), + ); + } finally { + session.retrying = false; + } + }; + + const restore = () => { + if (session.restored) return; + session.restored = true; + const plan = planRestore( + entries, + Date.now(), + (roomId) => mx.getRoom(roomId)?.getMyMembership() === KnownMembership.Join, + (entry) => { + const room = mx.getRoom(entry.roomId); + return !!room && isDelivered(room, entry.txnId); + }, + ); + plan.drop.forEach((entry) => persist(withoutEntry(entries, entry.txnId))); + plan.restore.forEach((entry) => { + const room = mx.getRoom(entry.roomId); + if (!room || live.has(entry.txnId)) return; + // Same shape as the SDK's own local echo (client.sendCompleteEvent). + const event = new MatrixEvent({ + type: entry.type, + content: entry.content, + event_id: `~${entry.roomId}:${entry.txnId}`, + sender: userId, + room_id: entry.roomId, + origin_server_ts: entry.ts, + }); + const thread = entry.threadId ? room.getThread(entry.threadId) : undefined; + if (thread) event.setThread(thread); + event.setTxnId(entry.txnId); + event.setStatus(EventStatus.NOT_SENT); + if (plan.autoSend.has(entry.txnId)) { + // Recent: treat like a send that lost the network (shown as Queued + // while offline, re-sent below and on reconnect). + // (The SDK types `error` as MatrixError but stores any send error there.) + event.error = new ConnectionError( + 'not sent before the app was closed', + ) as unknown as MatrixError; + } else { + // Old: back as "Failed to send"; the user decides (Retry / Cancel). + autoRetries.set(entry.txnId, OUTBOX_MAX_AUTO_RETRIES); + } + try { + room.addPendingEvent(event, entry.txnId); // emits LocalEchoUpdated → `live` + } catch { + // Already pending under this txnId: nothing to restore. + } + }); + if (!isOffline()) retryQueued(); + }; + + const updateOffline = () => { + const offline = isOffline(); + // Back online (sync recovered or the browser says so): send what's queued. + if (offlineNow && !offline && session.restored) retryQueued(); + offlineNow = offline; + setOffline(offline); + }; + + const onSync: ClientEventHandlerMap[ClientEvent.Sync] = (state, prevState) => { + updateOffline(); + if (!ONLINE.has(state)) return; + if (!session.restored) { + restore(); + return; + } + if (OFFLINE.has(prevState) || prevState === SyncState.Catchup) retryQueued(); + }; + const onBrowserOnline = () => updateOffline(); + + mx.on(RoomEvent.LocalEchoUpdated, onLocalEcho); + mx.on(ClientEvent.Sync, onSync); + window.addEventListener('online', onBrowserOnline); + window.addEventListener('offline', onBrowserOnline); + updateOffline(); + if (ONLINE.has(mx.getSyncState())) restore(); + + return () => { + disposed = true; + timers.forEach((t) => clearTimeout(t)); + mx.off(RoomEvent.LocalEchoUpdated, onLocalEcho); + mx.off(ClientEvent.Sync, onSync); + window.removeEventListener('online', onBrowserOnline); + window.removeEventListener('offline', onBrowserOnline); + }; + }, [mx, setOffline]); + + return null; +} diff --git a/src/app/features/room/message/Message.tsx b/src/app/features/room/message/Message.tsx index 129237f5f..ec7f52490 100644 --- a/src/app/features/room/message/Message.tsx +++ b/src/app/features/room/message/Message.tsx @@ -26,6 +26,7 @@ import { config, } from 'folds'; import React, { + CSSProperties, FormEventHandler, MouseEventHandler, ReactNode, @@ -38,7 +39,7 @@ import { useHover, useFocusWithin } from 'react-aria'; import { MatrixEvent, Room, EventStatus } from 'matrix-js-sdk'; import { Relations } from 'matrix-js-sdk/lib/models/relations'; import classNames from 'classnames'; -import { useAtom } from 'jotai'; +import { useAtom, useAtomValue } from 'jotai'; import { RoomPinnedEventsEventContent } from 'matrix-js-sdk/lib/types'; import { AvatarBase, @@ -98,27 +99,39 @@ import { useBookmarks } from '../../../hooks/useBookmarks'; import { PresenceRingAvatar } from '../../../components/presence'; import { useLotusShareBase } from '../../../hooks/useLotusLinkBase'; import { AvatarDecoration } from '../../../components/avatar-decoration/AvatarDecoration'; +import { sendOfflineAtom } from '../../../state/sendOffline'; +import { isRetryableSendError } from '../../../utils/outbox'; // Delivery status indicator for own messages function DeliveryStatus({ - status, + mEvent, + room, lotusTerminal, }: { - status: string | null; + mEvent: MatrixEvent; + room: Room; lotusTerminal: boolean; }) { + const mx = useMatrixClient(); + const offline = useAtomValue(sendOfflineAtom); + const { status } = mEvent; if (status === null) return null; // confirmed by server — read receipts take over let iconSrc: IconSrc; let label: string; let colorStyle: string; const isSending = status === EventStatus.SENDING || status === EventStatus.ENCRYPTING; - if (status === EventStatus.NOT_SENT || status === EventStatus.CANCELLED) { + // [Gitea #112] A network failure while offline is queued, not failed: the + // outbox sends it again when the connection is back. + const queued = status === EventStatus.NOT_SENT && offline && isRetryableSendError(mEvent.error); + const failed = !queued && (status === EventStatus.NOT_SENT || status === EventStatus.CANCELLED); + if (failed) { iconSrc = Icons.Cross; - label = 'Failed to send'; + label = status === EventStatus.NOT_SENT ? 'Failed to send. Click to retry' : 'Failed to send'; colorStyle = lotusTerminal ? 'var(--lt-accent-red)' : color.Critical.Main; - } else if (status === EventStatus.QUEUED || isSending) { - iconSrc = Icons.Send; - label = isSending ? 'Sending...' : 'Queued'; + } else if (queued || status === EventStatus.QUEUED || isSending) { + iconSrc = queued ? Icons.RecentClock : Icons.Send; + if (queued) label = "Queued. Will send when you're back online"; + else label = isSending ? 'Sending...' : 'Queued'; colorStyle = lotusTerminal ? 'color-mix(in srgb, var(--lt-accent-cyan) 60%, transparent)' : color.Secondary.Main; @@ -129,27 +142,49 @@ function DeliveryStatus({ ? 'color-mix(in srgb, var(--lt-accent-cyan) 70%, transparent)' : color.Secondary.Main; } + const retryable = failed && status === EventStatus.NOT_SENT; + const handleRetry: MouseEventHandler = (evt) => { + evt.stopPropagation(); + if (mEvent.status === EventStatus.NOT_SENT) mx.resendEvent(mEvent, room).catch(() => undefined); + }; + const style: CSSProperties = { + display: 'inline-flex', + alignItems: 'center', + marginTop: '2px', + lineHeight: 1, + color: colorStyle, + opacity: 0.85, + userSelect: 'none', + ...(lotusTerminal && failed ? { textShadow: 'var(--lt-glow-red)' } : {}), + }; + const icon = ( + + + + ); + if (retryable) { + return ( + + ); + } return ( - - - - + + {icon} ); } @@ -1102,7 +1137,7 @@ export const Message = React.memo( /> )} {isMine && !mEvent.isState() && readReceiptUsers.length === 0 && ( - + )} ); diff --git a/src/app/features/room/thread/ThreadTimeline.tsx b/src/app/features/room/thread/ThreadTimeline.tsx index b24101603..14995cac4 100644 --- a/src/app/features/room/thread/ThreadTimeline.tsx +++ b/src/app/features/room/thread/ThreadTimeline.tsx @@ -33,6 +33,8 @@ import { Badge, Box, Chip, Icon, Icons, Line, Scroll, Spinner, Text, color, conf import classNames from 'classnames'; import { Opts as LinkifyOpts } from 'linkifyjs'; import { isKeyHotkey } from 'is-hotkey'; +import { isRetryableSendError } from '../../../utils/outbox'; +import { sendOfflineAtom } from '../../../state/sendOffline'; import { eventWithShortcode, factoryEventSentBy } from '../../../utils/matrix'; import { useMatrixClient } from '../../../hooks/useMatrixClient'; import { useVirtualPaginator, ItemRange } from '../../../hooks/useVirtualPaginator'; @@ -253,6 +255,7 @@ export type ThreadTimelineProps = { export function ThreadTimeline({ room, thread, editor }: ThreadTimelineProps) { const mx = useMatrixClient(); + const sendOffline = useAtomValue(sendOfflineAtom); const alive = useAlive(); const useAuthentication = useMediaAuthentication(); @@ -992,8 +995,12 @@ export function ThreadTimeline({ room, thread, editor }: ThreadTimelineProps) { const showEmptyReplies = ready && thread.length === 0; const renderPendingEvent = (mEvent: MatrixEvent) => { + // [Gitea #112] Network failures while offline are queued, not failed. + const queued = + mEvent.status === EventStatus.NOT_SENT && sendOffline && isRetryableSendError(mEvent.error); const failed = - mEvent.status === EventStatus.NOT_SENT || mEvent.status === EventStatus.CANCELLED; + !queued && + (mEvent.status === EventStatus.NOT_SENT || mEvent.status === EventStatus.CANCELLED); return (
)} + {queued && ( + + + Queued. Will send when you're back online + + + )}
); }; diff --git a/src/app/pages/client/ClientNonUIFeatures.tsx b/src/app/pages/client/ClientNonUIFeatures.tsx index 95ea1bc44..ec09b460f 100644 --- a/src/app/pages/client/ClientNonUIFeatures.tsx +++ b/src/app/pages/client/ClientNonUIFeatures.tsx @@ -31,6 +31,7 @@ import { allInvitesAtom } from '../../state/room-list/inviteList'; import { useMatrixClient } from '../../hooks/useMatrixClient'; import { useClientConfig } from '../../hooks/useClientConfig'; import { useHydrateMsgDrafts } from '../../hooks/useHydrateMsgDrafts'; +import { OutboxFeature } from '../../features/outbox/OutboxFeature'; import { useSearchCacheInvalidation } from '../../utils/searchCacheInvalidation'; import { ClockSkewMonitor } from '../../utils/clockSkew'; import { clockSkewAtom } from '../../state/clockSkew'; @@ -1107,6 +1108,7 @@ export function ClientNonUIFeatures({ children }: ClientNonUIFeaturesProps) { + diff --git a/src/app/state/plaintextCaches.test.ts b/src/app/state/plaintextCaches.test.ts index aa0630771..8d0eaa24c 100644 --- a/src/app/state/plaintextCaches.test.ts +++ b/src/app/state/plaintextCaches.test.ts @@ -36,6 +36,7 @@ test('clearPlaintextCaches removes every plaintext/PII localStorage key', () => store.set('cinny_recent_forward_targets_v1', '[]'); store.set('cinny_recent_gifs_v1', '[]'); store.set('cinny_recent_stickers_v1', '[]'); + store.set('lotus_outbox_v1', '{"userId":"@me:server","entries":[]}'); removed.length = 0; clearPlaintextCaches(); @@ -47,6 +48,7 @@ test('clearPlaintextCaches removes every plaintext/PII localStorage key', () => 'cinny_recent_forward_targets_v1', 'cinny_recent_gifs_v1', 'cinny_recent_stickers_v1', + 'lotus_outbox_v1', // [Gitea #112] unsent message content ]) { assert.ok(removed.includes(key), `${key} cleared`); } diff --git a/src/app/state/plaintextCaches.ts b/src/app/state/plaintextCaches.ts index 8efa97e3a..83f640c0c 100644 --- a/src/app/state/plaintextCaches.ts +++ b/src/app/state/plaintextCaches.ts @@ -7,6 +7,7 @@ import { clearRecentStickers } from './recentStickers'; import { clearNavToActivePathStore } from './navToActivePath'; import { DRAFT_MSG_KEY_PREFIX } from '../utils/draft'; import { clearCallSession } from '../utils/callRejoin'; +import { clearOutbox } from '../utils/outbox'; /** * [Gitea #41] Wipe every persisted composer draft (`draft-msg-`). Drafts @@ -44,6 +45,7 @@ const clearMsgDrafts = (): void => { * - `cinny_recent_forward_targets_v1` — recent forward contact/room graph (PII) * - `cinny_recent_gifs_v1` / `cinny_recent_stickers_v1` — media the user sent * - `navToActivePath` — per-space last-visited room paths (needs userId) + * - `lotus_outbox_v1` — decrypted content of messages not yet sent ([Gitea #112]) * - `draft-msg-*` — unsent composer drafts (decrypted message text, unscoped by * user — see [Gitea #41]; previously deliberately preserved across logout * (N98), which let the next account on this device see/send a prior user's @@ -93,6 +95,7 @@ export const clearPlaintextCaches = (userId?: string): void => { clearRecentGifs(); clearRecentStickers(); clearMsgDrafts(); + clearOutbox(); clearCallSession(); clearStatusMessage(); if (userId) clearNavToActivePathStore(userId); diff --git a/src/app/state/sendOffline.ts b/src/app/state/sendOffline.ts new file mode 100644 index 000000000..416f9b115 --- /dev/null +++ b/src/app/state/sendOffline.ts @@ -0,0 +1,9 @@ +import { atom } from 'jotai'; + +/** + * [Gitea #112] True while the client can't reach the homeserver (sync is + * reconnecting / erroring, or the browser reports offline). Messages that + * failed for network reasons show "Queued" instead of "Failed to send" while + * this is set; the outbox sends them again once it clears. + */ +export const sendOfflineAtom = atom(false); diff --git a/src/app/utils/outbox.test.ts b/src/app/utils/outbox.test.ts new file mode 100644 index 000000000..873ed1b69 --- /dev/null +++ b/src/app/utils/outbox.test.ts @@ -0,0 +1,131 @@ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { + OUTBOX_AUTOSEND_MS, + OUTBOX_MAX_ENTRIES, + OutboxEntry, + isRetryableSendError, + outboxRetryDelayMs, + parseOutbox, + planRestore, + shouldKeepInOutbox, + withEntry, + withoutEntry, +} from './outbox'; + +const entry = (txnId: string, ts: number, roomId = '!a:hs'): OutboxEntry => ({ + txnId, + roomId, + threadId: null, + type: 'm.room.message', + content: { msgtype: 'm.text', body: txnId }, + ts, +}); + +test('keeps message-like events, not call signalling or redactions', () => { + assert.equal(shouldKeepInOutbox('m.room.message', { body: 'hi' }), true); + assert.equal(shouldKeepInOutbox('m.reaction', { 'm.relates_to': { event_id: '$x' } }), true); + assert.equal(shouldKeepInOutbox('org.matrix.msc3381.poll.start', {}), true); + assert.equal(shouldKeepInOutbox('org.matrix.msc4075.rtc.notification', {}), false); + assert.equal(shouldKeepInOutbox('m.call.invite', {}), false); + assert.equal(shouldKeepInOutbox('m.room.redaction', {}), false); +}); + +test('skips events that point at another unsent message (local id)', () => { + assert.equal( + shouldKeepInOutbox('m.reaction', { 'm.relates_to': { event_id: '~!a:hs:m123' } }), + false, + ); + assert.equal( + shouldKeepInOutbox('m.room.message', { + body: 'reply', + 'm.relates_to': { 'm.in_reply_to': { event_id: '~!a:hs:m1' } }, + }), + false, + ); + assert.equal( + shouldKeepInOutbox('m.room.message', { + body: 'reply', + 'm.relates_to': { 'm.in_reply_to': { event_id: '$real' } }, + }), + true, + ); +}); + +test('retryable: connection loss, timeouts, rate limits, server errors', () => { + const connectionError = new Error('fetch failed'); + Object.defineProperty(connectionError, 'name', { value: 'ConnectionError' }); + assert.equal(isRetryableSendError(connectionError), true); + assert.equal(isRetryableSendError({ httpStatus: 429 }), true); + assert.equal(isRetryableSendError({ httpStatus: 408 }), true); + assert.equal(isRetryableSendError({ httpStatus: 502 }), true); +}); + +test('not retryable: client errors, consent, unknown or missing errors', () => { + assert.equal(isRetryableSendError({ httpStatus: 403, errcode: 'M_FORBIDDEN' }), false); + assert.equal(isRetryableSendError({ httpStatus: 403, errcode: 'M_CONSENT_NOT_GIVEN' }), false); + assert.equal(isRetryableSendError({ httpStatus: 400 }), false); + assert.equal(isRetryableSendError(new Error('encryption failed')), false); + assert.equal(isRetryableSendError(undefined), false); + assert.equal(isRetryableSendError(null), false); +}); + +test('parse: only this user, only well-formed entries, junk tolerated', () => { + const good = entry('m1', 1); + const raw = JSON.stringify({ userId: '@me:hs', entries: [good, { txnId: 5 }, null] }); + assert.deepEqual(parseOutbox(raw, '@me:hs'), [good]); + assert.deepEqual(parseOutbox(raw, '@other:hs'), []); + assert.deepEqual(parseOutbox('{not json', '@me:hs'), []); + assert.deepEqual(parseOutbox(null, '@me:hs'), []); +}); + +test('withEntry: first write wins, sorted, capped (oldest dropped)', () => { + const a = entry('a', 2); + let list = withEntry([], a); + assert.equal(withEntry(list, { ...a, content: { body: 'changed' } }), list); + list = withEntry(list, entry('b', 1)); + assert.deepEqual( + list.map((e) => e.txnId), + ['b', 'a'], + ); + let many: OutboxEntry[] = []; + for (let i = 0; i < OUTBOX_MAX_ENTRIES + 5; i += 1) many = withEntry(many, entry(`t${i}`, i)); + assert.equal(many.length, OUTBOX_MAX_ENTRIES); + assert.equal(many[0].txnId, 't5'); +}); + +test('withoutEntry: removes, and returns the same list when absent', () => { + const list = [entry('a', 1), entry('b', 2)]; + assert.deepEqual( + withoutEntry(list, 'a').map((e) => e.txnId), + ['b'], + ); + assert.equal(withoutEntry(list, 'zzz'), list); +}); + +test('restore plan: drops delivered / left rooms, auto-sends only recent ones', () => { + const now = 10 * OUTBOX_AUTOSEND_MS; + const recent = entry('recent', now - 60_000); + const old = entry('old', now - OUTBOX_AUTOSEND_MS - 1); + const delivered = entry('delivered', now - 1000); + const left = entry('left', now - 1000, '!left:hs'); + const plan = planRestore( + [recent, left, old, delivered], + now, + (roomId) => roomId !== '!left:hs', + (e) => e.txnId === 'delivered', + ); + assert.deepEqual( + plan.restore.map((e) => e.txnId), + ['old', 'recent'], + ); + assert.deepEqual([...plan.autoSend], ['recent']); + assert.deepEqual(plan.drop.map((e) => e.txnId).sort(), ['delivered', 'left']); +}); + +test('blip retry delay backs off and is capped at a minute', () => { + assert.deepEqual( + [0, 1, 2, 3, 4, 9].map(outboxRetryDelayMs), + [5000, 10000, 20000, 40000, 60000, 60000], + ); +}); diff --git a/src/app/utils/outbox.ts b/src/app/utils/outbox.ts new file mode 100644 index 000000000..d0b381fc7 --- /dev/null +++ b/src/app/utils/outbox.ts @@ -0,0 +1,182 @@ +import { IContent } from 'matrix-js-sdk'; + +/** + * [Gitea #112] Offline outbox: unsent messages survive a reload and are + * retried when the connection comes back. + * + * The SDK keeps local echoes in memory only (chronological pending ordering), + * so a reload used to drop every message that hadn't reached the server. Each + * own message send is mirrored here (type + clear content + txnId) from its + * first local echo until the server confirms it or the user cancels it. On + * the next start it is put back as a failed local echo (Retry / Cancel as + * usual) and, if it's recent, sent again with the SAME txnId — the server + * deduplicates a transaction it already accepted. + * + * The content is the decrypted message, like a composer draft: it lives in + * localStorage until sent and is wiped on logout (clearPlaintextCaches). + */ + +export type OutboxEntry = { + txnId: string; + roomId: string; + threadId: string | null; + type: string; + content: IContent; + /** When the message was first sent (local echo timestamp, ms). */ + ts: number; +}; + +type Stored = { userId: string; entries: OutboxEntry[] }; + +const STORAGE_KEY = 'lotus_outbox_v1'; +/** Hard cap so a long offline stretch can't grow localStorage without bound. */ +export const OUTBOX_MAX_ENTRIES = 100; +/** + * Restored messages younger than this are sent again automatically after a + * reload; older ones come back as "Failed to send" and wait for the user + * (sending a message typed hours ago without asking would surprise people). + */ +export const OUTBOX_AUTOSEND_MS = 60 * 60 * 1000; +/** Auto-retries per message per session (one per reconnect). */ +export const OUTBOX_MAX_AUTO_RETRIES = 10; + +/** Delay before automatic retry number `attempt` (0-based) after a blip. */ +export const outboxRetryDelayMs = (attempt: number): number => + Math.min(5000 * 2 ** Math.max(0, attempt), 60_000); + +/** Message-like events worth keeping. Not call signalling, not redactions. */ +export const OUTBOX_EVENT_TYPES: ReadonlySet = new Set([ + 'm.room.message', + 'm.sticker', + 'm.reaction', + 'm.poll.start', + 'm.poll.response', + 'm.poll.end', + 'org.matrix.msc3381.poll.start', + 'org.matrix.msc3381.poll.response', + 'org.matrix.msc3381.poll.end', +]); + +const relatesToPendingEvent = (content: IContent): boolean => { + const rel = content['m.relates_to'] as + | { event_id?: unknown; 'm.in_reply_to'?: { event_id?: unknown } } + | undefined; + const ids = [rel?.event_id, rel?.['m.in_reply_to']?.event_id]; + // A local echo's id ("~!room:txn") means nothing after a reload. + return ids.some((id) => typeof id === 'string' && id.startsWith('~')); +}; + +/** Whether a send should be mirrored into the outbox. */ +export const shouldKeepInOutbox = (type: string, content: IContent): boolean => + OUTBOX_EVENT_TYPES.has(type) && !relatesToPendingEvent(content); + +/** + * A failure the network is to blame for: no connection (the SDK's + * ConnectionError), a timeout, rate limiting or a server error. Anything else + * (403, consent, bad request, encryption failure) needs the user. + */ +export const isRetryableSendError = (err: unknown): boolean => { + if (!err || typeof err !== 'object') return false; + const { name, httpStatus } = err as { name?: unknown; httpStatus?: unknown }; + if (name === 'ConnectionError') return true; + if (typeof httpStatus !== 'number') return false; + return httpStatus === 408 || httpStatus === 429 || httpStatus >= 500; +}; + +const isEntry = (e: unknown): e is OutboxEntry => { + if (!e || typeof e !== 'object') return false; + const o = e as Record; + return ( + typeof o.txnId === 'string' && + typeof o.roomId === 'string' && + (o.threadId === null || typeof o.threadId === 'string') && + typeof o.type === 'string' && + !!o.content && + typeof o.content === 'object' && + typeof o.ts === 'number' + ); +}; + +/** Parse the stored outbox, keeping only this user's well-formed entries. */ +export const parseOutbox = (raw: string | null, userId: string): OutboxEntry[] => { + if (!raw) return []; + try { + const parsed = JSON.parse(raw) as Partial; + if (parsed.userId !== userId || !Array.isArray(parsed.entries)) return []; + return parsed.entries.filter(isEntry); + } catch { + return []; + } +}; + +/** Add (or keep) an entry: first write wins, oldest dropped past the cap. */ +export const withEntry = (entries: OutboxEntry[], entry: OutboxEntry): OutboxEntry[] => { + if (entries.some((e) => e.txnId === entry.txnId)) return entries; + const next = [...entries, entry].sort((a, b) => a.ts - b.ts); + return next.length > OUTBOX_MAX_ENTRIES ? next.slice(next.length - OUTBOX_MAX_ENTRIES) : next; +}; + +export const withoutEntry = (entries: OutboxEntry[], txnId: string): OutboxEntry[] => + entries.some((e) => e.txnId === txnId) ? entries.filter((e) => e.txnId !== txnId) : entries; + +export type RestorePlan = { + /** Put back as failed local echoes, oldest first. */ + restore: OutboxEntry[]; + /** Subset of `restore` to send again right away. */ + autoSend: Set; + /** Already delivered or no longer sendable: forget them. */ + drop: OutboxEntry[]; +}; + +/** + * Decide what to do with the stored outbox on start. + * `canSend(roomId)`: the user is still joined; `delivered(entry)`: the + * server already has it (the transaction id came back down /sync). + */ +export const planRestore = ( + entries: OutboxEntry[], + now: number, + canSend: (roomId: string) => boolean, + delivered: (entry: OutboxEntry) => boolean, +): RestorePlan => { + const restore: OutboxEntry[] = []; + const autoSend = new Set(); + const drop: OutboxEntry[] = []; + [...entries] + .sort((a, b) => a.ts - b.ts) + .forEach((entry) => { + if (!canSend(entry.roomId) || delivered(entry)) { + drop.push(entry); + return; + } + restore.push(entry); + if (now - entry.ts < OUTBOX_AUTOSEND_MS) autoSend.add(entry.txnId); + }); + return { restore, autoSend, drop }; +}; + +export const loadOutbox = (userId: string): OutboxEntry[] => { + try { + return parseOutbox(localStorage.getItem(STORAGE_KEY), userId); + } catch { + return []; + } +}; + +export const saveOutbox = (userId: string, entries: OutboxEntry[]): void => { + try { + if (entries.length === 0) localStorage.removeItem(STORAGE_KEY); + else localStorage.setItem(STORAGE_KEY, JSON.stringify({ userId, entries } satisfies Stored)); + } catch { + // Storage full or blocked: the outbox is best-effort. + } +}; + +/** Wipe the outbox (logout): it holds decrypted message content. */ +export const clearOutbox = (): void => { + try { + localStorage.removeItem(STORAGE_KEY); + } catch { + /* localStorage unavailable — nothing to clear */ + } +};