大交通检索修复
This commit is contained in:
217
control-plane/src/arrangement-ledger.ts
Normal file
217
control-plane/src/arrangement-ledger.ts
Normal file
@@ -0,0 +1,217 @@
|
||||
import type pg from 'pg';
|
||||
import type { AppConfig } from './config.js';
|
||||
import { decryptText, encryptText, sha256Text } from './crypto.js';
|
||||
|
||||
export const ARRANGEMENT_LABELS = {
|
||||
arrangement_guide: '导游', arrangement_vehicle: '用车', arrangement_hotel: '酒店',
|
||||
arrangement_transport: '大交通', arrangement_other: '其他/备案'
|
||||
} as const;
|
||||
export type ArrangementAction = keyof typeof ARRANGEMENT_LABELS;
|
||||
type Json = Record<string, unknown>;
|
||||
type Client = Pick<pg.PoolClient, 'query'>;
|
||||
export interface ArrangementScope { organizationId: string; userId: string }
|
||||
export interface ArrangementEvent {
|
||||
version: 1;
|
||||
group_number: string;
|
||||
action: ArrangementAction;
|
||||
mode: 'create' | 'update' | 'clear';
|
||||
record_key: string;
|
||||
source_task_id: string;
|
||||
execution_id: string;
|
||||
effective_at: string;
|
||||
recorded_at: string;
|
||||
source_task_deleted_at?: string | null;
|
||||
details: Json;
|
||||
requested: Json;
|
||||
evidence: Json;
|
||||
}
|
||||
export interface ArrangementSummary {
|
||||
version: 1;
|
||||
scope: 'current_account';
|
||||
group_number: string;
|
||||
recorded_at: string;
|
||||
items: Array<{ action: ArrangementAction; label: string; arranged: boolean; count: number }>;
|
||||
}
|
||||
function object(value: unknown): Json {
|
||||
return value && typeof value === 'object' && !Array.isArray(value) ? value as Json : {};
|
||||
}
|
||||
function text(value: unknown): string { return typeof value === 'string' || typeof value === 'number' ? String(value).trim() : ''; }
|
||||
function first(...values: unknown[]): string { return values.map(text).find(Boolean) || ''; }
|
||||
export function arrangementGroupNumber(value: unknown): string {
|
||||
return text(value).normalize('NFKC').replace(/\s+/gu, '').toUpperCase();
|
||||
}
|
||||
export function isArrangementAction(value: unknown): value is ArrangementAction {
|
||||
return typeof value === 'string' && Object.hasOwn(ARRANGEMENT_LABELS, value);
|
||||
}
|
||||
function timestamp(value: unknown): string {
|
||||
const date = value instanceof Date ? value : new Date(String(value || ''));
|
||||
return Number.isFinite(date.getTime()) ? date.toISOString() : '';
|
||||
}
|
||||
function dateValue(value: unknown): string {
|
||||
return text(value).replace(/^(\d{4})-(\d{1,2})-(\d{1,2})$/, (_, y, m, d) => `${y}-${m.padStart(2, '0')}-${d.padStart(2, '0')}`);
|
||||
}
|
||||
|
||||
// Only confirmed, action-specific ERP evidence enters the ledger. Parser success,
|
||||
// bare completed flags and uncertain results are never arrangements.
|
||||
export function arrangementEventFromSuccess(input: {
|
||||
taskId: string; executionId: string; effectiveAt: unknown; recordedAt: unknown;
|
||||
operation: unknown; result: unknown; receipt: unknown;
|
||||
}): ArrangementEvent | null {
|
||||
const operation = object(input.operation);
|
||||
if (!isArrangementAction(operation.action)) return null;
|
||||
const result = object(input.result), report = object(result.report), receipt = object(input.receipt);
|
||||
const requery = object(report.requery || result.requery);
|
||||
const server = object(report.server_response || result.server_response);
|
||||
if (result.status !== 'completed' || receipt.success !== true || server.completed !== true || requery.matched !== true
|
||||
|| requery.skipped === true || result.uncertain === true || result.manual_review_required === true
|
||||
|| report.manual_review_required === true
|
||||
|| [result.blockers, report.blockers].some(value => Array.isArray(value) && value.length)) return null;
|
||||
const data = object(operation.data), arrangement = object(data.arrangement);
|
||||
const refs = object(report.resolved_refs || result.resolved_refs);
|
||||
const mode = text(arrangement.mode) || 'create';
|
||||
if (!['create', 'update', 'clear'].includes(mode)) return null;
|
||||
const group = arrangementGroupNumber(first(refs.identifier, refs.group_no, receipt.group_number, object(data.existing_refs).identifier));
|
||||
const suppliedGroup = arrangementGroupNumber(object(data.existing_refs).identifier);
|
||||
if (!group || (suppliedGroup && suppliedGroup !== group) || refs.kind === 'shared_child_order') return null;
|
||||
const recordedAt = timestamp(input.recordedAt), effectiveAt = timestamp(input.effectiveAt) || recordedAt;
|
||||
if (!recordedAt || !effectiveAt || !input.taskId || !input.executionId) return null;
|
||||
const subpage = object(requery.subpage);
|
||||
const target = object(subpage.arrangement_target || requery.arrangement_target || refs.arrangement_target);
|
||||
const checks = (Array.isArray(subpage.checks) ? subpage.checks : Array.isArray(requery.checks) ? requery.checks : []).map(object);
|
||||
const actual = (pattern: RegExp): string | undefined => {
|
||||
const matches = checks.filter(check => check.matched === true && pattern.test(text(check.field)) && check.actual != null);
|
||||
return matches.length === 1 ? text(matches[0].actual) : undefined;
|
||||
};
|
||||
const rowId = first(target.row_id, object(arrangement.target).row_id);
|
||||
// Guides are genuinely singleton. Indexed row position is not a stable identity.
|
||||
const recordKey = operation.action === 'arrangement_guide' ? 'guide'
|
||||
: rowId && rowId !== '0' ? `row:${rowId}` : `execution:${input.executionId}`;
|
||||
const changes = object(arrangement.changes);
|
||||
const quantity = first(target.quantity_value, actual(/^(?:jianshu|shuliang|tianshu)\d+$/u), changes.room_count, changes.quantity, arrangement.room_count, arrangement.quantity);
|
||||
const details: Json = { row_id: rowId };
|
||||
// Absence is different from a verified empty value. Narrow updates must not
|
||||
// discard details established by an earlier successful create/update.
|
||||
const set = (key: string, value: unknown, date = false): void => {
|
||||
if (value !== undefined && value !== null) details[key] = date ? dateValue(value) : value;
|
||||
};
|
||||
set('resource_id', target.resource_id);
|
||||
// Never label a requested search term as the actual selected resource.
|
||||
set('resource_name', actual(/^(?:danwei|daoyou)\d+$/u));
|
||||
set('item', target.item_value ?? actual(/^(?:fangxing|chehao|shuoming)\d+$/u));
|
||||
if (operation.action === 'arrangement_hotel' || operation.action === 'arrangement_vehicle') {
|
||||
set('start_date', actual(/^riqi\d+$/u) ?? arrangement.start_date, true);
|
||||
set('end_date', actual(/^riqis\d+$/u) ?? changes.end_date ?? arrangement.end_date, true);
|
||||
} else set('date', actual(/^riqi\d+$/u) ?? arrangement.date, true);
|
||||
if (quantity && Number.isFinite(Number(quantity))) details.quantity = Number(quantity);
|
||||
set('remark', actual(/^beizhu\d+$/u) ?? arrangement.remark);
|
||||
set('phone', actual(/^dianhua\d+$/u));
|
||||
set('grade', actual(/^dengji\d+$/u));
|
||||
if (operation.action === 'arrangement_other') {
|
||||
for (const [key, field] of Object.entries({ number: 'beianhao', date: 'beianri', entry_port: 'rujingkou', exit_port: 'chujingkou' })) {
|
||||
set(`filing_${key}`, actual(new RegExp(`^${field}$`, 'u')) ?? object(arrangement.filing)[key], key === 'date');
|
||||
}
|
||||
}
|
||||
return {
|
||||
version: 1, group_number: group, action: operation.action, mode: mode as ArrangementEvent['mode'], record_key: recordKey,
|
||||
source_task_id: input.taskId, execution_id: input.executionId, effective_at: effectiveAt, recorded_at: recordedAt,
|
||||
details, requested: arrangement,
|
||||
evidence: { receipt, resolved_refs: refs, arrangement_target: target, checks,
|
||||
server_response_completed: true, requery_matched: true }
|
||||
};
|
||||
}
|
||||
|
||||
export function currentArrangements(events: ArrangementEvent[]): ArrangementEvent[] {
|
||||
const current = new Map<string, ArrangementEvent>();
|
||||
const seen = new Set<string>();
|
||||
for (const event of [...events].sort((a, b) => a.effective_at.localeCompare(b.effective_at)
|
||||
|| a.recorded_at.localeCompare(b.recorded_at) || a.execution_id.localeCompare(b.execution_id))) {
|
||||
if (seen.has(event.execution_id)) continue;
|
||||
seen.add(event.execution_id);
|
||||
const key = `${event.group_number}:${event.action}:${event.record_key}`;
|
||||
if (event.mode === 'clear') current.delete(key);
|
||||
else current.set(key, event.mode === 'update' ? { ...event, details: { ...current.get(key)?.details, ...event.details } } : event);
|
||||
}
|
||||
return [...current.values()];
|
||||
}
|
||||
|
||||
export function arrangementSummary(events: ArrangementEvent[], groupNumber: string, recordedAt: string): ArrangementSummary {
|
||||
const group = arrangementGroupNumber(groupNumber);
|
||||
const current = currentArrangements(events.filter(event => event.group_number === group));
|
||||
return { version: 1, scope: 'current_account', group_number: group, recorded_at: recordedAt,
|
||||
items: (Object.keys(ARRANGEMENT_LABELS) as ArrangementAction[]).map(action => {
|
||||
const count = current.filter(event => event.action === action).length;
|
||||
return { action, label: ARRANGEMENT_LABELS[action], arranged: count > 0, count };
|
||||
}) };
|
||||
}
|
||||
|
||||
export function arrangementSummaryText(value: unknown): string {
|
||||
const summary = object(value);
|
||||
if (summary.version !== 1 || summary.scope !== 'current_account' || !Array.isArray(summary.items)) return '';
|
||||
const items = summary.items.map(object);
|
||||
if (items.length !== 5 || Object.keys(ARRANGEMENT_LABELS).some(action =>
|
||||
items.filter(item => item.action === action && typeof item.arranged === 'boolean').length !== 1)) return '';
|
||||
return ['本系统安排情况(当前账号):', ...Object.entries(ARRANGEMENT_LABELS).map(([action, label]) =>
|
||||
`${label}:${items.find(item => item.action === action)?.arranged ? '已安排' : '未安排'}`),
|
||||
'仅统计通过本系统成功执行的安排;已安排表示至少一条,不代表全行程已安排齐全。'].join('\n');
|
||||
}
|
||||
|
||||
export class ArrangementLedger {
|
||||
constructor(private readonly config: AppConfig) {}
|
||||
|
||||
// Serialize only this employee's ledger, including lazy historical backfill.
|
||||
// Callers already hold their task row lock; backfill never locks other tasks.
|
||||
async prepare(client: Client, scope: ArrangementScope, receiptFromResult: (value: unknown) => Json | null): Promise<void> {
|
||||
if (!scope.organizationId || !scope.userId) throw new Error('arrangement_ledger_scope_missing');
|
||||
await client.query(`INSERT INTO arrangement_ledger_owners (organization_id, owner_user_id)
|
||||
VALUES ($1, $2) ON CONFLICT DO NOTHING`, [scope.organizationId, scope.userId]);
|
||||
const lock = await client.query(`SELECT backfilled_at FROM arrangement_ledger_owners
|
||||
WHERE organization_id = $1 AND owner_user_id = $2 FOR UPDATE`, [scope.organizationId, scope.userId]);
|
||||
if (lock.rows[0]?.backfilled_at) return;
|
||||
let cursor = '';
|
||||
for (;;) {
|
||||
const history = await client.query(`SELECT t.task_id, t.operation_ciphertext, t.operation,
|
||||
t.execution_result_ciphertext, t.execution_result, t.success_receipt_at, t.updated_at,
|
||||
a.id AS execution_id, a.started_at AS effective_at
|
||||
FROM tasks t JOIN task_attempts a ON a.task_id = t.id AND a.phase = 'erp'
|
||||
WHERE t.organization_id = $1 AND t.assigned_user_id = $2 AND t.status = 'completed'
|
||||
AND t.task_id > $3 AND t.operation->>'action' = ANY($4::text[])
|
||||
ORDER BY t.task_id LIMIT 200`, [scope.organizationId, scope.userId, cursor, Object.keys(ARRANGEMENT_LABELS)]);
|
||||
for (const row of history.rows) {
|
||||
const decode = (encrypted: unknown, plain: unknown) => encrypted
|
||||
? JSON.parse(decryptText(this.config, text(encrypted))) : plain;
|
||||
const result = decode(row.execution_result_ciphertext, row.execution_result);
|
||||
const receipt = receiptFromResult(result);
|
||||
const event = arrangementEventFromSuccess({ taskId: row.task_id, executionId: row.execution_id,
|
||||
effectiveAt: row.effective_at, recordedAt: row.success_receipt_at || row.updated_at,
|
||||
operation: decode(row.operation_ciphertext, row.operation), result, receipt });
|
||||
if (event) await this.insert(client, scope, event);
|
||||
}
|
||||
if (history.rows.length < 200) break;
|
||||
cursor = String(history.rows.at(-1)?.task_id);
|
||||
}
|
||||
await client.query(`UPDATE arrangement_ledger_owners SET backfilled_at = now()
|
||||
WHERE organization_id = $1 AND owner_user_id = $2`, [scope.organizationId, scope.userId]);
|
||||
}
|
||||
|
||||
async insert(client: Client, scope: ArrangementScope, event: ArrangementEvent): Promise<void> {
|
||||
await client.query(`INSERT INTO arrangement_ledger_events
|
||||
(organization_id, owner_user_id, group_key, action, source_task_id, execution_id, effective_at, recorded_at, detail_ciphertext)
|
||||
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9) ON CONFLICT (organization_id, owner_user_id, execution_id) DO NOTHING`,
|
||||
[scope.organizationId, scope.userId, sha256Text(event.group_number), event.action, event.source_task_id,
|
||||
event.execution_id, event.effective_at, event.recorded_at, encryptText(this.config, JSON.stringify(event))]);
|
||||
}
|
||||
|
||||
async history(client: Client, scope: ArrangementScope, groupNumber: string): Promise<ArrangementEvent[]> {
|
||||
const result = await client.query(`SELECT detail_ciphertext, source_task_deleted_at FROM arrangement_ledger_events
|
||||
WHERE organization_id = $1 AND owner_user_id = $2 AND group_key = $3
|
||||
ORDER BY effective_at, recorded_at, id`, [scope.organizationId, scope.userId, sha256Text(arrangementGroupNumber(groupNumber))]);
|
||||
return result.rows.map(row => ({ ...JSON.parse(decryptText(this.config, row.detail_ciphertext)),
|
||||
source_task_deleted_at: row.source_task_deleted_at ? timestamp(row.source_task_deleted_at) : null }));
|
||||
}
|
||||
|
||||
async markTasksDeleted(client: Client, scope: ArrangementScope, taskIds: string[]): Promise<void> {
|
||||
await client.query(`UPDATE arrangement_ledger_events SET source_task_deleted_at = now()
|
||||
WHERE organization_id = $1 AND owner_user_id = $2 AND source_task_id = ANY($3::text[])
|
||||
AND source_task_deleted_at IS NULL`, [scope.organizationId, scope.userId, taskIds]);
|
||||
}
|
||||
}
|
||||
@@ -5,7 +5,7 @@ import { writeEmergencyDiagnostic } from './diagnostics.js';
|
||||
const { Pool } = pg;
|
||||
let pool: pg.Pool | null = null;
|
||||
|
||||
export const REQUIRED_SCHEMA_VERSION = '024_arrangement_update_routes';
|
||||
export const REQUIRED_SCHEMA_VERSION = '025_arrangement_ledger';
|
||||
|
||||
export interface DatabaseReadiness {
|
||||
ready: boolean;
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { normalizeTaskArtifactFileName } from './artifact-name.js';
|
||||
import { arrangementSummaryText, isArrangementAction } from './arrangement-ledger.js';
|
||||
|
||||
export const AGENTBUS_REPLY_CONTRACT_VERSION = 'agentbus-user-reply-v1';
|
||||
|
||||
@@ -561,12 +562,14 @@ export function buildCustomerSuccessReceipt(
|
||||
): Record<string, unknown> {
|
||||
const reply = buildCustomerSuccessReply(baseReceipt, resultValue, operationValue);
|
||||
if (!reply) return baseReceipt;
|
||||
const action = operationAction(operationValue, object(resultValue), resultContext(object(resultValue)));
|
||||
const summaryText = isArrangementAction(action) ? arrangementSummaryText(baseReceipt.arrangement_summary) : '';
|
||||
return {
|
||||
...baseReceipt,
|
||||
reply: {
|
||||
contract: reply.contract,
|
||||
profile: reply.profile,
|
||||
text: reply.text,
|
||||
text: summaryText ? `${reply.text}\n\n${summaryText}` : reply.text,
|
||||
fields: reply.fields
|
||||
},
|
||||
...(reply.attachments.length ? { reply_attachments: reply.attachments } : {})
|
||||
|
||||
@@ -31,6 +31,7 @@ import {
|
||||
} from './reply-contract.js';
|
||||
import { convertDocumentToPdf, convertDocumentToXlsx } from './document-converter.js';
|
||||
import { sanitizeOperationTiming } from './operation-timing.js';
|
||||
import { ArrangementLedger, arrangementEventFromSuccess, arrangementSummary } from './arrangement-ledger.js';
|
||||
import type { TaskInputAttachmentInput } from './input-attachment.js';
|
||||
import {
|
||||
diagnosticMetadataKeys,
|
||||
@@ -6616,7 +6617,9 @@ export class TaskService {
|
||||
const executionFailure = failureSummary(resultObject, '', status, stage);
|
||||
const operation = decryptedJson(this.config, row.operation_ciphertext) || row.operation || null;
|
||||
const baseSuccessReceipt = status === 'completed' ? successReceiptFromResult(resultObject) : null;
|
||||
const successReceipt = baseSuccessReceipt
|
||||
// This snapshot is exclusively derived by the server, never by a plugin.
|
||||
if (baseSuccessReceipt) delete baseSuccessReceipt.arrangement_summary;
|
||||
let successReceipt = baseSuccessReceipt
|
||||
? buildCustomerSuccessReceipt(baseSuccessReceipt, resultObject, operation)
|
||||
: null;
|
||||
const latestErrorSummary = isErrorSummaryStatus(status) ? errorSummaryPayload(executionFailure) : null;
|
||||
@@ -6652,6 +6655,21 @@ export class TaskService {
|
||||
if (text(row.lease_owner) !== leaseOwner && !canFinalizeTerminalReconciliation) {
|
||||
throw new TaskError('execution_owner_mismatch', '插件连接与任务领取记录不一致,回执已拒绝。');
|
||||
}
|
||||
const arrangementEvent = successReceipt ? arrangementEventFromSuccess({
|
||||
taskId, executionId, operation, result: resultObject, receipt: baseSuccessReceipt,
|
||||
effectiveAt: attemptRow.started_at, recordedAt: new Date().toISOString()
|
||||
}) : null;
|
||||
if (arrangementEvent) {
|
||||
const ledger = new ArrangementLedger(this.config);
|
||||
await ledger.prepare(client, context, successReceiptFromResult);
|
||||
await ledger.insert(client, context, arrangementEvent);
|
||||
const history = await ledger.history(client, context, arrangementEvent.group_number);
|
||||
// Freeze the account/group summary with this receipt in the SAME transaction.
|
||||
successReceipt = buildCustomerSuccessReceipt({
|
||||
...successReceipt,
|
||||
arrangement_summary: arrangementSummary(history, arrangementEvent.group_number, arrangementEvent.recorded_at)
|
||||
}, resultObject, operation);
|
||||
}
|
||||
const updated = await client.query(
|
||||
`UPDATE tasks
|
||||
SET execution_result_ciphertext = $1, execution_result = $2,
|
||||
@@ -7277,6 +7295,12 @@ export class TaskService {
|
||||
const missingTaskIds = normalizedTaskIds.filter((taskId) => !foundTaskIds.has(taskId));
|
||||
if (missingTaskIds.length) throw new TaskError('task_not_found', '任务不存在。', 404);
|
||||
|
||||
// Preserve historical successes before deleting their only original source.
|
||||
// Business records deliberately survive task cascades and retain snapshots.
|
||||
const ledger = new ArrangementLedger(this.config);
|
||||
await ledger.prepare(client, context, successReceiptFromResult);
|
||||
await ledger.markTasksDeleted(client, context, normalizedTaskIds);
|
||||
|
||||
// Hard delete intentionally has no status or handoff-state gate. The row
|
||||
// locks settle concurrent transitions; existing task foreign keys then
|
||||
// remove every task-owned record through ON DELETE CASCADE.
|
||||
|
||||
Reference in New Issue
Block a user