1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316 | /**
* Characterization scenarios for lib/mailbox's CONCURRENCY, driven through the
* `_` surface at the bottom of the file.
*
* These are not cross-implementation vectors: unlike vectors/mailbox.json (which
* came off the TypeScript and pins wire bytes), these record what the shipped
* ClojureScript actually does with time and re-entrancy, so that the reshaping
* Phase 6 has in mind for this file cannot change it silently. `mailbox` shipped
* two defects that took every room down — rooms would not open at all, and
* offline delivery had never once worked — and the suite was green either side
* of both fixes, because nothing in it could express "and then the retry timer
* fired".
*
* Each scenario returns ONE interleaved transcript: clock ops (from
* helpers/det-clock.mjs) and effects (from helpers/effect-log.mjs) share a sink,
* so a single deepEqual pins re-entrancy, side-effect ORDER and the backoff
* ladder at once. `["script", …]` rows are the harness narrating what it just
* did, so the transcript reads as a story rather than a list.
*
* Storage keys are deliberately the short literals "K.queue"/"K.cursor"/… and
* not `nsKey(…)`: the real key names are a storage contract already pinned by
* vectors/config.json and vectors/mailbox.json, and repeating them here would
* only make this transcript harder to read.
*/
import { installClock, makeSink } from "./det-clock.mjs";
import { effectLog, effectFn } from "./effect-log.mjs";
const KEYS = { cursor: "K.cursor", queue: "K.queue", sent: "K.sent", retry: "K.retry", ident: "K.ident" };
/** Two real dag-cbor CIDs, so `CID.parse(hash, base58btc)` inside build-payload
* sees exactly what OrbitDB produces. */
export async function fixtures() {
const { CID } = await import("multiformats/cid");
const { base58btc } = await import("multiformats/bases/base58");
const { sha256 } = await import("multiformats/hashes/sha2");
const mk = async (bytes) => {
const block = new Uint8Array(bytes);
const cid = CID.createV1(0x71, await sha256.digest(block));
return { block, cid, hash: cid.toString(base58btc) };
};
return { A: await mk([0xa1, 0x61, 0x61, 0x01]), B: await mk([0xa1, 0x61, 0x62, 0x02]) };
}
const b64 = (u8) => Buffer.from(u8).toString("base64");
const cat = (...parts) => {
const out = new Uint8Array(parts.reduce((n, p) => n + p.length, 0));
let off = 0;
for (const p of parts) { out.set(p, off); off += p.length; }
return out;
};
/** A no-op RoomLog double; each scenario overrides the handful it cares about. */
const baseLog = () => ({
async sealedEntryBytes() { return null; },
async rawBlock() { return null; },
myIdentityBlock() { return null; },
async imageDagBlocks() { return []; },
async hasEntry() { return false; },
async decodeSealedEntry() { return null; },
async putEntryBlock() {},
async putRawBlock() {},
async putIdentityBlock() {},
async ingestEntry() { return "rejected"; },
async entryMeta() { return null; },
async emitUnseen() {},
async pinImageIfLocal() { return true; },
});
/**
* Install the clock, wrap localStorage and the window listener registry, run
* `body`, and hand back the shared transcript.
*/
async function scenario(localStorageStub, body) {
const sink = makeSink();
localStorageStub.clear();
const clock = installClock({ t0: 1_700_000_000_000, sink });
const realStorage = globalThis.localStorage;
const realAdd = globalThis.addEventListener;
const realRemove = globalThis.removeEventListener;
const listeners = new Map();
globalThis.localStorage = effectLog(realStorage, "storage", sink, { only: ["setItem", "removeItem"] });
globalThis.addEventListener = (evt, fn) => listeners.set(evt, fn);
globalThis.removeEventListener = (evt) => listeners.delete(evt);
const say = (...what) => sink.push(["script", ...what]);
try {
const value = await body({ clock, sink, say, listeners, keys: KEYS });
return { transcript: sink, ...value };
} finally {
globalThis.localStorage = realStorage;
globalThis.addEventListener = realAdd;
globalThis.removeEventListener = realRemove;
clock.uninstall();
}
}
/**
* ONE: the `_flushing` guard.
*
* A second enqueue arriving while a deposit is in flight must not start a second
* deposit — the queue is drained by exactly one walker. At quiescence the queue
* must be empty AND no timer may be left pending: that is the honest way to test
* the tail re-check at the bottom of `flush!`. Testing the mechanism would mean
* reading `_flushing`; testing the CONSEQUENCE means asserting nothing was
* stranded.
*/
export async function flushReentrancy({ MailboxSync, localStorageStub, priv }) {
const { A, B } = await fixtures();
return scenario(localStorageStub, async ({ clock, say, keys }) => {
const sink = clock.sink;
const gates = [];
const release = () => gates.shift()();
const client = effectLog({
deposit(payload) {
return new Promise((r) => gates.push(r)).then(() => ({ ok: true, deduped: false, seq: `s${payload.length}` }));
},
async replay() { return { ok: false, kind: "net" }; },
}, "client", sink);
const log = effectLog({
...baseLog(),
async sealedEntryBytes(h) { return h === A.hash ? A.block : h === B.hash ? B.block : null; },
}, "log", sink, { only: ["sealedEntryBytes"] });
const sync = new MailboxSync({ log, client, cipher: null, maxMessageKb: 64, keys, hasIdentity: false });
say("enqueue A");
sync.enqueueEntry(A.hash);
await clock.drain();
say("enqueue B while A's deposit is in flight");
sync.enqueueEntry(B.hash);
await clock.drain();
say("call _flush() again, explicitly");
priv(sync, "flush")();
await clock.drain();
say("release A");
release();
await clock.drain();
say("release B");
release();
await clock.drain();
return {
queued: JSON.parse(localStorageStub.getItem(keys.queue)).length,
sent: JSON.parse(localStorageStub.getItem(keys.sent)),
pending: clock.pending(),
};
});
}
/**
* TWO: the crash-safety invariant.
*
* lib/mailbox says it in a comment and nothing pinned it: an entry goes onto the
* retry list from the moment its block is STORED, before the replay cursor
* advances past its page and long before any join is attempted. A crash between
* store and join must not strand the entry behind an advanced cursor — that is a
* message silently lost forever, which is precisely the failure mode this file
* already shipped once.
*
* The transcript makes it an index assertion: add-retry < cursor advance <
* ingestEntry.
*/
export async function crashSafetyOrder({ MailboxSync, localStorageStub, priv }) {
const { A } = await fixtures();
return scenario(localStorageStub, async ({ clock, say, keys }) => {
const sink = clock.sink;
const entry = { hash: A.hash };
const pages = [
{ ok: true, page: { messages: [{ seq: "0001", id: "m1", ts: 1, payload: b64(cat(new Uint8Array([1]), A.cid.bytes, A.block)) }], hasMore: false, oldestSeq: "0001", latestSeq: "0001" } },
{ ok: false, kind: "net" },
];
const client = effectLog({
async deposit() { return { ok: true, deduped: false, seq: "s1" }; },
async replay() { return pages.shift() ?? { ok: false, kind: "net" }; },
}, "client", sink);
const log = effectLog({
...baseLog(),
async sealedEntryBytes() { return A.block; },
async decodeSealedEntry() { return entry; },
// "deferred" = stored but not joinable yet (its ancestors expired), the
// exact state the retry list exists for
async ingestEntry() { return "deferred"; },
async entryMeta() { return { hash: A.hash, from: "F", clock: 1, op: { t: "msg" } }; },
}, "log", sink);
const sync = new MailboxSync({
log, client, cipher: null, maxMessageKb: 64, keys, hasIdentity: false,
onDisplayOnly: effectFn("onDisplayOnly", sink),
onNotice: effectFn("onNotice", sink),
onReplayed: effectFn("onReplayed", sink),
});
say("replay(force)");
await priv(sync, "replay")(true);
await clock.drain();
return {
retry: JSON.parse(localStorageStub.getItem(keys.retry)),
cursor: localStorageStub.getItem(keys.cursor),
pending: clock.pending(),
};
});
}
/**
* THREE: the backoff ladder, `_attempt` reset, and `start`'s kick.
*
* 1000, 2000, 4000 … 256000, then the 300000 cap twice. Then a deposit succeeds
* and the NEXT failure schedules 1000 again — the consequence of `_attempt`
* being reset, tested without reading `_attempt` (which is also asserted, since
* it is free).
*
* Then `start`'s wake kick: it must clear the live retry timer and leave no
* orphan. dom.mjs's `window.addEventListener` is a no-op, so the scenario
* captures the registry itself and fires "visibilitychange" by hand.
*/
export async function backoffLadder({ MailboxSync, localStorageStub }) {
const { A, B } = await fixtures();
return scenario(localStorageStub, async ({ clock, say, listeners, keys }) => {
const sink = clock.sink;
let ok = false;
const client = effectLog({
async deposit() { return ok ? { ok: true, deduped: false, seq: "s1" } : { ok: false, kind: "net" }; },
async replay() { return { ok: false, kind: "net" }; },
}, "client", sink);
const log = effectLog({ ...baseLog(), async sealedEntryBytes() { return A.block; } }, "log", sink,
{ only: ["sealedEntryBytes"] });
const sync = new MailboxSync({ log, client, cipher: null, maxMessageKb: 64, keys, hasIdentity: false });
say("enqueue — every deposit will fail with kind:net");
sync.enqueueEntry(A.hash);
await clock.drain();
// walk the ladder: each advance lands exactly on the pending deadline
const delays = [];
for (let i = 0; i < 11; i++) {
const next = clock.pending()[0];
delays.push(next.delay);
await clock.advance(next.delay);
}
say("deposits start succeeding");
ok = true;
const armed = clock.pending()[0];
await clock.advance(armed.delay);
const attemptAfterSuccess = sync._attempt;
say("enqueue again, deposits failing again — the ladder must restart at 1000");
ok = false;
sync.enqueueEntry(B.hash);
await clock.drain();
const restart = clock.pending()[0];
// start()'s own flush finds _retryTimer already live, so it must NOT arm a
// second one — the restart timer is still the only thing pending
say("start()");
sync.start();
await clock.drain();
const beforeKick = clock.pending();
say("wake kick: visibilitychange");
listeners.get("visibilitychange")();
await clock.drain();
const afterKick = clock.pending();
sync.stop();
return {
delays,
attemptAfterSuccess,
restartDelay: restart.delay,
clearedByKick: restart.id,
beforeKick,
afterKick,
};
});
}
/**
* FOUR: the replay throttle.
*
* `_replay(false)` twice inside REPLAY-THROTTLE-MS must replay once;
* advancing past the window lets the next one through; `_replay(true)` bypasses
* the window entirely. `start` depends on the bypass — if `force` stopped
* working, boot would silently skip the first replay and offline delivery would
* be broken again in exactly the way it already was.
*/
export async function replayThrottle({ MailboxSync, localStorageStub, priv }) {
return scenario(localStorageStub, async ({ clock, say, keys }) => {
const sink = clock.sink;
const client = effectLog({
async deposit() { return { ok: true, deduped: false, seq: "s1" }; },
async replay() { return { ok: false, kind: "net" }; },
}, "client", sink);
const log = effectLog(baseLog(), "log", sink, { only: [] });
const sync = new MailboxSync({ log, client, cipher: null, maxMessageKb: 64, keys, hasIdentity: false });
const replay = priv(sync, "replay");
say("_replay(false) — no previous replay, so it runs");
await replay(false);
say("_replay(false) again, inside the 30s window");
await replay(false);
say("_replay(true) — force bypasses the window");
await replay(true);
say("advance 30s");
await clock.advance(30000);
say("_replay(false) — the window has passed");
await replay(false);
await clock.drain();
return { pending: clock.pending() };
});
}
|