// Offline write queue with CID-aware retry. // // safePut / safeDelete try the network write immediately using the cached // swapRecord CID. On network failure the call is stashed in localStorage // (coalesced by collection+rkey) and replayed on `online` or at startup. // On InvalidSwap we refetch our PDS record, merge for CRDT collections, and // retry — so a stale queued write never regresses our repo. import { signal } from "@preact/signals"; import { putRecord, deleteRecord, getRecord, isSwapError } from "./atproto.js"; import { mergeRecords } from "./crdt.js"; import { session, ITEM_NSID } from "./state.js"; const queueKey = (did) => `atodo.queue.${did}`; const cids = signal(new Map()); export const queueSize = signal(0); const k = (collection, rkey) => `${collection}/${rkey}`; export const getCid = (collection, rkey) => cids.value.get(k(collection, rkey)); export function setCid(collection, rkey, cid) { const next = new Map(cids.value); if (cid) next.set(k(collection, rkey), cid); else next.delete(k(collection, rkey)); cids.value = next; } function loadQueue() { const did = session.value?.did; if (!did) return []; try { return JSON.parse(localStorage.getItem(queueKey(did))) ?? []; } catch { return []; } } function saveQueue(q) { const did = session.value?.did; if (!did) return; if (q.length === 0) localStorage.removeItem(queueKey(did)); else localStorage.setItem(queueKey(did), JSON.stringify(q)); queueSize.value = q.length; } export const getQueue = () => loadQueue(); function enqueue(entry) { const q = loadQueue(); // Preserve the expectedCid from the first offline edit — that's the CID // the PDS had before we started buffering. Later in-flight edits coalesce. const prior = q.find((e) => e.collection === entry.collection && e.rkey === entry.rkey); const merged = prior ? { ...entry, expectedCid: prior.expectedCid } : entry; saveQueue([ ...q.filter((e) => !(e.collection === entry.collection && e.rkey === entry.rkey)), merged, ]); } function dequeue(collection, rkey) { saveQueue(loadQueue().filter((e) => !(e.collection === collection && e.rkey === rkey))); } let forcedOffline = false; export function setForcedOffline(v) { forcedOffline = v; } const isNetwork = (err) => err instanceof TypeError; async function execute(op, collection, rkey, record, expectedCid) { if (op === "put") { const res = await putRecord(collection, rkey, record, expectedCid); setCid(collection, rkey, res.cid); return res.cid; } await deleteRecord(collection, rkey, expectedCid); setCid(collection, rkey, null); return null; } async function resolveSwap(op, collection, rkey, record) { const did = session.value.did; const fresh = await getRecord(did, collection, rkey).catch(() => null); if (!fresh) { if (op === "put") return execute("put", collection, rkey, record, null); setCid(collection, rkey, null); return null; } setCid(collection, rkey, fresh.cid); if (op === "delete") return execute("delete", collection, rkey, fresh.cid); const next = collection === ITEM_NSID ? mergeRecords([fresh.value, record]) : record; return execute("put", collection, rkey, next, fresh.cid); } export async function safePut(collection, rkey, record) { const expectedCid = getCid(collection, rkey); if (forcedOffline) { enqueue({ op: "put", collection, rkey, record, expectedCid }); return null; } try { return await execute("put", collection, rkey, record, expectedCid); } catch (err) { if (isNetwork(err)) { enqueue({ op: "put", collection, rkey, record, expectedCid }); return null; } if (!isSwapError(err)) throw err; try { return await resolveSwap("put", collection, rkey, record); } catch (err2) { if (!isNetwork(err2)) throw err2; enqueue({ op: "put", collection, rkey, record, expectedCid }); return null; } } } export async function safeDelete(collection, rkey) { const expectedCid = getCid(collection, rkey); if (forcedOffline) { enqueue({ op: "delete", collection, rkey, expectedCid }); return false; } try { await execute("delete", collection, rkey, expectedCid); return true; } catch (err) { if (isNetwork(err)) { enqueue({ op: "delete", collection, rkey, expectedCid }); return false; } if (!isSwapError(err)) throw err; try { await resolveSwap("delete", collection, rkey, null); return true; } catch (err2) { if (!isNetwork(err2)) throw err2; enqueue({ op: "delete", collection, rkey, expectedCid }); return false; } } } let flushing = false; export async function flushQueue() { if (flushing) return; flushing = true; try { while (true) { const q = loadQueue(); if (q.length === 0) break; const entry = q[0]; try { await execute(entry.op, entry.collection, entry.rkey, entry.record, entry.expectedCid); dequeue(entry.collection, entry.rkey); } catch (err) { if (isNetwork(err)) return; if (isSwapError(err)) { try { await resolveSwap(entry.op, entry.collection, entry.rkey, entry.record); dequeue(entry.collection, entry.rkey); } catch (err2) { if (isNetwork(err2)) return; console.error("flushQueue resolve", entry, err2); dequeue(entry.collection, entry.rkey); } } else { console.error("flushQueue", entry, err); dequeue(entry.collection, entry.rkey); } } } } finally { flushing = false; } } export function initQueue() { queueSize.value = loadQueue().length; } if (typeof window !== "undefined") { window.addEventListener("online", () => flushQueue()); }