Skip to content

Commit bb8494e

Browse files
authored
Merge pull request #47 from caido/ef-fix-early-event
Fix an issue where fast events would get missed
2 parents 8070409 + 5e0dbbd commit bb8494e

5 files changed

Lines changed: 169 additions & 41 deletions

File tree

.github/workflows/sdk-client-tests.yml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ jobs:
2525
caido-version:
2626
- 0.56.0
2727
- 0.57.0
28+
- 0.58.0
2829

2930
services:
3031
caido:
Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,92 @@
1+
import { describe, expect, it, vi } from "vitest";
2+
import { makeSubject, toAsyncIterable } from "wonka";
3+
4+
import { ReplaySDK } from "./replay.js";
5+
6+
import type { GraphQLClient } from "@/graphql/index.js";
7+
import { Version } from "@/version.js";
8+
9+
type FinishedEvent = {
10+
finishedTask: {
11+
task: {
12+
__typename: "ReplayTask";
13+
id: string;
14+
createdAt: string;
15+
replayEntry: { id: string };
16+
};
17+
status: "DONE";
18+
error: null;
19+
};
20+
};
21+
22+
describe("ReplaySDK.send", () => {
23+
it("resolves when the task finishes before the start mutation returns", async () => {
24+
const subject = makeSubject<FinishedEvent>();
25+
26+
const task = {
27+
__typename: "ReplayTask" as const,
28+
id: "task-1",
29+
createdAt: "2026-01-01T00:00:00.000Z",
30+
replayEntry: { id: "entry-1" },
31+
};
32+
33+
const graphql = {
34+
subscribe: vi.fn(() => toAsyncIterable(subject.source)),
35+
mutation: vi.fn(async () => {
36+
// Fast-finishing task: finished event arrives before mutation resolves
37+
subject.next({
38+
finishedTask: {
39+
task,
40+
status: "DONE",
41+
error: null,
42+
},
43+
});
44+
subject.complete();
45+
46+
return {
47+
startReplayTask: {
48+
task,
49+
error: null,
50+
},
51+
};
52+
}),
53+
query: vi.fn(async () => ({
54+
replayEntry: {
55+
id: "entry-1",
56+
createdAt: "2026-01-01T00:00:00.000Z",
57+
error: null,
58+
raw: undefined,
59+
connection: {
60+
__typename: "ConnectionInfo",
61+
host: "example.com",
62+
port: 80,
63+
isTLS: false,
64+
SNI: null,
65+
},
66+
request: null,
67+
session: { id: "session-1" },
68+
settings: { placeholders: [] },
69+
},
70+
})),
71+
} as unknown as GraphQLClient;
72+
73+
const replay = new ReplaySDK(graphql, Version.of("0.56.0"));
74+
75+
const result = await replay.send("session-1", {
76+
raw: "GET / HTTP/1.1\r\nHost: example.com\r\n\r\n",
77+
connection: {
78+
host: "example.com",
79+
port: 80,
80+
isTLS: false,
81+
},
82+
});
83+
84+
expect(result.status).toBe("DONE");
85+
expect(result.entry.id).toBe("entry-1");
86+
expect(graphql.subscribe).toHaveBeenCalled();
87+
expect(graphql.mutation).toHaveBeenCalled();
88+
expect(
89+
vi.mocked(graphql.subscribe).mock.invocationCallOrder[0]!,
90+
).toBeLessThan(vi.mocked(graphql.mutation).mock.invocationCallOrder[0]!);
91+
});
92+
});

packages/sdk-client/src/sdks/replay.ts

Lines changed: 39 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ import {
2020
type ReplaySendResult,
2121
TransportVersion,
2222
} from "@/types/index.js";
23+
import { bufferAsyncIterable } from "@/utils/asyncIterable.js";
2324
import { handleGraphQLError } from "@/utils/errors.js";
2425
import { isAbsent, isPresent } from "@/utils/optional.js";
2526
import type { Version } from "@/version.js";
@@ -143,19 +144,20 @@ export class ReplaySDK {
143144
});
144145
}
145146

146-
// Start the replay task
147-
const result = await this.graphql.mutation(LatestStartReplayTaskDocument, {
148-
sessionId: sessionId as string,
149-
});
150-
151-
const payload = result.startReplayTask;
152-
if (isPresent(payload.error)) {
153-
handleGraphQLError(payload.error);
154-
}
155-
const task = new ReplayTask(this.graphql, payload.task!);
147+
return this.waitForReplayTask(async () => {
148+
const result = await this.graphql.mutation(
149+
LatestStartReplayTaskDocument,
150+
{
151+
sessionId: sessionId as string,
152+
},
153+
);
156154

157-
// Wait for task to finish
158-
return this.waitForReplayTask(task);
155+
const payload = result.startReplayTask;
156+
if (isPresent(payload.error)) {
157+
handleGraphQLError(payload.error);
158+
}
159+
return new ReplayTask(this.graphql, payload.task!);
160+
});
159161
}
160162

