a proof of concept realtime collaborative todo list
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}