a proof of concept realtime collaborative todo list
0

Configure Feed

Select the types of activity you want to include in your feed.

atodo / offline.js
5.8 kB 187 lines
1// Offline write queue with CID-aware retry. 2// 3// safePut / safeDelete try the network write immediately using the cached 4// swapRecord CID. On network failure the call is stashed in localStorage 5// (coalesced by collection+rkey) and replayed on `online` or at startup. 6// On InvalidSwap we refetch our PDS record, merge for CRDT collections, and 7// retry — so a stale queued write never regresses our repo. 8 9import { signal } from "@preact/signals"; 10import { putRecord, deleteRecord, getRecord, isSwapError } from "./atproto.js"; 11import { mergeRecords } from "./crdt.js"; 12import { session, ITEM_NSID } from "./state.js"; 13 14const queueKey = (did) => `atodo.queue.${did}`; 15 16const cids = signal(new Map()); 17export const queueSize = signal(0); 18 19const k = (collection, rkey) => `${collection}/${rkey}`; 20 21export const getCid = (collection, rkey) => cids.value.get(k(collection, rkey)); 22 23export function setCid(collection, rkey, cid) { 24 const next = new Map(cids.value); 25 if (cid) next.set(k(collection, rkey), cid); 26 else next.delete(k(collection, rkey)); 27 cids.value = next; 28} 29 30function loadQueue() { 31 const did = session.value?.did; 32 if (!did) return []; 33 try { 34 return JSON.parse(localStorage.getItem(queueKey(did))) ?? []; 35 } catch { 36 return []; 37 } 38} 39 40function saveQueue(q) { 41 const did = session.value?.did; 42 if (!did) return; 43 if (q.length === 0) localStorage.removeItem(queueKey(did)); 44 else localStorage.setItem(queueKey(did), JSON.stringify(q)); 45 queueSize.value = q.length; 46} 47 48export const getQueue = () => loadQueue(); 49 50function enqueue(entry) { 51 const q = loadQueue(); 52 // Preserve the expectedCid from the first offline edit — that's the CID 53 // the PDS had before we started buffering. Later in-flight edits coalesce. 54 const prior = q.find((e) => e.collection === entry.collection && e.rkey === entry.rkey); 55 const merged = prior ? { ...entry, expectedCid: prior.expectedCid } : entry; 56 saveQueue([ 57 ...q.filter((e) => !(e.collection === entry.collection && e.rkey === entry.rkey)), 58 merged, 59 ]); 60} 61 62function dequeue(collection, rkey) { 63 saveQueue(loadQueue().filter((e) => !(e.collection === collection && e.rkey === rkey))); 64} 65 66let forcedOffline = false; 67export function setForcedOffline(v) { forcedOffline = v; } 68 69const isNetwork = (err) => err instanceof TypeError; 70 71async function execute(op, collection, rkey, record, expectedCid) { 72 if (op === "put") { 73 const res = await putRecord(collection, rkey, record, expectedCid); 74 setCid(collection, rkey, res.cid); 75 return res.cid; 76 } 77 await deleteRecord(collection, rkey, expectedCid); 78 setCid(collection, rkey, null); 79 return null; 80} 81 82async function resolveSwap(op, collection, rkey, record) { 83 const did = session.value.did; 84 const fresh = await getRecord(did, collection, rkey).catch(() => null); 85 if (!fresh) { 86 if (op === "put") return execute("put", collection, rkey, record, null); 87 setCid(collection, rkey, null); 88 return null; 89 } 90 setCid(collection, rkey, fresh.cid); 91 if (op === "delete") return execute("delete", collection, rkey, fresh.cid); 92 const next = collection === ITEM_NSID ? mergeRecords([fresh.value, record]) : record; 93 return execute("put", collection, rkey, next, fresh.cid); 94} 95 96export async function safePut(collection, rkey, record) { 97 const expectedCid = getCid(collection, rkey); 98 if (forcedOffline) { 99 enqueue({ op: "put", collection, rkey, record, expectedCid }); 100 return null; 101 } 102 try { 103 return await execute("put", collection, rkey, record, expectedCid); 104 } catch (err) { 105 if (isNetwork(err)) { 106 enqueue({ op: "put", collection, rkey, record, expectedCid }); 107 return null; 108 } 109 if (!isSwapError(err)) throw err; 110 try { 111 return await resolveSwap("put", collection, rkey, record); 112 } catch (err2) { 113 if (!isNetwork(err2)) throw err2; 114 enqueue({ op: "put", collection, rkey, record, expectedCid }); 115 return null; 116 } 117 } 118} 119 120export async function safeDelete(collection, rkey) { 121 const expectedCid = getCid(collection, rkey); 122 if (forcedOffline) { 123 enqueue({ op: "delete", collection, rkey, expectedCid }); 124 return false; 125 } 126 try { 127 await execute("delete", collection, rkey, expectedCid); 128 return true; 129 } catch (err) { 130 if (isNetwork(err)) { 131 enqueue({ op: "delete", collection, rkey, expectedCid }); 132 return false; 133 } 134 if (!isSwapError(err)) throw err; 135 try { 136 await resolveSwap("delete", collection, rkey, null); 137 return true; 138 } catch (err2) { 139 if (!isNetwork(err2)) throw err2; 140 enqueue({ op: "delete", collection, rkey, expectedCid }); 141 return false; 142 } 143 } 144} 145 146let flushing = false; 147 148export async function flushQueue() { 149 if (flushing) return; 150 flushing = true; 151 try { 152 while (true) { 153 const q = loadQueue(); 154 if (q.length === 0) break; 155 const entry = q[0]; 156 try { 157 await execute(entry.op, entry.collection, entry.rkey, entry.record, entry.expectedCid); 158 dequeue(entry.collection, entry.rkey); 159 } catch (err) { 160 if (isNetwork(err)) return; 161 if (isSwapError(err)) { 162 try { 163 await resolveSwap(entry.op, entry.collection, entry.rkey, entry.record); 164 dequeue(entry.collection, entry.rkey); 165 } catch (err2) { 166 if (isNetwork(err2)) return; 167 console.error("flushQueue resolve", entry, err2); 168 dequeue(entry.collection, entry.rkey); 169 } 170 } else { 171 console.error("flushQueue", entry, err); 172 dequeue(entry.collection, entry.rkey); 173 } 174 } 175 } 176 } finally { 177 flushing = false; 178 } 179} 180 181export function initQueue() { 182 queueSize.value = loadQueue().length; 183} 184 185if (typeof window !== "undefined") { 186 window.addEventListener("online", () => flushQueue()); 187}