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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all \u003cpre\u003e\u003ccode\u003e blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks"); } } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); } })(); (function(){ try { var __m = "github.com"; var __re = new RegExp('^' + "github\\.com" + '
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length \u003e 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Added
- [EE] Added guided reconnection for MCP connector authentication failures during Ask Sourcebot agent turns. [#1548](https://github.com/sourcebot-dev/sourcebot/pull/1548)

### Removed
- Removed the Langfuse integration. [#1536](https://github.com/sourcebot-dev/sourcebot/pull/1536)

Expand Down
2 changes: 1 addition & 1 deletion packages/web/src/app/api/(client)/client.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -320,7 +320,7 @@ export const getOffers = async (): Promise<OffersResponse | ServiceError> => {
return result as OffersResponse | ServiceError;
}

export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string }): Promise<ConnectMcpResponse | ServiceError> => {
export const connectMcpToAsk = async (body: { serverId: string; returnTo?: string; forceAuthorization?: boolean }): Promise<ConnectMcpResponse | ServiceError> => {
const result = await fetch('/api/ee/askmcp/connect', {
method: 'POST',
headers: {
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -48,7 +48,7 @@ vi.mock('@ai-sdk/mcp', () => ({
const { POST } = await import('./route');
const { getMcpOAuthReturnToFromState } = await import('@/ee/features/chat/mcp/mcpOAuthReturnTo');

function createRequest(body: { serverId: string; returnTo?: string } = { serverId: 'server-1' }) {
function createRequest(body: { serverId: string; returnTo?: string; forceAuthorization?: boolean } = { serverId: 'server-1' }) {
return new NextRequest('https://sourcebot.example.com/api/ee/askmcp/connect', {
method: 'POST',
headers: { 'content-type': 'application/json' },
Expand DownExpand Up@@ -229,6 +229,60 @@ describe('POST /api/ee/askmcp/connect', () => {
});
});

test('forces an interactive OAuth redirect for reconnect recovery', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
tx.userMcpServer.findUnique.mockResolvedValue({
tokens: 'encrypted:{"access_token":"stale-token"}',
codeVerifier: null,
state: null,
});
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockImplementation(async (provider) => {
await expect(provider.tokens()).resolves.toBeUndefined();
provider.authorizationUrl = 'https://oauth.example.com/authorize';
return 'REDIRECT';
});

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(await response.json()).toEqual({
authorizationUrl: 'https://oauth.example.com/authorize',
});
});

test('does not report a forced reconnect as successful without an OAuth redirect', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
mocks.authContext = {
org: { id: 1 },
user: { id: 'user-1' },
prisma,
};
mocks.unsafePrisma.$transaction.mockImplementation(async (callback, _options) => callback(tx));
mocks.mcpAuth.mockResolvedValue('AUTHORIZED');

const response = await POST(createRequest({
serverId: 'server-1',
returnTo: '/chat/abc123',
forceAuthorization: true,
}));

expect(response.status).toBe(502);
expect(await response.json()).toMatchObject({
message: 'Could not start connector reauthorization.',
});
});

test('ignores unsafe return paths', async () => {
const prisma = createPrismaMock();
const tx = createTransactionMock();
Expand Down
16 changes: 15 additions & 1 deletion packages/web/src/app/api/(server)/ee/askmcp/connect/route.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,6 +24,7 @@ import { getEnabledMcpOAuthScopeNames } from '@/ee/features/chat/mcp/oauthScopeU
const bodySchema = z.object({
serverId: z.string(),
returnTo: z.string().optional(),
forceAuthorization: z.boolean().optional().default(false),
});
const logger = createLogger('mcp-connect');
const MCP_AUTH_FETCH_TIMEOUT_MS = Math.min(env.SOURCEBOT_MCP_TOOL_CALL_TIMEOUT_MS, 30000);
Expand DownExpand Up@@ -146,6 +147,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
callbackReturnTo,
allowClientRegistration: true,
requestedOAuthScopes: getEnabledMcpOAuthScopeNames(mcpServer.oauthScopes),
forceAuthorization: parsed.data.forceAuthorization,
});

let authResult: Awaited<ReturnType<typeof mcpAuth>>;
Expand DownExpand Up@@ -200,7 +202,7 @@ export const POST = apiHandler(async (request: NextRequest) => {
throw error;
}

if (connectResult.authResult === 'AUTHORIZED') {
if (connectResult.authResult === 'AUTHORIZED' && !parsed.data.forceAuthorization) {
// Already has valid tokens (e.g., refreshed)
void captureEvent('ask_mcp_connector_connection_completed', {
...eventProperties,
Expand All@@ -209,6 +211,18 @@ export const POST = apiHandler(async (request: NextRequest) => {
return { authorizationUrl: null } satisfies ConnectMcpResponse;
}

if (connectResult.authResult === 'AUTHORIZED') {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
failureReason: 'missing_authorization_url',
});
throw new ServiceErrorException({
statusCode: StatusCodes.BAD_GATEWAY,
errorCode: ErrorCode.UNEXPECTED_ERROR,
message: 'Could not start connector reauthorization.',
});
}

if (!connectResult.authorizationUrl) {
void captureEvent('ask_mcp_connector_connection_failed', {
...eventProperties,
Expand Down
44 changes: 44 additions & 0 deletions packages/web/src/ee/features/chat/agent.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,6 +321,50 @@ describe('createMessageStream approval continuation', () => {
});
});

test('streams the connector ID when its tools fail to load', async () => {
const { getConnectedMcpClients } = await import('@/ee/features/chat/mcp/mcpClientFactory');
const { getMcpTools } = await import('@/ee/features/chat/mcp/mcpToolSets');
vi.mocked(getConnectedMcpClients).mockResolvedValueOnce([
{ serverId: 'server-linear', serverName: 'Linear' },
] as never);
vi.mocked(getMcpTools).mockResolvedValueOnce({
tools: {},
failedServers: [{ serverId: 'server-linear', serverName: 'Linear' }],
serverFaviconUrls: {},
toolDisplayNames: {},
cleanup: vi.fn(),
});
mockAi.streamText.mockReturnValue(createFakeStreamResult());

await createMessageStream({
chatId: 'chat-id',
messages: [createUserMessage()],
selectedRepos: [],
disabledMcpServerIds: [],
prisma: {},
model: {},
modelName: 'test-model',
promptCacheStrategy: noopStrategy,
onFinish: vi.fn(),
onError: () => 'error',
userId: 'user-id',
orgId: 1,
} as unknown as Parameters<typeof createMessageStream>[0]);

const execute = mockAi.latestCreateUIMessageStreamOptions?.execute;
if (!execute) {
throw new Error('Expected createUIMessageStream to capture execute callback.');
}

const write = vi.fn();
await execute({ writer: { merge: vi.fn(), write } });

expect(write).toHaveBeenCalledWith({
type: 'data-mcp-failed-server',
data: { serverId: 'server-linear', serverName: 'Linear' },
});
});

test.each([
['dynamic', dynamicApprovalRespondedPart],
['static', staticApprovalRespondedPart],
Expand Down
90 changes: 83 additions & 7 deletions packages/web/src/ee/features/chat/agent.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,12 @@ import { addLineNumbers, fileReferenceToString, formatAttachmentsForPrompt, getA
import { createTools } from "./tools";
import { getConnectedMcpClients } from "@/ee/features/chat/mcp/mcpClientFactory";
import { getMcpTools, McpToolsResult } from "@/ee/features/chat/mcp/mcpToolSets";
import {
createMcpAuthInterruptionDirective,
denyApprovedToolApprovalsForAuthInterruption,
getMcpAuthRequiredFailureFromAssistantMessage,
McpToolAuthFailure,
} from "@/ee/features/chat/mcp/mcpAuthFailure";
import { buildMcpToolRegistry, McpToolRegistryEntry } from "@/ee/features/chat/mcp/mcpToolRegistry";
import { PromptCacheStrategy, mergeProviderOptions, detectPromptCacheBreak, detectUnexpectedCacheMiss } from "./promptCaching";
import { hasEntitlement } from '@/lib/entitlements';
Expand DownExpand Up@@ -332,9 +338,21 @@ export const createMessageStream = async ({
? (lastMsg.metadata as SBChatMessageMetadata | undefined)
: undefined;

// When the response was interrupted by a reconnect-required authentication
// failure (detected via the safe tool error's marker text), the
// continuation must run its final step with tool use disabled. Any
// approval that was still approved is rewritten to a denial: once the
// response is authentication-terminal, later approval actions are invalid.
const priorMcpAuthFailure = hasApprovalContinuationReady
? getMcpAuthRequiredFailureFromAssistantMessage(lastMsg)
: undefined;

if (hasApprovalContinuationReady) {
const continuationMessage = priorMcpAuthFailure
? denyApprovedToolApprovalsForAuthInterruption(lastMsg, priorMcpAuthFailure.serverName)
: lastMsg;
const fullLastTurn = await convertToModelMessages(
[lastMsg],
[continuationMessage],
{ ignoreIncompleteToolCalls: true }
);
messageHistory = [...messageHistory, ...fullLastTurn];
Expand DownExpand Up@@ -375,12 +393,22 @@ export const createMessageStream = async ({
data: { modelToolName, rawToolName },
});
},
onMcpServerFailed: (serverName) => {
onMcpServerFailed: (server) => {
writer.write({
type: 'data-mcp-failed-server',
data: { serverName },
data: server,
});
},
onMcpAuthRequired: (failure) => {
// Transient: consumed live by the client to surface the
// connector reconnect UI, never folded into persisted parts.
writer.write({
type: 'data-mcp-auth-required',
data: failure,
transient: true,
});
},
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -508,7 +536,14 @@ interface AgentOptions {
onWriteSource: (source: Source) => void;
onMcpServerDiscovered: (sanitizedName: string, faviconUrl: string) => void;
onMcpToolDiscovered: (modelToolName: string, rawToolName: string) => void;
onMcpServerFailed: (serverName: string) => void;
onMcpServerFailed: (server: { serverId: string; serverName: string }) => void;
// Fired at most once per connector per response when a tool call fails
// with a reconnect-required authentication failure.
onMcpAuthRequired: (failure: McpToolAuthFailure) => void;
// Set when the incoming messages show this response was already
// interrupted by an authentication failure (approval continuation): the
// stream must run its final step with tool use disabled from step one.
priorMcpAuthFailure?: { serverName: string };
traceId: string;
chatId: string;
prisma: PrismaClient;
Expand All@@ -529,6 +564,8 @@ const createAgentStream = async ({
onMcpServerDiscovered,
onMcpToolDiscovered,
onMcpServerFailed,
onMcpAuthRequired,
priorMcpAuthFailure,
traceId,
chatId,
prisma,
Expand DownExpand Up@@ -564,6 +601,15 @@ const createAgentStream = async ({
}))
).filter((source) => source !== undefined);

// Mutable, response-scoped authentication failure state. `serverName` is
// the first failed connector's display name (V1 supports recovery for a
// single failed connector). Failures are deduplicated by connector so the
// client sees at most one transient event per connector per response.
const mcpAuthFailureState: { failure?: { serverName: string } } = {
...(priorMcpAuthFailure ? { failure: { serverName: priorMcpAuthFailure.serverName } } : {}),
};
const reportedMcpAuthFailureServerIds = new Set<string>();

let mcpToolSetsObj: McpToolsResult = { tools: {}, failedServers: [], serverFaviconUrls: {}, toolDisplayNames: {}, cleanup: async () => {} };
if (userId && orgId && await hasEntitlement('ask') && disabledMcpServerIds !== undefined) {
try {
Expand All@@ -573,6 +619,16 @@ const createAgentStream = async ({
chatId,
traceId,
source: 'sourcebot-ask-agent',
}, {
onAuthFailure: (failure) => {
if (!mcpAuthFailureState.failure) {
mcpAuthFailureState.failure = { serverName: failure.serverName };
}
if (!reportedMcpAuthFailureServerIds.has(failure.serverId)) {
reportedMcpAuthFailureServerIds.add(failure.serverId);
onMcpAuthRequired(failure);
}
},
});

for (const [sanitizedName, faviconUrl] of Object.entries(mcpToolSetsObj.serverFaviconUrls)) {
Expand All@@ -590,8 +646,8 @@ const createAgentStream = async ({
}
}

for (const serverName of mcpToolSetsObj.failedServers) {
onMcpServerFailed(serverName);
for (const server of mcpToolSetsObj.failedServers) {
onMcpServerFailed(server);
}

const mcpRegistry = buildMcpToolRegistry(mcpToolSetsObj.tools);
Expand DownExpand Up@@ -715,14 +771,34 @@ const createAgentStream = async ({
// rebuilds the step's messages each time as the original input plus
// its own accumulated response messages. Re-applying the moving tail marker
// to the new last message each step is safe and does not accumulate.
prepareStep: (tailMarker || hasMcpTools) ? ({ steps, messages }) => {
prepareStep: (tailMarker || hasMcpTools || mcpAuthFailureState.failure) ? ({ steps, messages }) => {
const stepMessages = (tailMarker && messages.length > 0)
? messages.map((message, index) =>
index === messages.length - 1
? { ...message, providerOptions: mergeProviderOptions(message.providerOptions, tailMarker) }
: message)
: undefined;

// Once a reconnect-required authentication failure occurs, the
// response is terminal for tool use: every remaining step runs
// with tool calling disabled (`toolChoice: 'none'` keeps the
// tool definitions byte-stable for prompt caching) plus an
// ephemeral directive to summarize completed work and prompt
// the user to reconnect. In-flight tool calls of the failing
// step have already run to completion by the time this fires.
if (mcpAuthFailureState.failure) {
return {
messages: [
...(stepMessages ?? messages),
{
role: 'user' as const,
content: createMcpAuthInterruptionDirective(mcpAuthFailureState.failure.serverName),
},
],
toolChoice: 'none' as const,
};
}

if (!hasMcpTools) {
return stepMessages ? { messages: stepMessages } : {};
}
Expand Down
Loading
Loading