Files
cinny/src/app/features/outbox/OutboxFeature.tsx
T

269 lines
9.1 KiB
TypeScript
Raw Normal View History

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;
}