mirror of
https://github.com/run-llama/workflows-ts.git
synced 2026-07-21 06:05:23 -04:00
fix: extends standard readable stream (#103)
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@llama-flow/core": patch
|
||||
---
|
||||
|
||||
fix: workflow stream extends standard readable stream
|
||||
@@ -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)) {
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user