Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion api/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@
"@azure/storage-blob": "^12.30.0",
"@google/genai": "^2.8.0",
"@keyv/redis": "5.1.6",
"@librechat/agents": "^3.9.6",
"@librechat/agents": "^3.9.7",
"@librechat/api": "*",
"@librechat/data-schemas": "*",
"@microsoft/microsoft-graph-client": "^3.0.7",
Expand Down
96 changes: 96 additions & 0 deletions api/server/controllers/agents/client.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -3172,6 +3172,102 @@ describe('AgentClient - startup telemetry', () => {
errorSpy.mockRestore();
});

it.each([
[
'closed',
'SocketError',
ErrorTypes.MODEL_STREAM_CLOSED,
'stream_closed',
'The model provider closed the connection before the response finished. Try again.',
],
[
'stalled',
'BodyTimeoutError',
ErrorTypes.MODEL_STREAM_STALLED,
'stream_stalled',
'The model provider stopped sending the response, and the request timed out. Try again.',
],
])(
'keeps partial content and a safe %s model error in the agent turn',
async (kind, causeName, type, errorType, prose) => {
jest.clearAllMocks();
const { logger } = require('@librechat/data-schemas');
const { errors: undiciErrors } = require('undici');
const privateValue = 'PRIVATE-TRANSPORT-DIAGNOSTIC';
const cause = new undiciErrors[causeName](privateValue);
const providerError = new TypeError('terminated', { cause });
const errorSpy = jest.spyOn(logger, 'error').mockImplementation(() => logger);
mockCreateRun.mockImplementation(async (options) => {
const tracker = options.modelCallbacks.find(
(callback) => callback.name === 'librechat-upstream-model-error-tracker',
);
return {
Graph: null,
processStream: jest.fn(async () => {
tracker.handleLLMError(providerError);
throw new Error('graph failed', { cause: providerError });
}),
getCalibrationRatio: jest.fn(() => 0),
};
});
mockIsHITLEnabled.mockReturnValue(false);
const partial = { type: ContentTypes.TEXT, [ContentTypes.TEXT]: 'Partial findings' };
const client = new AgentClient({
req: {
user: { id: 'user-123' },
body: {},
config: {
endpoints: { [EModelEndpoint.agents]: {} },
filters: { messages: { pii: {} } },
},
_resumableStreamId: `conversation-stream-${kind}`,
},
res: {},
agent: {
id: 'agent-123',
endpoint: EModelEndpoint.openAI,
provider: EModelEndpoint.openAI,
model_parameters: { model: 'gpt-4' },
hide_sequential_outputs: false,
},
endpointTokenConfig: {},
eventHandlers: {},
contentParts: [partial],
collectedUsage: [],
artifactPromises: [],
});
client.conversationId = `conversation-stream-${kind}`;
client.responseMessageId = `response-stream-${kind}`;
client.parentMessageId = `parent-stream-${kind}`;
client.recordCollectedUsage = jest.fn().mockResolvedValue();

try {
await client.chatCompletion({ payload: [] });

expect(client.contentParts).toEqual(
expect.arrayContaining([
partial,
{
type: ContentTypes.ERROR,
[ContentTypes.ERROR]: `${prose}\n${JSON.stringify({ type })}`,
},
]),
);
expect(JSON.stringify(client.contentParts)).not.toContain(privateValue);
expect(client.recordCollectedUsage).toHaveBeenCalledWith(
expect.objectContaining({ context: 'message' }),
);
expect(errorSpy).toHaveBeenCalledWith(
'[api/server/controllers/agents/client.js #sendCompletion] Upstream model error',
expect.objectContaining({ errorType }),
);
expect(JSON.stringify(errorSpy.mock.calls)).not.toContain(privateValue);
} finally {
errorSpy.mockRestore();
}
},
);

