Skip to content

Commit 66d49c0

Browse files
[world] Restructure stream interface, require run ID for all step and stream operations (#1293)
1 parent dc0c0dc commit 66d49c0

28 files changed

Lines changed: 1233 additions & 1211 deletions

File tree

‎.changeset/bright-pears-drum.md‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
---
2+
"@workflow/world": major
3+
"@workflow/world-local": major
4+
"@workflow/world-vercel": major
5+
"@workflow/world-postgres": major
6+
"@workflow/core": major
7+
"@workflow/cli": major
8+
"@workflow/web": major
9+
---
10+
11+
**BREAKING CHANGE**: Restructure stream methods on World interface to use `world.streams.*` namespace with `runId` as the first parameter. `writeToStream(name, runId, chunk)` → `streams.write(runId, name, chunk)`, `writeToStreamMulti` → `streams.writeMulti`, `closeStream` → `streams.close`, `readFromStream` → `streams.get(runId, name, startIndex?)`, `listStreamsByRunId` → `streams.list(runId)`.

‎.changeset/step-run-required.md‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
---
2+
"@workflow/world-postgres": major
3+
"@workflow/world-vercel": major
4+
"@workflow/world-local": major
5+
"@workflow/world": major
6+
"@workflow/core": major
7+
"@workflow/cli": major
8+
"@workflow/web": major
9+
---
10+
11+
Require `runId` argument for `world.steps.get`.

‎docs/content/docs/api-reference/workflow-api/world/streams.mdx‎

Lines changed: 35 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -2,26 +2,26 @@
22
title: Streams
33
description: Read, write, and manage real-time data streams for workflow runs.
44
type: reference
5-
summary: "Methods: writeToStream(), writeToStreamMulti(), readFromStream(), closeStream(), listStreamsByRunId(), getStreamChunks(), getStreamInfo(). Stream methods live directly on the world object."
5+
summary: "Methods: streams.write(), streams.writeMulti(), streams.get(), streams.close(), streams.list(), streams.getChunks(), streams.getInfo(). Stream methods live on world.streams."
66
prerequisites:
77
- /docs/api-reference/workflow-api/get-world
88
related:
99
- /docs/foundations/streaming
1010
- /docs/api-reference/workflow/get-writable
1111
keywords:
12-
- writeToStream
13-
- writeToStreamMulti
14-
- readFromStream
15-
- closeStream
16-
- listStreamsByRunId
17-
- getStreamChunks
18-
- getStreamInfo
12+
- streams.write
13+
- streams.writeMulti
14+
- streams.get
15+
- streams.close
16+
- streams.list
17+
- streams.getChunks
18+
- streams.getInfo
1919
- Streamer interface
2020
- real-time streaming
2121
- stream lifecycle
2222
---
2323

24-
Stream methods live directly on the `world` object returned by `getWorld()`. Use them to write chunks, read streams, and manage stream lifecycle outside of the standard `getWritable()` pattern.
24+
Stream methods live on `world.streams` (the `streams` sub-object of the `world` object returned by `getWorld()`). Use them to write chunks, read streams, and manage stream lifecycle outside of the standard `getWritable()` pattern.
2525

2626
<Callout type="info">
2727
For most streaming use cases, use [`getWritable()`](/docs/api-reference/workflow/get-writable) inside steps. Direct stream methods are for advanced scenarios like building custom stream consumers or managing streams from outside a workflow.
@@ -33,81 +33,82 @@ Stream methods live directly on the `world` object returned by `getWorld()`. Use
3333
import { getWorld } from "workflow/runtime";
3434

3535
const world = getWorld(); // [!code highlight]
36-
// Stream methods are called directly on world — e.g. world.writeToStream()
36+
// Stream methods are called on world.streams — e.g. world.streams.write()
3737
```
3838

3939
## Methods
4040

41-
### writeToStream()
41+
### write()
4242

4343
Write a data chunk to a named stream.
4444

4545
```typescript lineNumbers
46-
await world.writeToStream("default", runId, chunk); // [!code highlight]
46+
await world.streams.write(runId, "default", chunk); // [!code highlight]
4747
```
4848

4949
**Parameters:**
5050

5151
| Parameter | Type | Description |
5252
|-----------|------|-------------|
53-
| `name` | `string` | The stream name |
5453
| `runId` | `string` | The workflow run ID |
54+
| `name` | `string` | The stream name |
5555
| `chunk` | `string \| Uint8Array` | Data to write |
5656

57-
### writeToStreamMulti()
57+
### writeMulti()
5858

59-
Write multiple chunks in a single operation. Optional optimization — not all World implementations support it. Falls back to sequential `writeToStream()` calls if unavailable.
59+
Write multiple chunks in a single operation. Optional optimization — not all World implementations support it. Falls back to sequential `write()` calls if unavailable.
6060

6161
```typescript lineNumbers
62-
await world.writeToStreamMulti?.("default", runId, [chunk1, chunk2]); // [!code highlight]
62+
await world.streams.writeMulti?.(runId, "default", [chunk1, chunk2]); // [!code highlight]
6363
```
6464

6565
**Parameters:**
6666

6767
| Parameter | Type | Description |
6868
|-----------|------|-------------|
69-
| `name` | `string` | The stream name |
7069
| `runId` | `string` | The workflow run ID |
70+
| `name` | `string` | The stream name |
7171
| `chunks` | `(string \| Uint8Array)[]` | Chunks to write, in order |
7272

73-
### readFromStream()
73+
### get()
7474

7575
Read data from a named stream as a live `ReadableStream` that waits for new chunks in real time.
7676

7777
```typescript lineNumbers
78-
const readable = await world.readFromStream("default"); // [!code highlight]
78+
const readable = await world.streams.get(runId, "default"); // [!code highlight]
7979
```
8080

8181
**Parameters:**
8282

8383
| Parameter | Type | Description |
8484
|-----------|------|-------------|
85+
| `runId` | `string` | The workflow run ID |
8586
| `name` | `string` | The stream name |
8687
| `startIndex` | `number` | Optional. Positive values skip chunks from the start (0-based). Negative values read from the tail (e.g. `-3` starts 3 chunks from the end). Clamped to 0. |
8788

8889
**Returns:** `ReadableStream<Uint8Array>`
8990

90-
### closeStream()
91+
### close()
9192

9293
Close a stream when done writing.
9394

9495
```typescript lineNumbers
95-
await world.closeStream("default", runId); // [!code highlight]
96+
await world.streams.close(runId, "default"); // [!code highlight]
9697
```
9798

9899
**Parameters:**
99100

100101
| Parameter | Type | Description |
101102
|-----------|------|-------------|
102-
| `name` | `string` | The stream name |
103103
| `runId` | `string` | The workflow run ID |
104+
| `name` | `string` | The stream name |
104105

105-
### listStreamsByRunId()
106+
### list()
106107

107108
List all stream names associated with a workflow run.
108109

109110
```typescript lineNumbers
110-
const streamNames = await world.listStreamsByRunId(runId); // [!code highlight]
111+
const streamNames = await world.streams.list(runId); // [!code highlight]
111112
```
112113

113114
**Parameters:**
@@ -118,12 +119,12 @@ const streamNames = await world.listStreamsByRunId(runId); // [!code highlight]
118119

119120
**Returns:** `string[]`
120121

121-
### getStreamChunks()
122+
### getChunks()
122123

123-
Fetch stream chunks with cursor-based pagination. Unlike `readFromStream()` (which returns a live `ReadableStream`), this returns a snapshot of currently available chunks.
124+
Fetch stream chunks with cursor-based pagination. Unlike `get()` (which returns a live `ReadableStream`), this returns a snapshot of currently available chunks.
124125

125126
```typescript lineNumbers
126-
const result = await world.getStreamChunks("default", runId, { // [!code highlight]
127+
const result = await world.streams.getChunks(runId, "default", { // [!code highlight]
127128
limit: 50,
128129
}); // [!code highlight]
129130
// result.data: StreamChunk[], result.cursor, result.hasMore, result.done
@@ -133,8 +134,8 @@ const result = await world.getStreamChunks("default", runId, { // [!code highlig
133134

134135
| Parameter | Type | Description |
135136
|-----------|------|-------------|
136-
| `name` | `string` | The stream name |
137137
| `runId` | `string` | The workflow run ID |
138+
| `name` | `string` | The stream name |
138139
| `options.limit` | `number` | Max chunks per page (default: 100, max: 1000) |
139140
| `options.cursor` | `string` | Cursor from a previous response |
140141

@@ -147,21 +148,21 @@ const result = await world.getStreamChunks("default", runId, { // [!code highlig
147148
| `hasMore` | `boolean` | Whether more pages of already-written chunks exist |
148149
| `done` | `boolean` | Whether the stream is fully closed. When `false`, new chunks may appear in future requests even after `hasMore` is `false`. |
149150

150-
### getStreamInfo()
151+
### getInfo()
151152

152153
Retrieve lightweight metadata about a stream without fetching chunks.
153154

154155
```typescript lineNumbers
155-
const info = await world.getStreamInfo("default", runId); // [!code highlight]
156+
const info = await world.streams.getInfo(runId, "default"); // [!code highlight]
156157
// info.tailIndex: last chunk index (-1 if empty), info.done: whether stream is closed
157158
```
158159

159160
**Parameters:**
160161

161162
| Parameter | Type | Description |
162163
|-----------|------|-------------|
163-
| `name` | `string` | The stream name |
164164
| `runId` | `string` | The workflow run ID |
165+
| `name` | `string` | The stream name |
165166

166167
**Returns:** `StreamInfoResponse`
167168

@@ -181,8 +182,9 @@ import { getWorld } from "workflow/runtime";
181182
export async function GET(req: Request) {
182183
const url = new URL(req.url);
183184
const streamName = url.searchParams.get("name") ?? "default";
185+
const runId = url.searchParams.get("runId")!;
184186
const world = getWorld();
185-
const readable = await world.readFromStream(streamName); // [!code highlight]
187+
const readable = await world.streams.get(runId, streamName); // [!code highlight]
186188

187189
return new Response(readable, {
188190
headers: { "Content-Type": "application/octet-stream" },
@@ -199,7 +201,7 @@ const world = getWorld();
199201
let cursor: string | undefined;
200202

201203
do {
202-
const result = await world.getStreamChunks("default", runId, { cursor }); // [!code highlight]
204+
const result = await world.streams.getChunks(runId, "default", { cursor }); // [!code highlight]
203205
for (const chunk of result.data) {
204206
console.log(`Chunk ${chunk.index}:`, chunk.data);
205207
}

‎docs/content/docs/deploying/building-a-world.mdx‎

Lines changed: 45 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -166,54 +166,56 @@ The Streamer interface enables real-time data streaming:
166166
{/* @skip-typecheck - interface definition, not runnable code */}
167167
```typescript
168168
interface Streamer {
169-
writeToStream(
170-
name: string,
171-
runId: string,
172-
chunk: string | Uint8Array
173-
): Promise<void>;
174-
175-
writeToStreamMulti?(
176-
name: string,
177-
runId: string,
178-
chunks: (string | Uint8Array)[]
179-
): Promise<void>;
180-
181-
closeStream(
182-
name: string,
183-
runId: string
184-
): Promise<void>;
185-
186-
readFromStream(
187-
name: string,
188-
startIndex?: number
189-
): Promise<ReadableStream<Uint8Array>>;
190-
191-
listStreamsByRunId(runId: string): Promise<string[]>;
192-
193-
/** Paginated snapshot of stream chunks. */
194-
getStreamChunks(
195-
name: string,
196-
runId: string,
197-
options?: { limit?: number; cursor?: string }
198-
): Promise<{
199-
data: { index: number; data: Uint8Array }[];
200-
cursor: string | null;
201-
hasMore: boolean;
202-
done: boolean;
203-
}>;
204-
205-
/** Lightweight metadata: tail index and completion flag. */
206-
getStreamInfo(
207-
name: string,
208-
runId: string
209-
): Promise<{ tailIndex: number; done: boolean }>;
169+
streamFlushIntervalMs?: number;
170+
171+
streams: {
172+
write(
173+
runId: string,
174+
name: string,
175+
chunk: string | Uint8Array
176+
): Promise<void>;
177+
178+
writeMulti?(
179+
runId: string,
180+
name: string,
181+
chunks: (string | Uint8Array)[]
182+
): Promise<void>;
183+
184+
close(runId: string, name: string): Promise<void>;
185+
186+
get(
187+
runId: string,
188+
name: string,
189+
startIndex?: number
190+
): Promise<ReadableStream<Uint8Array>>;
191+
192+
list(runId: string): Promise<string[]>;
193+
194+
/** Paginated snapshot of stream chunks. */
195+
getChunks(
196+
runId: string,
197+
name: string,
198+
options?: { limit?: number; cursor?: string }
199+
): Promise<{
200+
data: { index: number; data: Uint8Array }[];
201+
cursor: string | null;
202+
hasMore: boolean;
203+
done: boolean;
204+
}>;
205+
206+
/** Lightweight metadata: tail index and completion flag. */
207+
getInfo(
208+
runId: string,
209+
name: string
210+
): Promise<{ tailIndex: number; done: boolean }>;
211+
};
210212
}
211213
```
212214

213215
Streams are identified by a combination of `runId` and `name`. Each workflow run can have multiple named streams.
214-
`writeToStreamMulti()` is an optional optimization for batching multiple writes.
216+
`writeMulti()` is an optional optimization for batching multiple writes.
215217

216-
`getStreamChunks` returns a paginated snapshot of currently available chunks (unlike `readFromStream` which returns a live `ReadableStream` that waits for new chunks). `getStreamInfo` returns the tail index (last chunk index, 0-based, or `-1` when empty) and whether the stream is complete — useful for resolving negative `startIndex` values into absolute positions.
218+
`getChunks` returns a paginated snapshot of currently available chunks (unlike `get` which returns a live `ReadableStream` that waits for new chunks). `getInfo` returns the tail index (last chunk index, 0-based, or `-1` when empty) and whether the stream is complete — useful for resolving negative `startIndex` values into absolute positions.
217219

218220
## Reference Implementations
219221

‎packages/cli/src/lib/inspect/output.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -849,7 +849,10 @@ export const showStream = async (
849849
'Filtering by step-id is not supported when showing a stream, ignoring filter.'
850850
);
851851
}
852-
const rawStream = await world.readFromStream(streamId);
852+
if (!opts.runId) {
853+
throw new Error('--run is required when showing a stream');
854+
}
855+
const rawStream = await world.streams.get(opts.runId, streamId);
853856

854857
// Only resolve the encryption key when --decrypt is passed and --run is provided.
855858
// We fetch the full WorkflowRun object so that getEncryptionKeyForRun has
@@ -921,7 +924,7 @@ export const listStreamsByRunId = async (
921924
}
922925

923926
try {
924-
const streamIds = await world.listStreamsByRunId(runId);
927+
const streamIds = await world.streams.list(runId);
925928
const matchingStreams = streamIds.map((streamId) => ({
926929
runId,
927930
streamId,

‎packages/core/e2e/e2e.test.ts‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -718,7 +718,7 @@ describe('e2e', () => {
718718
});
719719

720720
describe.skipIf(isLocalDeployment())(
721-
'outputStreamWorkflow - getTailIndex and getStreamChunks',
721+
'outputStreamWorkflow - getTailIndex and getChunks',
722722
() => {
723723
test(
724724
'getTailIndex returns correct index after stream completes',
@@ -755,7 +755,7 @@ describe('e2e', () => {
755755
);
756756

757757
test(
758-
'getStreamChunks returns same content as reading the stream',
758+
'getChunks returns same content as reading the stream',
759759
{
760760
timeout: 60_000,
761761
},
@@ -772,13 +772,13 @@ describe('e2e', () => {
772772
streamChunks.push(value);
773773
}
774774

775-
// Read all chunks via getStreamChunks pagination
775+
// Read all chunks via getChunks pagination
776776
const world = getWorld();
777777
const streamName = `${run.runId.replace('wrun_', 'strm_')}_user`;
778778
const paginatedChunks: Uint8Array[] = [];
779779
let cursor: string | null = null;
780780
do {
781-
const page = await world.getStreamChunks(streamName, run.runId, {
781+
const page = await world.streams.getChunks(run.runId, streamName, {
782782
limit: 1, // small page size to exercise pagination
783783
...(cursor ? { cursor } : {}),
784784
});

0 commit comments

Comments
 (0)