回执增加内容

This commit is contained in:
andy committed 2026-09-20 14:59:19 +08:00
1 parent ea9c4b34c6
commit 40230ce33d
17 files changed
+729 -9

No files matched your search

+10
View File
@@ -3,6 +3,16 @@ export const AGENTBUS_ROSTER_WAITING_TEXT = '已识别名单业务,等待 .xls
export const AGENTBUS_ROSTER_ATTACHMENT_RECEIVED_TEXT = '名单附件已收到,正在校验并处理。';
export const AGENTBUS_FINAL_REPLY_ELIGIBLE_FIELD = '_final_reply_eligible';
// Preserve the complete receipt on the platform; never silently cut a long
// multi-team/detail receipt in the middle of a field on the channel.
export function boundedAgentBusReplyText(value: string): string {
if (value.length <= 20_000) return value;
const suffix = '\n\n回执明细较长,此处仅展示部分;完整内容请在平台任务回执中查看。';
const prefix = value.slice(0, 20_000 - suffix.length);
const line = prefix.lastIndexOf('\n');
return (line > 0 ? prefix.slice(0, line) : prefix.replace(/[\uD800-\uDBFF]$/, '')) + suffix;
}
export interface AgentBusAcceptedDeliveryOptions {
text?: string;
finalReplyEligible?: boolean;
+3 -2
View File
@@ -21,6 +21,7 @@ import {
} from './input-attachment.js';
import { diagnosticDurationMs, diagnosticError } from './diagnostics.js';
import {
boundedAgentBusReplyText,
AGENTBUS_GENERIC_ACCEPTED_TEXT,
AGENTBUS_ROSTER_ATTACHMENT_RECEIVED_TEXT,
publicAgentBusDeliveryPayload
@@ -482,7 +483,7 @@ export function createTaskResultFrame(
payload: {
event: 'task.result',
status,
text: resultText.slice(0, 20_000),
text: boundedAgentBusReplyText(resultText),
...(attachments.length ? { attachments } : {}),
...(importantMessage ? { important_message: channelImportantMessage(importantMessage) } : {})
}
@@ -1149,7 +1150,7 @@ export class AgentBusListener {
event: 'task.result',
status: taskResultStatus(task),
task_id: task.task_id,
text: taskResultText(task).slice(0, 20_000),
text: boundedAgentBusReplyText(taskResultText(task)),
...(attachmentRefs.length ? { attachments: attachmentRefs } : {}),
...(importantMessage ? { important_message: importantMessage } : {})
};
+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 = '025_arrangement_ledger';
export const REQUIRED_SCHEMA_VERSION = '026_team_progress';
export interface DatabaseReadiness {
ready: boolean;
+3 -1
View File
@@ -1,5 +1,6 @@
import { normalizeTaskArtifactFileName } from './artifact-name.js';
import { arrangementSummaryText, isArrangementAction } from './arrangement-ledger.js';
import { teamProgressText } from './team-progress.js';
export const AGENTBUS_REPLY_CONTRACT_VERSION = 'agentbus-user-reply-v1';
@@ -563,7 +564,8 @@ export function buildCustomerSuccessReceipt(
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) : '';
const summaryText = teamProgressText(baseReceipt.team_progress)
|| (isArrangementAction(action) ? arrangementSummaryText(baseReceipt.arrangement_summary) : '');
return {
...baseReceipt,
reply: {
+22 -1
View File
@@ -32,6 +32,7 @@ import {
import { convertDocumentToPdf, convertDocumentToXlsx } from './document-converter.js';
import { sanitizeOperationTiming } from './operation-timing.js';
import { ArrangementLedger, arrangementEventFromSuccess, arrangementSummary } from './arrangement-ledger.js';
import { TeamProgressLedger, progressEventsFromSuccess, partialBatchProgressEvents } from './team-progress.js';
import type { TaskInputAttachmentInput } from './input-attachment.js';
import {
diagnosticMetadataKeys,
@@ -6618,7 +6619,10 @@ export class TaskService {
const operation = decryptedJson(this.config, row.operation_ciphertext) || row.operation || null;
const baseSuccessReceipt = status === 'completed' ? successReceiptFromResult(resultObject) : null;
// This snapshot is exclusively derived by the server, never by a plugin.
if (baseSuccessReceipt) delete baseSuccessReceipt.arrangement_summary;
if (baseSuccessReceipt) {
delete baseSuccessReceipt.arrangement_summary;
delete baseSuccessReceipt.team_progress;
}
let successReceipt = baseSuccessReceipt
? buildCustomerSuccessReceipt(baseSuccessReceipt, resultObject, operation)
: null;
@@ -6670,6 +6674,20 @@ export class TaskService {
arrangement_summary: arrangementSummary(history, arrangementEvent.group_number, arrangementEvent.recorded_at)
}, resultObject, operation);
}
const progressInput = { taskId, executionId, operation, result: resultObject,
receipt: baseSuccessReceipt, effectiveAt: attemptRow.started_at, recordedAt: new Date().toISOString() };
const progressEvents = [...progressEventsFromSuccess(progressInput), ...partialBatchProgressEvents(progressInput)];
if (successReceipt || progressEvents.length) {
const ledger = new ArrangementLedger(this.config);
await ledger.prepare(client, context, successReceiptFromResult);
const progress = new TeamProgressLedger(this.config);
await progress.prepare(client, context, successReceiptFromResult);
for (const event of progressEvents) await progress.insert(client, context, event);
if (successReceipt) successReceipt = buildCustomerSuccessReceipt({ ...successReceipt,
team_progress: await progress.snapshot(client, context, progressEvents, progressInput.recordedAt,
groupNumber => ledger.history(client, context, groupNumber))
}, resultObject, operation);
}
const updated = await client.query(
`UPDATE tasks
SET execution_result_ciphertext = $1, execution_result = $2,
@@ -7300,6 +7318,9 @@ export class TaskService {
const ledger = new ArrangementLedger(this.config);
await ledger.prepare(client, context, successReceiptFromResult);
await ledger.markTasksDeleted(client, context, normalizedTaskIds);
const progress = new TeamProgressLedger(this.config);
await progress.prepare(client, context, successReceiptFromResult);
await progress.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
+291
View File
@@ -0,0 +1,291 @@
import type pg from 'pg';
import type { AppConfig } from './config.js';
import { decryptText, encryptText, sha256Text } from './crypto.js';
import { ARRANGEMENT_LABELS, arrangementGroupNumber, currentArrangements,
type ArrangementEvent, type ArrangementScope } from './arrangement-ledger.js';
type Json = Record<string, unknown>;
type Client = Pick<pg.PoolClient, 'query'>;
const obj = (v: unknown): Json => v && typeof v === 'object' && !Array.isArray(v) ? v as Json : {};
const txt = (v: unknown): string => typeof v === 'string' || typeof v === 'number' ? String(v).trim() : '';
const first = (...v: unknown[]): string => v.map(txt).find(Boolean) || '';
const list = (v: unknown): Json[] => Array.isArray(v) ? v.map(obj) : [];
const strings = (...values: unknown[]): string[] => [...new Set(values.flatMap(v => Array.isArray(v) ? v.map(txt) : [txt(v)]).filter(Boolean))];
const date = (v: unknown): string => { const d = v instanceof Date ? v : new Date(txt(v)); return Number.isFinite(d.getTime()) ? d.toISOString() : ''; };
const group = arrangementGroupNumber;
// Lifecycle receipts may fall back to internal tid/ddid values. Those are not
// group numbers and must never make unrelated business records look connected.
const businessGroup = (v: unknown): string => {
const value = group(v);
return /^\d+$/.test(value) || /^D\d+$/.test(value) ? '' : value;
};
const ACTIONS = new Set(['team_order_create', 'team_order_batch_create', 'shared_plan_create', 'shared_child_order_create',
'shared_child_order_batch_create', 'passenger_list_import', 'order_update_independent', 'order_update_shared_plan',
'order_update_shared_child', 'order_cancel', 'order_restore', 'confirmation_export', ...Object.keys(ARRANGEMENT_LABELS)]);
const FILES: Record<string, string> = { 'xingyou-confirm': '星游客户确认单', 'liantai-confirm': '联泰客户确认单',
'job-order': '团队申请书', 'visitor-list': '游客名单', 'guide-confirm': '导游通知书', 'hotel-preorder': '酒店预订单',
'transport-preorder': '团体票申请书', 'filing-current': '备案确认书', 'filing-history': '历史备案确认书', 'pickup-sign': '接机牌' };
const UPDATE_LABELS: Record<string, string> = { 'rooms.SGL': '单间', 'rooms.TWN': '标间', 'rooms.DBL': '大床房', 'rooms.TRP': '三人间',
twin_room_count: '标间', planned_capacity: '计划收客数', lodging_note: '订房说明', 'pax.adult': '成人',
'pax.child_bed': '占床儿童', 'pax.child_no_bed': '不占床儿童', 'pax.leader': '领队人数', departure_date: '出发日期', remark: '备注' };
export interface ProgressEvent {
version: 1; group_number: string; subject_kind: 'team' | 'child'; subject_number: string;
action: string; source_task_id: string; execution_id: string; effective_at: string; recorded_at: string;
source_task_deleted_at?: string | null; details: Json;
}
export interface ProgressSnapshot {
version: 1; scope: 'current_account'; recorded_at: string;
groups: Array<{ group_number: string; events: ProgressEvent[]; arrangements: Array<{ action: string; details: Json }> }>;
unlinked_subjects: string[];
}
export interface ProgressInput {
taskId: string; executionId: string; effectiveAt: unknown; recordedAt: unknown;
operation: unknown; result: unknown; receipt: unknown;
}
// This is a projection of accepted success evidence, never a copy of operations,
// passenger TSV, attachment contents, URLs or arbitrary plugin metadata.
export function progressEventsFromSuccess(input: ProgressInput): ProgressEvent[] {
const op = obj(input.operation), data = obj(op.data), result = obj(input.result), report = obj(result.report), receipt = obj(input.receipt);
const action = txt(op.action), refs = obj(report.resolved_refs || result.resolved_refs);
if (!ACTIONS.has(action) || result.status !== 'completed' || !Object.keys(receipt).length || receipt.success === false
|| result.uncertain === true || result.manual_review_required === true || report.manual_review_required === true
|| [result.blockers, report.blockers].some(v => Array.isArray(v) && v.length)) return [];
const recordedAt = date(input.recordedAt), effectiveAt = date(input.effectiveAt) || recordedAt;
if (!recordedAt || !input.taskId || !input.executionId) return [];
const context = obj(result.business_reply_context || report.business_reply_context || report.business_receipt);
const details: Json = {};
// Search keywords are displayed as input conditions, never as verified names.
for (const [key, value] of Object.entries({ customer: context.customer_name, product: context.product_name,
customer_condition: obj(data.customer).name || obj(data.customer).keyword,
product_condition: obj(data.product).name || obj(data.product).keyword })) if (txt(value)) details[key] = txt(value);
if (Array.isArray(data.departure_dates) && data.departure_dates.length === 1) details.departure_date = txt(data.departure_dates[0]);
if (action === 'passenger_list_import') {
const n = context.imported_count ?? context.row_count ?? result.imported_count ?? obj(data.passenger_list).row_count;
if (n !== undefined && n !== null && txt(n) !== '' && Number.isInteger(Number(n)) && Number(n) >= 0) details.roster_count = Number(n);
}
if (action.startsWith('order_update_')) details.updates = list(obj(data.updates).actions)
.filter(u => Object.hasOwn(UPDATE_LABELS, txt(u.target)))
.map(u => ({ label: UPDATE_LABELS[txt(u.target)], value: txt(u.value), mode: txt(u.operation) }));
if (action === 'order_cancel' || action === 'order_restore') details.order_status = action === 'order_cancel'
? '已取消' : first(context.to_status, obj(data.transition).to_status, report.to_status) || '已恢复';
if (action === 'confirmation_export') details.files = [...new Set(list(report.artifacts)
.filter(a => a.agentbus_visible !== false).map(a => txt(a.type)).filter(t => Object.hasOwn(FILES, t)))];
if (refs.kind === 'shared_plan' || action === 'shared_plan_create' || action === 'order_update_shared_plan') details.team_kind = 'shared_plan';
const targets: Array<{ group: string; child: string; departure?: string; createdChild?: boolean }> = [];
const batches = list(result.batch_results).length ? list(result.batch_results) : list(report.batch_results);
if (action === 'shared_child_order_batch_create') {
for (const b of batches) if (b.status === 'completed' && !list(b.blockers).length && txt(b.child_order_no))
targets.push({ group: businessGroup(b.parent_group_no), child: group(b.child_order_no), departure: txt(b.departure_date) });
} else if (action === 'shared_child_order_create' || refs.kind === 'shared_child_order') {
const children = [...new Set(strings(receipt.order_number, refs.child_order_no,
refs.kind === 'shared_child_order' ? refs.identifier : '', result.order_number).map(group).filter(Boolean))];
const parents = [...new Set(strings(refs.parent_group_no, receipt.parent_group_no,
obj(report.erp_receipt).parent_group_no).map(businessGroup).filter(Boolean))];
if (children.length > 1 || parents.length > 1) return [];
if (children.length === 1) targets.push({ group: parents[0] || '', child: children[0] });
} else {
const groups = strings(receipt.group_numbers, receipt.group_number, refs.identifier, refs.group_no).map(businessGroup);
const unique = [...new Set(groups.filter(Boolean))];
// A single operation resolving to conflicting group identities is not merged.
const batch = ['team_order_batch_create', 'shared_plan_create'].includes(action);
if (!batch && unique.length > 1) return [];
for (const g of unique) {
const matching = batches.filter(b => strings(b.group_number, obj(b.erp_receipt).group_number).map(group).includes(g));
targets.push({ group: g, child: '', ...(matching.length === 1 ? { departure: first(matching[0].date, matching[0].departure_date) } : {}) });
}
// Split-plan creation can prove concrete child identities per parent/date.
for (const row of list(obj(obj(report.erp_receipt || result.erp_receipt).split_order_probe).per_date)) {
if (row.facts_determined !== true || !unique.includes(group(row.parent_group_no))) continue;
for (const child of list(row.child_refs)) if (txt(child.child_order_no)) targets.push({ group: group(row.parent_group_no),
child: group(child.child_order_no), departure: txt(row.departure_date), createdChild: true });
}
const requery = obj(report.requery || result.requery);
if (action === 'order_cancel' && refs.kind === 'shared_plan' && unique.length === 1 && requery.child_status_propagation_matched === true) {
for (const child of list(requery.child_refs)) if (/^\d+$/.test(txt(child.ddid))) targets.push({ group: unique[0], child: `D${txt(child.ddid)}` });
}
}
const seen = new Set<string>();
return targets.filter(t => { const key = `${t.group}/${t.child}`; if (seen.has(key)) return false; seen.add(key); return true; })
.map(t => ({ version: 1, group_number: t.group, subject_kind: t.child ? 'child' : 'team', subject_number: t.child || t.group,
action: t.createdChild ? 'shared_child_order_create' : action, source_task_id: input.taskId, execution_id: input.executionId,
effective_at: effectiveAt, recorded_at: recordedAt, details: { ...details, ...(t.departure ? { departure_date: t.departure } : {}) } }));
}
// An interrupted batch is not a successful task. Preserve only independently
// confirmed rows from the existing adapter's per-item completion contract.
export function partialBatchProgressEvents(input: ProgressInput): ProgressEvent[] {
const op = obj(input.operation), result = obj(input.result), report = obj(result.report);
if (result.status === 'completed' || !['blocked', 'failed', 'reconciliation_pending'].includes(txt(result.status))) return [];
const batches = list(result.batch_results).length ? list(result.batch_results) : list(report.batch_results);
if (op.action === 'shared_child_order_batch_create' && report.status === 'split_child_batch_stopped_after_write') {
const verified = batches.filter(b => b.status === 'completed' && b.adapter_status === 'split_child_completed'
&& txt(b.parent_group_no) && txt(b.child_order_no) && (!Array.isArray(b.blockers) || !b.blockers.length));
return progressEventsFromSuccess({ ...input, receipt: { success: true }, result: { status: 'completed', batch_results: verified } });
}
if (op.action === 'team_order_batch_create' && report.status === 'batch_fallback_incomplete') {
return batches.flatMap(b => {
const receipt = obj(b.erp_receipt);
if (b.status !== 'completed' || receipt.success === false || !txt(receipt.group_number)) return [];
return progressEventsFromSuccess({ ...input, operation: { ...op,
data: { ...obj(op.data), departure_dates: txt(b.date) ? [txt(b.date)] : [] } },
result: { status: 'completed' }, receipt });
});
}
return [];
}
function ordered(events: ProgressEvent[]): ProgressEvent[] {
return [...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));
}
function latest(events: ProgressEvent[], match: (e: ProgressEvent) => boolean): ProgressEvent | undefined { return events.filter(match).at(-1); }
function displayTime(value: string): string {
return new Intl.DateTimeFormat('sv-SE', { timeZone: 'Asia/Shanghai', year: 'numeric', month: '2-digit', day: '2-digit',
hour: '2-digit', minute: '2-digit' }).format(new Date(value));
}
function arrangementLine(action: string, details: Json): string {
const fields = [txt(details.resource_name), txt(details.item),
details.start_date ? `${txt(details.start_date)} 至 ${txt(details.end_date) || '离店/结束日未记录'}` : txt(details.date),
details.quantity === undefined ? '' : `${txt(details.quantity)}${action === 'arrangement_hotel' ? '间' : '(数量)'}`,
details.phone ? `${action === 'arrangement_vehicle' ? '司机电话' : '电话'}:${txt(details.phone)}` : '',
details.vehicle_number ? `车号:${txt(details.vehicle_number)}` : '', details.grade ? `等级:${txt(details.grade)}` : '',
details.filing_number ? `备案号:${txt(details.filing_number)}` : '', details.filing_date ? `备案日期:${txt(details.filing_date)}` : '',
details.filing_entry_port ? `入境口岸:${txt(details.filing_entry_port)}` : '', details.filing_exit_port ? `出境口岸:${txt(details.filing_exit_port)}` : '',
details.remark ? `备注:${txt(details.remark)}` : ''].filter(Boolean);
if (!details.resource_name) fields.unshift('资源名称未记录');
return fields.length > 1 ? fields.join('|') : '已安排,明细未记录';
}
export function teamProgressText(value: unknown): string {
const snap = value as ProgressSnapshot;
if (!snap || snap.version !== 1 || snap.scope !== 'current_account' || !Array.isArray(snap.groups)) return '';
const lines = ['本系统团队办理进度(当前账号)'];
for (const entry of snap.groups) {
const events = ordered(entry.events), team = events.filter(e => e.subject_kind === 'team');
lines.push(`\n团号:${entry.group_number}`);
const facts = Object.assign({}, ...team.map(e => e.details)) as Json;
for (const [key, label] of [['customer', '客户'], ['product', '产品'], ['departure_date', '出发日期']] as const)
if (txt(facts[key])) lines.push(`${label}:${txt(facts[key])}`);
if (!facts.customer && facts.customer_condition) lines.push(`客户条件:${txt(facts.customer_condition)}`);
if (!facts.product && facts.product_condition) lines.push(`产品条件:${txt(facts.product_condition)}`);
const status = latest(team, e => e.action === 'order_cancel' || e.action === 'order_restore'
|| ['team_order_create','team_order_batch_create','shared_plan_create'].includes(e.action));
lines.push(`订单状态:${status ? txt(status.details.order_status) || '已创建' : '暂无本系统状态记录'}`);
const subjects = [...new Set(events.filter(e => e.subject_kind === 'child').map(e => e.subject_number))].sort();
const shared = facts.team_kind === 'shared_plan' || subjects.length > 0;
const creation = latest(team, e => ['team_order_create','team_order_batch_create','shared_plan_create'].includes(e.action));
lines.push(`${shared ? '建计划' : '下单'}:${creation ? `已完成(${displayTime(creation.effective_at)})` : '暂无本系统成功记录'}`);
const subjectProgress = (subset: ProgressEvent[], label: string, parent = false): void => {
const roster = latest(subset, e => e.action === 'passenger_list_import');
if (!parent) lines.push(`${label}名单:${roster ? `已导入,最近一次${roster.details.roster_count === undefined ? '人数未记录' : `${roster.details.roster_count}人`}(${displayTime(roster.effective_at)})` : '暂无本系统导入记录'}`);
const update = latest(subset, e => e.action.startsWith('order_update_'));
lines.push(`${label}最近修改:${update ? `${list(update.details.updates).map(u => `${txt(u.label)}${u.mode === 'append' ? '追加' : '改为'}${txt(u.value)}`).join(';') || '修改成功,明细未记录'}(${displayTime(update.effective_at)})` : '暂无本系统修改记录'}`);
const exports = new Map<string, ProgressEvent>();
for (const e of subset) if (e.action === 'confirmation_export') for (const f of strings(e.details.files)) exports.set(f, e);
const unknownExport = latest(subset, e => e.action === 'confirmation_export' && !strings(e.details.files).length);
if (!exports.size && !unknownExport) lines.push(`${label}文件:暂无本系统导出记录`);
if (unknownExport) lines.push(`${label}文件:已导出,文件类型未记录(${displayTime(unknownExport.effective_at)})`);
for (const [file, e] of exports) {
const changed = events.some(change => change.action !== 'confirmation_export' && change.effective_at > e.effective_at
&& (change.subject_kind === 'team' || subset.includes(change) || (file === 'visitor-list' && shared)));
lines.push(`${label}${parent && file === 'visitor-list' ? '整团游客信息' : FILES[file] || '团队文件'}:已导出(${displayTime(e.effective_at)})${changed ? ';之后有业务变更,可能需要重新导出' : ''}`);
}
};
if (!shared) subjectProgress(team, '');
else {
lines.push(`本系统记录子单:${subjects.length}个(不代表母团全部子单)`);
// Parent modifications and whole-group exports belong to the parent.
subjectProgress(team, '母团', true);
for (const child of subjects) {
const subset = events.filter(e => e.subject_kind === 'child' && e.subject_number === child);
const created = subset.some(e => ['shared_child_order_create','shared_child_order_batch_create'].includes(e.action));
const state = latest(subset, e => ['order_cancel','order_restore','shared_child_order_create','shared_child_order_batch_create'].includes(e.action));
lines.push(`子单 ${child}:${created ? '已新增' : '暂无本系统新增记录'};状态:${state ? txt(state.details.order_status) || '已创建' : '未记录'}`);
subjectProgress(subset, ` ${child} `);
}
}
lines.push('【安排情况与明细】');
for (const [action, label] of Object.entries(ARRANGEMENT_LABELS)) {
const items = entry.arrangements.filter(a => a.action === action);
lines.push(`${label}:${items.length ? `已安排,共${items.length}条` : '未安排'}`);
items.forEach((a,i) => lines.push(` ${i+1}. ${arrangementLine(action, a.details)}`));
}
}
if (!snap.groups.length || snap.unlinked_subjects.length) lines.push(`\n${snap.unlinked_subjects.length ? `对象:${snap.unlinked_subjects.join('、')}。` : ''}暂无法关联团号,未推断安排状态。`);
lines.push('\n仅统计当前账号通过本系统成功完成的操作;未记录不代表 ERP 中没有。已安排表示至少一条,不代表全行程已安排齐全。');
return lines.join('\n');
}
export class TeamProgressLedger {
constructor(private readonly config: AppConfig) {}
async prepare(client: Client, scope: ArrangementScope, receiptFromResult: (v: unknown) => Json | null): Promise<void> {
if (!scope.organizationId || !scope.userId) throw new Error('team_progress_scope_missing');
await client.query('INSERT INTO team_progress_owners (organization_id, owner_user_id) VALUES ($1,$2) ON CONFLICT DO NOTHING', [scope.organizationId, scope.userId]);
const locked = await client.query('SELECT backfilled_at FROM team_progress_owners WHERE organization_id = $1 AND owner_user_id = $2 FOR UPDATE', [scope.organizationId, scope.userId]);
if (locked.rows[0]?.backfilled_at) return;
let cursor = '';
for (;;) {
const rows = 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'
AND a.id::text = COALESCE(t.execution_result->>'execution_id', '')
WHERE t.organization_id = $1 AND t.assigned_user_id = $2
AND t.status IN ('completed','blocked','failed','reconciliation_pending') AND t.task_id > $3
ORDER BY t.task_id LIMIT 200`, [scope.organizationId, scope.userId, cursor]);
for (const row of rows.rows) {
const decode = (cipher: unknown, plain: unknown): unknown => cipher ? JSON.parse(decryptText(this.config, txt(cipher))) : plain;
const result = decode(row.execution_result_ciphertext, row.execution_result);
const input = { 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: receiptFromResult(result) };
const events = [...progressEventsFromSuccess(input), ...partialBatchProgressEvents(input)];
for (const e of events) await this.insert(client, scope, e);
}
if (rows.rows.length < 200) break;
cursor = txt(rows.rows.at(-1)?.task_id);
}
await client.query('UPDATE team_progress_owners SET backfilled_at = now() WHERE organization_id = $1 AND owner_user_id = $2', [scope.organizationId, scope.userId]);
}
async insert(client: Client, scope: ArrangementScope, e: ProgressEvent): Promise<void> {
await client.query(`INSERT INTO team_progress_events
(organization_id, owner_user_id, group_key, subject_key, subject_kind, event_key, source_task_id, execution_id, effective_at, recorded_at, detail_ciphertext)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) ON CONFLICT (organization_id, owner_user_id, execution_id, event_key) DO NOTHING`,
[scope.organizationId, scope.userId, e.group_number ? sha256Text(e.group_number) : '', sha256Text(e.subject_number), e.subject_kind,
sha256Text(`${e.action}/${e.group_number}/${e.subject_kind}/${e.subject_number}`), e.source_task_id, e.execution_id,
e.effective_at, e.recorded_at, encryptText(this.config, JSON.stringify(e))]);
}
async knownGroups(client: Client, scope: ArrangementScope, child: string): Promise<string[]> {
const rows = await client.query(`SELECT detail_ciphertext FROM team_progress_events WHERE organization_id = $1 AND owner_user_id = $2
AND subject_key = $3 AND subject_kind = 'child' AND group_key <> ''`, [scope.organizationId, scope.userId, sha256Text(child)]);
return strings(...rows.rows.map(r => obj(JSON.parse(decryptText(this.config, r.detail_ciphertext))).group_number));
}
async history(client: Client, scope: ArrangementScope, groupNumber: string): Promise<ProgressEvent[]> {
const rows = await client.query(`SELECT e.detail_ciphertext, e.source_task_deleted_at FROM team_progress_events e
WHERE e.organization_id = $1 AND e.owner_user_id = $2 AND (e.group_key = $3 OR (e.group_key = '' AND e.subject_kind = 'child'
AND e.subject_key IN (SELECT subject_key FROM team_progress_events WHERE organization_id = $1 AND owner_user_id = $2 AND group_key = $3)
AND NOT EXISTS (SELECT 1 FROM team_progress_events conflict WHERE conflict.organization_id = $1 AND conflict.owner_user_id = $2
AND conflict.subject_key = e.subject_key AND conflict.group_key <> '' AND conflict.group_key <> $3)))
ORDER BY e.effective_at, e.recorded_at, e.id`, [scope.organizationId, scope.userId, sha256Text(group(groupNumber))]);
return rows.rows.map(r => ({ ...JSON.parse(decryptText(this.config, r.detail_ciphertext)), source_task_deleted_at: date(r.source_task_deleted_at) || null }));
}
async snapshot(client: Client, scope: ArrangementScope, events: ProgressEvent[], recordedAt: string,
arrangements: (g: string) => Promise<ArrangementEvent[]>): Promise<ProgressSnapshot> {
const groups = new Set<string>(), unlinked = new Set<string>();
for (const e of events) {
if (e.group_number) groups.add(e.group_number);
else {
const parents = await this.knownGroups(client, scope, e.subject_number);
if (parents.length === 1) groups.add(parents[0]); else unlinked.add(e.subject_number);
}
}
const snapshots: ProgressSnapshot['groups'] = [];
for (const g of groups) snapshots.push({ group_number: g, events: await this.history(client, scope, g),
arrangements: currentArrangements(await arrangements(g)).map(a => ({ action: a.action, details: a.details })) });
return { version: 1, scope: 'current_account', recorded_at: recordedAt, groups: snapshots, unlinked_subjects: [...unlinked] };
}
async markTasksDeleted(client: Client, scope: ArrangementScope, taskIds: string[]): Promise<void> {
await client.query(`UPDATE team_progress_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]);
}
}