feat(sdk): allow enqueuing runs (#1818)

This commit is contained in:
David Duong
2025-12-15 14:46:05 +01:00
committed by GitHub
parent b1ed761610
commit 1497df973f
4 changed files with 560 additions and 288 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@langchain/langgraph-sdk": minor
---
feat(sdk): add support for enqueuing `useStream().submit(...)` calls while the agent is still running
+494 -263
View File
@@ -1,6 +1,6 @@
import "@testing-library/jest-dom/vitest";
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
import { it, expect, beforeEach, afterEach, vi } from "vitest";
import { render, screen, waitFor, within } from "@testing-library/react";
import { userEvent } from "@testing-library/user-event";
import { setupServer } from "msw/node";
@@ -11,11 +11,12 @@ import { Client, type Message } from "@langchain/langgraph-sdk";
import {
StateGraph,
MessagesAnnotation,
START,
Runtime,
interrupt,
END,
pushMessage,
START,
END,
type Runtime,
type Pregel,
} from "@langchain/langgraph";
import { MemorySaver } from "@langchain/langgraph-checkpoint";
import { FakeStreamingChatModel } from "@langchain/core/utils/testing";
@@ -55,6 +56,7 @@ const model = new FakeStreamingChatModel({
responses: [new AIMessage("Hey")],
sleep: 100,
});
const agent = new StateGraph(MessagesAnnotation)
.addNode(
"agent",
@@ -68,111 +70,108 @@ const agent = new StateGraph(MessagesAnnotation)
.addEdge(START, "agent")
.compile();
const parentAgent = new StateGraph(MessagesAnnotation)
.addNode("child", agent, { subgraphs: [agent] })
.addEdge(START, "child")
.compile();
type AnyPregel = Pregel<any, any, any, any, any>;
const interruptAgent = new StateGraph(MessagesAnnotation)
.addNode("beforeInterrupt", async () => {
return { messages: [new AIMessage("Before interrupt")] };
})
.addNode("agent", async () => {
const resume = interrupt({ nodeName: "agent" });
return { messages: [new AIMessage(`Hey: ${resume}`)] };
})
.addNode("afterInterrupt", async () => {
return { messages: [new AIMessage("After interrupt")] };
})
.addEdge(START, "beforeInterrupt")
.addEdge("beforeInterrupt", "agent")
.addEdge("agent", "afterInterrupt")
.addEdge("afterInterrupt", END)
.compile();
const removeMessageAgent = new StateGraph(MessagesAnnotation)
.addSequence({
step1: () => ({ messages: [new AIMessage("Step 1: To Remove")] }),
step2: async (state) => {
const messages: BaseMessage[] = [
...state.messages
.filter((m) => AIMessage.isInstance(m))
.map((m) => new RemoveMessage({ id: m.id! })),
new AIMessage({ id: randomUUID(), content: "Step 2: To Keep" }),
];
for (const message of messages) {
pushMessage(message, { stateKey: null });
}
return { messages };
},
step3: () => ({ messages: [new AIMessage("Step 3: To Keep")] }),
})
.addEdge(START, "step1")
.compile();
const app = createEmbedServer({
graph: { agent, parentAgent, interruptAgent, removeMessageAgent },
checkpointer,
threads,
});
const server = setupServer(http.all("*", (ctx) => app.fetch(ctx.request)));
function TestChatComponent() {
const { messages, isLoading, error, submit, stop } = useStream({
assistantId: "agent",
apiKey: "test-api-key",
});
return (
<div>
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
<div data-testid="loading">
{isLoading ? "Loading..." : "Not loading"}
</div>
{error ? <div data-testid="error">{String(error)}</div> : null}
<button
data-testid="submit"
onClick={() =>
submit({ messages: [{ content: "Hello", type: "human" }] })
}
>
Send
</button>
<button data-testid="stop" onClick={stop}>
Stop
</button>
</div>
);
function getHttpHandler(graph: Record<string, AnyPregel>) {
const app = createEmbedServer({ graph, checkpointer, threads });
return http.all("*", (ctx) => app.fetch(ctx.request));
}
describe("useStream", () => {
beforeEach(() => server.listen());
const server = setupServer();
afterEach(() => {
server.resetHandlers();
server.close();
vi.clearAllMocks();
});
beforeEach(() => server.listen());
afterEach(() => server.close());
it(
"renders initial state correctly",
server.boundary(() => {
server.use(getHttpHandler({ agent }));
function TestChatComponent() {
const { messages, isLoading, error, submit, stop } = useStream({
assistantId: "agent",
apiKey: "test-api-key",
});
return (
<div>
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
<div data-testid="loading">
{isLoading ? "Loading..." : "Not loading"}
</div>
{error ? <div data-testid="error">{String(error)}</div> : null}
<button
data-testid="submit"
onClick={() =>
submit({ messages: [{ content: "Hello", type: "human" }] })
}
>
Send
</button>
<button data-testid="stop" onClick={stop}>
Stop
</button>
</div>
);
}
it("renders initial state correctly", () => {
render(<TestChatComponent />);
expect(screen.getByTestId("loading")).toHaveTextContent("Not loading");
expect(screen.getByTestId("messages")).toBeEmptyDOMElement();
expect(screen.queryByTestId("error")).not.toBeInTheDocument();
});
})
);
it(
"handles message submission and streaming",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
function TestChatComponent() {
const { messages, isLoading, error, submit, stop } = useStream({
assistantId: "agent",
apiKey: "test-api-key",
});
return (
<div>
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
<div data-testid="loading">
{isLoading ? "Loading..." : "Not loading"}
</div>
{error ? <div data-testid="error">{String(error)}</div> : null}
<button
data-testid="submit"
onClick={() =>
submit({ messages: [{ content: "Hello", type: "human" }] })
}
>
Send
</button>
<button data-testid="stop" onClick={stop}>
Stop
</button>
</div>
);
}
it("handles message submission and streaming", async () => {
const user = userEvent.setup();
render(<TestChatComponent />);
@@ -189,9 +188,50 @@ describe("useStream", () => {
// Check final state
expect(screen.getByTestId("loading")).toHaveTextContent("Not loading");
});
})
);
it(
"handles stop functionality",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
function TestChatComponent() {
const { messages, isLoading, error, submit, stop } = useStream({
assistantId: "agent",
apiKey: "test-api-key",
});
return (
<div>
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
<div data-testid="loading">
{isLoading ? "Loading..." : "Not loading"}
</div>
{error ? <div data-testid="error">{String(error)}</div> : null}
<button
data-testid="submit"
onClick={() =>
submit({ messages: [{ content: "Hello", type: "human" }] })
}
>
Send
</button>
<button data-testid="stop" onClick={stop}>
Stop
</button>
</div>
);
}
it("handles stop functionality", async () => {
const user = userEvent.setup();
render(<TestChatComponent />);
@@ -203,9 +243,14 @@ describe("useStream", () => {
await waitFor(() => {
expect(screen.getByTestId("loading")).toHaveTextContent("Not loading");
});
});
})
);
it(
"displays initial values immediately and clears them when submitting",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("displays initial values immediately and clears them when submitting", async () => {
const user = userEvent.setup();
function TestCachedComponent() {
@@ -276,9 +321,14 @@ describe("useStream", () => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("message-1")).toHaveTextContent("Hey");
});
});
})
);
it(
"accepts newThreadId option without errors",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("accepts newThreadId option without errors", async () => {
const user = userEvent.setup();
const spy = vi.fn();
@@ -328,9 +378,14 @@ describe("useStream", () => {
assistant_id: "agent",
},
});
});
})
);
it(
"onStop does not clear stream values",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("onStop does not clear stream values", async () => {
const user = userEvent.setup();
function TestComponent() {
@@ -390,9 +445,14 @@ describe("useStream", () => {
});
expect(screen.getByTestId("message-1")).toHaveTextContent("H");
});
})
);
it(
"onStop callback is called when stop is called",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("onStop callback is called when stop is called", async () => {
const user = userEvent.setup();
const onStopCallback = vi.fn();
@@ -428,9 +488,14 @@ describe("useStream", () => {
mutate: expect.any(Function),
})
);
});
})
);
it(
"onStop mutate function updates stream values immediately",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("onStop mutate function updates stream values immediately", async () => {
const user = userEvent.setup();
function TestComponent() {
@@ -489,9 +554,14 @@ describe("useStream", () => {
"Stream stopped"
);
});
});
})
);
it(
"onStop handles functional updates correctly",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("onStop handles functional updates correctly", async () => {
const user = userEvent.setup();
function TestComponent() {
@@ -542,9 +612,14 @@ describe("useStream", () => {
"item1, item2, stopped"
);
});
});
})
);
it(
"onStop is not called when stream completes naturally",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("onStop is not called when stream completes naturally", async () => {
const user = userEvent.setup();
const onStopCallback = vi.fn();
@@ -574,9 +649,14 @@ describe("useStream", () => {
await waitFor(() => {
expect(onStopCallback).not.toHaveBeenCalled();
});
});
})
);
it(
"make sure to pass metadata to the thread",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("make sure to pass metadata to the thread", async () => {
const user = userEvent.setup();
const onStopCallback = vi.fn();
@@ -628,9 +708,14 @@ describe("useStream", () => {
const thread = await client.threads.get(threadId);
expect(thread.metadata).toMatchObject({ random: "123" });
});
})
);
it(
"branching",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("branching", async () => {
const user = userEvent.setup();
function BranchControls(props: {
@@ -841,62 +926,76 @@ describe("useStream", () => {
within(screen.getByTestId("message-1")).getByRole("navigation")
).toHaveTextContent("1 / 2");
});
});
})
);
it.each([false, { limit: 2 }])(
"fetchStateHistory: %s",
async (fetchStateHistory) => {
const user = userEvent.setup();
it.each([false, { limit: 2 }])(
"fetchStateHistory: %s",
server.boundary(async (fetchStateHistory) => {
server.use(getHttpHandler({ agent }));
function TestComponent() {
const { submit, messages, isLoading } = useStream({
assistantId: "agent",
apiKey: "test-api-key",
fetchStateHistory,
});
const user = userEvent.setup();
return (
<div>
<div data-testid="loading">
{isLoading ? "Loading..." : "Not loading"}
</div>
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
<button
data-testid="submit"
onClick={() =>
submit({ messages: [{ content: "Hello", type: "human" }] })
}
>
Send
</button>
function TestComponent() {
const { submit, messages, isLoading } = useStream({
assistantId: "agent",
apiKey: "test-api-key",
fetchStateHistory,
});
return (
<div>
<div data-testid="loading">
{isLoading ? "Loading..." : "Not loading"}
</div>
);
}
render(<TestComponent />);
await user.click(screen.getByTestId("submit"));
await waitFor(() => {
expect(screen.getByTestId("loading")).toHaveTextContent("Loading...");
});
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("message-1")).toHaveTextContent("Hey");
expect(screen.getByTestId("loading")).toHaveTextContent("Not loading");
});
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
<button
data-testid="submit"
onClick={() =>
submit({ messages: [{ content: "Hello", type: "human" }] })
}
>
Send
</button>
</div>
);
}
);
it("streamSubgraphs: true", async () => {
render(<TestComponent />);
await user.click(screen.getByTestId("submit"));
await waitFor(() => {
expect(screen.getByTestId("loading")).toHaveTextContent("Loading...");
});
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("message-1")).toHaveTextContent("Hey");
expect(screen.getByTestId("loading")).toHaveTextContent("Not loading");
});
})
);
it(
"streamSubgraphs: true",
server.boundary(async () => {
server.use(
getHttpHandler({
parentAgent: new StateGraph(MessagesAnnotation)
.addNode("child", agent, { subgraphs: [agent] })
.addEdge(START, "child")
.compile(),
})
);
const user = userEvent.setup();
const onCheckpointEvent = vi.fn();
@@ -1008,9 +1107,14 @@ describe("useStream", () => {
expect(onCustomEvent.mock.calls).toMatchObject([
["Custom events", { namespace: [expect.any(String)] }],
]);
});
})
);
it(
"streamMetadata",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("streamMetadata", async () => {
const user = userEvent.setup();
function TestComponent() {
@@ -1060,9 +1164,14 @@ describe("useStream", () => {
expect(screen.getByTestId("message-1")).toHaveTextContent("Hey");
expect(screen.getByTestId("stream-metadata")).toHaveTextContent("agent");
});
});
})
);
it(
"onRequest gets called when a request is made",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
it("onRequest gets called when a request is made", async () => {
const user = userEvent.setup();
const onRequestCallback = vi.fn();
@@ -1128,101 +1237,151 @@ describe("useStream", () => {
},
],
]);
});
})
);
it.each([[{ fetchStateHistory: false }], [{ fetchStateHistory: true }]])(
"interrupts (%s)",
async ({ fetchStateHistory }) => {
const user = userEvent.setup();
it.each([[{ fetchStateHistory: false }], [{ fetchStateHistory: true }]])(
"interrupts (%s)",
server.boundary(async ({ fetchStateHistory }) => {
server.use(
getHttpHandler({
interruptAgent: new StateGraph(MessagesAnnotation)
.addNode("beforeInterrupt", async () => {
return { messages: [new AIMessage("Before interrupt")] };
})
.addNode("agent", async () => {
const resume = interrupt({ nodeName: "agent" });
return { messages: [new AIMessage(`Hey: ${resume}`)] };
})
.addNode("afterInterrupt", async () => {
return { messages: [new AIMessage("After interrupt")] };
})
.addEdge(START, "beforeInterrupt")
.addEdge("beforeInterrupt", "agent")
.addEdge("agent", "afterInterrupt")
.addEdge("afterInterrupt", END)
.compile(),
})
);
function TestComponent() {
const { submit, interrupt, messages } = useStream<
{ messages: Message[] },
{ InterruptType: { nodeName: string } }
>({
assistantId: "interruptAgent",
apiKey: "test-api-key",
fetchStateHistory,
});
const user = userEvent.setup();
return (
<div>
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
{interrupt ? (
<>
<div data-testid="interrupt">
{interrupt.when ?? interrupt.value?.nodeName}
</div>
<button
data-testid="resume"
onClick={() =>
submit(null, { command: { resume: "Resuming" } })
}
>
Resume
</button>
</>
) : null}
<button
data-testid="submit"
onClick={() =>
submit(
{ messages: [{ content: "Hello", type: "human" }] },
{ interruptBefore: ["beforeInterrupt"] }
)
}
>
Send
</button>
function TestComponent() {
const { submit, interrupt, messages } = useStream<
{ messages: Message[] },
{ InterruptType: { nodeName: string } }
>({
assistantId: "interruptAgent",
apiKey: "test-api-key",
fetchStateHistory,
});
return (
<div>
<div data-testid="messages">
{messages.map((msg, i) => (
<div key={msg.id ?? i} data-testid={`message-${i}`}>
{typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)}
</div>
))}
</div>
);
}
render(<TestComponent />);
await user.click(screen.getByTestId("submit"));
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("interrupt")).toHaveTextContent("breakpoint");
});
await user.click(screen.getByTestId("resume"));
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("message-1")).toHaveTextContent(
"Before interrupt"
);
expect(screen.getByTestId("interrupt")).toHaveTextContent("agent");
});
await user.click(screen.getByTestId("resume"));
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("message-1")).toHaveTextContent(
"Before interrupt"
);
expect(screen.getByTestId("message-2")).toHaveTextContent(
"Hey: Resuming"
);
expect(screen.getByTestId("message-3")).toHaveTextContent(
"After interrupt"
);
});
{interrupt ? (
<>
<div data-testid="interrupt">
{interrupt.when ?? interrupt.value?.nodeName}
</div>
<button
data-testid="resume"
onClick={() =>
submit(null, { command: { resume: "Resuming" } })
}
>
Resume
</button>
</>
) : null}
<button
data-testid="submit"
onClick={() =>
submit(
{ messages: [{ content: "Hello", type: "human" }] },
{ interruptBefore: ["beforeInterrupt"] }
)
}
>
Send
</button>
</div>
);
}
);
it("handle message removal", async () => {
render(<TestComponent />);
await user.click(screen.getByTestId("submit"));
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("interrupt")).toHaveTextContent("breakpoint");
});
await user.click(screen.getByTestId("resume"));
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("message-1")).toHaveTextContent(
"Before interrupt"
);
expect(screen.getByTestId("interrupt")).toHaveTextContent("agent");
});
await user.click(screen.getByTestId("resume"));
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello");
expect(screen.getByTestId("message-1")).toHaveTextContent(
"Before interrupt"
);
expect(screen.getByTestId("message-2")).toHaveTextContent(
"Hey: Resuming"
);
expect(screen.getByTestId("message-3")).toHaveTextContent(
"After interrupt"
);
});
})
);
it(
"handle message removal",
server.boundary(async () => {
server.use(
getHttpHandler({
removeMessageAgent: new StateGraph(MessagesAnnotation)
.addSequence({
step1: () => ({ messages: [new AIMessage("Step 1: To Remove")] }),
step2: async (state) => {
const messages: BaseMessage[] = [
...state.messages
.filter((m) => AIMessage.isInstance(m))
.map((m) => new RemoveMessage({ id: m.id! })),
new AIMessage({ id: randomUUID(), content: "Step 2: To Keep" }),
];
for (const message of messages) {
pushMessage(message, { stateKey: null });
}
return { messages };
},
step3: () => ({ messages: [new AIMessage("Step 3: To Keep")] }),
})
.addEdge(START, "step1")
.compile(),
})
);
const user = userEvent.setup();
const messagesValues = new Set<string>();
@@ -1292,5 +1451,77 @@ describe("useStream", () => {
["human: Hello", "ai: Step 2: To Keep", "ai: Step 3: To Keep"],
].map((msg) => msg.join("\n"))
);
});
});
})
);
it(
"enqueue multiple .submit() calls",
server.boundary(async () => {
server.use(getHttpHandler({ agent }));
const user = userEvent.setup();
const messagesValues = new Set<string>();
function TestComponent() {
const { submit, messages, isLoading } = useStream({
assistantId: "agent",
apiKey: "test-api-key",
});
const rawMessages = messages.map((msg, i) => ({
id: msg.id ?? i,
content: `${msg.type}: ${
typeof msg.content === "string"
? msg.content
: JSON.stringify(msg.content)
}`,
}));
messagesValues.add(rawMessages.map((msg) => msg.content).join("\n"));
return (
<div>
<div data-testid="loading">
{isLoading ? "Loading..." : "Not loading"}
</div>
<div data-testid="messages">
{rawMessages.map((msg, i) => (
<div key={msg.id} data-testid={`message-${i}`}>
<span>{msg.content}</span>
</div>
))}
</div>
<button
data-testid="submit-first"
onClick={() =>
submit({ messages: [{ content: "Hello (1)", type: "human" }] })
}
>
Send First
</button>
<button
data-testid="submit-second"
onClick={() =>
submit({ messages: [{ content: "Hello (2)", type: "human" }] })
}
>
Send Second
</button>
</div>
);
}
render(<TestComponent />);
await user.click(screen.getByTestId("submit-first"));
await user.click(screen.getByTestId("submit-second"));
await waitFor(() => {
expect(screen.getByTestId("message-0")).toHaveTextContent("Hello (1)");
expect(screen.getByTestId("message-1")).toHaveTextContent("Hey");
expect(screen.getByTestId("message-2")).toHaveTextContent("Hello (2)");
expect(screen.getByTestId("message-3")).toHaveTextContent("Hey");
});
})
);
+14 -20
View File
@@ -245,9 +245,6 @@ export function useStreamLGP<
hasTaskListener,
]);
const clearCallbackRef = useRef<() => void>(null!);
clearCallbackRef.current = stream.clear;
const threadIdRef = useRef<string | null>(threadId);
const threadIdStreamingRef = useRef<string | null>(null);
@@ -384,29 +381,26 @@ export function useStreamLGP<
// to ensure we're not accidentally submitting to a wrong branch
includeImplicitBranch;
stream.setStreamValues(() => {
const prev = shouldRefetch
? historyValues
: { ...historyValues, ...stream.values };
if (submitOptions?.optimisticValues != null) {
return {
...prev,
...(typeof submitOptions.optimisticValues === "function"
? submitOptions.optimisticValues(prev)
: submitOptions.optimisticValues),
};
}
return { ...prev };
});
let callbackMeta: RunCallbackMeta | undefined;
let rejoinKey: `lg:stream:${string}` | undefined;
let usableThreadId = threadId;
await stream.start(
async (signal: AbortSignal) => {
stream.setStreamValues((values) => {
const prev = { ...historyValues, ...(values ?? {}) };
if (submitOptions?.optimisticValues != null) {
return {
...prev,
...(typeof submitOptions.optimisticValues === "function"
? submitOptions.optimisticValues(prev)
: submitOptions.optimisticValues),
};
}
return { ...prev };
});
if (!usableThreadId) {
const thread = await client.threads.create({
threadId: submitOptions?.threadId,
+47 -5
View File
@@ -106,6 +106,10 @@ export class StreamManager<
private throttle: number | boolean;
private queue: Promise<unknown> = Promise.resolve();
private queueSize: number = 0;
private state: {
isLoading: boolean;
values: [values: StateType, kind: "stream" | "stop"] | null;
@@ -223,7 +227,7 @@ export class StreamManager<
return expected === actual || actual.startsWith(`${expected}|`);
};
start = async (
protected enqueue = async (
action: (
signal: AbortSignal
) => Promise<
@@ -255,10 +259,9 @@ export class StreamManager<
onFinish?: () => void;
}
): Promise<void> => {
if (this.state.isLoading) return;
) => {
try {
this.queueSize = Math.max(0, this.queueSize - 1);
this.setState({ isLoading: true, error: undefined });
this.abortRef = new AbortController();
@@ -342,7 +345,9 @@ export class StreamManager<
if (streamError != null) throw streamError;
const values = await options.onSuccess?.();
if (typeof values !== "undefined") this.setStreamValues(values);
if (typeof values !== "undefined" && this.queueSize === 0) {
this.setStreamValues(values);
}
} catch (error) {
if (
!(
@@ -361,6 +366,43 @@ export class StreamManager<
}
};
start = async (
action: (
signal: AbortSignal
) => Promise<
AsyncGenerator<
EventStreamEvent<
StateType,
GetUpdateType<Bag, StateType>,
GetCustomEventType<Bag>
>
>
>,
options: {
getMessages: (values: StateType) => Message[];
setMessages: (current: StateType, messages: Message[]) => StateType;
initialValues: StateType;
callbacks: StreamManagerEventCallbacks<StateType, Bag>;
onSuccess: () =>
| StateType
| null
| undefined
| void
| Promise<StateType | null | undefined | void>;
onError: (error: unknown) => void | Promise<void>;
onFinish?: () => void;
}
): Promise<void> => {
this.queueSize += 1;
this.queue = this.queue.then(() => this.enqueue(action, options));
};
stop = async (
historyValues: StateType,
options: {