feat: stream.on API (#89)

This commit is contained in:
Alex Yang
2025-04-25 18:37:35 -07:00
committed by GitHub
parent 4342f0204e
commit 2a18aca125
4 changed files with 131 additions and 2 deletions
+18
View File
@@ -0,0 +1,18 @@
---
"@llama-flow/core": patch
---
feat: `stream.on` API
```ts
workflow.handle([startEvent], () => {
const { sendEvent } = getContext();
sendEvent(messageEvent.with("Hello World"));
});
const { stream, sendEvent } = workflow.createContext();
const unsubscribe = stream.on(messageEvent, (event) => {
expect(event.data).toBe("Hello World");
});
sendEvent(startEvent.with());
```
+10 -2
View File
@@ -8,6 +8,7 @@ import {
type Subscribable,
} from "./utils";
import { createAsyncContext } from "@llama-flow/core/async-context";
import { WorkflowStream } from "./stream";
export type Handler<
AcceptEvents extends WorkflowEvent<any>[],
@@ -53,7 +54,7 @@ export type ContextNext = (
) => void;
export type WorkflowContext = {
get stream(): ReadableStream<WorkflowEventData<any>>;
get stream(): WorkflowStream<WorkflowEvent<any>>;
get signal(): AbortSignal;
sendEvent: (...events: WorkflowEventData<any>[]) => void;
@@ -198,7 +199,7 @@ export const createContext = ({
): WorkflowContext => ({
get stream() {
let unsubscribe: () => void;
return new ReadableStream({
const stream = new ReadableStream({
start: async (controller) => {
unsubscribe =
rootWorkflowContext.__internal__call_send_event.subscribe(
@@ -220,6 +221,13 @@ export const createContext = ({
}
},
});
return new WorkflowStream(
rootWorkflowContext.__internal__call_send_event as unknown as Subscribable<
[event: WorkflowEventData<any>],
void
>,
stream,
);
},
get signal() {
return handlerContext.abortController.signal;
+78
View File
@@ -0,0 +1,78 @@
import { type WorkflowEvent, type WorkflowEventData } from "./event";
import type { Subscribable } from "./utils";
export class WorkflowStream<Event extends WorkflowEvent<any>> {
#subscribable: Subscribable<[event: WorkflowEventData<any>], void>;
#stream: ReadableStream<ReturnType<Event["with"]>>;
on<T extends Event>(
event: T,
handler: (event: ReturnType<T["with"]>) => void,
): () => void {
return this.#subscribable.subscribe((ev) => {
if (event.include(ev)) {
handler(ev as ReturnType<T["with"]>);
}
});
}
constructor(
subscribable: Subscribable<[event: WorkflowEventData<any>], void>,
stream: ReadableStream<ReturnType<Event["with"]>>,
) {
this.#subscribable = subscribable;
this.#stream = stream;
}
get locked() {
return this.#stream.locked;
}
[Symbol.asyncIterator](): ReadableStreamAsyncIterator<
ReturnType<Event["with"]>
> {
return this.#stream[Symbol.asyncIterator]();
}
cancel(reason?: any): Promise<void> {
return this.#stream.cancel(reason);
}
// make type compatible with Web ReadableStream API
getReader(options: { mode: "byob" }): ReadableStreamBYOBReader;
getReader(): ReadableStreamDefaultReader<ReturnType<Event["with"]>>;
getReader(
options?: ReadableStreamGetReaderOptions,
): ReadableStreamReader<ReturnType<Event["with"]>>;
getReader(): any {
return this.#stream.getReader();
}
pipeThrough<T = WorkflowEventData<any>>(
transform: ReadableWritablePair<T, ReturnType<Event["with"]>>,
options?: StreamPipeOptions,
): ReadableStream<T> {
return this.#stream.pipeThrough(transform, options);
}
pipeTo(
destination: WritableStream<ReturnType<Event["with"]>>,
options?: StreamPipeOptions,
): Promise<void> {
return this.#stream.pipeTo(destination, options);
}
tee(): [WorkflowStream<Event>, WorkflowStream<Event>] {
const [l, r] = this.#stream.tee();
return [
new WorkflowStream(this.#subscribable, l),
new WorkflowStream(this.#subscribable, r),
];
}
values(
options?: ReadableStreamIteratorOptions,
): ReadableStreamAsyncIterator<ReturnType<Event["with"]>> {
return this.#stream.values(options);
}
}
+25
View File
@@ -0,0 +1,25 @@
import { describe, expect, test, vi } from "vitest";
import { createWorkflow, getContext, workflowEvent } from "@llama-flow/core";
describe("workflow listener api", () => {
test("can listen message event", () => {
const workflow = createWorkflow();
const startEvent = workflowEvent();
const messageEvent = workflowEvent<string>({});
workflow.handle([startEvent], () => {
const { sendEvent } = getContext();
sendEvent(messageEvent.with("Hello World"));
});
const { stream, sendEvent } = workflow.createContext();
const callback = vi.fn((event: ReturnType<typeof messageEvent.with>) => {
expect(event.data).toBe("Hello World");
});
const unsubscribe = stream.on(messageEvent, callback);
sendEvent(startEvent.with());
expect(callback).toHaveBeenCalledTimes(1);
unsubscribe();
sendEvent(startEvent.with());
expect(callback).toHaveBeenCalledTimes(1);
});
});