fix: extends standard readable stream (#103)

This commit is contained in:
Alex Yang
2025-05-02 15:57:16 -07:00
committed by GitHub
parent 4201aabb0d
commit 4402a6dc15
4 changed files with 12 additions and 8 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@llama-flow/core": patch
---
fix: workflow stream extends standard readable stream
+3 -2
View File
@@ -8,10 +8,10 @@ import {
import { createSubscribable, type Subscribable } from "./utils";
export class WorkflowStream<R = any>
implements AsyncIterable<R>, ReadableStream<R>
extends ReadableStream<R>
implements AsyncIterable<R>
{
#stream: ReadableStream<R>;
#subscribable: Subscribable<[data: R], void>;
on(
@@ -40,6 +40,7 @@ export class WorkflowStream<R = any>
"Either subscribable or root stream must be provided",
);
}
super();
if (!subscribable) {
this.#subscribable = createSubscribable<[data: R], void>();
this.#stream = rootStream!.pipeThrough(
@@ -56,8 +56,7 @@ describe("stream-chain", () => {
sendEvent(startEvent.with());
const outputs: WorkflowEventData<any>[] = [];
const pipeline = chain([
// fixme: upstream should be treat it as a stream
stream.until(stopEvent)[Symbol.asyncIterator],
stream.until(stopEvent),
new TransformStream({
transform: (event: WorkflowEventData<any>, controller) => {
if (messageEvent.include(event)) {
+3 -4
View File
@@ -5,7 +5,6 @@ import {
workflowEvent,
WorkflowStream,
} from "@llama-flow/core";
import { collect } from "@llama-flow/core/stream/consumer";
import { createStatefulMiddleware } from "@llama-flow/core/middleware/state";
type Handler<
@@ -179,7 +178,7 @@ export class Workflow<ContextData, Start, Stop> {
Object.assign(result, {
then: async (resolve: any, reject: any) => {
try {
const events = await collect(result);
const events = await result.toArray();
resolve(events.at(-1)!);
} catch (error) {
reject(error);
@@ -187,14 +186,14 @@ export class Workflow<ContextData, Start, Stop> {
},
catch: async (reject: any) => {
try {
await collect(result);
await result.toArray();
} catch (error) {
reject(error);
}
},
finally: async (resolve: any) => {
try {
await collect(result);
await result.toArray();
} finally {
resolve();
}