161163
private async sendV056(
@@ -183,25 +185,35 @@ export class ReplaySDK {
183185
},
184186
} satisfies StartReplayTaskInput;
185187

186-
const result = await this.graphql.mutation(V056StartReplayTaskDocument, {
187-
sessionId: sessionId as string,
188-
input,
188+
return this.waitForReplayTask(async () => {
189+
const result = await this.graphql.mutation(V056StartReplayTaskDocument, {
190+
sessionId: sessionId as string,
191+
input,
192+
});
193+
194+
const payload = result.startReplayTask;
195+
if (isPresent(payload.error)) {
196+
handleGraphQLError(payload.error);
197+
}
198+
return new ReplayTask(this.graphql, payload.task!);
189199
});
200+
}
190201

191-
const payload = result.startReplayTask;
192-
if (isPresent(payload.error)) {
193-
handleGraphQLError(payload.error);
194-
}
195-
const task = new ReplayTask(this.graphql, payload.task!);
202+
/**
203+
* Open the finished-task subscription before running `start`, then resolve
204+
* when the started task finishes.
205+
*/
206+
private async waitForReplayTask(
207+
start: () => Promise<ReplayTask>,
208+
): Promise<ReplaySendResult> {
209+
const finished = bufferAsyncIterable(this.tasks.finished());
210+
const task = await start();
196211

197-
// Wait for task to finish
198-
return this.waitForReplayTask(task);
199-
}
212+
for await (const result of finished) {
213+
if (result.task.id !== task.id) {
214+
continue;
215+
}
200216

201-
private async waitForReplayTask(task: ReplayTask): Promise<ReplaySendResult> {
202-
for await (const result of this.tasks.finished(
203-
(finished) => finished.task.id === task.id,
204-
)) {
205217
const entry = await this.entries.get(task.replayEntryId);
206218
if (isAbsent(entry)) {
207219
throw new OtherUserError("INTERNAL", "Replay entry not found");

packages/sdk-client/src/sdks/task.ts

Lines changed: 16 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -42,21 +42,23 @@ export class TaskSDK {
4242
}
4343

4444
finished(
45-
filter: (taskResult: TaskResult) => boolean,
45+
filter?: (taskResult: TaskResult) => boolean,
4646
): AsyncIterable<TaskResult> {
47-
return filterAsyncIterable(
48-
filter,
49-
mapAsyncIterable((event) => {
50-
const task = new Task(this.graphql, event.finishedTask.task);
51-
return {
52-
task,
53-
status: event.finishedTask.status as TaskStatus,
54-
error: event.finishedTask.error
55-
? { code: event.finishedTask.error.code }
56-
: undefined,
57-
};
58-
}, this.graphql.subscribe(FinishedTaskDocument)),
59-
);
47+
const results = mapAsyncIterable((event) => {
48+
const task = new Task(this.graphql, event.finishedTask.task);
49+
return {
50+
task,
51+
status: event.finishedTask.status as TaskStatus,
52+
error: event.finishedTask.error
53+
? { code: event.finishedTask.error.code }
54+
: undefined,
55+
};
56+
}, this.graphql.subscribe(FinishedTaskDocument));
57+
58+
if (filter === undefined) {
59+
return results;
60+
}
61+
return filterAsyncIterable(filter, results);
6062
}
6163
}
6264

packages/sdk-client/src/utils/asyncIterable.ts

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,3 +18,24 @@ export async function* filterAsyncIterable<T>(
1818
}
1919
}
2020
}
21+
22+
export function bufferAsyncIterable<T>(
23+
source: AsyncIterable<T>,
24+
): AsyncIterable<T> {
25+
const iterator = source[Symbol.asyncIterator]();
26+
const first = iterator.next();
27+
28+
return {
29+
async *[Symbol.asyncIterator]() {
30+
try {
31+
let result = await first;
32+
while (result.done !== true) {
33+
yield result.value;
34+
result = await iterator.next();
35+
}
36+
} finally {
37+
await iterator.return?.(undefined);
38+
}
39+
},
40+
};
41+
}

0 commit comments

Comments
 (0)