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
1 change: 1 addition & 0 deletions eslint.config.js
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ export default tseslint.config(
'.integration-tests/**',
'packages/**/.integration-test/**',
'dist/**',
'demo/**/dist/**',
'docs-site/.next/**',
'docs-site/out/**',
'.qwen/**',
Expand Down
112 changes: 112 additions & 0 deletions packages/acp-bridge/src/bridge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -384,6 +384,118 @@ describe('createAcpSessionBridge', () => {
await bridge.shutdown();
});

it('refreshes extensions across live sessions and broadcasts merged results', async () => {
const handles: ChannelHandle[] = [];
const bridge = makeBridge({
channelFactory: async () => {
const h = makeChannel({
extMethodImpl: (method, params) => {
if (
method === 'qwen/control/workspace/extensions/refresh' &&
String(params['sessionId']).endsWith('#2')
) {
throw new Error('refresh failed');
}
return {};
},
});
handles.push(h);
return h.channel;
},
});
const first = await bridge.spawnOrAttach({
workspaceCwd: WS_A,
sessionScope: 'thread',
});
const second = await bridge.spawnOrAttach({
workspaceCwd: WS_A,
sessionScope: 'thread',
});
const abort = new AbortController();
const iter = bridge.subscribeEvents(first.sessionId, {
signal: abort.signal,
});
const nextEvent = iter[Symbol.asyncIterator]().next();

const result = await bridge.refreshExtensionsForAllSessions({
status: 'updated',
name: 'test-ext',
});

expect(result).toEqual({ refreshed: 1, failed: 1 });
expect(handles[0]?.agent.extMethodCalls).toEqual([
{
method: 'qwen/control/workspace/extensions/refresh',
params: { sessionId: first.sessionId },
},
{
method: 'qwen/control/workspace/extensions/refresh',
params: { sessionId: second.sessionId },
},
]);
const event = await nextEvent;
expect(event.value).toMatchObject({
type: 'extensions_changed',
data: {
status: 'updated',
name: 'test-ext',
refreshed: 1,
failed: 1,
},
});
abort.abort();
await bridge.shutdown();
});

it('does not refresh or broadcast extensions when no sessions are live', async () => {
const bridge = makeBridge();

await expect(bridge.refreshExtensionsForAllSessions()).resolves.toEqual({
refreshed: 0,
failed: 0,
});

await bridge.shutdown();
});

it('skips dying sessions when refreshing extensions', async () => {
let releaseKill: (() => void) | undefined;
const handles: ChannelHandle[] = [];
const bridge = makeBridge({
channelFactory: async () => {
const h = makeChannel();
const originalKill = h.channel.kill;
h.channel.kill = async () => {
await new Promise<void>((resolve) => {
releaseKill = resolve;
});
await originalKill();
};
handles.push(h);
return h.channel;
},
});
const session = await bridge.spawnOrAttach({ workspaceCwd: WS_A });
const killPromise = bridge.killSession(session.sessionId);
await vi.waitFor(() => {
expect(releaseKill).toBeDefined();
});

await expect(bridge.refreshExtensionsForAllSessions()).resolves.toEqual({
refreshed: 0,
failed: 0,
});
expect(
handles[0]?.agent.extMethodCalls.filter(
(call) => call.method === 'qwen/control/workspace/extensions/refresh',
),
).toEqual([]);

releaseKill?.();
await killPromise;
await bridge.shutdown();
});

it('rejects session status requests for unknown sessions', async () => {
const bridge = makeBridge({
channelFactory: async () => makeChannel().channel,
Expand Down
55 changes: 55 additions & 0 deletions packages/acp-bridge/src/bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3788,6 +3788,61 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge {
);
},

async refreshExtensionsForAllSessions(data) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Suggestion] refreshExtensionsForAllSessions iterates all sessions across all workspaces (Array.from(byId.values())). In a multi-workspace daemon, an extension mutation in workspace A triggers unnecessary workspaceExtensionsRefresh calls and extensions_changed broadcasts for workspace B sessions.

Consider accepting a workspaceCwd filter parameter and skipping sessions whose entry.workspaceCwd does not match the workspace where the mutation occurred.

— qwen3.7-max via Qwen Code /review

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Valid optimization, but I am leaving it for follow-up. Scoping refresh by workspace would require carrying workspace identity through bridge/facade/event callers and updating the surrounding tests. The current PR keeps the existing broadcast model and focuses on making extension mutations refresh loaded sessions correctly.

const sessions = Array.from(byId.values());

const results = await Promise.all(
sessions.map(async (entry) => {
const info = channelInfoForEntry(entry);
if (!info || info.isDying) {
return { refreshed: 0, failed: 0 };
}
try {
await Promise.race([
withTimeout(
entry.connection.extMethod(
SERVE_CONTROL_EXT_METHODS.workspaceExtensionsRefresh,
{ sessionId: entry.sessionId },
),
30_000,
SERVE_CONTROL_EXT_METHODS.workspaceExtensionsRefresh,
),
getTransportClosedReject(entry),
]);
return { refreshed: 1, failed: 0 };
} catch (err) {
writeServeDebugLine(
`refreshExtensions: session ${entry.sessionId} failed: ` +
`${err instanceof Error ? err.message : String(err)}`,
);
return { refreshed: 0, failed: 1 };
}
}),
);

const refreshed = results.reduce(
(sum, result) => sum + result.refreshed,
0,
);
const failed = results.reduce((sum, result) => sum + result.failed, 0);

if (refreshed > 0 || failed > 0 || data?.status !== undefined) {
broadcastWorkspaceEvent({
type: 'extensions_changed',
data: { ...data, refreshed, failed },
});
}

return { refreshed, failed };
},

broadcastExtensionsChanged(data) {
broadcastWorkspaceEvent({
type: 'extensions_changed',
data,
});
},

async setSessionModel(sessionId, req, context) {
const entry = byId.get(sessionId);
if (!entry) throw new SessionNotFoundError(sessionId);
Expand Down
30 changes: 30 additions & 0 deletions packages/acp-bridge/src/bridgeTypes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,22 @@ export interface BridgeDaemonStatusSnapshot {
sessions: BridgeDaemonSessionDiagnostic[];
}

export interface BridgeExtensionsChangedData {
refreshed: number;
failed: number;
status?:
| 'installed'
| 'enabled'
| 'disabled'
| 'updated'
| 'uninstalled'
| 'failed';
source?: string;
name?: string;
version?: string;
error?: string;
}

export interface AcpSessionBridge {
/** Read-only daemon diagnostics for status endpoints. */
getDaemonStatusSnapshot(): BridgeDaemonStatusSnapshot;
Expand Down Expand Up @@ -472,6 +488,20 @@ export interface AcpSessionBridge {
/** Read workspace-level installed extension status. */
getWorkspaceExtensionsStatus(): Promise<ServeWorkspaceExtensionsStatus>;

/**
* Broadcast extension refresh to all active sessions and emit an
* `extensions_changed` workspace event when complete.
*/
refreshExtensionsForAllSessions(
data?: Omit<BridgeExtensionsChangedData, 'refreshed' | 'failed'>,
): Promise<{
refreshed: number;
failed: number;
}>;

/** Emit an extension lifecycle event without refreshing sessions. */
broadcastExtensionsChanged(data: BridgeExtensionsChangedData): void;

/**
* Switch the active model service for a session. Throws
* `SessionNotFoundError` for unknown ids.
Expand Down
23 changes: 23 additions & 0 deletions packages/acp-bridge/src/status.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,7 @@ export const SERVE_CONTROL_EXT_METHODS = {
workspaceMcpRuntimeAdd: 'qwen/control/workspace/mcp/runtime-add',
workspaceMcpRuntimeRemove: 'qwen/control/workspace/mcp/runtime-remove',
workspaceReload: 'qwen/control/workspace/reload',
workspaceExtensionsRefresh: 'qwen/control/workspace/extensions/refresh',
} as const;

export type ServeStatus =
Expand Down Expand Up @@ -872,6 +873,26 @@ export interface ServeExtensionCapabilities {
hasSettings: boolean;
}

export type ServeExtensionUpdateState =
| 'checking for updates'
| 'updated, needs restart'
| 'updating'
| 'updated'
| 'update available'
| 'up to date'
| 'error'
| 'not updatable'
| 'unknown';

export interface ServeExtensionDetails {
mcpServers: string[];
commands: string[];
skills: string[];
agents: string[];
contextFiles: string[];
settings: string[];
}

export interface ServeExtensionEntry {
kind: 'extension';
id: string;
Expand All @@ -885,7 +906,9 @@ export interface ServeExtensionEntry {
originSource?: ServeExtensionOriginSource;
ref?: string;
autoUpdate?: boolean;
updateState?: ServeExtensionUpdateState;
capabilities: ServeExtensionCapabilities;
details?: ServeExtensionDetails;
}

export interface ServeWorkspaceExtensionsStatus {
Expand Down
112 changes: 112 additions & 0 deletions packages/cli/src/acp-integration/acpAgent.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6139,4 +6139,116 @@ describe('sessionLanguage multi-session propagation', () => {
mockConnectionState.resolve();
await agentPromise;
});

it('refreshes extension commands for the live session', async () => {
const extensionManager = {
refreshCache: vi.fn().mockResolvedValue(undefined),
refreshTools: vi.fn().mockResolvedValue(undefined),
};
const cfg = makeConfig({
getSessionId: vi.fn().mockReturnValue('s-ext'),
getExtensionManager: vi.fn().mockReturnValue(extensionManager),
});
const sendAvailableCommandsUpdate = vi.fn().mockResolvedValue(undefined);

vi.mocked(loadSettings).mockReturnValue({
merged: { mcpServers: {} },
getUserHooks: vi.fn().mockReturnValue({}),
getProjectHooks: vi.fn().mockReturnValue({}),
} as unknown as LoadedSettings);
vi.mocked(loadCliConfig).mockResolvedValue(cfg as unknown as Config);
vi.mocked(Session).mockImplementation(
() =>
({
getId: vi.fn().mockReturnValue('s-ext'),
getConfig: vi.fn().mockReturnValue(cfg),
sendAvailableCommandsUpdate,
installRewriter: vi.fn(),
startCronScheduler: vi.fn(),
dispose: vi.fn(),
}) as unknown as InstanceType<typeof Session>,
);

const agentPromise = runAcpAgent(
makeConfig() as unknown as Config,
{ merged: { mcpServers: {} } } as unknown as LoadedSettings,
mockArgv,
);
await vi.waitFor(() => expect(capturedAgentFactory).toBeDefined());
const agent = capturedAgentFactory!({
get closed() {
return mockConnectionState.promise;
},
});

await agent.newSession({ cwd: '/ext', mcpServers: [] });
await expect(
agent.extMethod(SERVE_CONTROL_EXT_METHODS.workspaceExtensionsRefresh, {
sessionId: 's-ext',
}),
).resolves.toEqual({ ok: true });

expect(extensionManager.refreshCache).toHaveBeenCalledOnce();
expect(extensionManager.refreshTools).toHaveBeenCalledOnce();
expect(sendAvailableCommandsUpdate).toHaveBeenCalledOnce();

mockConnectionState.resolve();
await agentPromise;
});

it('still sends available commands update when extension tool refresh fails', async () => {
const extensionManager = {
refreshCache: vi.fn().mockResolvedValue(undefined),
refreshTools: vi.fn().mockRejectedValue(new Error('bad tool schema')),
};
const cfg = makeConfig({
getSessionId: vi.fn().mockReturnValue('s-ext'),
getExtensionManager: vi.fn().mockReturnValue(extensionManager),
});
const sendAvailableCommandsUpdate = vi.fn().mockResolvedValue(undefined);

vi.mocked(loadSettings).mockReturnValue({
merged: { mcpServers: {} },
getUserHooks: vi.fn().mockReturnValue({}),
getProjectHooks: vi.fn().mockReturnValue({}),
} as unknown as LoadedSettings);
vi.mocked(loadCliConfig).mockResolvedValue(cfg as unknown as Config);
vi.mocked(Session).mockImplementation(
() =>
({
getId: vi.fn().mockReturnValue('s-ext'),
getConfig: vi.fn().mockReturnValue(cfg),
sendAvailableCommandsUpdate,
installRewriter: vi.fn(),
startCronScheduler: vi.fn(),
dispose: vi.fn(),
}) as unknown as InstanceType<typeof Session>,
);

const agentPromise = runAcpAgent(
makeConfig() as unknown as Config,
{ merged: { mcpServers: {} } } as unknown as LoadedSettings,
mockArgv,
);
await vi.waitFor(() => expect(capturedAgentFactory).toBeDefined());
const agent = capturedAgentFactory!({
get closed() {
return mockConnectionState.promise;
},
});

await agent.newSession({ cwd: '/ext', mcpServers: [] });
await expect(
agent.extMethod(SERVE_CONTROL_EXT_METHODS.workspaceExtensionsRefresh, {
sessionId: 's-ext',
}),
).resolves.toEqual({ ok: true });

expect(extensionManager.refreshCache).toHaveBeenCalledOnce();
expect(extensionManager.refreshTools).toHaveBeenCalledOnce();
expect(sendAvailableCommandsUpdate).toHaveBeenCalledOnce();

mockConnectionState.resolve();
await agentPromise;
});
});
Loading
Loading