Skip to content
File

Blob: src/workerd/api/tests/jsrpc-timing-test.js

javascript63 lines
1// Copyright (c) 2017-2024 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5// Test to validate timing semantics for JSRPC invocations with streaming responses.
6// This test verifies that the Return event is emitted at the correct time relative
7// to Onset and Outcome events.
8//
9// Expected timeline:
10// T=0: Onset (invocation starts)
11// T=~500ms: Return (handler returns the stream)
12// T=~950ms: Outcome (stream fully consumed)
13 
14import { WorkerEntrypoint } from 'cloudflare:workers';
15 
16export class StreamingService extends WorkerEntrypoint {
17 async getStreamWithDelays() {
18 // Sleep 500ms before returning the stream
19 await scheduler.wait(500);
20 
21 // Return a ReadableStream that takes ~450ms to drain
22 const encoder = new TextEncoder();
23 let chunksSent = 0;
24 
25 const stream = new ReadableStream({
26 async pull(controller) {
27 if (chunksSent >= 3) {
28 controller.close();
29 return;
30 }
31 // Sleep 150ms between chunks (3 chunks = ~450ms to drain)
32 await scheduler.wait(150);
33 controller.enqueue(encoder.encode(`chunk${chunksSent++}\n`));
34 },
35 });
36 
37 return stream;
38 }
39}
40 
41export default {
42 async test(controller, env, ctx) {
43 // Call the streaming RPC method
44 const stream = await env.StreamingService.getStreamWithDelays();
45 
46 // Consume the stream fully
47 const reader = stream.getReader();
48 let chunks = [];
49 while (true) {
50 const { done, value } = await reader.read();
51 if (done) break;
52 chunks.push(new TextDecoder().decode(value));
53 }
54 
55 // Verify we got all chunks
56 if (chunks.length !== 3) {
57 throw new Error(`Expected 3 chunks, got ${chunks.length}`);
58 }
59 
60 // The actual timing validation is done in the tail worker's test() handler
61 },
62};