merge: integrate leader summary webhook delivery

# Conflicts:
#	control-plane/README.md
#	control-plane/src/db.ts
#	control-plane/test/control-plane.test.ts
This commit is contained in:
inman committed 2026-09-09 17:25:23 +08:00
commit 295409bff0
21 files changed
+1260 -1269

No files matched your search

-2
View File
@@ -655,7 +655,6 @@ export interface AgentBusManagerOptions {
scheduleParseQueue: () => Promise<void>;
logger?: AgentBusLogger;
socketFactory?: AgentBusSocketFactory;
leaderNotifications?: AgentBusListenerOptions['leaderNotifications'];
}
export class AgentBusManager {
@@ -723,7 +722,6 @@ export class AgentBusManager {
scheduleParseQueue: this.options.scheduleParseQueue,
socketFactory: this.options.socketFactory,
logger: this.options.logger,
leaderNotifications: this.options.leaderNotifications,
channel: {
id: channel.id,
displayName: channel.display_name,
-147
View File
@@ -25,7 +25,6 @@ import {
AGENTBUS_ROSTER_ATTACHMENT_RECEIVED_TEXT,
publicAgentBusDeliveryPayload
} from './agentbus-delivery.js';
import type { LeaderTaskSummaryDelivery } from './leader-notification-service.js';
const OPEN_READY_STATE = 1;
const MAX_COMPLETED_TASK_IDS = 2_048;
@@ -134,24 +133,6 @@ export interface AgentBusTaskGateway {
releaseAgentBusDeliveries?(channelId: string, leaseOwner: string): Promise<void>;
}
export interface LeaderNotificationGateway {
observeLeaderRoute?(input: {
organizationId: string;
leaderUserId: string;
channelId: string;
recipientAddress: string;
conversationId?: string;
}): Promise<void>;
claimDeliveries(
channelId: string,
leaseOwner: string,
limit?: number
): Promise<LeaderTaskSummaryDelivery[]>;
markDeliveryDelivered(deliveryId: string): Promise<void>;
markDeliveryFailed(deliveryId: string, errorMessage: string): Promise<void>;
releaseDeliveries(channelId: string, leaseOwner: string): Promise<void>;
}
export interface AgentBusChannelConnection {
id: string;
displayName: string;
@@ -170,7 +151,6 @@ export interface AgentBusListenerOptions {
socketFactory?: AgentBusSocketFactory;
logger?: AgentBusLogger;
channel?: AgentBusChannelConnection;
leaderNotifications?: LeaderNotificationGateway;
onStatusChange?: (
status: 'disabled' | 'connecting' | 'connected' | 'error',
errorMessage: string | null,
@@ -538,22 +518,6 @@ function createDurableDeliveryFrame(
return frame;
}
export function createLeaderTaskSummaryFrame(
delivery: LeaderTaskSummaryDelivery,
session: AgentBusSession
): AgentBusFrame {
return {
id: `leader-summary-${delivery.id}`,
type: 'event',
from: session.address,
to: delivery.recipient_address,
session_id: session.id,
epoch: session.epoch,
conversation_id: delivery.conversation_id,
payload: { ...delivery.payload }
};
}
export function taskResultStatus(task: PublicTask): 'completed' | 'failed' {
return FAILED_TASK_STATUSES.has(text(task.status).toLowerCase()) ? 'failed' : 'completed';
}
@@ -653,7 +617,6 @@ export class AgentBusListener {
private readonly socketFactory: AgentBusSocketFactory;
private readonly logger: AgentBusLogger;
private readonly channel: AgentBusChannelConnection | null;
private readonly leaderNotifications: LeaderNotificationGateway | null;
private readonly onStatusChange: AgentBusListenerOptions['onStatusChange'];
private readonly leaseOwner: string;
private socket: AgentBusSocket | null = null;
@@ -665,7 +628,6 @@ export class AgentBusListener {
private readonly pendingFinalReplies = new Map<string, PendingFinalReply>();
private deliveryTimer: NodeJS.Timeout | null = null;
private durableFlushInFlight: Promise<void> | null = null;
private leaderFlushInFlight: Promise<void> | null = null;
constructor(options: AgentBusListenerOptions) {
this.config = options.config;
@@ -675,7 +637,6 @@ export class AgentBusListener {
this.socketFactory = options.socketFactory || defaultSocketFactory;
this.logger = options.logger || noopLogger;
this.channel = options.channel || null;
this.leaderNotifications = options.leaderNotifications || null;
this.onStatusChange = options.onStatusChange;
this.leaseOwner = `agentbus:${this.channel?.id || 'legacy'}`;
}
@@ -794,9 +755,6 @@ export class AgentBusListener {
if (this.channel && this.tasks.releaseAgentBusDeliveries) {
void this.tasks.releaseAgentBusDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined);
}
if (this.channel && this.leaderNotifications) {
void this.leaderNotifications.releaseDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined);
}
if (persistStatus) this.setRuntimeStatus('disabled', null, sessionEpoch);
}
@@ -877,9 +835,6 @@ export class AgentBusListener {
if (this.channel && this.tasks.releaseAgentBusDeliveries) {
void this.tasks.releaseAgentBusDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined);
}
if (this.channel && this.leaderNotifications) {
void this.leaderNotifications.releaseDeliveries(this.channel.id, this.leaseOwner).catch(() => undefined);
}
for (const pending of this.pendingFinalReplies.values()) pending.sending = false;
this.logger.warn({
agentbus_event: 'socket_close',
@@ -1066,25 +1021,6 @@ export class AgentBusListener {
...(this.channel ? { channelId: this.channel.id } : {})
};
const result = await this.tasks.ingestMessage(context, input);
if (
this.channel?.ownerRole === 'team_lead'
&& this.leaderNotifications?.observeLeaderRoute
) {
void this.leaderNotifications.observeLeaderRoute({
organizationId: this.organizationId,
leaderUserId: this.channel.ownerUserId,
channelId: this.channel.id,
recipientAddress: text(frame.from),
conversationId
}).catch((error) => {
this.logger.warn({
agentbus_event: 'leader_summary_route_observation_failed',
channel_id: this.channel?.id,
owner_user_id: this.channel?.ownerUserId,
...diagnosticError(error, 'leader_summary_route_observation_failed')
}, 'AgentBus leader summary route observation failed');
});
}
this.logger.info({
agentbus_event: 'task_ingested',
inbound_frame_id: taskId,
@@ -1249,90 +1185,7 @@ export class AgentBusListener {
}
private async flushOutboundDeliveries(): Promise<void> {
// Employee replies always get the first claim/send opportunity on a tick.
// Leader summaries use a separate, smaller queue and cannot delay task replies.
await this.flushDurableDeliveries();
await this.flushLeaderDeliveries();
}
private async flushLeaderDeliveries(): Promise<void> {
if (!this.channel || !this.leaderNotifications) return;
if (this.leaderFlushInFlight) return this.leaderFlushInFlight;
if (!this.socket || this.socket.readyState !== OPEN_READY_STATE || !this.session) return;
this.leaderFlushInFlight = (async () => {
const deliveries = await this.leaderNotifications!.claimDeliveries(
this.channel!.id,
this.leaseOwner,
5
);
for (const delivery of deliveries) void this.sendLeaderDelivery(delivery);
})()
.catch((error) => {
this.logger.warn({
agentbus_event: 'leader_summary_delivery_flush_failed',
channel_id: this.channel?.id || null,
...diagnosticError(error, 'leader_summary_delivery_flush_failed')
}, 'AgentBus leader task summary flush failed');
})
.finally(() => {
this.leaderFlushInFlight = null;
});
return this.leaderFlushInFlight;
}
private async sendLeaderDelivery(delivery: LeaderTaskSummaryDelivery): Promise<void> {
if (!this.socket || this.socket.readyState !== OPEN_READY_STATE || !this.session || !this.leaderNotifications) {
return;
}
const socket = this.socket;
const session = this.session;
const notifications = this.leaderNotifications;
const frame = createLeaderTaskSummaryFrame(delivery, session);
this.logger.info({
agentbus_event: 'leader_summary_delivery_sending',
channel_id: delivery.channel_id,
delivery_id: delivery.id,
task_id: delivery.task_id,
recipient_fingerprint: delivery.recipient_fingerprint,
conversation_fingerprint: delivery.conversation_fingerprint,
attempt_count: delivery.attempt_count
}, 'AgentBus leader task summary sending');
const fail = (error: unknown) => {
const errorMetadata = diagnosticError(error, 'leader_summary_delivery_failed');
void notifications.markDeliveryFailed(
delivery.id,
`${String(errorMetadata.error_code)}:${String(errorMetadata.error_fingerprint)}`
).catch(() => undefined);
this.logger.warn({
agentbus_event: 'leader_summary_delivery_failed',
channel_id: delivery.channel_id,
delivery_id: delivery.id,
task_id: delivery.task_id,
recipient_fingerprint: delivery.recipient_fingerprint,
conversation_fingerprint: delivery.conversation_fingerprint,
...errorMetadata
}, 'AgentBus leader task summary failed');
};
try {
socket.send(JSON.stringify(frame), (error) => {
if (error) {
fail(error);
return;
}
void notifications.markDeliveryDelivered(delivery.id)
.then(() => this.logger.info({
agentbus_event: 'leader_summary_delivery_sent',
channel_id: delivery.channel_id,
delivery_id: delivery.id,
task_id: delivery.task_id,
recipient_fingerprint: delivery.recipient_fingerprint,
conversation_fingerprint: delivery.conversation_fingerprint
}, 'AgentBus leader task summary sent'))
.catch(fail);
});
} catch (error) {
fail(error);
}
}
private async sendDurableDelivery(
+40
View File
@@ -9,6 +9,14 @@ const optionalString = z.preprocess(
z.string().optional()
);
const optionalExactString = z.preprocess(
(value) => {
if (value == null || value === '') return undefined;
return String(value);
},
z.string().optional()
);
const optionalUrl = z.preprocess(
(value) => {
const normalized = String(value ?? '').trim();
@@ -86,6 +94,10 @@ const envSchema = z.object({
AGENTBUS_WS_RECONNECT_DELAY: z.string().default('5s').transform(parseDurationMs),
AGENTBUS_TASK_TIMEOUT_MS: z.coerce.number().int().positive().default(300_000),
AGENTBUS_LOG_PAYLOADS: z.enum(['true', 'false']).default('false').transform((value) => value === 'true'),
// Webhook configuration is validated separately so a typo in this optional
// notification feature cannot prevent normal task or AgentBus startup.
WEBHOOK_SEND_URL: optionalString,
WEBHOOK_EXTERNAL_TOKEN: optionalExactString,
ARTIFACT_STORAGE_BACKEND: z.enum(['database', 'oss']).default('database'),
ARTIFACT_MAX_BYTES: z.coerce.number().int().positive().max(50_000_000).default(10_485_760),
DOCUMENT_CONVERTER_PATH: z.string().trim().min(1).default('soffice'),
@@ -105,6 +117,8 @@ const envSchema = z.object({
export type AppConfig = z.infer<typeof envSchema> & {
fieldEncryptionKey: Buffer;
agentBusEnabled: boolean;
leaderSummaryWebhookEnabled: boolean;
leaderSummaryWebhookConfigurationError: string | null;
ossRegion: string;
};
@@ -143,6 +157,30 @@ export function loadConfig(env: NodeJS.ProcessEnv = process.env): AppConfig {
throw new Error(`AgentBus is enabled but missing configuration: ${missing.join(', ')}`);
}
}
const webhookCredentialsPresent = Boolean(parsed.WEBHOOK_SEND_URL || parsed.WEBHOOK_EXTERNAL_TOKEN);
const webhookConfigurationComplete = Boolean(parsed.WEBHOOK_SEND_URL && parsed.WEBHOOK_EXTERNAL_TOKEN);
let leaderSummaryWebhookConfigurationError: string | null = null;
if (webhookCredentialsPresent && !webhookConfigurationComplete) {
leaderSummaryWebhookConfigurationError = 'webhook_configuration_incomplete';
} else if (parsed.WEBHOOK_EXTERNAL_TOKEN && parsed.WEBHOOK_EXTERNAL_TOKEN.length !== 32) {
leaderSummaryWebhookConfigurationError = 'webhook_token_length_invalid';
} else if (parsed.WEBHOOK_EXTERNAL_TOKEN && /\s/u.test(parsed.WEBHOOK_EXTERNAL_TOKEN)) {
leaderSummaryWebhookConfigurationError = 'webhook_token_whitespace_invalid';
} else if (parsed.WEBHOOK_SEND_URL) {
let webhookUrl: URL | null = null;
try {
webhookUrl = new URL(parsed.WEBHOOK_SEND_URL);
} catch {
leaderSummaryWebhookConfigurationError = 'webhook_url_invalid';
}
if (webhookUrl && !webhookUrl.pathname.endsWith('/webhookInfo/sendMessageByOut')) {
leaderSummaryWebhookConfigurationError = 'webhook_route_unconfirmed';
} else if (webhookUrl && parsed.NODE_ENV === 'production' && webhookUrl.protocol !== 'https:') {
leaderSummaryWebhookConfigurationError = 'webhook_https_required';
}
}
const leaderSummaryWebhookEnabled = webhookConfigurationComplete
&& !leaderSummaryWebhookConfigurationError;
if (parsed.ARTIFACT_STORAGE_BACKEND === 'oss') {
const missing = [
['OSS_ACCESS_KEY_ID', parsed.OSS_ACCESS_KEY_ID],
@@ -164,6 +202,8 @@ export function loadConfig(env: NodeJS.ProcessEnv = process.env): AppConfig {
...parsed,
fieldEncryptionKey: decodeEncryptionKey(parsed.FIELD_ENCRYPTION_KEY, parsed.NODE_ENV),
agentBusEnabled,
leaderSummaryWebhookEnabled,
leaderSummaryWebhookConfigurationError,
ossRegion
};
}
+1 -1
View File
@@ -5,7 +5,7 @@ import { writeEmergencyDiagnostic } from './diagnostics.js';
const { Pool } = pg;
let pool: pg.Pool | null = null;
export const REQUIRED_SCHEMA_VERSION = '022_shared_child_order_batch_create';
export const REQUIRED_SCHEMA_VERSION = '023_leader_summary_webhook_delivery';
export interface DatabaseReadiness {
ready: boolean;
@@ -0,0 +1,178 @@
export const EXTERNAL_WEBHOOK_CONFIG_ID = '9999';
export const EXTERNAL_WEBHOOK_TIMEOUT_MS = 15_000;
const MAX_RESPONSE_BYTES = 65_536;
export type ExternalWebhookSendResult =
| {
outcome: 'accepted';
httpStatus: 200;
errorCode: null;
}
| {
outcome: 'rejected' | 'uncertain';
httpStatus: number | null;
errorCode: string;
};
export type ExternalWebhookFetch = (
input: string | URL | Request,
init?: RequestInit
) => Promise<Response>;
export interface ExternalWebhookClientOptions {
url: string;
token: string;
timeoutMs?: number;
fetchImpl?: ExternalWebhookFetch;
}
function jsonObject(value: unknown): Record<string, unknown> {
return value && typeof value === 'object' && !Array.isArray(value)
? value as Record<string, unknown>
: {};
}
async function readBoundedResponse(response: Response): Promise<{
text: string;
tooLarge: boolean;
}> {
const declaredLength = Number(response.headers.get('content-length'));
if (Number.isFinite(declaredLength) && declaredLength > MAX_RESPONSE_BYTES) {
await response.body?.cancel().catch(() => undefined);
return { text: '', tooLarge: true };
}
if (!response.body) return { text: '', tooLarge: false };
const reader = response.body.getReader();
const chunks: Uint8Array[] = [];
let total = 0;
while (true) {
const { done, value } = await reader.read();
if (done) break;
if (!value) continue;
total += value.byteLength;
if (total > MAX_RESPONSE_BYTES) {
await reader.cancel().catch(() => undefined);
return { text: '', tooLarge: true };
}
chunks.push(value);
}
const bytes = new Uint8Array(total);
let offset = 0;
for (const chunk of chunks) {
bytes.set(chunk, offset);
offset += chunk.byteLength;
}
return { text: new TextDecoder().decode(bytes), tooLarge: false };
}
function rejectedErrorCode(status: number, payload: Record<string, unknown>): string {
const message = String(payload.msg ?? '').trim();
if (status === 401 && message === 'Invalid token') return 'webhook_invalid_token';
if (status === 400 && message === 'id and content must not be blank') return 'webhook_invalid_request';
if (status === 404 && message === 'Webhook configuration not found') {
return 'webhook_configuration_not_found';
}
if (status === 404) return 'webhook_route_not_found';
if (status >= 400 && status < 500) return `webhook_http_${status}`;
return 'webhook_business_rejected';
}
export class ExternalWebhookClient {
private readonly timeoutMs: number;
private readonly fetchImpl: ExternalWebhookFetch;
constructor(private readonly options: ExternalWebhookClientOptions) {
this.timeoutMs = options.timeoutMs || EXTERNAL_WEBHOOK_TIMEOUT_MS;
this.fetchImpl = options.fetchImpl || globalThis.fetch;
}
async send(content: string): Promise<ExternalWebhookSendResult> {
const normalizedContent = String(content ?? '').trim();
if (!normalizedContent) {
throw new Error('Webhook content must not be blank.');
}
if (this.options.token.length !== 32 || /\s/u.test(this.options.token)) {
throw new Error('Webhook token must contain exactly 32 non-whitespace characters.');
}
let response: Response;
try {
response = await this.fetchImpl(this.options.url, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'x-token': this.options.token
},
body: JSON.stringify({
id: EXTERNAL_WEBHOOK_CONFIG_ID,
content: normalizedContent
}),
signal: AbortSignal.timeout(this.timeoutMs)
});
} catch (error) {
return {
outcome: 'uncertain',
httpStatus: null,
errorCode: error instanceof DOMException && error.name === 'TimeoutError'
? 'webhook_request_timeout'
: 'webhook_request_uncertain'
};
}
let responseText = '';
try {
const bounded = await readBoundedResponse(response);
if (bounded.tooLarge) {
return {
outcome: 'uncertain',
httpStatus: response.status,
errorCode: 'webhook_response_too_large'
};
}
responseText = bounded.text;
} catch {
return {
outcome: 'uncertain',
httpStatus: response.status,
errorCode: 'webhook_response_read_failed'
};
}
let payload: Record<string, unknown> = {};
try {
payload = jsonObject(JSON.parse(responseText));
} catch {
if (response.status >= 400 && response.status < 500) {
return {
outcome: 'rejected',
httpStatus: response.status,
errorCode: rejectedErrorCode(response.status, payload)
};
}
return {
outcome: 'uncertain',
httpStatus: response.status,
errorCode: 'webhook_response_invalid_json'
};
}
if (response.status === 200 && payload.code === 0 && payload.data === true) {
return { outcome: 'accepted', httpStatus: 200, errorCode: null };
}
if (response.status >= 500) {
return {
outcome: 'uncertain',
httpStatus: response.status,
errorCode: `webhook_http_${response.status}_uncertain`
};
}
return {
outcome: 'rejected',
httpStatus: response.status,
errorCode: rejectedErrorCode(response.status, payload)
};
}
}
File diff suppressed because it is too large. Load diff
+27 -7
View File
@@ -443,6 +443,8 @@ export async function buildServer({
database_ssl: config.DATABASE_SSL,
artifact_storage_backend: config.ARTIFACT_STORAGE_BACKEND,
agentbus_enabled: config.agentBusEnabled,
leader_summary_webhook_enabled: config.leaderSummaryWebhookEnabled,
leader_summary_webhook_configuration_error: config.leaderSummaryWebhookConfigurationError,
parser_loop_enabled: startParserLoop,
data_retention_enabled: config.DATA_RETENTION_ENABLED,
raw_payload_logging: config.AGENTBUS_LOG_PAYLOADS
@@ -539,9 +541,9 @@ export async function buildServer({
error: (metadata, message) => app.log.error(agentBusDiagnosticMetadata(metadata), message)
});
const leaderNotificationService = new LeaderNotificationService(config, {
info: (metadata, message) => app.log.info(agentBusDiagnosticMetadata(metadata), message),
warn: (metadata, message) => app.log.warn(agentBusDiagnosticMetadata(metadata), message),
error: (metadata, message) => app.log.error(agentBusDiagnosticMetadata(metadata), message)
info: (metadata, message) => app.log.info(metadata, message),
warn: (metadata, message) => app.log.warn(metadata, message),
error: (metadata, message) => app.log.error(metadata, message)
});
let agentBus: AgentBusManager | null = null;
@@ -844,7 +846,6 @@ export async function buildServer({
tasks,
organizationId: organization.id,
scheduleParseQueue,
leaderNotifications: leaderNotificationService,
logger: {
info: (metadata, message) => app.log.info(agentBusDiagnosticMetadata(metadata), message),
warn: (metadata, message) => app.log.warn(agentBusDiagnosticMetadata(metadata), message),
@@ -852,7 +853,26 @@ export async function buildServer({
}
});
await agentBus.start();
leaderNotificationService.startProjector(organization.id);
}
if (config.leaderSummaryWebhookEnabled) {
try {
const organization = await auth.getOrganization();
if (organization) {
leaderNotificationService.startProjector(organization.id);
} else {
app.log.warn({
diagnostic_event: 'leader_summary.webhook_initialization_skipped',
notification_event: 'webhook_initialization_skipped',
error_code: 'leader_summary_webhook_organization_not_found'
}, 'Leader summary webhook initialization skipped; normal task services remain available');
}
} catch (error) {
app.log.warn({
diagnostic_event: 'leader_summary.webhook_initialization_failed',
notification_event: 'webhook_initialization_failed',
...diagnosticError(error, 'leader_summary_webhook_initialization_failed')
}, 'Leader summary webhook initialization failed; normal task services remain available');
}
}
async function getAiProbe(): Promise<unknown> {
@@ -1206,11 +1226,11 @@ export async function buildServer({
return { ok: true, ...result };
});
app.get('/api/settings/leader-summary-subscriptions', async (request) => {
app.get('/api/settings/leader-summary-webhook', async (request) => {
const session = await requireAdminSession(request);
return {
ok: true,
subscriptions: await leaderNotificationService.listSubscriptions(session.user.organizationId)
status: await leaderNotificationService.getWebhookStatus(session.user.organizationId)
};
});