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 | /**
* Characterization scenarios for lib/media's perfect-negotiation races.
*
* Before 3d this file had no `:testlib` export, no `_` seam and not one test —
* 445 lines of glare handling, transceiver reuse and ICE buffering, the largest
* single hole in the suite. Everything below drives the real class through the
* `_` vtable added in the same commit; the transcripts are the shared
* clock+effect sink helpers/det-clock.mjs and helpers/effect-log.mjs produce.
*
* Which seam each scenario uses, and why it must be that one:
*
* · the MUTEX goes through the PUBLIC frame handler (the callback the
* constructor registers with `net.handleFrames`), because the serialization
* lives in `on-sig-frame`'s `peer.queue = peer.queue.then(…)` — not in
* `handle-signal`. Two `_handleSignal` calls would interleave by design and
* would pin the wrong thing.
* · the `settingRemoteAnswer` WINDOW goes through `_handleSignal`, because the
* queue is exactly what closes that window: through the public path the
* second signal waits and the flag is false again by the time it is read.
*
* Politeness is `myId < peerId` by JS string comparison, so the peer-id pair is
* chosen per scenario: "AAA…" is polite against "ZZZ…", and the reverse.
*/
import { installClock, makeSink } from "./det-clock.mjs";
import { effectFn, summarize } from "./effect-log.mjs";
import { installRtc, fakeSigNet, FakeMediaStream, track } from "./rtc-fakes.mjs";
export const CALL_SIG = "/sueta/2/call-sig/1.0";
const SECRET = "s".repeat(43);
const SALT = "chat.example/v2";
const LOW = "12D3KooAAAApolite";
const HIGH = "12D3KooZZZZimpolite";
const te = new TextEncoder();
/**
* Drain to quiescence with the REAL macrotask timer: WebCrypto's seal/open take
* an unspecified number of turns, so a fixed count would be a fixture that
* depends on how fast node's threadpool is. Stopping when the transcript stops
* growing is the same rule helpers/app-fakes.mjs's `settle` uses.
*/
async function idle(sink, clock, max = 400) {
let prev = -1;
let stable = 0;
for (let i = 0; i < max; i++) {
await clock.drain();
if (sink.length === prev) {
// a generous margin: a seal that took longer than the quiescence window
// would let the next `say()` land BEFORE its sendFrame, and the fixture
// would encode how fast this machine's threadpool is
if (++stable > 12) return;
} else {
stable = 0;
prev = sink.length;
}
}
throw new Error("media scenario never quiesced");
}
/**
* Build a Media wired to fakes. `myId`/`peerId` decide politeness.
*/
async function harness({ Media, RoomCrypto }, myId, peerId) {
const sink = makeSink();
const clock = installClock({ t0: 1_700_000_000_000, sink });
const rtc = installRtc(sink);
const net = fakeSigNet(myId, sink);
const rcMe = await RoomCrypto.create(SECRET, SALT, myId);
const rcPeer = await RoomCrypto.create(SECRET, SALT, peerId);
const media = new Media({
crypto: rcMe,
ice: { get: () => Promise.resolve({ iceServers: [] }) },
net: net.node,
myName: () => "me",
events: { track: effectFn("ev.track", sink), mediaState: effectFn("ev.mediaState", sink) },
});
// lib/media swallows negotiation errors into console.warn — which is exactly
// where a real fault would hide. Route them into the transcript instead: the
// swallow becomes an assertable event rather than noise scrolling past.
const realWarn = console.warn;
console.warn = (...a) => sink.push(["console.warn", ...a.map(summarize)]);
const say = (...what) => sink.push(["script", ...what]);
/** A sealed frame from the peer, exactly as the wire delivers it. */
const frame = async (payload) => te.encode(JSON.stringify(await rcPeer.sealSignal(myId, payload)));
/** Deliver a frame through the PUBLIC handler, without awaiting it. */
const deliver = (bytes) => net.handlers.get(CALL_SIG)(peerId, bytes);
return {
sink, clock, rtc, net, media, say, frame, deliver, peerId,
idle: () => idle(sink, clock),
/** Decrypt everything lib/media sent, so the transcript's byte counts have
* a readable companion. */
async outgoing() {
const out = [];
for (const bytes of net.sent) {
const env = JSON.parse(new TextDecoder().decode(bytes));
const p = await rcPeer.openSignal(myId, env);
out.push(p === null ? null : { kind: p.kind, sdp: p.sdp ?? null, candidate: p.candidate ?? null });
}
return out;
},
teardown() {
console.warn = realWarn;
rtc.uninstall();
clock.uninstall();
},
};
}
/** attachMedia + let setup-pc! finish; returns the pc the code created. */
async function attached(h, kinds = ["audio", "video"]) {
h.say("attachMedia");
h.media.attachMedia(h.peerId, new FakeMediaStream(kinds.map((k) => track(k))), undefined);
await h.idle();
return h.rtc.created[0];
}
/**
* ONE and TWO: glare, both sides.
*
* A remote offer arriving while our own offer is outstanding is a collision.
* The POLITE side (smaller peer id) rolls back and accepts it; the IMPOLITE side
* sets `ignoreOffer` and drops it on the floor — and "drops it" has to mean the
* transcript contains NO setRemoteDescription at all, not merely that the end
* state looks right.
*/
export async function glare(M, { polite }) {
const [me, them] = polite ? [LOW, HIGH] : [HIGH, LOW];
const h = await harness(M, me, them);
try {
const pc = await attached(h);
h.say("our own offer goes out (onnegotiationneeded)");
pc.onnegotiationneeded();
await h.idle();
h.say(`a remote offer arrives while signalingState is ${pc.signalingState} — glare`);
const mark = h.sink.length;
h.deliver(await h.frame({ kind: "offer", sdp: "their-offer", name: "them" }));
await h.idle();
const afterGlare = h.sink.slice(mark);
return {
transcript: h.sink,
polite: h.media._peer(them).polite,
ignoreOffer: h.media._peer(them).ignoreOffer,
signalingState: pc.signalingState,
setRemoteAfterGlare: afterGlare.filter((e) => e[0] === "pc.setRemoteDescription").length,
outgoing: await h.outgoing(),
};
} finally {
h.teardown();
}
}
/**
* THREE: the `settingRemoteAnswer` window.
*
* `readyForOffer` is true when the state is "stable" OR we are mid-way through
* applying a remote ANSWER. That second clause is the whole reason an impolite
* peer can still accept an offer that arrives in "have-local-offer": the answer
* it is applying is about to make the state stable anyway. Close the window and
* the identical offer is dropped.
*
* Only reachable through `_handleSignal`: `on-sig-frame`'s per-peer queue is
* precisely what serializes the two signals and closes the window.
*/
export async function settingRemoteAnswerWindow(M) {
const h = await harness(M, HIGH, LOW); // impolite: the drop is the default
try {
const pc = await attached(h);
h.say("our own offer goes out");
pc.onnegotiationneeded();
await h.idle();
const peer = h.media._peer(LOW);
h.say("an ANSWER starts being applied, and hangs inside setRemoteDescription");
const release = h.rtc.gate("setRemoteDescription");
h.media._handleSignal(peer, { kind: "answer", sdp: "their-answer", name: "them" });
await h.idle();
const inWindow = { settingRemoteAnswer: peer.settingRemoteAnswer, signalingState: pc.signalingState };
h.say("an OFFER arrives inside that window — impolite, but accepted");
const mark = h.sink.length;
h.media._handleSignal(peer, { kind: "offer", sdp: "their-offer", name: "them" });
await h.idle();
const acceptedInWindow = h.sink.slice(mark).filter((e) => e[0] === "pc.setRemoteDescription").length;
const ignoreInWindow = peer.ignoreOffer;
h.say("release");
release();
await h.idle();
h.say("window closed: back in have-local-offer with settingRemoteAnswer false");
pc.signalingState = "have-local-offer";
const mark2 = h.sink.length;
await h.media._handleSignal(peer, { kind: "offer", sdp: "their-offer-2", name: "them" });
await h.idle();
return {
transcript: h.sink,
inWindow,
acceptedInWindow,
ignoreInWindow,
droppedOutsideWindow: h.sink.slice(mark2).filter((e) => e[0] === "pc.setRemoteDescription").length,
ignoreOutsideWindow: peer.ignoreOffer,
outgoing: await h.outgoing(),
};
} finally {
h.teardown();
}
}
/**
* FOUR: ICE candidates buffered before the remote description, flushed in
* ARRIVAL order.
*
* `addIceCandidate` before a remote description throws in the browser, so
* candidates queue on `pendingCandidates` and are spliced out the moment the
* description lands. Order matters: ICE is a priority-ordered list.
*/
export async function candidateBuffering(M) {
const h = await harness(M, LOW, HIGH);
try {
await attached(h);
const peer = h.media._peer(HIGH);
h.say("three candidates arrive before any remote description");
const mark = h.sink.length;
for (const c of ["c1", "c2", "c3"]) {
await h.media._handleSignal(peer, { kind: "ice", candidate: { candidate: c }, name: "them" });
}
await h.idle();
const addedEarly = h.sink.slice(mark).filter((e) => e[0] === "pc.addIceCandidate").length;
const buffered = peer.pendingCandidates.map((c) => c.candidate);
h.say("the offer arrives — the buffer flushes in arrival order");
await h.media._handleSignal(peer, { kind: "offer", sdp: "their-offer", name: "them" });
await h.idle();
h.say("a fourth candidate, now that remoteDescription is set");
await h.media._handleSignal(peer, { kind: "ice", candidate: { candidate: "c4" }, name: "them" });
await h.idle();
h.say("a null candidate — end-of-candidates, passed as undefined");
await h.media._handleSignal(peer, { kind: "ice", candidate: null, name: "them" });
await h.idle();
return {
transcript: h.sink,
addedEarly,
buffered,
flushOrder: h.sink.filter((e) => e[0] === "pc.addIceCandidate").map((e) => e[1]),
pendingAfter: peer.pendingCandidates.length,
outgoing: await h.outgoing(),
};
} finally {
h.teardown();
}
}
/**
* FIVE: THE MUTEX.
*
* Two frames delivered back to back, without awaiting either. `on-sig-frame`
* chains the second onto `peer.queue`, so A's whole sequence must appear before
* B's first step — never interleaved. This is the one scenario that catches
* "someone replaced a promise-chain mutex with `p/let`", which is exactly what
* the idiomatic rewrite tempts you into: `p/let` preserves the ordering WITHIN
* one handler and destroys the ordering BETWEEN two.
*/
export async function signalMutex(M) {
const h = await harness(M, LOW, HIGH);
try {
await attached(h);
const release = h.rtc.gate("setRemoteDescription");
const a = await h.frame({ kind: "offer", sdp: "A", name: "them" });
const b = await h.frame({ kind: "offer", sdp: "B", name: "them" });
// A is delivered and never awaited; the harness only lets it reach the
// gate. B is then delivered while A's handler is still suspended inside
// setRemoteDescription. Deliberately NOT both in the same turn: two
// concurrent open-signal calls would race in WebCrypto, and a fixture that
// depended on which decrypt finished first would be recording node's
// threadpool, not lib/media's mutex.
h.say("deliver A; it reaches setRemoteDescription and hangs there");
h.deliver(a);
await h.idle();
h.say("deliver B while A is still in flight");
h.deliver(b);
await h.idle();
const whileGated = h.sink.filter((e) => e[0] === "pc.setRemoteDescription").map((e) => e[2]);
h.say("release A — only now may B start");
release();
await h.idle();
const steps = h.sink
.filter((e) => e[0].startsWith("pc.") || e[0] === "net.sendFrame")
.map((e) => e.slice(0, 3));
return { transcript: h.sink, whileGated, steps, outgoing: await h.outgoing() };
} finally {
h.teardown();
}
}
|