mirror of
https://github.com/run-llama/workflows-ts.git
synced 2026-07-21 06:05:23 -04:00
feat: move third party to top-level (#71)
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@llama-flow/core": patch
|
||||
---
|
||||
|
||||
feat: move third party to top-level
|
||||
@@ -40,21 +40,17 @@
|
||||
"default": "./async-context/index.js"
|
||||
}
|
||||
},
|
||||
"./interrupter/hono": {
|
||||
"types": "./interrupter/hono.d.ts",
|
||||
"default": "./interrupter/hono.js"
|
||||
"./hono": {
|
||||
"types": "./hono.d.ts",
|
||||
"default": "./hono.js"
|
||||
},
|
||||
"./interrupter/mcp": {
|
||||
"types": "./interrupter/mcp.d.ts",
|
||||
"default": "./interrupter/mcp.js"
|
||||
"./mcp": {
|
||||
"types": "./mcp.d.ts",
|
||||
"default": "./mcp.js"
|
||||
},
|
||||
"./interrupter/next": {
|
||||
"types": "./interrupter/next.d.ts",
|
||||
"default": "./interrupter/next.js"
|
||||
},
|
||||
"./interrupter/promise": {
|
||||
"types": "./interrupter/promise.d.ts",
|
||||
"default": "./interrupter/promise.js"
|
||||
"./next": {
|
||||
"types": "./next.d.ts",
|
||||
"default": "./next.js"
|
||||
},
|
||||
"./middleware/store": {
|
||||
"types": "./middleware/store.d.ts",
|
||||
|
||||
@@ -4,7 +4,7 @@ import type {
|
||||
WorkflowEvent,
|
||||
WorkflowEventData,
|
||||
} from "@llama-flow/core";
|
||||
import { promiseHandler } from "./promise";
|
||||
import { runWorkflow } from "./stream/run";
|
||||
|
||||
export const createHonoHandler = <Start, Stop>(
|
||||
workflow: Workflow,
|
||||
@@ -20,7 +20,7 @@ export const createHonoHandler = <Start, Stop>(
|
||||
};
|
||||
}
|
||||
return async (c) => {
|
||||
const stop = await promiseHandler(workflow, await getStart(c), stopEvent);
|
||||
const stop = await runWorkflow(workflow, await getStart(c), stopEvent);
|
||||
return wrapStopEvent(c, stop.data);
|
||||
};
|
||||
};
|
||||
@@ -1,34 +0,0 @@
|
||||
import type {
|
||||
WorkflowContext,
|
||||
WorkflowEvent,
|
||||
WorkflowEventData,
|
||||
} from "@llama-flow/core";
|
||||
import { collect } from "@llama-flow/core/stream/consumer";
|
||||
import { until } from "@llama-flow/core/stream/until";
|
||||
|
||||
/**
|
||||
* Interrupter that wraps a workflow in a promise.
|
||||
*
|
||||
* Resolves when the workflow reads the stop event.
|
||||
* reject if the workflow throws an error.
|
||||
*/
|
||||
export async function promiseHandler<
|
||||
Start,
|
||||
Stop,
|
||||
WorkflowLike extends {
|
||||
createContext(): WorkflowContext;
|
||||
},
|
||||
>(
|
||||
workflow: WorkflowLike,
|
||||
start: WorkflowEventData<Start>,
|
||||
stop: WorkflowEvent<Stop>,
|
||||
): Promise<WorkflowEventData<Stop>> {
|
||||
const { stream, sendEvent } = workflow.createContext();
|
||||
sendEvent(start);
|
||||
const events = await collect(until(stream, stop));
|
||||
const stopEvent = events.reverse().find((e) => stop.include(e));
|
||||
if (stopEvent) {
|
||||
return stopEvent;
|
||||
}
|
||||
throw new Error("Workflow did not return a stop event");
|
||||
}
|
||||
@@ -1,7 +1,7 @@
|
||||
import { createAsyncContext } from "@llama-flow/core/async-context";
|
||||
import { z, type ZodRawShape, type ZodTypeAny } from "zod";
|
||||
import type { Workflow, WorkflowEvent } from "@llama-flow/core";
|
||||
import { promiseHandler } from "./promise";
|
||||
import { runWorkflow } from "./stream/run";
|
||||
import type { RequestHandlerExtra } from "@modelcontextprotocol/sdk/shared/protocol.js";
|
||||
import type { CallToolResult } from "@modelcontextprotocol/sdk/types.js";
|
||||
|
||||
@@ -30,7 +30,7 @@ export function mcpTool<
|
||||
) => CallToolResult | Promise<CallToolResult> {
|
||||
return async (args, extra) =>
|
||||
requestHandlerExtraAsyncLocalStorage.run(extra, async () => {
|
||||
const { data } = await promiseHandler(workflow, start.with(args), stop);
|
||||
const { data } = await runWorkflow(workflow, start.with(args), stop);
|
||||
return data;
|
||||
});
|
||||
}
|
||||
@@ -4,7 +4,7 @@ import type {
|
||||
WorkflowEventData,
|
||||
WorkflowEvent,
|
||||
} from "@llama-flow/core";
|
||||
import { promiseHandler } from "./promise";
|
||||
import { runWorkflow } from "./stream/run";
|
||||
|
||||
type WorkflowAPI = {
|
||||
GET: (request: NextRequest) => Promise<Response>;
|
||||
@@ -19,11 +19,7 @@ export const createNextHandler = <Start, Stop>(
|
||||
): WorkflowAPI => {
|
||||
return {
|
||||
GET: async (request) => {
|
||||
const result = await promiseHandler(
|
||||
workflow,
|
||||
await getStart(request),
|
||||
stop,
|
||||
);
|
||||
const result = await runWorkflow(workflow, await getStart(request), stop);
|
||||
return Response.json(result.data);
|
||||
},
|
||||
};
|
||||
@@ -5,7 +5,7 @@ import {
|
||||
getContext,
|
||||
workflowEvent,
|
||||
} from "@llama-flow/core";
|
||||
import { promiseHandler } from "@llama-flow/core/interrupter/promise";
|
||||
import { runWorkflow } from "@llama-flow/core/stream/run";
|
||||
import { collect } from "@llama-flow/core/stream/consumer";
|
||||
import { until } from "@llama-flow/core/stream/until";
|
||||
|
||||
@@ -23,9 +23,9 @@ describe("sub workflow", () => {
|
||||
sendEvent(stopEvent.with());
|
||||
});
|
||||
await Promise.all([
|
||||
promiseHandler(subWorkflow, startEvent.with(), stopEvent),
|
||||
promiseHandler(subWorkflow, startEvent.with(), stopEvent),
|
||||
promiseHandler(subWorkflow, startEvent.with(), stopEvent),
|
||||
runWorkflow(subWorkflow, startEvent.with(), stopEvent),
|
||||
runWorkflow(subWorkflow, startEvent.with(), stopEvent),
|
||||
runWorkflow(subWorkflow, startEvent.with(), stopEvent),
|
||||
]).then((evt) => sendEvent(...evt));
|
||||
sendEvent(haltEvent.with());
|
||||
});
|
||||
|
||||
@@ -2,7 +2,7 @@ import { describe, expect, test, vi, afterEach } from "vitest";
|
||||
import { createWorkflow, workflowEvent } from "@llama-flow/core";
|
||||
import { withValidation } from "@llama-flow/core/middleware/validation";
|
||||
import { find } from "@llama-flow/core/stream/find";
|
||||
import { promiseHandler } from "@llama-flow/core/interrupter/promise";
|
||||
import { runWorkflow } from "@llama-flow/core/stream/run";
|
||||
|
||||
describe("with directed graph", () => {
|
||||
const consoleWarnMock = vi
|
||||
@@ -69,7 +69,7 @@ describe("with directed graph", () => {
|
||||
sendEvent(stopEvent.with(1));
|
||||
});
|
||||
|
||||
const result = await promiseHandler(workflow, startEvent.with(), stopEvent);
|
||||
const result = await runWorkflow(workflow, startEvent.with(), stopEvent);
|
||||
expect(result.data).toBe(1);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -31,8 +31,8 @@ export default defineConfig([
|
||||
},
|
||||
// Interrupter APIs
|
||||
{
|
||||
entry: ["src/interrupter/*.ts"],
|
||||
outDir: "interrupter",
|
||||
entry: ["src/*.ts"],
|
||||
outDir: "dist",
|
||||
format: ["esm"],
|
||||
external: [
|
||||
"next",
|
||||
|
||||
Reference in New Issue
Block a user