basant307/AI_Governance_Project
048
1import * as path from 'node:path';2import { pathToFileURL } from 'node:url';3import { Storage } from '@qwen-code/qwen-code-core';4import type {5 SessionRouter,6 ChannelAgentBridge,7 ChannelBase,8 ChannelBaseOptions,9 ChannelPlugin,10 PermissionRequestEvent,11 PermissionResolvedEvent,12 ToolCallEvent,13} from '@qwen-code/channel-base';14import { sanitizeLogText } from '@qwen-code/channel-base';15import { loadSettings, type LoadedSettings } from '../../config/settings.js';16import { writeStderrLine, writeStdoutLine } from '../../utils/stdioHelpers.js';17import { getExtensionManager } from '../extensions/utils.js';18import { getPlugin, registerPlugin } from './channel-registry.js';19import { parseChannelConfig } from './config-utils.js';20 21export type ParsedChannelConfig = Awaited<22 ReturnType<typeof parseChannelConfig>23>;24 25export interface ParsedChannel {26 name: string;27 config: ParsedChannelConfig;28}29 30export function sessionsPath(): string {31 return path.join(Storage.getGlobalQwenDir(), 'channels', 'sessions.json');32}33 34export function channelLoopPath(): string {35 return path.join(Storage.getGlobalQwenDir(), 'channels', 'cron.json');36}37 38export function loadChannelsConfig(39 cwd: string = process.cwd(),40 settings: LoadedSettings = loadSettings(cwd),41): Record<string, unknown> {42 const channels = (43 settings.merged as unknown as { channels?: Record<string, unknown> }44 ).channels;45 return channels || {};46}47 48export function resolveExtensionChannelEntrySpecifier(49 extensionPath: string,50 entry: string,51): string {52 return pathToFileURL(path.join(extensionPath, entry)).href;53}54 55/**56 * Load channel plugins from active extensions.57 * Extensions declare channels in their qwen-extension.json manifest.58 */59export async function loadChannelsFromExtensions(): Promise<number> {60 let loaded = 0;61 try {62 const extensionManager = await getExtensionManager();63 const extensions = extensionManager64 .getLoadedExtensions()65 .filter((e) => e.isActive && e.channels);66 67 for (const ext of extensions) {68 for (const [channelType, channelDef] of Object.entries(ext.channels!)) {69 if (await getPlugin(channelType)) {70 writeStderrLine(71 `[Extensions] Skipping channel "${channelType}" from "${ext.name}": type already registered`,72 );73 continue;74 }75 76 const entrySpecifier = resolveExtensionChannelEntrySpecifier(77 ext.path,78 channelDef.entry,79 );80 try {81 const module = (await import(entrySpecifier)) as {82 plugin?: ChannelPlugin;83 };84 const plugin = module.plugin;85 86 if (!plugin || typeof plugin.createChannel !== 'function') {87 writeStderrLine(88 `[Extensions] "${ext.name}": channel entry point does not export a valid plugin object`,89 );90 continue;91 }92 93 if (plugin.channelType !== channelType) {94 writeStderrLine(95 `[Extensions] "${ext.name}": channelType mismatch — manifest says "${channelType}", plugin says "${plugin.channelType}"`,96 );97 continue;98 }99 100 registerPlugin(plugin);101 loaded++;102 writeStdoutLine(103 `[Extensions] Loaded channel "${channelType}" from "${ext.name}"`,104 );105 } catch (err) {106 writeStderrLine(107 `[Extensions] Failed to load channel "${channelType}" from "${ext.name}": ${err instanceof Error ? err.message : String(err)}`,108 );109 }110 }111 }112 } catch (err) {113 writeStderrLine(114 `[Extensions] Failed to load extensions: ${err instanceof Error ? err.message : String(err)}`,115 );116 }117 return loaded;118}119 120export async function createChannel(121 name: string,122 config: ParsedChannelConfig,123 bridge: ChannelAgentBridge,124 options?: ChannelBaseOptions,125): Promise<ChannelBase> {126 const channelPlugin = await getPlugin(config.type);127 if (!channelPlugin) {128 throw new Error(`Unknown channel type: "${config.type}".`);129 }130 return channelPlugin.createChannel(name, config, bridge, options);131}132 133export function selectFirstModel(134 parsed: ParsedChannel[],135 bridgeLabel: string,136): string | undefined {137 const models = [138 ...new Set(139 parsed140 .map((channel) => channel.config.model)141 .filter((model): model is string => Boolean(model)),142 ),143 ];144 if (models.length > 1) {145 writeStderrLine(146 `[Channel] Warning: Multiple models configured (${models.join(', ')}). ` +147 `${bridgeLabel} will use "${models[0]}".`,148 );149 }150 return models[0];151}152 153export function registerToolCallDispatch(154 bridge: ChannelAgentBridge,155 router: SessionRouter,156 channels: Map<string, ChannelBase>,157): void {158 bridge.on('toolCall', (event: ToolCallEvent) => {159 const target = router.getTarget(event.sessionId);160 if (target) {161 const channel = channels.get(target.channelName);162 if (channel) {163 channel.dispatchToolCall(event);164 }165 }166 });167}168 169function cancelPermissionRequest(170 bridge: ChannelAgentBridge,171 requestId: string,172): void {173 if (!bridge.respondToPermission) {174 return;175 }176 void bridge177 .respondToPermission(requestId, { outcome: { outcome: 'cancelled' } })178 .catch((err: unknown) => {179 writeStderrLine(180 `[Channel] Permission cancellation failed for ${sanitizeLogText(requestId, 128)}: ${err instanceof Error ? sanitizeLogText(err.message, 512) : sanitizeLogText(String(err), 512)}`,181 );182 });183}184 185export function registerPermissionRelay(186 bridge: ChannelAgentBridge,187 router: SessionRouter,188 channels: Map<string, ChannelBase>,189): void {190 bridge.on('permissionRequest', (event: PermissionRequestEvent) => {191 const target = router.getTarget(event.sessionId);192 if (!target) {193 writeStderrLine(194 `[Channel] No route for session ${sanitizeLogText(event.sessionId, 128)}; cancelling permission ${sanitizeLogText(event.requestId, 128)}`,195 );196 cancelPermissionRequest(bridge, event.requestId);197 return;198 }199 const channel = channels.get(target.channelName);200 if (!channel) {201 writeStderrLine(202 `[Channel] No channel "${sanitizeLogText(target.channelName, 64)}" for session ${sanitizeLogText(event.sessionId, 128)}; cancelling permission ${sanitizeLogText(event.requestId, 128)}`,203 );204 cancelPermissionRequest(bridge, event.requestId);205 return;206 }207 channel.dispatchPermissionRequest(event).catch((err: unknown) => {208 writeStderrLine(209 `[Channel] Permission relay failed for ${sanitizeLogText(event.requestId, 128)}: ${err instanceof Error ? sanitizeLogText(err.message, 512) : sanitizeLogText(String(err), 512)}`,210 );211 cancelPermissionRequest(bridge, event.requestId);212 });213 });214 215 bridge.on('permissionResolved', (event: PermissionResolvedEvent) => {216 for (const channel of channels.values()) {217 channel.dispatchPermissionResolved(event);218 }219 });220}221 222export function registerSessionCleanup(223 bridge: ChannelAgentBridge,224 router: SessionRouter,225 channels: Map<string, ChannelBase>,226): void {227 bridge.on('sessionDied', (event: { sessionId: string; reason?: string }) => {228 const safeId = sanitizeLogText(event.sessionId, 128);229 const safeReason = event.reason ? sanitizeLogText(event.reason, 512) : '';230 writeStderrLine(231 `[Channel] Session ${safeId} died${safeReason ? ` (${safeReason})` : ''}, removing routing state`,232 );233 const target = router.getTarget(event.sessionId);234 const channel = target ? channels.get(target.channelName) : undefined;235 if (channel) {236 channel.onSessionDied(event.sessionId);237 } else {238 router.removeSessionId(event.sessionId);239 }240 });241}242 243export async function parseConfiguredChannels(244 channelsConfig: Record<string, unknown>,245 selectedNames: string[],246 opts: { defaultCwd?: string } = {},247): Promise<ParsedChannel[]> {248 const parsed: ParsedChannel[] = [];249 for (const name of selectedNames) {250 const rawConfig = channelsConfig[name];251 if (!rawConfig || typeof rawConfig !== 'object') {252 throw new Error(253 `Error in channel "${name}": channel is not configured. Add a "${name}" entry under "channels" in settings.json.`,254 );255 }256 try {257 parsed.push({258 name,259 config: await parseChannelConfig(260 name,261 rawConfig as Record<string, unknown>,262 opts.defaultCwd,263 { resolveEnvVars: 'available' },264 ),265 });266 } catch (err) {267 throw new Error(268 `Error in channel "${name}": ${err instanceof Error ? err.message : String(err)}`,269 );270 }271 }272 return parsed;273}274 