fix(api): apply namespace to messages-tuple stream mode (#1386)

This commit is contained in:
David Duong
2025-07-14 14:03:05 +02:00
committed by GitHub
parent ab58412bb1
commit d172de38b9
3 changed files with 29 additions and 9 deletions
+7
View File
@@ -0,0 +1,7 @@
---
"@langchain/langgraph-api": patch
"@langchain/langgraph-cli": patch
"@langchain/langgraph-ui": patch
---
Fix apply namespace to messages-tuple stream mode
+5 -1
View File
@@ -271,7 +271,11 @@ export async function* streamState(
if (mode === "messages") {
if (userStreamMode.includes("messages-tuple")) {
yield { event: "messages", data };
if (kwargs.subgraphs && ns?.length) {
yield { event: `messages|${ns.join("|")}`, data };
} else {
yield { event: "messages", data };
}
}
} else if (userStreamMode.includes(mode)) {
if (kwargs.subgraphs && ns?.length) {
+17 -8
View File
@@ -3,7 +3,10 @@ import type {
BaseMessageLike,
MessageType,
} from "@langchain/core/messages";
import { Client, type FeedbackStreamEvent } from "@langchain/langgraph-sdk";
import {
Client,
type MessagesTupleStreamEvent,
} from "@langchain/langgraph-sdk";
import { RemoteGraph } from "@langchain/langgraph/remote";
import postgres from "postgres";
import { afterAll, beforeAll, beforeEach, describe, expect, it } from "vitest";
@@ -936,18 +939,24 @@ describe("runs", () => {
const stream = await client.runs.stream(
thread.thread_id,
assistant.assistant_id,
{ input, streamMode: "messages-tuple", config: globalConfig }
{
input,
streamMode: "messages-tuple",
config: globalConfig,
streamSubgraphs: true,
}
);
const chunks = await gatherIterator(stream);
const runId = findLast(
chunks,
(i): i is FeedbackStreamEvent => i.event === "metadata"
)?.data.run_id;
const runId = findLast(chunks, (i) => i.event === "metadata")?.data.run_id;
expect(runId).not.toBeUndefined();
expect(runId).not.toBeNull();
const messages = chunks
.filter((i) => i.event === "messages")
.filter(
(i): i is MessagesTupleStreamEvent =>
i.event.startsWith("messages|") || i.event === "messages"
)
.map((i) => i.data[0]);
expect(messages).toHaveLength("begin".length + "end".length + 1);
@@ -957,7 +966,7 @@ describe("runs", () => {
..."end".split("").map((c) => ({ content: c })),
]);
const seenEventTypes = new Set(chunks.map((i) => i.event));
const seenEventTypes = new Set(chunks.map((i) => i.event.split("|")[0]));
expect(seenEventTypes).toEqual(new Set(["metadata", "messages"]));
const run = await client.runs.get(thread.thread_id, runId as string);