mirror of
https://github.com/run-llama/workflows-ts.git
synced 2026-07-21 06:05:23 -04:00
feat: add rxjs binding (#73)
This commit is contained in:
@@ -0,0 +1,6 @@
|
||||
---
|
||||
"demo": patch
|
||||
"@llama-flow/core": patch
|
||||
---
|
||||
|
||||
feat: add rxjs binding
|
||||
@@ -3,18 +3,21 @@ import {
|
||||
messageEvent,
|
||||
startEvent,
|
||||
} from "../workflows/file-parse-agent.js";
|
||||
import { from, filter } from "rxjs";
|
||||
import { filter, map } from "rxjs";
|
||||
import { eventSource } from "@llama-flow/core";
|
||||
import type { WorkflowEventData } from "@llama-flow/core";
|
||||
import { toObservable } from "@llama-flow/core/observable";
|
||||
|
||||
const directory = "..";
|
||||
|
||||
const { stream, sendEvent } = fileParseWorkflow.createContext();
|
||||
|
||||
from(stream as unknown as AsyncIterable<WorkflowEventData<any>>)
|
||||
.pipe(filter((ev) => eventSource(ev) === messageEvent))
|
||||
.subscribe((ev) => {
|
||||
console.log(ev.data);
|
||||
toObservable(stream)
|
||||
.pipe(
|
||||
filter((ev) => eventSource(ev) === messageEvent),
|
||||
map((ev) => ev.data),
|
||||
)
|
||||
.subscribe((data) => {
|
||||
console.log(data);
|
||||
});
|
||||
|
||||
sendEvent(startEvent.with(directory));
|
||||
|
||||
+1
-1
@@ -3,7 +3,7 @@
|
||||
"target": "esnext",
|
||||
"module": "esnext",
|
||||
"esModuleInterop": true,
|
||||
"moduleResolution": "node16",
|
||||
"moduleResolution": "bundler",
|
||||
"types": ["DOM", "DOM.AsyncIterable", "DOM.Iterable", "esnext"],
|
||||
"forceConsistentCasingInFileNames": true,
|
||||
"strict": true,
|
||||
|
||||
@@ -41,16 +41,20 @@
|
||||
}
|
||||
},
|
||||
"./hono": {
|
||||
"types": "./hono.d.ts",
|
||||
"default": "./hono.js"
|
||||
"types": "./dist/hono.d.ts",
|
||||
"default": "./dist/hono.js"
|
||||
},
|
||||
"./mcp": {
|
||||
"types": "./mcp.d.ts",
|
||||
"default": "./mcp.js"
|
||||
"types": "./dist/mcp.d.ts",
|
||||
"default": "./dist/mcp.js"
|
||||
},
|
||||
"./next": {
|
||||
"types": "./next.d.ts",
|
||||
"default": "./next.js"
|
||||
"types": "./dist/next.d.ts",
|
||||
"default": "./dist/next.js"
|
||||
},
|
||||
"./observable": {
|
||||
"types": "./dist/observable.d.ts",
|
||||
"default": "./dist/observable.js"
|
||||
},
|
||||
"./middleware/store": {
|
||||
"types": "./middleware/store.d.ts",
|
||||
@@ -111,6 +115,7 @@
|
||||
"next": "^15.3.1",
|
||||
"p-retry": "^6.2.1",
|
||||
"rimraf": "^6.0.1",
|
||||
"rxjs": "^7.8.2",
|
||||
"stream-chain": "^3.4.0",
|
||||
"tsdown": "^0.9.1",
|
||||
"typescript": "^5.8.3",
|
||||
@@ -121,6 +126,7 @@
|
||||
"hono": "^4.7.4",
|
||||
"next": "^15.2.2",
|
||||
"p-retry": "^6.2.1",
|
||||
"rxjs": "^7.8.2",
|
||||
"zod": "^3.24.2"
|
||||
},
|
||||
"license": "MIT",
|
||||
@@ -137,6 +143,9 @@
|
||||
"p-retry": {
|
||||
"optional": true
|
||||
},
|
||||
"rxjs": {
|
||||
"optional": true
|
||||
},
|
||||
"zod": {
|
||||
"optional": true
|
||||
}
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
import { Observable } from "rxjs";
|
||||
import type { WorkflowEventData } from "@llama-flow/core";
|
||||
|
||||
export const toObservable = (
|
||||
stream: ReadableStream<WorkflowEventData<any>>,
|
||||
): Observable<WorkflowEventData<any>> => {
|
||||
return new Observable((subscriber) => {
|
||||
const reader = stream.getReader();
|
||||
|
||||
const read = async () => {
|
||||
try {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) {
|
||||
subscriber.complete();
|
||||
} else {
|
||||
subscriber.next(value);
|
||||
read();
|
||||
}
|
||||
} catch (error) {
|
||||
subscriber.error(error);
|
||||
}
|
||||
};
|
||||
|
||||
read().catch(subscriber.error);
|
||||
|
||||
return () => {
|
||||
reader.cancel().catch(subscriber.error);
|
||||
};
|
||||
});
|
||||
};
|
||||
Generated
+3
@@ -148,6 +148,9 @@ importers:
|
||||
rimraf:
|
||||
specifier: ^6.0.1
|
||||
version: 6.0.1
|
||||
rxjs:
|
||||
specifier: ^7.8.2
|
||||
version: 7.8.2
|
||||
stream-chain:
|
||||
specifier: ^3.4.0
|
||||
version: 3.4.0
|
||||
|
||||
Reference in New Issue
Block a user