mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-02-06 17:51:50 +01:00
📬 feat: Implement Delta Buffering System for Out-of-Order SSE Events (#11643)
* ✨ test: Add MCP tool definitions tests for server name variants
- Introduced new test cases for loading MCP tools with underscored and hyphenated server names, ensuring correct functionality and handling of tool definitions.
- Validated that the tool definitions are loaded accurately based on different server name formats, enhancing test coverage for the MCP tool integration.
- Included assertions to verify the expected behavior and properties of the loaded tools, improving reliability and maintainability of the tests.
* refactor: useStepHandler to support additional delta events and buffer management
- Added support for Agents.ReasoningDeltaEvent and Agents.RunStepDeltaEvent in the TStepEvent type.
- Introduced a pendingDeltaBuffer to store deltas that arrive before their corresponding run step, ensuring they are processed in the correct order.
- Updated event handling to buffer deltas when no corresponding run step is found, improving the reliability of message processing.
- Cleared the pendingDeltaBuffer during cleanup to prevent memory leaks.
This commit is contained in:
parent
c8e4257342
commit
1ba5bf87b0
3 changed files with 1297 additions and 3 deletions
1085
client/src/hooks/SSE/__tests__/useStepHandler.spec.ts
Normal file
1085
client/src/hooks/SSE/__tests__/useStepHandler.spec.ts
Normal file
File diff suppressed because it is too large
Load diff
|
|
@ -31,6 +31,8 @@ type TStepEvent = {
|
|||
event: string;
|
||||
data:
|
||||
| Agents.MessageDeltaEvent
|
||||
| Agents.ReasoningDeltaEvent
|
||||
| Agents.RunStepDeltaEvent
|
||||
| Agents.AgentUpdate
|
||||
| Agents.RunStep
|
||||
| Agents.ToolEndEvent
|
||||
|
|
@ -61,6 +63,8 @@ export default function useStepHandler({
|
|||
const toolCallIdMap = useRef(new Map<string, string | undefined>());
|
||||
const messageMap = useRef(new Map<string, TMessage>());
|
||||
const stepMap = useRef(new Map<string, Agents.RunStep>());
|
||||
/** Buffer for deltas that arrive before their corresponding run step */
|
||||
const pendingDeltaBuffer = useRef(new Map<string, TStepEvent[]>());
|
||||
|
||||
/**
|
||||
* Calculate content index for a run step.
|
||||
|
|
@ -350,6 +354,14 @@ export default function useStepHandler({
|
|||
|
||||
setMessages(updatedMessages);
|
||||
}
|
||||
|
||||
const bufferedDeltas = pendingDeltaBuffer.current.get(runStep.id);
|
||||
if (bufferedDeltas && bufferedDeltas.length > 0) {
|
||||
pendingDeltaBuffer.current.delete(runStep.id);
|
||||
for (const bufferedDelta of bufferedDeltas) {
|
||||
stepHandler({ event: bufferedDelta.event, data: bufferedDelta.data }, submission);
|
||||
}
|
||||
}
|
||||
} else if (event === 'on_agent_update') {
|
||||
const { agent_update } = data as Agents.AgentUpdate;
|
||||
let responseMessageId = agent_update.runId || '';
|
||||
|
|
@ -391,7 +403,9 @@ export default function useStepHandler({
|
|||
}
|
||||
|
||||
if (!runStep || !responseMessageId) {
|
||||
console.warn('No run step or runId found for message delta event');
|
||||
const buffer = pendingDeltaBuffer.current.get(messageDelta.id) ?? [];
|
||||
buffer.push({ event: 'on_message_delta', data: messageDelta });
|
||||
pendingDeltaBuffer.current.set(messageDelta.id, buffer);
|
||||
return;
|
||||
}
|
||||
|
||||
|
|
@ -432,7 +446,9 @@ export default function useStepHandler({
|
|||
}
|
||||
|
||||
if (!runStep || !responseMessageId) {
|
||||
console.warn('No run step or runId found for reasoning delta event');
|
||||
const buffer = pendingDeltaBuffer.current.get(reasoningDelta.id) ?? [];
|
||||
buffer.push({ event: 'on_reasoning_delta', data: reasoningDelta });
|
||||
pendingDeltaBuffer.current.set(reasoningDelta.id, buffer);
|
||||
return;
|
||||
}
|
||||
|
||||
|
|
@ -473,7 +489,9 @@ export default function useStepHandler({
|
|||
}
|
||||
|
||||
if (!runStep || !responseMessageId) {
|
||||
console.warn('No run step or runId found for run step delta event');
|
||||
const buffer = pendingDeltaBuffer.current.get(runStepDelta.id) ?? [];
|
||||
buffer.push({ event: 'on_run_step_delta', data: runStepDelta });
|
||||
pendingDeltaBuffer.current.set(runStepDelta.id, buffer);
|
||||
return;
|
||||
}
|
||||
|
||||
|
|
@ -578,6 +596,7 @@ export default function useStepHandler({
|
|||
toolCallIdMap.current.clear();
|
||||
messageMap.current.clear();
|
||||
stepMap.current.clear();
|
||||
pendingDeltaBuffer.current.clear();
|
||||
}, []);
|
||||
|
||||
/**
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue