导入问题修复

This commit is contained in:
andy committed 2026-09-15 11:34:59 +08:00
1 parent 0282c30a2a
commit ede1d47c47
26 files changed
+934 -45

No files matched your search

@@ -13,7 +13,7 @@ const ABSOLUTE_MAX_DATA_ROWS = 5_000;
const MAX_PASSENGER_SEQUENCE = 5_000;
const MAX_HEADER_ROW = 100;
export const PASSENGER_ROSTER_WORKBOOK_VERSION = 'ltjt-passenger-roster-workbook-v1.5.0';
export const PASSENGER_ROSTER_WORKBOOK_VERSION = 'ltjt-passenger-roster-workbook-v1.5.1';
const REQUIRED_SOURCE_FIELDS = [
'序号', '姓名', '英文姓名', '性别', '出生日期', '出生地', '护照号码',
@@ -36,12 +36,14 @@ const HEADER_FIELD_ENTRIES: ReadonlyArray<readonly [string, RequiredSourceField]
['性别', '性别'],
['出生日期', '出生日期'],
['出生地', '出生地'],
['出生地点', '出生地'],
['护照号码', '护照号码'],
['证件号码', '护照号码'],
['签发地', '签发地'],
['签发日期', '签发日期'],
['签发日', '签发日期'],
['有效期', '有效期'],
['有效期至', '有效期'],
['电话', '电话'],
['备注', '备注']
];
+81 -4
View File
@@ -1699,6 +1699,11 @@ function executionSummary(value: unknown): Record<string, unknown> {
no_erp_write: result.no_erp_write === true,
write_attempted: result.write_attempted === true,
execution_id: result.execution_id || '',
...(result.execution_phase === 'dispatch_uncertain' ? {
execution_phase: 'dispatch_uncertain',
reconciliation_required: result.reconciliation_required === true,
reconciliation_stale: result.reconciliation_stale === true
} : {}),
...(operationTiming ? { operation_timing: operationTiming } : {})
};
}
@@ -2351,6 +2356,58 @@ export function hasPrewriteNoErpEvidence(value: unknown): boolean {
&& ajaxRecords.every((record) => record.prevented_from_network === true || record.network_request_attempted === false);
}
/** Only a dispatch acknowledgement timeout may be closed by a late, proven
* read-only lookup failure. This is result ingestion, never permission to run
* the task again. The caller still enforces actor, task and execution access. */
export function canResolveDispatchNoWrite(input: {
taskStatus: string; attemptStatus: string; previous: unknown; incoming: unknown;
executionId: string; connectionId: string; claimedConnectionId: string;
}): boolean {
const previous = jsonObject(input.previous);
const result = jsonObject(input.incoming);
const report = jsonObject(result.report);
const preflight = jsonObject(report.preflight);
const timedOut = input.taskStatus === 'failed'
&& previous.reconciliation_required === true && previous.reconciliation_stale === true;
const sameExecution = Boolean(input.executionId)
&& text(previous.execution_id) === input.executionId && text(result.execution_id) === input.executionId;
const sameConnection = Boolean(input.connectionId)
&& input.connectionId === input.claimedConnectionId;
return (input.taskStatus === 'reconciliation_pending' || timedOut)
&& input.attemptStatus === 'reconciliation_pending' && sameExecution && sameConnection
&& text(previous.execution_phase) === 'dispatch_uncertain'
&& ['dispatch', 'reconciliation_timeout'].includes(text(previous.stage))
&& !previous.report && !previous.native_request && !previous.server_response
&& !previous.erp_receipt && !previous.success_receipt && !hasCompletionEvidence(previous)
&& text(result.status) === 'blocked' && text(result.stage) === 'erp_resolution'
&& text(result.failure_source) === 'plugin_executor'
&& result.no_erp_write === true && result.write_attempted === false && result.erp_write_started !== true
&& !/submit|write_started|erp_write/i.test(text(result.execution_phase))
&& text(report.status) === 'erp_resolution_blocked'
&& text(preflight.stage) === 'erp_readonly_resolution'
&& report.no_erp_write === true && report.write_attempted === false
&& preflight.no_erp_write === true && preflight.write_attempted === false
&& report.erp_write_started !== true && preflight.erp_write_started !== true
&& jsonObject(report.submit_safety).live_submit_attempted !== true
&& !report.native_request && !report.server_response && !report.requery
&& !result.native_request && !result.server_response && !result.requery
&& !result.erp_receipt && !result.success_receipt && !report.erp_receipt && !report.success_receipt
&& !hasCompletionEvidence(result)
&& Array.isArray(report.blockers) && report.blockers.length > 0
&& Array.isArray(report.side_effects)
&& report.side_effects.every((item) => ['no_procurement', 'no_payment', 'no_notification', 'no_external_send'].includes(text(item)))
&& (report.ajax_records == null || (Array.isArray(report.ajax_records) && report.ajax_records.length === 0));
}
export function shouldIgnoreDispatchTimeout(previousValue: unknown, incomingValue: unknown): boolean {
const previous = jsonObject(previousValue);
const incoming = jsonObject(incomingValue);
return text(incoming.execution_phase) === 'dispatch_uncertain' && text(incoming.stage) === 'dispatch'
&& Boolean(text(incoming.execution_id)) && text(previous.execution_id) === text(incoming.execution_id)
&& Boolean(text(previous.stage)) && text(previous.stage) !== 'dispatch'
&& ['accepted', 'running', 'completed', 'blocked', 'failed', 'cancelled', 'dry_run', 'reconciliation_pending'].includes(text(previous.status));
}
function publicShadowDifference(row: Record<string, unknown>, config: AppConfig): PublicParserShadowDifference {
const programCandidate = decryptedJson(config, row.program_result_ciphertext);
const aiCandidate = decryptedJson(config, row.ai_result_ciphertext);
@@ -6505,6 +6562,10 @@ export class TaskService {
const previousExecutionResult = jsonObject(
decryptedJson(this.config, row.execution_result_ciphertext) || row.execution_result
);
if (text(jsonObject(attemptRow.details).connection_id) === connectionId
&& shouldIgnoreDispatchTimeout(previousExecutionResult, preparedResult)) {
return { row, event: null };
}
let resultObject = preparedResult;
if (prepared.required && !prepared.failureCode) {
const stored: StoredTaskArtifact[] = [];
@@ -6536,6 +6597,23 @@ export class TaskService {
const stage = text(resultObject.stage || 'erp').slice(0, 160);
if (rawStatus !== status) resultObject.result_status = rawStatus;
resultObject.status = status;
const attemptDetails = jsonObject(attemptRow.details);
const canFinalizeDispatchNoWrite = canResolveDispatchNoWrite({
taskStatus: text(row.status), attemptStatus: text(attemptRow.status),
previous: previousExecutionResult, incoming: resultObject, executionId,
connectionId, claimedConnectionId: text(attemptDetails.connection_id)
});
const repeatsResolvedDispatchResult = text(row.status) === 'blocked'
&& previousExecutionResult.reconciliation_resolution === 'late_dispatch_readonly_failure'
&& text(previousExecutionResult.execution_id) === executionId
&& text(attemptDetails.connection_id) === connectionId;
// Restore server-derived metadata for an exact duplicate's hash. Any
// changed result still fails the immutable-state check below.
if (canFinalizeDispatchNoWrite || repeatsResolvedDispatchResult) {
resultObject.reconciliation_resolved = true;
resultObject.reconciliation_required = false;
resultObject.reconciliation_resolution = 'late_dispatch_readonly_failure';
}
const lifecycle = executionLifecycleFacts(resultObject, status, true);
Object.assign(resultObject, lifecycle);
const message = text(resultObject.message || (uncertain ? 'ERP 结果不确定,请勿重复提交同一任务,等待任务所属账号只读核验。' : '插件已返回状态。')).slice(0, 2_000);
@@ -6550,7 +6628,6 @@ export class TaskService {
if (text(attemptRow.response_hash) === resultHash) {
return { row, event: null };
}
const attemptDetails = jsonObject(attemptRow.details);
const sameConnection = !text(attemptDetails.connection_id) || text(attemptDetails.connection_id) === connectionId;
const failedReconciliation = text(row.status) === 'failed'
&& previousExecutionResult.reconciliation_required === true;
@@ -6569,7 +6646,7 @@ export class TaskService {
&& hasPrewriteNoErpEvidence(previousExecutionResult)
&& hasPrewriteNoErpEvidence(resultObject)
&& sameConnection;
const canFinalizeTerminalReconciliation = canFinalizeReconciliation || canFinalizeNoWriteReconciliation;
const canFinalizeTerminalReconciliation = canFinalizeReconciliation || canFinalizeNoWriteReconciliation || canFinalizeDispatchNoWrite;
if (executionStatusImmutable(text(row.status)) && !canFinalizeTerminalReconciliation) {
throw new TaskError('result_after_terminal', `任务已处于不可变状态 ${row.status},迟到结果已忽略。`);
}
@@ -6652,7 +6729,7 @@ export class TaskService {
execution_phase: text(resultObject.execution_phase)
}
}, context.userId);
await this.audit(client, context, canFinalizeNoWriteReconciliation
await this.audit(client, context, canFinalizeNoWriteReconciliation || canFinalizeDispatchNoWrite
? 'task.reconciliation_no_write_resolved'
: canFinalizeReconciliation
? 'task.reconciliation_resolved'
@@ -6662,7 +6739,7 @@ export class TaskService {
connection_id: connectionId,
terminal,
reconciliation_resolved: canFinalizeTerminalReconciliation,
no_erp_write_resolution: canFinalizeNoWriteReconciliation
no_erp_write_resolution: canFinalizeNoWriteReconciliation || canFinalizeDispatchNoWrite
});
return { row: updatedRow, event };
});
@@ -0,0 +1,189 @@
import test from 'node:test';
import assert from 'node:assert/strict';
import { TaskService, canResolveDispatchNoWrite, shouldIgnoreDispatchTimeout } from '../src/task-service.js';
import { loadConfig } from '../src/config.js';
import { getPool, closePool } from '../src/db.js';
import { decryptText } from '../src/crypto.js';
function fixture() {
return {
taskStatus: 'reconciliation_pending', attemptStatus: 'reconciliation_pending',
executionId: 'execution-a', connectionId: 'connection-a', claimedConnectionId: 'connection-a',
previous: {
execution_id: 'execution-a', stage: 'dispatch', execution_phase: 'dispatch_uncertain',
status: 'reconciliation_pending', write_attempted: true, erp_write_started: true, no_erp_write: false
} as Record<string, unknown>,
incoming: {
execution_id: 'execution-a', status: 'blocked', stage: 'erp_resolution', failure_source: 'plugin_executor',
no_erp_write: true, write_attempted: false,
report: {
status: 'erp_resolution_blocked', no_erp_write: true, write_attempted: false,
preflight: { stage: 'erp_readonly_resolution', no_erp_write: true, write_attempted: false },
native_request: null, server_response: null, requery: null,
blockers: ['erp_resolution_navigation_failed:native_list_search_not_ready'],
side_effects: ['no_procurement', 'no_payment', 'no_notification', 'no_external_send']
} as Record<string, unknown>
} as Record<string, unknown>
};
}
test('same execution read-only failure resolves provisional dispatch uncertainty', () => {
assert.equal(canResolveDispatchNoWrite(fixture()), true);
});
test('maintenance-expired dispatch uncertainty still accepts the original explicit no-write result', () => {
const f = fixture(); f.taskStatus = 'failed';
Object.assign(f.previous, { stage: 'reconciliation_timeout', reconciliation_required: true, reconciliation_stale: true });
assert.equal(canResolveDispatchNoWrite(f), true);
delete f.previous.reconciliation_stale;
assert.equal(canResolveDispatchNoWrite(f), false);
});
test('late no-write resolution enforces execution, original connection and immutable state boundaries', () => {
for (const [key, values] of Object.entries({
taskStatus: ['completed', 'cancelled', 'blocked', 'running', 'failed'],
attemptStatus: ['completed', 'cancelled', 'blocked', 'running', 'failed'],
executionId: ['', 'another'], connectionId: ['', 'another'], claimedConnectionId: ['', 'another']
})) for (const value of values) {
assert.equal(canResolveDispatchNoWrite({ ...fixture(), [key]: value }), false, `${key}=${value}`);
}
for (const key of ['incoming', 'previous'] as const) {
const f = fixture(); f[key].execution_id = 'another';
assert.equal(canResolveDispatchNoWrite(f), false, key);
}
});
test('true post-write uncertainty and prior native evidence cannot be overwritten', () => {
for (const patch of [
{ execution_phase: 'write_started' }, { execution_phase: '' }, { stage: 'lifecycle_live' },
{ report: { native_request: { url: 'fixture' } } }, { native_request: {} }, { server_response: {} },
{ erp_receipt: { group_numbers: ['TEST-001'] } }
]) {
const f = fixture(); Object.assign(f.previous, patch);
assert.equal(canResolveDispatchNoWrite(f), false, JSON.stringify(patch));
}
});
test('bare no-write flags, contradictory write evidence and non-lookup results are rejected', () => {
for (const patch of [
{ report: null }, { status: 'running' }, { status: 'failed' }, { stage: 'lifecycle_live' },
{ failure_source: 'program_parser' }, { no_erp_write: false }, { write_attempted: true },
{ write_attempted: undefined }, { erp_write_started: true }, { execution_phase: 'submit_started' },
{ native_request: {} }, { server_response: {} }, { requery: {} }, { erp_receipt: { group_numbers: ['TEST-001'] } }
]) {
const f = fixture(); Object.assign(f.incoming, patch);
assert.equal(canResolveDispatchNoWrite(f), false, JSON.stringify(patch));
}
});
test('page-level no-write evidence must be complete and internally consistent', () => {
for (const patch of [
{ status: 'live_submit_blocked' }, { no_erp_write: false }, { write_attempted: true },
{ native_request: {} }, { server_response: {} }, { requery: {} },
{ blockers: [] }, { blockers: null }, { side_effects: ['payment'] }, { side_effects: {} },
{ ajax_records: {} }, { ajax_records: [{ network_request_attempted: true }] },
{ preflight: {} }, { preflight: { stage: 'erp_readonly_resolution', no_erp_write: true, write_attempted: true } }
]) {
const f = fixture(); Object.assign(f.incoming.report as object, patch);
assert.equal(canResolveDispatchNoWrite(f), false, JSON.stringify(patch));
}
});
test('provisional timeout cannot replace already-received plugin progress or a final outcome', () => {
const incoming = fixture().previous;
for (const status of ['accepted', 'running', 'completed', 'blocked', 'failed', 'cancelled', 'dry_run', 'reconciliation_pending']) {
assert.equal(shouldIgnoreDispatchTimeout({ execution_id: 'execution-a', stage: 'erp_resolution', status }, incoming), true, status);
}
for (const previous of [{}, { execution_id: 'another', stage: 'erp_resolution', status: 'running' },
{ execution_id: 'execution-a', stage: 'dispatch', status: 'reconciliation_pending' }]) {
assert.equal(shouldIgnoreDispatchTimeout(previous, incoming), false);
}
assert.equal(shouldIgnoreDispatchTimeout({ execution_id: 'execution-a', stage: 'erp_resolution', status: 'running' }, { ...incoming, execution_phase: 'write_started' }), false);
});
test('real result-ingestion transaction closes only eligible dispatch failures and emits consistent facts', async t => {
const config = loadConfig({ NODE_ENV: 'test', FIELD_ENCRYPTION_KEY: Buffer.alloc(32, 27).toString('base64'),
DATABASE_URL: 'postgresql://fixture:fixture@127.0.0.1:1/fixture' });
const pool = getPool(config);
t.after(async () => { t.mock.restoreAll(); await closePool(); });
let row: Record<string, unknown>;
let attempt: Record<string, unknown>;
let events: Array<Record<string, unknown>>;
let queries: string[];
const client = { release() {}, async query(sql: string, params: unknown[] = []) {
queries.push(sql.trim().split(/\s+/).slice(0, 2).join(' '));
if (/SELECT \* FROM task_attempts/.test(sql)) return { rowCount: params[0] === 'execution-a' ? 1 : 0, rows: [attempt] };
if (/UPDATE tasks\s/.test(sql)) {
row.execution_result_ciphertext = params[0]; row.execution_result = params[1];
row.status = params[2]; row.stage = params[3]; row.message = params[4];
if (params[8]) row.lease_owner = null;
return { rowCount: 1, rows: [row] };
}
if (/UPDATE task_attempts/.test(sql)) {
attempt.status = params[0]; attempt.response_hash = params[1];
return { rowCount: 1, rows: [] };
}
if (!['BEGIN', 'COMMIT', 'ROLLBACK'].includes(sql)) throw new Error(`unexpected SQL ${sql}`);
return { rowCount: 0, rows: [] };
} };
t.mock.method(pool, 'connect', async () => client);
const service = new TaskService(config);
Object.assign(service, {
lockTaskForAccess: async () => row,
requireBrowserConnection: async () => {},
getTask: async () => row,
emitEvent: async (_client: unknown, _row: unknown, event: Record<string, unknown>) => { events.push(event); return event; },
audit: async () => {}, notify: () => {}
});
const context = { organizationId: 'org-a', userId: 'user-a', role: 'user' as const, requestId: 'fixture' };
for (const scenario of ['late', 'expired', 'wrong_connection', 'wrong_execution', 'wrong_actor', 'completed', 'postwrite', 'bare_flags', 'progress_first']) {
const f = fixture();
if (scenario === 'expired') {
f.taskStatus = 'failed';
Object.assign(f.previous, { stage: 'reconciliation_timeout', reconciliation_required: true, reconciliation_stale: true });
}
if (scenario === 'postwrite') Object.assign(f.previous, { stage: 'lifecycle_live', execution_phase: 'submitted' });
if (scenario === 'bare_flags') delete f.incoming.report;
if (scenario === 'progress_first') {
f.taskStatus = 'running';
Object.assign(f.previous, { stage: 'erp_resolution', status: 'running', execution_phase: '' });
}
row = { id: 'row-a', task_id: 'task-a', assigned_user_id: 'user-a', organization_id: 'org-a',
status: scenario === 'completed' ? 'completed' : f.taskStatus, execution_result: f.previous,
lease_owner: scenario === 'progress_first' ? 'browser:connection-a' : null };
attempt = { id: 'execution-a', status: f.attemptStatus, details: { connection_id: 'connection-a' } };
events = []; queries = [];
const operation = service.recordExecutionResult(
scenario === 'wrong_actor' ? { ...context, userId: 'user-b' } : context,
'task-a', scenario === 'progress_first' ? fixture().previous : f.incoming,
scenario === 'wrong_connection' ? 'connection-b' : 'connection-a',
scenario === 'wrong_execution' ? 'execution-b' : 'execution-a', { skipArtifactPersistence: true }
);
if (['late', 'expired'].includes(scenario)) {
await operation;
const stored = JSON.parse(decryptText(config, String(row.execution_result_ciphertext)));
assert.equal(row.status, 'blocked', scenario);
assert.equal(attempt.status, 'blocked');
assert.equal(stored.reconciliation_resolved, true);
assert.equal(stored.no_erp_write, true);
assert.equal(stored.erp_write_started, false);
assert.equal(stored.write_attempted, false);
assert.equal(events.length, 1);
assert.equal((events[0].payload as Record<string, unknown>).reconciliation_resolved, true);
assert.ok(queries.includes('COMMIT'));
await service.recordExecutionResult(context, 'task-a', structuredClone(fixture().incoming), 'connection-a', 'execution-a', { skipArtifactPersistence: true });
assert.equal(events.length, 1, 'a duplicate read must not emit another terminal event');
await assert.rejects(service.recordExecutionResult(context, 'task-a', { ...fixture().incoming, message: 'changed' }, 'connection-a', 'execution-a', { skipArtifactPersistence: true }), /不可变状态/);
} else if (scenario === 'progress_first') {
await operation;
assert.equal(row.status, 'running');
assert.equal(events.length, 0);
assert.ok(!queries.includes('UPDATE tasks'));
} else {
await assert.rejects(operation, Error, scenario);
assert.equal(events.length, 0);
assert.ok(!queries.includes('UPDATE tasks'), scenario);
assert.ok(queries.includes('ROLLBACK') || scenario === 'wrong_execution');
}
}
});
@@ -246,6 +246,39 @@ test('the approved attachment aliases and reverse issue-date formula are normali
].join('\n'));
});
test('birthplace and expiry aliases preserve the canonical passenger fields and values', async () => {
const normalize = (content: Buffer) => normalizePassengerRosterWorkbook({
content, fileName: 'synthetic.xlsx', contentType: 'application/octet-stream'
});
const original = await normalize(await syntheticWorkbook());
const aliased = await normalize(await syntheticWorkbook(workbook => {
const sheet = workbook.getWorksheet('Sheet1')!;
sheet.getCell('H2').value = '出生地点';
sheet.getCell('L2').value = '有效期至';
}));
assert.equal(aliased.headerRow, original.headerRow);
assert.equal(aliased.rowCount, original.rowCount);
assert.equal(aliased.canonicalTsv, original.canonicalTsv);
});
test('birthplace and expiry aliases cannot silently select duplicate semantic columns', async () => {
for (const [column, alias] of [['H2', '出生地点'], ['L2', '有效期至']] as const) {
for (const duplicate of ['canonical', 'alias'] as const) {
const content = await syntheticWorkbook(workbook => {
const sheet = workbook.getWorksheet('Sheet1')!;
const original = sheet.getCell(column).value;
sheet.getCell(column).value = alias;
sheet.getCell('O2').value = duplicate === 'canonical' ? original : alias;
sheet.getCell('O3').value = '合成冲突值';
});
await expectRosterError(
() => normalizePassengerRosterWorkbook({ content, fileName: 'synthetic.xlsx', contentType: 'application/octet-stream' }),
'roster_workbook_header_not_found'
);
}
}
});
test('source header order may vary, but every required label must appear exactly once', async () => {
const missingHeader = await syntheticWorkbook((workbook) => {
workbook.getWorksheet('Sheet1')!.getCell('N2').value = '备注说明';