/** A compaction's only record of having been one is the marker on the part it
* produced, and Compact runs on whatever leaf the branch ends with. Without
* the marker on the failure, a compaction that failed on a user leaf keeps a
Expand Down
24 changes: 24 additions & 0 deletions api/server/services/Files/Code/process.js
Original file line number Diff line number Diff line change
Expand Up @@ -1088,10 +1088,14 @@ async function readWorkspaceFile({
req,
signal,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
}) {
return executeWorkspaceTool({
baseURL: codeApiBaseUrl,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
/** Minted per admission attempt: a queued call outlives one token TTL. */
authHeaders: async () => ({
...(await getCodeApiAuthHeaders(req, bridgeWorkerId)),
Expand Down Expand Up @@ -1137,10 +1141,14 @@ async function searchWorkspace({
req,
signal,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
}) {
return executeWorkspaceTool({
baseURL: codeApiBaseUrl,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
/** Minted per admission attempt: a queued call outlives one token TTL. */
authHeaders: async () => ({
...(await getCodeApiAuthHeaders(req, bridgeWorkerId)),
Expand Down Expand Up @@ -1186,10 +1194,14 @@ async function listWorkspaceFiles({
req,
signal,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
}) {
return executeWorkspaceTool({
baseURL: codeApiBaseUrl,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
/** Minted per admission attempt: a queued call outlives one token TTL. */
authHeaders: async () => ({
...(await getCodeApiAuthHeaders(req, bridgeWorkerId)),
Expand Down Expand Up @@ -1221,10 +1233,14 @@ async function writeWorkspaceFile({
req,
signal,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
}) {
return executeWorkspaceTool({
baseURL: codeApiBaseUrl,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
/** Minted per admission attempt: a queued call outlives one token TTL. */
authHeaders: async () => ({
...(await getCodeApiAuthHeaders(req, bridgeWorkerId)),
Expand Down Expand Up @@ -1256,10 +1272,14 @@ async function editWorkspaceFile({
req,
signal,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
}) {
return executeWorkspaceTool({
baseURL: codeApiBaseUrl,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
/** Minted per admission attempt: a queued call outlives one token TTL. */
authHeaders: async () => ({
...(await getCodeApiAuthHeaders(req, bridgeWorkerId)),
Expand Down Expand Up @@ -1290,10 +1310,14 @@ async function previewWorkspaceEdit({
req,
signal,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
}) {
return executeWorkspaceTool({
baseURL: codeApiBaseUrl,
maxQueueWaitMs,
maxRequestTimeoutMs,
deadlineAtMs,
/** Minted per admission attempt: a queued call outlives one token TTL. */
authHeaders: async () => ({
...(await getCodeApiAuthHeaders(req, bridgeWorkerId)),
Expand Down
18 changes: 18 additions & 0 deletions api/server/services/Files/Code/process.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -2061,6 +2061,8 @@ describe('Code Process', () => {
req: mockReq,
signal: controller.signal,
maxQueueWaitMs: 0,
maxRequestTimeoutMs: 125_000,
deadlineAtMs: 160_000,
}),
).resolves.toBe(result);

Expand All @@ -2081,6 +2083,8 @@ describe('Code Process', () => {
baseURL: 'https://attached-code.example.com/v1',
authHeaders: expect.any(Function),
maxQueueWaitMs: 0,
maxRequestTimeoutMs: 125_000,
deadlineAtMs: 160_000,
request: {
protocolVersion: 1,
operation: 'read_file',
Expand Down Expand Up @@ -2140,6 +2144,8 @@ describe('Code Process', () => {
baseURL: 'https://attached-code.example.com/v1',
authHeaders: expect.any(Function),
maxQueueWaitMs: 0,
maxRequestTimeoutMs: undefined,
deadlineAtMs: undefined,
request: {
protocolVersion: 1,
operation: 'search_text',
Expand Down Expand Up @@ -2198,6 +2204,8 @@ describe('Code Process', () => {
baseURL: 'https://attached-code.example.com/v1',
authHeaders: expect.any(Function),
maxQueueWaitMs: 0,
maxRequestTimeoutMs: undefined,
deadlineAtMs: undefined,
request: {
protocolVersion: 1,
operation: 'list_files',
Expand Down Expand Up @@ -2258,6 +2266,8 @@ describe('Code Process', () => {
baseURL: 'https://attached-code.example.com/v1',
authHeaders: expect.any(Function),
maxQueueWaitMs: 0,
maxRequestTimeoutMs: undefined,
deadlineAtMs: undefined,
request: {
protocolVersion: 1,
operation: 'write_file',
Expand Down Expand Up @@ -2296,11 +2306,15 @@ describe('Code Process', () => {
bridgeWorkerId: 'worker-user-1',
req: mockReq,
expected_base_sha256: 'a'.repeat(64),
maxRequestTimeoutMs: 125_000,
deadlineAtMs: 160_000,
}),
).resolves.toBe(result);

expect(mockExecuteWorkspaceTool).toHaveBeenCalledWith(
expect.objectContaining({
maxRequestTimeoutMs: 125_000,
deadlineAtMs: 160_000,
request: {
protocolVersion: 1,
operation: 'edit_file',
Expand Down Expand Up @@ -2338,11 +2352,15 @@ describe('Code Process', () => {
codeApiBaseUrl: 'https://attached-code.example.com/v1',
executionProfile: 'stateful',
req: mockReq,
maxRequestTimeoutMs: 125_000,
deadlineAtMs: 160_000,
}),
).resolves.toBe(result);

expect(mockExecuteWorkspaceTool).toHaveBeenCalledWith(
expect.objectContaining({
maxRequestTimeoutMs: 125_000,
deadlineAtMs: 160_000,
request: {
protocolVersion: 1,
operation: 'preview_edit',
Expand Down
12 changes: 12 additions & 0 deletions api/server/services/MCP.js
Original file line number Diff line number Diff line change
Expand Up @@ -842,6 +842,8 @@ async function reconnectServer({
upstreamTokenProviderResolver,
recoveryPolicy,
oboIdentityContext,
streamId,
jobCreatedAt,
forceNew: true,
returnOnOAuth: false,
connectionTimeout: Time.THIRTY_SECONDS,
Expand Down Expand Up @@ -1386,6 +1388,16 @@ function createToolInstance({
);
}

// The schedule service checks the typed cause and verified job identity before
// recording a durable tool failure; other tool errors are a cheap no-op.
await require('~/server/services/Schedules').recordMCPToolAuthFailure({
error,
streamId,
jobCreatedAt,
userId,
serverName,
});

/** Carries the actionable re-auth message; the substring heuristic below would misreport it as an OAuth configuration problem */
if (
error instanceof OpenIDReauthRequiredError ||
Expand Down
50 changes: 50 additions & 0 deletions api/server/services/MCP.spec.js
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,10 @@ jest.mock('./Tools/mcp', () => ({
reinitMCPServer: jest.fn(),
}));

jest.mock('~/server/services/Schedules', () => ({
recordMCPToolAuthFailure: jest.fn(async () => true),
}));

jest.mock('./GraphTokenService', () => ({
getGraphApiToken: jest.fn(),
}));
Expand Down Expand Up @@ -1544,6 +1548,44 @@ describe('User parameter passing tests', () => {
});

describe('createMCPTool', () => {
it('records a typed OBO failure against the scheduled generation before returning an error', async () => {
const user = { id: 'scheduled-owner', role: 'USER' };
const missing = new Error('Unattended provider not configured');
const error = Object.assign(new Error('MCP tool error'), { cause: missing });
const receipt = require('~/server/services/Schedules').recordMCPToolAuthFailure;
require('~/models').getRoleByName.mockResolvedValue({
permissions: { [PermissionTypes.MCP_SERVERS]: { [Permissions.USE]: true } },
});
mockGetMCPManager.mockReturnValue({ callTool: jest.fn().mockRejectedValue(error) });
const tool = await createMCPTool({
user,
toolKey: `test-tool${D}test-server`,
provider: 'openai',
streamId: 'scheduled-conversation',
jobCreatedAt: 42,
config: { type: 'streamable-http', url: 'https://mcp.example.com' },
availableTools: {
[`test-tool${D}test-server`]: {
function: { description: 'Test MCP', parameters: { type: 'object', properties: {} } },
},
},
});
await expect(
tool.func({}, undefined, {
configurable: { user },
metadata: { provider: 'openai', thread_id: 'scheduled-conversation', run_id: 'run-1' },
toolCall: {},
}),
).rejects.toThrow();
expect(receipt).toHaveBeenCalledWith({
error,
streamId: 'scheduled-conversation',
jobCreatedAt: 42,
userId: 'scheduled-owner',
serverName: 'test-server',
});
});

it('keeps shared OAuth recovery alive when one tool caller aborts', async () => {
const mockUser = { id: 'shared-recovery-user', role: 'USER' };
const mockRes = { write: jest.fn(), flush: jest.fn() };
Expand Down Expand Up @@ -2783,13 +2825,17 @@ describe('User parameter passing tests', () => {
serverName: 'server1',
provider: 'anthropic',
userMCPAuthMap: {},
streamId: 'scheduled-stream',
jobCreatedAt: 42,
});

// Verify all calls to reinitMCPServer had the user
expect(reinitCalls.length).toBeGreaterThan(0);
reinitCalls.forEach((call) => {
expect(call.user).toBe(mockUser);
expect(call.user.id).toBe('user-001');
expect(call.streamId).toBe('scheduled-stream');
expect(call.jobCreatedAt).toBe(42);
});
});

Expand Down Expand Up @@ -2817,12 +2863,16 @@ describe('User parameter passing tests', () => {
provider: 'google',
userMCPAuthMap: {},
availableTools: undefined, // Force reinit
streamId: 'resumed-stream',
jobCreatedAt: 43,
});

// Verify the call to reinitMCPServer had the user
expect(reinitCalls.length).toBe(1);
expect(reinitCalls[0].user).toBe(mockUser);
expect(reinitCalls[0].user.id).toBe('user-002');
expect(reinitCalls[0].streamId).toBe('resumed-stream');
expect(reinitCalls[0].jobCreatedAt).toBe(43);
});
});

Expand Down
Loading
Loading