mirror of
https://github.com/run-llama/workflows-ts.git
synced 2026-07-21 14:15:24 -04:00
feat: stream.on API (#89)
This commit is contained in:
@@ -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());
|
||||
```
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user