feat: offline outbox — unsent messages survive reload and retry (#112)
CI / Build & Quality Checks (pull_request) Successful in 6m4s
CI / Trigger Desktop Build (pull_request) Skipped
CI / Secret scan (gitleaks) (pull_request) Successful in 27s
CI / Docker image build & smoke test (pull_request) Skipped
CI / Playwright smoke (e2e) (pull_request) Successful in 10m28s
CI / Build & Quality Checks (pull_request) Successful in 6m4s
CI / Trigger Desktop Build (pull_request) Skipped
CI / Secret scan (gitleaks) (pull_request) Successful in 27s
CI / Docker image build & smoke test (pull_request) Skipped
CI / Playwright smoke (e2e) (pull_request) Successful in 10m28s
Until now a send that failed (offline, homeserver down, a blip) went straight to "Failed to send": nothing retried it, and a reload dropped it without trace (chronological pending ordering keeps local echoes in memory only). - Outbox (utils/outbox.ts + features/outbox/OutboxFeature): own message sends (text, stickers, reactions, polls; not call signalling or redactions) are mirrored to localStorage from their first local echo until the server confirms them or the user cancels. - After a reload they come back as local echoes via room.addPendingEvent, same shape as the SDK's own. Recent ones (< 1 h) are sent again with the same txnId; older ones come back as "Failed to send" for the user to retry or cancel. Ones the server already has (transaction id seen in /sync) are dropped, so no duplicates. - Retries: network failures (ConnectionError, 408/429/5xx) are re-sent when the connection returns (sync recovers or the browser goes back online), and after a blip while online (5 s, backing off, max 10 per message). Oldest first, in order per room. 4xx / consent / encryption failures are left to the user. - UI: a network failure while offline shows a clock, "Queued. Will send when you're back online" (thread view too), not the red ✕. The ✕ is now a button: click to retry. - Logout wipes the outbox with the other plaintext caches (the content is decrypted, like drafts). Tested end to end against a local Synapse (Chromium): offline → queued → sent once on reconnect; homeserver unreachable → queued → sent once; failed send → reload → sent once and shown once; server accepted but response lost → reload → no duplicate; 2 h old entry → failed, not sent, click ✕ → sent; cancel → gone after reload; encrypted room → restored message goes out as m.room.encrypted with no plaintext and decrypts; one-off failure retried by itself in ~5 s; no page errors. Unit tests for the pure parts; Playwright 20 passed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01PPmy3tPq869XDW4njjVaKA
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
c0c93213c1
commit
02d86caeb0
@@ -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<EventStatus | null> = new Set([
|
||||
EventStatus.SENDING,
|
||||
EventStatus.ENCRYPTING,
|
||||
EventStatus.QUEUED,
|
||||
EventStatus.NOT_SENT,
|
||||
]);
|
||||
const ONLINE: ReadonlySet<SyncState | null> = new Set([
|
||||
SyncState.Prepared,
|
||||
SyncState.Syncing,
|
||||
SyncState.Catchup,
|
||||
]);
|
||||
const OFFLINE: ReadonlySet<SyncState | null> = 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<string, { event: MatrixEvent; room: Room }>;
|
||||
autoRetries: Map<string, number>;
|
||||
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<MatrixClient, OutboxSession>();
|
||||
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<ReturnType<typeof setTimeout>>();
|
||||
|
||||
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<Room, MatrixEvent[]>();
|
||||
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;
|
||||
}
|
||||
Reference in New Issue
Block a user