fix: auto-enable leader summary routing
This commit is contained in:
1 parent
8fb066f345
commit
0b3aa5c42d
14 files changed
+618
-623
No files matched your search
@@ -377,11 +377,6 @@ export class AgentBusChannelService {
|
||||
const displayName = input.displayName === undefined ? text(current.display_name) : text(input.displayName).slice(0, 120);
|
||||
if (!displayName) throw new TaskError('channel_name_required', '渠道名称不能为空。', 400);
|
||||
const currentExternalUserRef = text(current.external_user_ref) || null;
|
||||
const externalUserRef = currentExternalUserRef === LEGACY_CHANNEL_REF
|
||||
? LEGACY_CHANNEL_REF
|
||||
: input.externalUserRef === undefined
|
||||
? currentExternalUserRef
|
||||
: (text(input.externalUserRef).slice(0, 200) || null);
|
||||
const enabled = input.enabled === undefined
|
||||
? current.enabled === true || text(current.enabled) === 'true'
|
||||
: input.enabled;
|
||||
@@ -391,6 +386,14 @@ export class AgentBusChannelService {
|
||||
? text(current.owner_user_id)
|
||||
: text(input.ownerUserId);
|
||||
const owner = await this.requireAssignableOwner(client, context.organizationId, ownerUserId);
|
||||
const ownerChanged = ownerUserId !== text(current.owner_user_id);
|
||||
const externalUserRef = currentExternalUserRef === LEGACY_CHANNEL_REF
|
||||
? LEGACY_CHANNEL_REF
|
||||
: input.externalUserRef !== undefined
|
||||
? (text(input.externalUserRef).slice(0, 200) || null)
|
||||
: ownerChanged
|
||||
? null
|
||||
: currentExternalUserRef;
|
||||
if (enabled) {
|
||||
const keyResult = await client.query(
|
||||
`SELECT agentbus_ws_token_ciphertext
|
||||
@@ -436,7 +439,9 @@ export class AgentBusChannelService {
|
||||
await this.audit(client, context, 'agentbus_channel.updated', channelId, {
|
||||
enabled,
|
||||
display_name: displayName,
|
||||
owner_user_id: text(owner.id)
|
||||
owner_user_id: text(owner.id),
|
||||
owner_changed: ownerChanged,
|
||||
external_user_ref_present: Boolean(externalUserRef)
|
||||
});
|
||||
return text(updated.rows[0].id);
|
||||
});
|
||||
|
||||
@@ -135,6 +135,13 @@ export interface AgentBusTaskGateway {
|
||||
}
|
||||
|
||||
export interface LeaderNotificationGateway {
|
||||
observeLeaderRoute?(input: {
|
||||
organizationId: string;
|
||||
leaderUserId: string;
|
||||
channelId: string;
|
||||
recipientAddress: string;
|
||||
conversationId?: string;
|
||||
}): Promise<void>;
|
||||
claimDeliveries(
|
||||
channelId: string,
|
||||
leaseOwner: string,
|
||||
@@ -930,6 +937,13 @@ export class AgentBusListener {
|
||||
this.logFrame('frame_ignored', frame, 'AgentBus frame ignored', { ignore_reason: ignoreReason });
|
||||
return;
|
||||
}
|
||||
if (!this.session) {
|
||||
this.logger.warn({
|
||||
agentbus_event: 'task_before_session_ready',
|
||||
...frameLogData(frame, this.config.AGENTBUS_LOG_PAYLOADS)
|
||||
}, 'Ignoring AgentBus task before session.ready');
|
||||
return;
|
||||
}
|
||||
const taskId = text(frame.id);
|
||||
if (this.inFlightTaskIds.has(taskId) || this.completedTaskIds.has(taskId) || this.pendingFinalReplies.has(taskId)) {
|
||||
this.logFrame('task_duplicate_ignored', frame, 'AgentBus duplicate task ignored', {
|
||||
@@ -939,13 +953,6 @@ export class AgentBusListener {
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (!this.session) {
|
||||
this.logger.warn({
|
||||
agentbus_event: 'task_before_session_ready',
|
||||
...frameLogData(frame, this.config.AGENTBUS_LOG_PAYLOADS)
|
||||
}, 'Ignoring AgentBus task before session.ready');
|
||||
return;
|
||||
}
|
||||
this.logFrame('inbound_task_accepted', frame, 'AgentBus inbound task accepted');
|
||||
const processingStartedAt = process.hrtime.bigint();
|
||||
const processing = this.processInboundTask(frame)
|
||||
@@ -1059,6 +1066,25 @@ 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,
|
||||
|
||||
@@ -7,7 +7,6 @@ import {
|
||||
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;
|
||||
@@ -25,9 +24,11 @@ export interface PublicLeaderTaskSummarySubscription {
|
||||
scope: 'organization';
|
||||
include_manual: boolean;
|
||||
include_agentbus: boolean;
|
||||
automatic: true;
|
||||
enabled: boolean;
|
||||
eligible: boolean;
|
||||
eligibility_message: string;
|
||||
route_ready: boolean;
|
||||
target_verified: boolean;
|
||||
recipient_fingerprint: string;
|
||||
conversation_fingerprint: string;
|
||||
@@ -41,16 +42,12 @@ export interface PublicLeaderTaskSummarySubscription {
|
||||
updated_at: string;
|
||||
}
|
||||
|
||||
export interface LeaderTaskSummarySubscriptionInput {
|
||||
export interface ObservedLeaderAgentBusRoute {
|
||||
organizationId: string;
|
||||
leaderUserId: string;
|
||||
channelId: string;
|
||||
recipientAddress?: string;
|
||||
recipientAddress: string;
|
||||
conversationId?: string;
|
||||
includeManual: boolean;
|
||||
includeAgentBus: boolean;
|
||||
enabled: boolean;
|
||||
targetVerified?: boolean;
|
||||
expectedRevision?: number;
|
||||
}
|
||||
|
||||
export interface LeaderTaskSummaryDelivery {
|
||||
@@ -93,6 +90,17 @@ const NEEDS_REVIEW_STATUSES = [
|
||||
'uncertain'
|
||||
] as const;
|
||||
|
||||
const LEGACY_CHANNEL_REF = 'legacy-env';
|
||||
|
||||
interface AutomaticLeaderRoute {
|
||||
organizationId: string;
|
||||
leaderUserId: string;
|
||||
channelId: string;
|
||||
recipientAddress: string;
|
||||
conversationId: string;
|
||||
source: 'agentbus_inbound' | 'channel_external_user_ref';
|
||||
}
|
||||
|
||||
const noopLogger: LeaderNotificationLogger = {
|
||||
info: () => undefined,
|
||||
warn: () => undefined,
|
||||
@@ -135,10 +143,16 @@ function eligibility(
|
||||
if (!booleanValue(row.channel_enabled)) {
|
||||
return { eligible: false, message: 'AgentBus 渠道已停用,通知不会发送。' };
|
||||
}
|
||||
if (!row.target_verified_at) {
|
||||
return { eligible: false, message: '主动投递目标尚未验证,通知不会发送。' };
|
||||
if (!text(row.leader_erp_account)) {
|
||||
return { eligible: false, message: '组长渠道尚未满足账号绑定条件,通知不会发送。' };
|
||||
}
|
||||
return { eligible: true, message: '投递条件已就绪。' };
|
||||
if (!row.target_verified_at) {
|
||||
return { eligible: false, message: 'AgentBus 尚未提供可用路由;收到组长消息后会自动同步。' };
|
||||
}
|
||||
if (!booleanValue(row.enabled)) {
|
||||
return { eligible: false, message: '自动订阅正在同步,暂时不会发送。' };
|
||||
}
|
||||
return { eligible: true, message: '组长身份和 AgentBus 路由已就绪,摘要会自动推送。' };
|
||||
}
|
||||
|
||||
function publicSubscription(
|
||||
@@ -156,9 +170,11 @@ function publicSubscription(
|
||||
scope: 'organization',
|
||||
include_manual: booleanValue(row.include_manual),
|
||||
include_agentbus: booleanValue(row.include_agentbus),
|
||||
automatic: true,
|
||||
enabled: booleanValue(row.enabled),
|
||||
eligible: state.eligible,
|
||||
eligibility_message: state.message,
|
||||
route_ready: Boolean(row.target_verified_at),
|
||||
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),
|
||||
@@ -212,7 +228,17 @@ export class LeaderNotificationService {
|
||||
private async projectTick(): Promise<void> {
|
||||
if (!this.projectorOrganizationId) return;
|
||||
if (this.projectorInFlight) return this.projectorInFlight;
|
||||
this.projectorInFlight = this.projectPending(this.projectorOrganizationId, 100)
|
||||
this.projectorInFlight = this.reconcileAutomaticSubscriptions(this.projectorOrganizationId)
|
||||
.then((changed) => {
|
||||
if (changed > 0) {
|
||||
this.log('info', {
|
||||
agentbus_event: 'leader_summary_automatic_subscriptions_reconciled',
|
||||
organization_id: this.projectorOrganizationId,
|
||||
changed_count: changed
|
||||
}, 'Automatic leader task summary subscriptions reconciled');
|
||||
}
|
||||
return this.projectPending(this.projectorOrganizationId, 100);
|
||||
})
|
||||
.then((count) => {
|
||||
if (count > 0) {
|
||||
this.log('info', {
|
||||
@@ -241,6 +267,7 @@ export class LeaderNotificationService {
|
||||
leader.username AS leader_username,
|
||||
leader.role AS leader_role,
|
||||
leader.is_active AS leader_is_active,
|
||||
leader.erp_account AS leader_erp_account,
|
||||
channel.display_name AS channel_name,
|
||||
channel.owner_user_id AS channel_owner_user_id,
|
||||
channel.enabled AS channel_enabled,
|
||||
@@ -277,242 +304,332 @@ export class LeaderNotificationService {
|
||||
.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);
|
||||
private automaticRoute(
|
||||
input: Omit<AutomaticLeaderRoute, 'recipientAddress' | 'conversationId'> & {
|
||||
recipientAddress: unknown;
|
||||
conversationId?: unknown;
|
||||
}
|
||||
): AutomaticLeaderRoute | null {
|
||||
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 (!recipientAddress || recipientAddress === LEGACY_CHANNEL_REF) return null;
|
||||
const conversationId = (
|
||||
text(input.conversationId)
|
||||
|| `agentbus:${recipientAddress}`
|
||||
).slice(0, 500);
|
||||
if (!conversationId) return null;
|
||||
return { ...input, recipientAddress, conversationId };
|
||||
}
|
||||
|
||||
private async auditAutomaticChange(
|
||||
client: import('pg').PoolClient,
|
||||
input: {
|
||||
organizationId: string;
|
||||
subscriptionId: string;
|
||||
eventType: string;
|
||||
metadata: Record<string, unknown>;
|
||||
}
|
||||
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
|
||||
): Promise<void> {
|
||||
await client.query(
|
||||
`INSERT INTO audit_events
|
||||
(organization_id, actor_user_id, event_type, entity_type, entity_id, request_id, metadata)
|
||||
VALUES ($1, NULL, $2, 'leader_task_summary_subscription', $3, NULL, $4)`,
|
||||
[input.organizationId, input.eventType, input.subscriptionId, input.metadata]
|
||||
);
|
||||
this.log('info', {
|
||||
agentbus_event: 'leader_summary_automatic_subscription_audit_staged',
|
||||
entity_id: input.subscriptionId,
|
||||
metadata_keys: diagnosticMetadataKeys(input.metadata)
|
||||
}, 'Automatic leader summary subscription audit event staged');
|
||||
}
|
||||
|
||||
private async syncAutomaticSubscription(
|
||||
client: import('pg').PoolClient,
|
||||
route: AutomaticLeaderRoute,
|
||||
existing?: Record<string, unknown>
|
||||
): Promise<boolean> {
|
||||
const recipientFingerprint = sha256Text(route.recipientAddress);
|
||||
const conversationFingerprint = sha256Text(route.conversationId);
|
||||
const routeChanged = !existing
|
||||
|| route.channelId !== text(existing.channel_id)
|
||||
|| recipientFingerprint !== text(existing.recipient_address_fingerprint)
|
||||
|| conversationFingerprint !== text(existing.conversation_id_fingerprint);
|
||||
const policyChanged = !existing
|
||||
|| !booleanValue(existing.enabled)
|
||||
|| !booleanValue(existing.include_manual)
|
||||
|| !booleanValue(existing.include_agentbus)
|
||||
|| !existing.target_verified_at;
|
||||
if (!routeChanged && !policyChanged) return false;
|
||||
|
||||
const nextRevision = existing ? Number(existing.revision || 0) + 1 : 0;
|
||||
let subscriptionId: string;
|
||||
if (existing) {
|
||||
const updated = await client.query(
|
||||
`UPDATE leader_task_summary_subscriptions
|
||||
SET channel_id = $1,
|
||||
include_manual = true,
|
||||
include_agentbus = true,
|
||||
enabled = true,
|
||||
starts_at = now(),
|
||||
revision = $2,
|
||||
recipient_address_ciphertext = $3,
|
||||
recipient_address_fingerprint = $4,
|
||||
conversation_id_ciphertext = $5,
|
||||
conversation_id_fingerprint = $6,
|
||||
target_verified_at = now(),
|
||||
target_verified_by = NULL,
|
||||
updated_at = now()
|
||||
WHERE id = $7
|
||||
RETURNING id`,
|
||||
[
|
||||
route.channelId,
|
||||
nextRevision,
|
||||
encryptText(this.config, route.recipientAddress),
|
||||
recipientFingerprint,
|
||||
encryptText(this.config, route.conversationId),
|
||||
conversationFingerprint,
|
||||
existing.id
|
||||
]
|
||||
);
|
||||
subscriptionId = text(updated.rows[0]?.id);
|
||||
} 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, true, true, true, now(), 0,
|
||||
$4, $5, $6, $7, now(), NULL, NULL)
|
||||
RETURNING id`,
|
||||
[
|
||||
route.organizationId,
|
||||
route.leaderUserId,
|
||||
route.channelId,
|
||||
encryptText(this.config, route.recipientAddress),
|
||||
recipientFingerprint,
|
||||
encryptText(this.config, route.conversationId),
|
||||
conversationFingerprint
|
||||
]
|
||||
);
|
||||
subscriptionId = text(inserted.rows[0]?.id);
|
||||
}
|
||||
|
||||
const subscriptionId = await withTransaction(this.config, async (client) => {
|
||||
if (existing) {
|
||||
await client.query(
|
||||
`SELECT pg_advisory_xact_lock(hashtextextended($1::text || ':leader-summary:' || $2::text, 0))`,
|
||||
[context.organizationId, leaderUserId]
|
||||
`UPDATE leader_task_summary_deliveries
|
||||
SET delivery_status = 'cancelled',
|
||||
last_error = '组长身份或 AgentBus 路由已变化,旧路由待发送摘要已取消。',
|
||||
updated_at = now()
|
||||
WHERE subscription_id = $1
|
||||
AND subscription_revision <> $2
|
||||
AND delivery_status IN ('pending', 'sending', 'failed')`,
|
||||
[subscriptionId, nextRevision]
|
||||
);
|
||||
}
|
||||
await this.auditAutomaticChange(client, {
|
||||
organizationId: route.organizationId,
|
||||
subscriptionId,
|
||||
eventType: existing && routeChanged
|
||||
? 'leader_summary_subscription.automatic_route_updated'
|
||||
: 'leader_summary_subscription.automatic_enabled',
|
||||
metadata: {
|
||||
leader_user_id: route.leaderUserId,
|
||||
channel_id: route.channelId,
|
||||
automatic: true,
|
||||
include_manual: true,
|
||||
include_agentbus: true,
|
||||
enabled: true,
|
||||
route_source: route.source,
|
||||
recipient_fingerprint: recipientFingerprint.slice(0, 12),
|
||||
conversation_fingerprint: conversationFingerprint.slice(0, 12),
|
||||
revision: nextRevision,
|
||||
historical_backfill: false
|
||||
}
|
||||
});
|
||||
return true;
|
||||
}
|
||||
|
||||
private async disableAutomaticSubscription(
|
||||
client: import('pg').PoolClient,
|
||||
existing: Record<string, unknown>,
|
||||
reason: 'channel_disabled' | 'route_unavailable' | 'role_or_channel_invalid'
|
||||
): Promise<boolean> {
|
||||
if (!booleanValue(existing.enabled) && !existing.target_verified_at) return false;
|
||||
const nextRevision = Number(existing.revision || 0) + 1;
|
||||
const updated = await client.query(
|
||||
`UPDATE leader_task_summary_subscriptions
|
||||
SET enabled = false,
|
||||
starts_at = now(),
|
||||
revision = $2,
|
||||
target_verified_at = NULL,
|
||||
target_verified_by = NULL,
|
||||
updated_at = now()
|
||||
WHERE id = $1
|
||||
RETURNING id`,
|
||||
[existing.id, nextRevision]
|
||||
);
|
||||
const subscriptionId = text(updated.rows[0]?.id);
|
||||
await client.query(
|
||||
`UPDATE leader_task_summary_deliveries
|
||||
SET delivery_status = 'cancelled',
|
||||
last_error = '组长身份或 AgentBus 路由已失效,待发送摘要已取消。',
|
||||
updated_at = now()
|
||||
WHERE subscription_id = $1
|
||||
AND delivery_status IN ('pending', 'sending', 'failed')`,
|
||||
[subscriptionId]
|
||||
);
|
||||
await this.auditAutomaticChange(client, {
|
||||
organizationId: text(existing.organization_id),
|
||||
subscriptionId,
|
||||
eventType: 'leader_summary_subscription.automatic_disabled',
|
||||
metadata: {
|
||||
leader_user_id: text(existing.leader_user_id),
|
||||
channel_id: text(existing.channel_id),
|
||||
automatic: true,
|
||||
enabled: false,
|
||||
reason,
|
||||
revision: nextRevision,
|
||||
historical_backfill: false
|
||||
}
|
||||
});
|
||||
return true;
|
||||
}
|
||||
|
||||
async reconcileAutomaticSubscriptions(organizationId: string): Promise<number> {
|
||||
return withTransaction(this.config, async (client) => {
|
||||
await client.query(
|
||||
`SELECT pg_advisory_xact_lock(
|
||||
hashtextextended($1::text || ':leader-summary-automatic', 0)
|
||||
)`,
|
||||
[organizationId]
|
||||
);
|
||||
const existingResult = await client.query(
|
||||
`SELECT *
|
||||
FROM leader_task_summary_subscriptions
|
||||
WHERE organization_id = $1
|
||||
FOR UPDATE`,
|
||||
[organizationId]
|
||||
);
|
||||
const existingByLeader = new Map<string, Record<string, unknown>>(
|
||||
(existingResult.rows as Record<string, unknown>[])
|
||||
.map((row) => [text(row.leader_user_id), row])
|
||||
);
|
||||
const candidates = await client.query(
|
||||
`SELECT leader.id AS leader_user_id,
|
||||
channel.id AS channel_id,
|
||||
channel.enabled AS channel_enabled,
|
||||
channel.external_user_ref,
|
||||
latest_route.inbound_from,
|
||||
latest_route.conversation_id
|
||||
FROM users leader
|
||||
JOIN user_channels channel
|
||||
ON channel.organization_id = leader.organization_id
|
||||
AND channel.owner_user_id = leader.id
|
||||
LEFT JOIN LATERAL (
|
||||
SELECT delivery.inbound_from, delivery.conversation_id
|
||||
FROM agentbus_deliveries delivery
|
||||
JOIN tasks route_task
|
||||
ON route_task.id = delivery.task_id
|
||||
AND route_task.organization_id = delivery.organization_id
|
||||
AND route_task.assigned_user_id = leader.id
|
||||
WHERE delivery.organization_id = leader.organization_id
|
||||
AND delivery.channel_id = channel.id
|
||||
AND btrim(delivery.inbound_from) <> ''
|
||||
ORDER BY delivery.created_at DESC, delivery.id DESC
|
||||
LIMIT 1
|
||||
) latest_route ON true
|
||||
WHERE leader.organization_id = $1
|
||||
AND leader.role = 'team_lead'
|
||||
AND leader.is_active = true
|
||||
AND leader.erp_account IS NOT NULL
|
||||
ORDER BY leader.id`,
|
||||
[organizationId]
|
||||
);
|
||||
|
||||
const currentLeaders = new Set<string>();
|
||||
let changed = 0;
|
||||
for (const row of candidates.rows as Record<string, unknown>[]) {
|
||||
const leaderUserId = text(row.leader_user_id);
|
||||
const channelId = text(row.channel_id);
|
||||
currentLeaders.add(leaderUserId);
|
||||
const existing = existingByLeader.get(leaderUserId);
|
||||
if (!booleanValue(row.channel_enabled)) {
|
||||
if (existing && await this.disableAutomaticSubscription(client, existing, 'channel_disabled')) changed += 1;
|
||||
continue;
|
||||
}
|
||||
const inboundRecipient = text(row.inbound_from);
|
||||
const externalRecipient = text(row.external_user_ref);
|
||||
const route = this.automaticRoute({
|
||||
organizationId,
|
||||
leaderUserId,
|
||||
channelId,
|
||||
recipientAddress: inboundRecipient || externalRecipient,
|
||||
conversationId: inboundRecipient ? row.conversation_id : undefined,
|
||||
source: inboundRecipient ? 'agentbus_inbound' : 'channel_external_user_ref'
|
||||
});
|
||||
if (!route) {
|
||||
if (existing && await this.disableAutomaticSubscription(client, existing, 'route_unavailable')) changed += 1;
|
||||
continue;
|
||||
}
|
||||
if (await this.syncAutomaticSubscription(client, route, existing)) changed += 1;
|
||||
}
|
||||
|
||||
for (const [leaderUserId, existing] of existingByLeader) {
|
||||
if (currentLeaders.has(leaderUserId)) continue;
|
||||
if (await this.disableAutomaticSubscription(client, existing, 'role_or_channel_invalid')) changed += 1;
|
||||
}
|
||||
return changed;
|
||||
});
|
||||
}
|
||||
|
||||
async observeLeaderRoute(input: ObservedLeaderAgentBusRoute): Promise<void> {
|
||||
const route = this.automaticRoute({
|
||||
organizationId: text(input.organizationId),
|
||||
leaderUserId: text(input.leaderUserId),
|
||||
channelId: text(input.channelId),
|
||||
recipientAddress: input.recipientAddress,
|
||||
conversationId: input.conversationId,
|
||||
source: 'agentbus_inbound'
|
||||
});
|
||||
if (!route || !route.organizationId || !route.leaderUserId || !route.channelId) return;
|
||||
await withTransaction(this.config, async (client) => {
|
||||
await client.query(
|
||||
`SELECT pg_advisory_xact_lock(
|
||||
hashtextextended($1::text || ':leader-summary-automatic', 0)
|
||||
)`,
|
||||
[route.organizationId]
|
||||
);
|
||||
const eligible = await client.query(
|
||||
`SELECT channel.id
|
||||
FROM user_channels channel
|
||||
JOIN users leader
|
||||
ON leader.id = channel.owner_user_id
|
||||
AND leader.organization_id = channel.organization_id
|
||||
AND leader.role = 'team_lead'
|
||||
AND leader.is_active = true
|
||||
AND leader.erp_account IS NOT NULL
|
||||
WHERE channel.organization_id = $1
|
||||
AND channel.id = $2
|
||||
AND channel.owner_user_id = $3
|
||||
AND channel.enabled = true
|
||||
FOR SHARE OF channel, leader`,
|
||||
[route.organizationId, route.channelId, route.leaderUserId]
|
||||
);
|
||||
if (!eligible.rowCount) return;
|
||||
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]
|
||||
[route.organizationId, route.leaderUserId]
|
||||
);
|
||||
const existing = existingResult.rows[0] as Record<string, unknown> | undefined;
|
||||
const safeDisable = Boolean(
|
||||
existing
|
||||
&& !input.enabled
|
||||
&& channelId === text(existing.channel_id)
|
||||
await this.syncAutomaticSubscription(
|
||||
client,
|
||||
route,
|
||||
existingResult.rows[0] as Record<string, unknown> | undefined
|
||||
);
|
||||
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> {
|
||||
@@ -550,6 +667,7 @@ export class LeaderNotificationService {
|
||||
AND leader.organization_id = subscription.organization_id
|
||||
AND leader.role = 'team_lead'
|
||||
AND leader.is_active = true
|
||||
AND leader.erp_account IS NOT NULL
|
||||
JOIN user_channels channel
|
||||
ON channel.id = subscription.channel_id
|
||||
AND channel.organization_id = subscription.organization_id
|
||||
@@ -714,6 +832,7 @@ export class LeaderNotificationService {
|
||||
AND leader.organization_id = subscription.organization_id
|
||||
AND leader.role = 'team_lead'
|
||||
AND leader.is_active = true
|
||||
AND leader.erp_account IS NOT NULL
|
||||
JOIN user_channels channel
|
||||
ON channel.id = subscription.channel_id
|
||||
AND channel.organization_id = subscription.organization_id
|
||||
@@ -752,6 +871,7 @@ export class LeaderNotificationService {
|
||||
AND leader.organization_id = subscription.organization_id
|
||||
AND leader.role = 'team_lead'
|
||||
AND leader.is_active = true
|
||||
AND leader.erp_account IS NOT NULL
|
||||
JOIN user_channels channel
|
||||
ON channel.id = subscription.channel_id
|
||||
AND channel.organization_id = subscription.organization_id
|
||||
|
||||
@@ -157,16 +157,6 @@ const channelUpdateSchema = z.object({
|
||||
enabled: z.boolean().optional()
|
||||
});
|
||||
const channelRotateKeySchema = z.object({ agentbus_key: z.string().min(1).max(4_000) });
|
||||
const leaderSummarySubscriptionSchema = z.object({
|
||||
channel_id: z.string().uuid(),
|
||||
recipient_address: z.string().trim().min(1).max(500).optional(),
|
||||
conversation_id: z.string().trim().min(1).max(500).optional(),
|
||||
include_manual: z.boolean(),
|
||||
include_agentbus: z.boolean(),
|
||||
enabled: z.boolean(),
|
||||
target_verified: z.boolean().optional(),
|
||||
expected_revision: z.number().int().min(0).optional()
|
||||
});
|
||||
const listTasksQuerySchema = z.object({
|
||||
status: z.string().max(80).optional(),
|
||||
search: z.string().max(200).optional(),
|
||||
@@ -1191,27 +1181,6 @@ export async function buildServer({
|
||||
};
|
||||
});
|
||||
|
||||
app.put('/api/settings/leader-summary-subscriptions/:leaderUserId', async (request) => {
|
||||
const session = await requireAdminMutationSession(request);
|
||||
const params = request.params as { leaderUserId: string };
|
||||
const body = leaderSummarySubscriptionSchema.parse(request.body);
|
||||
const subscription = await leaderNotificationService.upsertSubscription(
|
||||
contextFor(session, request),
|
||||
{
|
||||
leaderUserId: params.leaderUserId,
|
||||
channelId: body.channel_id,
|
||||
recipientAddress: body.recipient_address,
|
||||
conversationId: body.conversation_id,
|
||||
includeManual: body.include_manual,
|
||||
includeAgentBus: body.include_agentbus,
|
||||
enabled: body.enabled,
|
||||
targetVerified: body.target_verified,
|
||||
expectedRevision: body.expected_revision
|
||||
}
|
||||
);
|
||||
return { ok: true, subscription };
|
||||
});
|
||||
|
||||
app.post('/api/auth/logout', async (request, reply) => {
|
||||
setAuthNoStore(reply);
|
||||
const session = await requireAuthenticatedMutationSession(request);
|
||||
|
||||
Reference in new issue
Block a user