feat: add leader AgentBus task summaries

This commit is contained in:
inman committed 2026-09-07 13:01:01 +08:00
1 parent f4664997a8
commit 1a3ab63700
18 files changed
+2158 -19

No files matched your search

@@ -0,0 +1,873 @@
import type { AppConfig } from './config.js';
import { decryptText, encryptText, sha256Text } from './crypto.js';
import { getPool, withTransaction } from './db.js';
import { diagnosticError, diagnosticMetadataKeys } from './diagnostics.js';
import {
buildLeaderTaskSummary,
isLeaderTaskSummaryStatus,
type LeaderTaskSummaryStatus
} from './leadership-task-summary.js';
import { TaskError, type TaskContext } from './task-service.js';
export interface LeaderNotificationLogger {
info(metadata: Record<string, unknown>, message?: string): void;
warn(metadata: Record<string, unknown>, message?: string): void;
error(metadata: Record<string, unknown>, message?: string): void;
}
export interface PublicLeaderTaskSummarySubscription {
id: string;
organization_id: string;
leader_user_id: string;
leader_username: string;
channel_id: string;
channel_name: string;
scope: 'organization';
include_manual: boolean;
include_agentbus: boolean;
enabled: boolean;
eligible: boolean;
eligibility_message: string;
target_verified: boolean;
recipient_fingerprint: string;
conversation_fingerprint: string;
starts_at: string;
revision: number;
pending_count: number;
failed_count: number;
delivered_count: number;
last_delivered_at: string | null;
created_at: string;
updated_at: string;
}
export interface LeaderTaskSummarySubscriptionInput {
leaderUserId: string;
channelId: string;
recipientAddress?: string;
conversationId?: string;
includeManual: boolean;
includeAgentBus: boolean;
enabled: boolean;
targetVerified?: boolean;
expectedRevision?: number;
}
export interface LeaderTaskSummaryDelivery {
id: string;
channel_id: string;
task_id: string;
recipient_address: string;
recipient_fingerprint: string;
conversation_id: string;
conversation_fingerprint: string;
payload: {
event: 'task.summary';
status: LeaderTaskSummaryStatus;
task_id: string;
text: string;
};
attempt_count: number;
}
const PROJECTABLE_STATUSES = [
'completed',
'dry_run',
'failed',
'parse_failed',
'parse_blocked',
'agent_parse_blocked',
'blocked',
'operation_blocked',
'cancelled',
'reconciliation_pending',
'saved_unverified',
'execution_uncertain',
'uncertain'
] as const;
const NEEDS_REVIEW_STATUSES = [
'reconciliation_pending',
'saved_unverified',
'execution_uncertain',
'uncertain'
] as const;
const noopLogger: LeaderNotificationLogger = {
info: () => undefined,
warn: () => undefined,
error: () => undefined
};
function text(value: unknown): string {
return value == null ? '' : String(value).trim();
}
function booleanValue(value: unknown): boolean {
return value === true || text(value) === 'true';
}
function iso(value: unknown): string | null {
if (!value) return null;
const date = new Date(String(value));
return Number.isNaN(date.getTime()) ? null : date.toISOString();
}
function jsonObject(value: unknown): Record<string, unknown> {
return value && typeof value === 'object' && !Array.isArray(value)
? value as Record<string, unknown>
: {};
}
function eligibility(
row: Record<string, unknown>,
agentBusEnabled: boolean
): { eligible: boolean; message: string } {
if (!agentBusEnabled) {
return { eligible: false, message: 'AgentBus 全局开关未启用,通知不会发送。' };
}
if (!booleanValue(row.leader_is_active) || text(row.leader_role) !== 'team_lead') {
return { eligible: false, message: '组长账号已停用或角色已变化,通知不会发送。' };
}
if (text(row.channel_owner_user_id) !== text(row.leader_user_id)) {
return { eligible: false, message: 'AgentBus 渠道已不再归属该组长,通知不会发送。' };
}
if (!booleanValue(row.channel_enabled)) {
return { eligible: false, message: 'AgentBus 渠道已停用,通知不会发送。' };
}
if (!row.target_verified_at) {
return { eligible: false, message: '主动投递目标尚未验证,通知不会发送。' };
}
return { eligible: true, message: '投递条件已就绪。' };
}
function publicSubscription(
row: Record<string, unknown>,
agentBusEnabled: boolean
): PublicLeaderTaskSummarySubscription {
const state = eligibility(row, agentBusEnabled);
return {
id: text(row.id),
organization_id: text(row.organization_id),
leader_user_id: text(row.leader_user_id),
leader_username: text(row.leader_username) || '未知组长',
channel_id: text(row.channel_id),
channel_name: text(row.channel_name) || '未命名渠道',
scope: 'organization',
include_manual: booleanValue(row.include_manual),
include_agentbus: booleanValue(row.include_agentbus),
enabled: booleanValue(row.enabled),
eligible: state.eligible,
eligibility_message: state.message,
target_verified: Boolean(row.target_verified_at),
recipient_fingerprint: text(row.recipient_address_fingerprint).slice(0, 12),
conversation_fingerprint: text(row.conversation_id_fingerprint).slice(0, 12),
starts_at: iso(row.starts_at) || new Date(0).toISOString(),
revision: Number(row.revision || 0),
pending_count: Number(row.pending_count || 0),
failed_count: Number(row.failed_count || 0),
delivered_count: Number(row.delivered_count || 0),
last_delivered_at: iso(row.last_delivered_at),
created_at: iso(row.created_at) || new Date(0).toISOString(),
updated_at: iso(row.updated_at) || new Date(0).toISOString()
};
}
export class LeaderNotificationService {
private projectorTimer: NodeJS.Timeout | null = null;
private projectorInFlight: Promise<void> | null = null;
private projectorOrganizationId = '';
constructor(
private readonly config: AppConfig,
private readonly logger: LeaderNotificationLogger = noopLogger
) {}
private log(
level: 'info' | 'warn' | 'error',
metadata: Record<string, unknown>,
message: string
): void {
try {
this.logger[level](metadata, message);
} catch {
// Notification persistence and task state never depend on logging.
}
}
startProjector(organizationId: string): void {
if (this.projectorTimer) return;
this.projectorOrganizationId = organizationId;
this.projectorTimer = setInterval(() => void this.projectTick(), 2_000);
void this.projectTick();
}
async stopProjector(): Promise<void> {
if (this.projectorTimer) clearInterval(this.projectorTimer);
this.projectorTimer = null;
if (this.projectorInFlight) await this.projectorInFlight.catch(() => undefined);
this.projectorOrganizationId = '';
}
private async projectTick(): Promise<void> {
if (!this.projectorOrganizationId) return;
if (this.projectorInFlight) return this.projectorInFlight;
this.projectorInFlight = this.projectPending(this.projectorOrganizationId, 100)
.then((count) => {
if (count > 0) {
this.log('info', {
agentbus_event: 'leader_summary_projected',
organization_id: this.projectorOrganizationId,
projected_count: count
}, 'Leader task summaries projected');
}
})
.catch((error) => {
this.log('warn', {
agentbus_event: 'leader_summary_projection_failed',
organization_id: this.projectorOrganizationId,
...diagnosticError(error, 'leader_summary_projection_failed')
}, 'Leader task summary projection failed');
})
.finally(() => {
this.projectorInFlight = null;
});
return this.projectorInFlight;
}
async listSubscriptions(organizationId: string): Promise<PublicLeaderTaskSummarySubscription[]> {
const result = await getPool(this.config).query(
`SELECT subscription.*,
leader.username AS leader_username,
leader.role AS leader_role,
leader.is_active AS leader_is_active,
channel.display_name AS channel_name,
channel.owner_user_id AS channel_owner_user_id,
channel.enabled AS channel_enabled,
COALESCE(delivery.pending_count, 0)::int AS pending_count,
COALESCE(delivery.failed_count, 0)::int AS failed_count,
COALESCE(delivery.delivered_count, 0)::int AS delivered_count,
delivery.last_delivered_at
FROM leader_task_summary_subscriptions subscription
JOIN users leader
ON leader.id = subscription.leader_user_id
AND leader.organization_id = subscription.organization_id
JOIN user_channels channel
ON channel.id = subscription.channel_id
AND channel.organization_id = subscription.organization_id
LEFT JOIN LATERAL (
SELECT count(*) FILTER (
WHERE d.delivery_status IN ('pending', 'sending')
AND d.subscription_revision = subscription.revision
) AS pending_count,
count(*) FILTER (
WHERE d.delivery_status = 'failed'
AND d.subscription_revision = subscription.revision
) AS failed_count,
count(*) FILTER (WHERE d.delivery_status = 'delivered') AS delivered_count,
max(d.delivered_at) AS last_delivered_at
FROM leader_task_summary_deliveries d
WHERE d.subscription_id = subscription.id
) delivery ON true
WHERE subscription.organization_id = $1
ORDER BY leader.username, subscription.id`,
[organizationId]
);
return (result.rows as Record<string, unknown>[])
.map((row) => publicSubscription(row, this.config.agentBusEnabled));
}
async upsertSubscription(
context: TaskContext,
input: LeaderTaskSummarySubscriptionInput
): Promise<PublicLeaderTaskSummarySubscription> {
const leaderUserId = text(input.leaderUserId);
const channelId = text(input.channelId);
const recipientAddress = text(input.recipientAddress).slice(0, 500);
const conversationId = text(input.conversationId).slice(0, 500);
if (!leaderUserId || !channelId) {
throw new TaskError('leader_summary_target_required', '请选择组长及其 AgentBus 渠道。', 400);
}
if (!input.includeManual && !input.includeAgentBus) {
throw new TaskError('leader_summary_source_required', '人工任务和 AgentBus 任务至少选择一种。', 400);
}
if (input.enabled && !this.config.agentBusEnabled) {
throw new TaskError(
'leader_summary_agentbus_disabled',
'AgentBus 全局开关未启用,暂时不能启用组长摘要抄送。',
409
);
}
const subscriptionId = await withTransaction(this.config, async (client) => {
await client.query(
`SELECT pg_advisory_xact_lock(hashtextextended($1::text || ':leader-summary:' || $2::text, 0))`,
[context.organizationId, leaderUserId]
);
const existingResult = await client.query(
`SELECT *
FROM leader_task_summary_subscriptions
WHERE organization_id = $1 AND leader_user_id = $2
FOR UPDATE`,
[context.organizationId, leaderUserId]
);
const existing = existingResult.rows[0] as Record<string, unknown> | undefined;
const safeDisable = Boolean(
existing
&& !input.enabled
&& channelId === text(existing.channel_id)
);
const leaderResult = await client.query(
`SELECT id, role, is_active
FROM users
WHERE organization_id = $1 AND id = $2
FOR SHARE`,
[context.organizationId, leaderUserId]
);
if (!leaderResult.rowCount) {
throw new TaskError('leader_summary_leader_not_found', '组长账号不存在。', 404);
}
const leader = leaderResult.rows[0] as Record<string, unknown>;
if ((!booleanValue(leader.is_active) || text(leader.role) !== 'team_lead') && !safeDisable) {
throw new TaskError('leader_summary_leader_invalid', '只有有效的组长账号可以接收任务摘要。', 409);
}
const channelResult = await client.query(
`SELECT id, owner_user_id, enabled
FROM user_channels
WHERE organization_id = $1 AND id = $2
FOR SHARE`,
[context.organizationId, channelId]
);
if (!channelResult.rowCount) {
throw new TaskError('leader_summary_channel_not_found', 'AgentBus 渠道不存在。', 404);
}
const channel = channelResult.rows[0] as Record<string, unknown>;
if (text(channel.owner_user_id) !== leaderUserId && !safeDisable) {
throw new TaskError('leader_summary_channel_owner_mismatch', '所选 AgentBus 渠道不属于该组长。', 409);
}
if (input.enabled && !booleanValue(channel.enabled)) {
throw new TaskError('leader_summary_channel_disabled', '请先启用该组长的 AgentBus 渠道。', 409);
}
const currentRevision = Number(existing?.revision || 0);
if (existing && input.expectedRevision !== undefined && input.expectedRevision !== currentRevision) {
throw new TaskError(
'leader_summary_revision_conflict',
'组长摘要设置已被其他管理员修改,请刷新后重试。',
409,
{ current_revision: currentRevision }
);
}
if (!existing && (!recipientAddress || !conversationId)) {
throw new TaskError(
'leader_summary_route_required',
'首次配置必须填写 AgentBus 收件地址和微信会话 ID。',
400
);
}
const recipientFingerprint = recipientAddress
? sha256Text(recipientAddress)
: text(existing?.recipient_address_fingerprint);
const conversationFingerprint = conversationId
? sha256Text(conversationId)
: text(existing?.conversation_id_fingerprint);
const recipientCiphertext = recipientAddress
? encryptText(this.config, recipientAddress)
: text(existing?.recipient_address_ciphertext);
const conversationCiphertext = conversationId
? encryptText(this.config, conversationId)
: text(existing?.conversation_id_ciphertext);
const routeChanged = !existing
|| channelId !== text(existing.channel_id)
|| recipientFingerprint !== text(existing.recipient_address_fingerprint)
|| conversationFingerprint !== text(existing.conversation_id_fingerprint);
const targetVerified = input.targetVerified === undefined
? !routeChanged && Boolean(existing?.target_verified_at)
: input.targetVerified === true;
if (input.enabled && !targetVerified) {
throw new TaskError(
'leader_summary_target_unverified',
'启用前必须确认该收件地址和微信会话已经过主动投递验证。',
409
);
}
const nextRevision = existing ? currentRevision + 1 : 0;
let id: string;
if (existing) {
const updated = await client.query(
`UPDATE leader_task_summary_subscriptions
SET channel_id = $1,
include_manual = $2,
include_agentbus = $3,
enabled = $4,
starts_at = now(),
revision = $5,
recipient_address_ciphertext = $6,
recipient_address_fingerprint = $7,
conversation_id_ciphertext = $8,
conversation_id_fingerprint = $9,
target_verified_at = CASE WHEN $10 THEN now() ELSE NULL END,
target_verified_by = CASE WHEN $10 THEN $11::uuid ELSE NULL END,
updated_at = now()
WHERE id = $12
RETURNING id`,
[
channelId,
input.includeManual,
input.includeAgentBus,
input.enabled,
nextRevision,
recipientCiphertext,
recipientFingerprint,
conversationCiphertext,
conversationFingerprint,
targetVerified,
context.userId || null,
existing.id
]
);
id = text(updated.rows[0].id);
await client.query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'cancelled',
last_error = '订阅设置已变化,旧目标待发送摘要已取消。',
updated_at = now()
WHERE subscription_id = $1
AND subscription_revision <> $2
AND delivery_status IN ('pending', 'sending', 'failed')`,
[id, nextRevision]
);
} else {
const inserted = await client.query(
`INSERT INTO leader_task_summary_subscriptions
(organization_id, leader_user_id, channel_id, include_manual,
include_agentbus, enabled, starts_at, revision,
recipient_address_ciphertext, recipient_address_fingerprint,
conversation_id_ciphertext, conversation_id_fingerprint,
target_verified_at, target_verified_by, created_by)
VALUES ($1, $2, $3, $4, $5, $6, now(), 0,
$7, $8, $9, $10,
CASE WHEN $11 THEN now() ELSE NULL END,
CASE WHEN $11 THEN $12::uuid ELSE NULL END,
$12::uuid)
RETURNING id`,
[
context.organizationId,
leaderUserId,
channelId,
input.includeManual,
input.includeAgentBus,
input.enabled,
recipientCiphertext,
recipientFingerprint,
conversationCiphertext,
conversationFingerprint,
targetVerified,
context.userId || null
]
);
id = text(inserted.rows[0].id);
}
await client.query(
`INSERT INTO audit_events
(organization_id, actor_user_id, event_type, entity_type, entity_id, request_id, metadata)
VALUES ($1, $2, $3, 'leader_task_summary_subscription', $4, $5, $6)`,
[
context.organizationId,
context.userId || null,
input.enabled ? 'leader_summary_subscription.enabled' : 'leader_summary_subscription.disabled',
id,
context.requestId,
{
leader_user_id: leaderUserId,
channel_id: channelId,
include_manual: input.includeManual,
include_agentbus: input.includeAgentBus,
enabled: input.enabled,
target_verified: targetVerified,
recipient_fingerprint: recipientFingerprint.slice(0, 12),
conversation_fingerprint: conversationFingerprint.slice(0, 12),
revision: nextRevision,
historical_backfill: false
}
]
);
this.log('info', {
agentbus_event: 'leader_summary_subscription_audit_staged',
request_id: context.requestId,
entity_id: id,
metadata_keys: diagnosticMetadataKeys({
leader_user_id: leaderUserId,
channel_id: channelId,
enabled: input.enabled,
revision: nextRevision
})
}, 'Leader summary subscription audit event staged');
return id;
});
const subscriptions = await this.listSubscriptions(context.organizationId);
const subscription = subscriptions.find((item) => item.id === subscriptionId);
if (!subscription) throw new TaskError('leader_summary_subscription_not_found', '组长摘要设置不存在。', 404);
return subscription;
}
async projectPending(organizationId: string, limit = 100): Promise<number> {
const boundedLimit = Math.max(1, Math.min(500, Math.trunc(limit)));
return withTransaction(this.config, async (client) => {
const lock = await client.query(
`SELECT pg_try_advisory_xact_lock(
hashtextextended($1::text || ':leader-summary-projector', 0)
) AS acquired`,
[organizationId]
);
if (!booleanValue(lock.rows[0]?.acquired)) return 0;
const candidates = await client.query(
`SELECT subscription.id AS subscription_id,
subscription.revision AS subscription_revision,
subscription.leader_user_id,
subscription.channel_id,
subscription.recipient_address_ciphertext,
subscription.recipient_address_fingerprint,
subscription.conversation_id_ciphertext,
subscription.conversation_id_fingerprint,
task.id AS task_row_id,
task.task_id,
task.status,
task.business_route_id,
task.success_receipt,
task.created_at,
assignee.username AS assignee_username,
source_event.id AS source_outbox_event_id,
COALESCE(delivery_state.has_needs_review, false) AS has_needs_review,
COALESCE(delivery_state.has_exposed_needs_review, false) AS has_exposed_needs_review
FROM leader_task_summary_subscriptions subscription
JOIN users leader
ON leader.id = subscription.leader_user_id
AND leader.organization_id = subscription.organization_id
AND leader.role = 'team_lead'
AND leader.is_active = true
JOIN user_channels channel
ON channel.id = subscription.channel_id
AND channel.organization_id = subscription.organization_id
AND channel.owner_user_id = subscription.leader_user_id
AND channel.enabled = true
JOIN tasks task
ON task.organization_id = subscription.organization_id
AND task.assigned_user_id IS NOT NULL
AND task.assigned_user_id <> subscription.leader_user_id
AND task.source IN ('manual', 'agentbus')
AND ((task.source = 'manual' AND subscription.include_manual)
OR (task.source = 'agentbus' AND subscription.include_agentbus))
JOIN users assignee
ON assignee.id = task.assigned_user_id
AND assignee.organization_id = task.organization_id
AND assignee.role <> 'admin'
JOIN LATERAL (
SELECT event.id
FROM outbox_events event
WHERE event.organization_id = subscription.organization_id
AND event.topic = 'task.updated'
AND event.aggregate_type = 'task'
AND event.aggregate_id = task.task_id
AND event.created_at >= subscription.starts_at
AND event.payload ->> 'status' = ANY($2::text[])
AND (event.payload ->> 'archived') IS DISTINCT FROM 'true'
AND (event.payload ->> 'restored') IS DISTINCT FROM 'true'
ORDER BY event.id DESC
LIMIT 1
) source_event ON true
LEFT JOIN LATERAL (
SELECT bool_or(delivery.milestone = 'needs_review' AND delivery.delivery_status <> 'cancelled') AS has_needs_review,
bool_or(delivery.milestone = 'needs_review' AND delivery.delivery_status IN ('sending', 'delivered'))
AS has_exposed_needs_review,
bool_or(delivery.milestone = 'final' AND delivery.delivery_status <> 'cancelled') AS has_final,
bool_or(delivery.milestone = 'resolved' AND delivery.delivery_status <> 'cancelled') AS has_resolved
FROM leader_task_summary_deliveries delivery
WHERE delivery.subscription_id = subscription.id
AND delivery.subscription_revision = subscription.revision
AND delivery.task_id = task.id
) delivery_state ON true
WHERE subscription.organization_id = $1
AND subscription.enabled = true
AND subscription.target_verified_at IS NOT NULL
AND task.status = ANY($2::text[])
AND (
(task.status = ANY($3::text[]) AND NOT COALESCE(delivery_state.has_needs_review, false))
OR
(NOT (task.status = ANY($3::text[])) AND (
(COALESCE(delivery_state.has_exposed_needs_review, false) AND NOT COALESCE(delivery_state.has_resolved, false))
OR
(NOT COALESCE(delivery_state.has_exposed_needs_review, false) AND NOT COALESCE(delivery_state.has_final, false))
))
)
ORDER BY source_event.id ASC, subscription.id, task.id
LIMIT $4`,
[organizationId, [...PROJECTABLE_STATUSES], [...NEEDS_REVIEW_STATUSES], boundedLimit]
);
let projected = 0;
for (const row of candidates.rows as Record<string, unknown>[]) {
if (!isLeaderTaskSummaryStatus(row.status)) continue;
const currentNeedsReview = (NEEDS_REVIEW_STATUSES as readonly string[]).includes(text(row.status));
let hadExposedNeedsReview = booleanValue(row.has_exposed_needs_review);
if (!currentNeedsReview) {
await client.query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'cancelled',
last_error = '任务已形成确定结果,未发送的旧核验提醒已取消。',
updated_at = now()
WHERE subscription_id = $1
AND subscription_revision = $2
AND task_id = $3
AND milestone = 'needs_review'
AND delivery_status IN ('pending', 'failed')`,
[row.subscription_id, row.subscription_revision, row.task_row_id]
);
const exposure = await client.query(
`SELECT EXISTS (
SELECT 1
FROM leader_task_summary_deliveries
WHERE subscription_id = $1
AND subscription_revision = $2
AND task_id = $3
AND milestone = 'needs_review'
AND delivery_status IN ('sending', 'delivered')
) AS exposed`,
[row.subscription_id, row.subscription_revision, row.task_row_id]
);
hadExposedNeedsReview = booleanValue(exposure.rows[0]?.exposed);
}
const projection = buildLeaderTaskSummary({
taskId: text(row.task_id),
status: text(row.status),
businessRouteId: text(row.business_route_id) || null,
assigneeUsername: text(row.assignee_username),
createdAt: String(row.created_at),
successReceipt: jsonObject(row.success_receipt),
hadNeedsReview: hadExposedNeedsReview
});
if (!projection) continue;
const payload = {
event: 'task.summary' as const,
status: projection.deliveryStatus,
task_id: text(row.task_id),
text: projection.messageText
};
const payloadText = JSON.stringify(payload);
const inserted = await client.query(
`INSERT INTO leader_task_summary_deliveries
(organization_id, subscription_id, subscription_revision, task_id,
leader_user_id, channel_id, source_outbox_event_id, milestone,
recipient_address_ciphertext, recipient_address_fingerprint,
conversation_id_ciphertext, conversation_id_fingerprint,
payload_ciphertext, payload_fingerprint)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8,
$9, $10, $11, $12, $13, $14)
ON CONFLICT (subscription_id, subscription_revision, task_id, milestone) DO NOTHING
RETURNING id`,
[
organizationId,
row.subscription_id,
row.subscription_revision,
row.task_row_id,
row.leader_user_id,
row.channel_id,
row.source_outbox_event_id,
projection.milestone,
row.recipient_address_ciphertext,
row.recipient_address_fingerprint,
row.conversation_id_ciphertext,
row.conversation_id_fingerprint,
encryptText(this.config, payloadText),
sha256Text(payloadText)
]
);
projected += Number(inserted.rowCount || 0);
}
return projected;
});
}
async claimDeliveries(
channelId: string,
leaseOwner: string,
limit = 5
): Promise<LeaderTaskSummaryDelivery[]> {
const boundedLimit = Math.max(1, Math.min(10, Math.trunc(limit)));
return withTransaction(this.config, async (client) => {
await client.query(
`UPDATE leader_task_summary_deliveries delivery
SET delivery_status = 'cancelled',
last_error = '订阅或目标已失效,摘要未改投其他渠道。',
updated_at = now()
WHERE delivery.channel_id = $1
AND delivery.delivery_status IN ('pending', 'sending', 'failed')
AND NOT EXISTS (
SELECT 1
FROM leader_task_summary_subscriptions subscription
JOIN users leader
ON leader.id = subscription.leader_user_id
AND leader.organization_id = subscription.organization_id
AND leader.role = 'team_lead'
AND leader.is_active = true
JOIN user_channels channel
ON channel.id = subscription.channel_id
AND channel.organization_id = subscription.organization_id
AND channel.owner_user_id = subscription.leader_user_id
AND channel.enabled = true
WHERE subscription.id = delivery.subscription_id
AND subscription.enabled = true
AND subscription.target_verified_at IS NOT NULL
AND subscription.revision = delivery.subscription_revision
AND subscription.channel_id = delivery.channel_id
)`,
[channelId]
);
await client.query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'pending',
last_error = COALESCE(last_error, 'delivery lease expired'),
updated_at = now()
WHERE channel_id = $1
AND delivery_status = 'sending'
AND updated_at < now() - interval '1 minute'`,
[channelId]
);
const pending = await client.query(
`SELECT delivery.*,
task.task_id AS public_task_id
FROM leader_task_summary_deliveries delivery
JOIN leader_task_summary_subscriptions subscription
ON subscription.id = delivery.subscription_id
AND subscription.enabled = true
AND subscription.target_verified_at IS NOT NULL
AND subscription.revision = delivery.subscription_revision
AND subscription.channel_id = delivery.channel_id
JOIN users leader
ON leader.id = subscription.leader_user_id
AND leader.organization_id = subscription.organization_id
AND leader.role = 'team_lead'
AND leader.is_active = true
JOIN user_channels channel
ON channel.id = subscription.channel_id
AND channel.organization_id = subscription.organization_id
AND channel.owner_user_id = subscription.leader_user_id
AND channel.enabled = true
JOIN tasks task
ON task.id = delivery.task_id
AND task.organization_id = delivery.organization_id
WHERE delivery.channel_id = $1
AND delivery.delivery_status IN ('pending', 'failed')
AND delivery.next_attempt_at <= now()
ORDER BY delivery.created_at ASC, delivery.id ASC
FOR UPDATE OF delivery SKIP LOCKED
LIMIT $2`,
[channelId, boundedLimit]
);
const deliveries: LeaderTaskSummaryDelivery[] = [];
for (const row of pending.rows as Record<string, unknown>[]) {
try {
const payloadText = decryptText(this.config, text(row.payload_ciphertext));
const payload = JSON.parse(payloadText) as LeaderTaskSummaryDelivery['payload'];
const recipientAddress = decryptText(this.config, text(row.recipient_address_ciphertext)).trim();
const conversationId = decryptText(this.config, text(row.conversation_id_ciphertext)).trim();
if (
payload.event !== 'task.summary'
|| !['completed', 'failed', 'cancelled', 'needs_review'].includes(text(payload.status))
|| text(payload.task_id) !== text(row.public_task_id)
|| !text(payload.text)
|| !recipientAddress
|| !conversationId
|| sha256Text(payloadText) !== text(row.payload_fingerprint)
|| sha256Text(recipientAddress) !== text(row.recipient_address_fingerprint)
|| sha256Text(conversationId) !== text(row.conversation_id_fingerprint)
) {
throw new Error('invalid leader summary delivery payload');
}
const updated = await client.query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'sending',
attempt_count = attempt_count + 1,
updated_at = now(),
last_error = $2
WHERE id = $1
RETURNING attempt_count`,
[row.id, `sending:${leaseOwner}`]
);
deliveries.push({
id: text(row.id),
channel_id: text(row.channel_id),
task_id: text(row.public_task_id),
recipient_address: recipientAddress,
recipient_fingerprint: text(row.recipient_address_fingerprint).slice(0, 12),
conversation_id: conversationId,
conversation_fingerprint: text(row.conversation_id_fingerprint).slice(0, 12),
payload,
attempt_count: Number(updated.rows[0]?.attempt_count || 1)
});
} catch (error) {
await client.query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'cancelled',
last_error = '摘要投递密文或结构无效,已停止重试。',
updated_at = now()
WHERE id = $1`,
[row.id]
);
this.log('error', {
agentbus_event: 'leader_summary_delivery_invalid',
delivery_id: text(row.id),
channel_id: text(row.channel_id),
task_id: text(row.public_task_id),
...diagnosticError(error, 'leader_summary_delivery_invalid')
}, 'Leader task summary delivery is invalid');
}
}
return deliveries;
});
}
async markDeliveryDelivered(deliveryId: string): Promise<void> {
await getPool(this.config).query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'delivered',
delivered_at = now(),
updated_at = now(),
last_error = NULL
WHERE id = $1 AND delivery_status = 'sending'`,
[deliveryId]
);
}
async markDeliveryFailed(deliveryId: string, errorMessage: string): Promise<void> {
const normalized = text(errorMessage).slice(0, 500) || 'leader_summary_delivery_failed';
await getPool(this.config).query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'failed',
next_attempt_at = now()
+ LEAST(300, GREATEST(5, power(2, LEAST(attempt_count, 8)))) * interval '1 second',
last_error = $2,
updated_at = now()
WHERE id = $1 AND delivery_status = 'sending'`,
[deliveryId, normalized]
);
}
async releaseDeliveries(channelId: string, leaseOwner: string): Promise<void> {
await getPool(this.config).query(
`UPDATE leader_task_summary_deliveries
SET delivery_status = 'pending',
next_attempt_at = now(),
last_error = 'AgentBus 连接已断开,等待重发。',
updated_at = now()
WHERE channel_id = $1
AND delivery_status = 'sending'
AND last_error = $2`,
[channelId, `sending:${leaseOwner}`]
);
}
}