feat: move third party to top-level (#71)

This commit is contained in:
Alex Yang
2025-04-23 00:59:55 -07:00
committed by GitHub
parent 9c4bbd9c1d
commit 78b3141bc2
9 changed files with 28 additions and 65 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@llama-flow/core": patch
---
feat: move third party to top-level
+9 -13
View File
@@ -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);
};
};
-34
View File
@@ -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);
});
});
+2 -2
View File
@@ -31,8 +31,8 @@ export default defineConfig([
},
// Interrupter APIs
{
entry: ["src/interrupter/*.ts"],
outDir: "interrupter",
entry: ["src/*.ts"],
outDir: "dist",
format: ["esm"],
external: [
"next",