feat: ship deterministic parser and lifecycle release

This commit is contained in:
inman committed 2026-08-28 17:13:58 +08:00
1 parent 208c434b31
commit 7b5d855b09
176 files changed
+20248 -1362

No files matched your search

+103 -36
View File
@@ -19,6 +19,7 @@ export interface PublicAgentBusChannel {
enabled: boolean;
status: 'disabled' | 'connecting' | 'connected' | 'error';
key_configured: true;
deletable: boolean;
last_connected_at: string | null;
last_error: string | null;
created_at: string;
@@ -57,8 +58,13 @@ function iso(value: unknown): string | null {
return value ? new Date(String(value)).toISOString() : null;
}
function publicChannel(row: Record<string, unknown>): PublicAgentBusChannel {
function publicChannel(
row: Record<string, unknown>,
legacyEnvironmentManaged = false
): PublicAgentBusChannel {
const status = text(row.status);
const environmentManaged = text(row.external_user_ref) === LEGACY_CHANNEL_REF
&& legacyEnvironmentManaged;
return {
id: text(row.id),
organization_id: text(row.organization_id),
@@ -71,6 +77,7 @@ function publicChannel(row: Record<string, unknown>): PublicAgentBusChannel {
? status as PublicAgentBusChannel['status']
: 'error',
key_configured: true,
deletable: !environmentManaged,
last_connected_at: iso(row.last_connected_at),
last_error: text(row.last_error) || null,
created_at: new Date(String(row.created_at)).toISOString(),
@@ -115,6 +122,13 @@ export class AgentBusChannelService {
constructor(private readonly config: AppConfig) {}
private legacyEnvironmentManaged(): boolean {
return this.config.agentBusEnabled
&& Boolean(text(this.config.AGENTBUS_WS_URL))
&& Boolean(text(this.config.AGENTBUS_WS_TOKEN))
&& Boolean(text(this.config.AGENTBUS_BOT_ADDRESS));
}
async list(organizationId: string): Promise<PublicAgentBusChannel[]> {
const result = await getPool(this.config).query(
`WITH canonical_legacy AS (
@@ -133,7 +147,9 @@ export class AgentBusChannelService {
ORDER BY channel.created_at ASC, channel.id ASC`,
[organizationId]
);
return (result.rows as Record<string, unknown>[]).map(publicChannel);
const legacyEnvironmentManaged = this.legacyEnvironmentManaged();
return (result.rows as Record<string, unknown>[])
.map((row) => publicChannel(row, legacyEnvironmentManaged));
}
async listEnabledSecrets(organizationId: string): Promise<AgentBusChannelSecret[]> {
@@ -171,7 +187,7 @@ export class AgentBusChannelService {
continue;
}
channels.push({
...publicChannel(row),
...publicChannel(row, this.legacyEnvironmentManaged()),
ws_url: wsUrl,
ws_token: wsToken,
bot_address: botAddress
@@ -217,7 +233,7 @@ export class AgentBusChannelService {
});
return inserted.rows[0] as Record<string, unknown>;
});
return publicChannel(result);
return publicChannel(result, this.legacyEnvironmentManaged());
} catch (error) {
if (error && typeof error === 'object' && String((error as { code?: unknown }).code || '') === '23505') {
throw new TaskError('channel_name_conflict', '同一组织下的渠道名称已存在。', 409);
@@ -276,7 +292,7 @@ export class AgentBusChannelService {
});
return updated.rows[0] as Record<string, unknown>;
});
return publicChannel(result);
return publicChannel(result, this.legacyEnvironmentManaged());
} catch (error) {
if (error && typeof error === 'object' && String((error as { code?: unknown }).code || '') === '23505') {
throw new TaskError('channel_name_conflict', '同一组织下的渠道名称已存在。', 409);
@@ -305,7 +321,51 @@ export class AgentBusChannelService {
await this.audit(client, context, 'agentbus_channel.key_rotated', channelId, {});
return updated.rows[0] as Record<string, unknown>;
});
return publicChannel(result);
return publicChannel(result, this.legacyEnvironmentManaged());
}
async delete(
context: TaskContext,
channelId: string
): Promise<{ channel_id: string; deleted: true }> {
return withTransaction(this.config, async (client) => {
const lookup = await client.query(
`SELECT id, display_name, external_user_ref, enabled
FROM user_channels
WHERE organization_id = $1 AND id = $2
FOR UPDATE`,
[context.organizationId, channelId]
);
if (!lookup.rowCount) throw new TaskError('channel_not_found', '用户渠道不存在。', 404);
const row = lookup.rows[0] as Record<string, unknown>;
if (text(row.external_user_ref) === LEGACY_CHANNEL_REF && this.legacyEnvironmentManaged()) {
throw new TaskError(
'channel_managed_by_environment',
'该兼容渠道仍由旧 AgentBus 环境变量托管;请先移除环境配置并重启服务,再删除数据库渠道。',
409
);
}
const related = await client.query(
`SELECT (SELECT count(*)::int FROM tasks WHERE channel_id = $1) AS linked_task_count,
(SELECT count(*)::int FROM agentbus_deliveries WHERE channel_id = $1) AS delivery_count`,
[channelId]
);
const relatedRow = related.rows[0] as Record<string, unknown>;
await this.audit(client, context, 'agentbus_channel.deleted', channelId, {
display_name: text(row.display_name),
enabled: booleanValue(row.enabled),
linked_task_count: Number(relatedRow.linked_task_count || 0),
removed_delivery_count: Number(relatedRow.delivery_count || 0)
});
const deleted = await client.query(
'DELETE FROM user_channels WHERE organization_id = $1 AND id = $2 RETURNING id',
[context.organizationId, channelId]
);
if (!deleted.rowCount) throw new TaskError('channel_not_found', '用户渠道不存在。', 404);
return { channel_id: text(deleted.rows[0].id), deleted: true };
});
}
async ensureLegacyChannel(organizationId: string): Promise<void> {
@@ -429,6 +489,7 @@ export class AgentBusManager {
private readonly options: AgentBusManagerOptions;
private started = false;
private reloadInFlight: Promise<void> | null = null;
private reloadRequested = false;
constructor(options: AgentBusManagerOptions) {
this.options = options;
@@ -447,37 +508,42 @@ export class AgentBusManager {
}
async reload(): Promise<void> {
if (!this.started || this.reloadInFlight) return this.reloadInFlight || Promise.resolve();
if (!this.started) return;
this.reloadRequested = true;
if (this.reloadInFlight) return this.reloadInFlight;
this.reloadInFlight = (async () => {
await this.channels.ensureLegacyChannel(this.options.organizationId);
for (const listener of this.listeners.values()) listener.stop(false);
this.listeners.clear();
const secrets = await this.channels.listEnabledSecrets(this.options.organizationId);
for (const channel of secrets) {
const listener = new AgentBusListener({
config: this.options.config,
tasks: this.options.tasks,
organizationId: this.options.organizationId,
scheduleParseQueue: this.options.scheduleParseQueue,
socketFactory: this.options.socketFactory,
logger: this.options.logger,
channel: {
id: channel.id,
displayName: channel.display_name,
wsUrl: channel.ws_url,
wsToken: channel.ws_token,
botAddress: channel.bot_address
},
onStatusChange: (status, error, epoch) => this.channels.setRuntimeStatus(
this.options.organizationId,
channel.id,
status,
error,
epoch
).catch(() => undefined)
});
this.listeners.set(channel.id, listener);
listener.start();
while (this.started && this.reloadRequested) {
this.reloadRequested = false;
await this.channels.ensureLegacyChannel(this.options.organizationId);
for (const listener of this.listeners.values()) listener.stop(false);
this.listeners.clear();
const secrets = await this.channels.listEnabledSecrets(this.options.organizationId);
for (const channel of secrets) {
const listener = new AgentBusListener({
config: this.options.config,
tasks: this.options.tasks,
organizationId: this.options.organizationId,
scheduleParseQueue: this.options.scheduleParseQueue,
socketFactory: this.options.socketFactory,
logger: this.options.logger,
channel: {
id: channel.id,
displayName: channel.display_name,
wsUrl: channel.ws_url,
wsToken: channel.ws_token,
botAddress: channel.bot_address
},
onStatusChange: (status, error, epoch) => this.channels.setRuntimeStatus(
this.options.organizationId,
channel.id,
status,
error,
epoch
).catch(() => undefined)
});
this.listeners.set(channel.id, listener);
listener.start();
}
}
})().finally(() => {
this.reloadInFlight = null;
@@ -487,6 +553,7 @@ export class AgentBusManager {
async stop(): Promise<void> {
this.started = false;
this.reloadRequested = false;
if (this.reloadInFlight) await this.reloadInFlight.catch(() => undefined);
// Process shutdown is not an administrative channel disable. Do not let
// an older process write `disabled` after a newer process has connected.