feat: add rxjs binding (#73)

This commit is contained in:
Alex Yang
2025-04-23 09:28:03 -07:00
committed by GitHub
parent 78b3141bc2
commit e2f8e23869
6 changed files with 64 additions and 13 deletions
+6
View File
@@ -0,0 +1,6 @@
---
"demo": patch
"@llama-flow/core": patch
---
feat: add rxjs binding
+9 -6
View File
@@ -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
View File
@@ -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,
+15 -6
View File
@@ -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
}
+30
View File
@@ -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);
};
});
};
+3
View File
@@ -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