From 976a4fa7b2bc445e327f412531c6ccc5e7ddc99a Mon Sep 17 00:00:00 2001 From: Wyndham ARR Date: Fri, 9 Oct 2026 00:12:06 +0800 Subject: [PATCH] Automatically resume empty saved source reviews --- .../tasks/20261008-production-review-9e7b.md | 13 + arr_web/DIRECT_DATA_ENTRY.md | 5 + arr_web/arr_data_executor.py | 20 +- arr_web/arr_data_review.py | 75 ++-- arr_web/arr_downloads.py | 30 ++ arr_web/local_ohip.py | 4 + arr_web/static/app.js | 8 +- deploy/OHIP_RELEASE_HANDOVER.md | 1 + tests/javascript/arr_data_review.cjs | 72 ++++ tests/javascript/helpers/arr_ui_harness.cjs | 7 +- tests/test_arr_empty_review_recovery.py | 379 ++++++++++++++++++ 11 files changed, 585 insertions(+), 29 deletions(-) create mode 100644 tests/test_arr_empty_review_recovery.py diff --git a/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md b/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md index b4154b2..1e6d43f 100644 --- a/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md +++ b/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md @@ -330,3 +330,16 @@ Read: memory-index, project-positioning, current-state latest September sections - Authenticated existing request04df0edb6743fbaa576d25063008c127 remains `needs_data_review`, job_id=null; source acquisition is complete106/106/100%. Review revision4 is editing, pending_count=0,total_count=0,can_finalize=true, with18 cancellations and11 PM exclusions. No user field input is still required. `DataFieldReviews.prepare` explicitly retains reanalyzed/status-evidence legacy reviews until finalized even when all gaps disappear; ordinary new complete input without those legacy interpretation checkpoints skips manual review. DirectARRExecutor returns needs_data_review until this retained review has a frozen payload. This is a legacy confirmation hold, not an additional field/rule failure or unfinished acquisition. - Outcome/boundaries: explain that corrected optional association interpretation and approved exclusions cleared the historical issues, but the old task deliberately retained its final continue/generate action. The empty panel currently misrepresents that action as manual checking. Its enabled “确认并生成日报” continues the same saved request/source without refetching; subsequent shared processing may still identify price review if required. No Oracle request, actual review/value/finalization, job generation, report/monthly mutation, UI/code change, restart or primary edit performed. No test run needed for read-only diagnosis. Documentation structure/ownership/whitespace gates apply. - Product follow-up/promotion candidate: integration-owned product flow should distinguish a legacy0-item ready-to-continue state from actual pending field/price issues, rather than showing an empty manual-check table. Current user question establishes confusion, not authorization to silently finalize an existing saved-source interpretation or auto-generate a report. Preserve original task/audit and use existing explicit continue control until the intended transition is changed. + + +## Same-task Follow-up: Automatically Continue Empty Legacy Reviews + +- User explicitly authorizes fixing the empty9/17 manual-check hold and carrying out the fix. This supersedes the preceding read-only question's no-finalization boundary for9/17 and the product behavior for0-item source reviews; it does not authorize invented prices, bypassing pending issues, changing shared classification/pricing or resetting published reports. Concurrent Task Gate Passed for resumed20261008-production-review-9e7b/feature/codex/owned checkout/codex/arr-production-review/base2417b1a, no peers, initially clean. Project Context Loaded: retained required memory/positioning/current-state/ADR004/006/007/architecture/domain/evidence/reflection/commitments/stale context, refreshed entry/planning/task records and inspected current source-review/executor/queue contracts. Canonical memory remains stale/integration-owned; the new direct user authorization governs this narrowly scoped continuation. Planning Gate Passed. +- Goal/plan: a validated0-item legacy source review must freeze its audited derived source and continue the same saved request automatically, rather than showing an empty human-check table. Preserve explicit final confirmation for genuine human decisions and keep all unresolved source/price items pending. Add a narrowly guarded automatic freeze/continuation path with system actor, durable retry/restart recovery and no mutation on ordinary GET reads; resume already stranded0-item tasks on safe service startup without a new source query. Regression-check no issues versus saved-but-confirmable manual decisions, source context/integrity failure, crash/retry and independent days. Safely back up/restart the local instance only with no active tasks; verify the actual9/17 progresses using saved data, retains audit/exclusions and stops at any genuine price review. No real price entry, source modification, reset/reprocess of completed10/7, new Oracle acquisition, schema change or primary-checkout edit. Update operational handoff, preserve before/after evidence and publish the completed fix to the already authorized repository. + +- Implementation: source-review recovery recomputes current shared rules under the existing revision lock, requires total_count=0/pending_count=0/can_finalize and collected source status, verifies exact request/context/checkpoint/original hash, then uses the existing durable freeze with actor `system:arr-empty-source-review`. Frozen audit preserves source interpretations, original observations and prior decisions, including historical decisions on now-excluded records. A real nonempty manual list still requires explicit final confirmation even when all entries are saved. Ordinary GET reads do not initiate a freeze; the worker checks stranded reviews on startup, resumes the same request with compare-and-set, isolates invalid reviews and reuses a completed freeze after interrupted queue publication. Missing/changed collection checkpoints cannot cause refetching. Ordinary complete input still skips source review. No business processing/price/schema/config identity changed. +- Frontend: a strictly matched0-item source-review state keeps polling the original request with GET only; processing, real price-review and terminal transitions clear the obsolete empty source panel. No client-side confirmation or value submission was added. Updated direct-entry and release handoff guides. +- Regression evidence:89 Python cases passed across new recovery, existing source/direct executor, queue, cancellation and PM suites, with disposable PostgreSQL enabled. The final checkpoint guard added1 case; all16 recovery cases were rerun and passed after it, for90 distinct backend cases (not105). Root separately passed13 existing static UI cases and all88 JavaScript cases. Coverage includes nonempty-but-fully-saved manual lists, unresolved/failed/ambiguous source, excluded historical decisions, integrity/context/checkpoint conflicts, interrupted freeze/queue update, startup without browser access, repeated restart and exactly one synthetic Finance version. Node syntax, Git whitespace, project-doc structure and ownership checks passed. No synthetic regression used the private real database or Oracle. +- Actual activation on2026-10-09 local time: queue had0 active tasks, service was stopped gracefully and PostgreSQL exit confirmed before a cold backup of the private instance to `ARR2.0/empty-review-service-backup-20261008-160608` (only runtime dependencies and prior nested backups excluded). Restarted the existing8875 LaunchAgent from this owned checkout. Same9/17 request04df0edb6743fbaa576d25063008c127 moved automatically from needs_data_review to needs_review, attempts1→2, with acquisition106/106 unchanged. Source revision4→5 is finalized, exactly one system finalization event; original source, prior events and decisions are unchanged. Shared processor4.4.0 created jobarrbatch-2651d152601cc7955e587459c5fc0e2b271f12b16e5d7707 and stopped at exactly1 real price item: HONGTAI/LBKB/Opera1000,2 records/2 rooms/4 room-nights. Reference candidate Opera900→1800 does not match1000 and was not applied. Actual price is still null, case0/1 open, no9/17 Finance version or report published. This supersedes the earlier9/17 explicit-continuation-only state in this task's historical sections. +- Preservation evidence: all3669 original source files byte/hash-identical; full9/15 source review and9/16 price review identical; all other5 queue rows/IDs unchanged; published10/7 trace identical apart from its per-read refreshed_at timestamp. No new Oracle acquisition or invented/changed price. Private before/after, cold-backup metadata and screenshot are in `production-validation-20261007/empty-review-auto-{before-20261008,after-20261009,backup}.json` and `empty-review-auto-local-20261009.png`, excluded from Git. Browser verified9/17 now opens the genuine price review, both other dates remain available, manual input blank and generation disabled; retained result tab12. Existing user tabs were not reloaded and no actual save/finalize button was clicked. +- Integration/promotion: promote automatic0-item continuation and the separation of completed acquisition versus pending price decision to canonical product/data-flow/current-state at the integration gate; canonical files remain untouched in this feature task. Keep the owned worktree/service running. Real9/17 completion remains intentionally pending the user's nightly processing price; separately deferred universal XML/API population alignment and production deployment are not claimed resolved. Push this narrow fix to the already authorized arr0918/main without rewriting history. diff --git a/arr_web/DIRECT_DATA_ENTRY.md b/arr_web/DIRECT_DATA_ENTRY.md index 7549c09..0419a9a 100644 --- a/arr_web/DIRECT_DATA_ENTRY.md +++ b/arr_web/DIRECT_DATA_ENTRY.md @@ -46,6 +46,11 @@ 后续从已保存来源继续,不重新获取订单、不改写原始捕获。冻结后不能再修改字段,重复确认或恢复不会重复提交日报。 正常完整输入无需人工步骤,直接进入原处理链路。 +旧任务若因已核实的来源解释或取消/PM排除而清空整张字段清单,系统会重新验证来源依据, +以系统身份保存处理记录并自动继续同一请求,不再要求对0项字段点击确认。服务重启会恢复这种已滞留的任务, +仍使用已保存的原始数据,不重新取数。清单仍有人工核对项的任务,即使已全部填写,也保留“确认并生成日报”; +后续若缺少处理价,仍进入正常价格核对,系统不代填价格。 + 字段完善后端接口归属于ARR自身,沿用网页登录权限与POST的CSRF保护;它们不向Oracle写入业务资料: | 操作 | ARR接口 | 正文 | diff --git a/arr_web/arr_data_executor.py b/arr_web/arr_data_executor.py index cf98d24..175f5e5 100644 --- a/arr_web/arr_data_executor.py +++ b/arr_web/arr_data_executor.py @@ -3,7 +3,7 @@ from datetime import date from pathlib import Path from arr_web.arr_download_handoff import outcome_from_handoff -from arr_web.arr_downloads import DownloadOutcome, validate_request_id +from arr_web.arr_downloads import DownloadOutcome, validate_report_date, validate_request_id from integrations.ohip.arr_data import ARRDataSource, VERSION as SOURCE_VERSION from integrations.ohip import processing_handoff as handoff from integrations.ohip.audit_arr_capture import protected_read @@ -34,6 +34,23 @@ class DirectARRExecutor: def finalize_data_review(self, request_id, revision, actor): return self.data_reviews.finalize(request_id, revision, actor) + def recover_data_review(self, request_id): + # Recovery may resume only an already collected, exactly bound local + # request. A missing/conflicting checkpoint must not start another fetch. + validate_request_id(request_id) + directory = self.root / "requests" / request_id + stored = strict_json(protected_read(directory / "request.json", 65536)) + report_date = validate_report_date(stored.get("report_date")) + identity = {**self.source_identity(), "request_id": request_id, "report_date": report_date} + require(fingerprint(stored) == fingerprint(identity), "data_executor_request_conflict") + checkpoint = strict_json(protected_read(directory / "collected.json", 65536)) + require(fingerprint(checkpoint["identity"]) == fingerprint(identity), "data_collection_checkpoint_conflict") + payload, manifest = self.data_reviews.original(request_id) + require(handoff._hash(payload) == checkpoint["data_sha256"] + and manifest == checkpoint["manifest_sha256"] + and strict_json(payload).get("report_date") == report_date, "data_collection_checkpoint_changed") + return self.data_reviews.recover_ready(request_id) + def source_identity(self): return {"version": "arr-direct-data-executor/v1", "source_version": SOURCE_VERSION, "hotel_id": self.source.hotel_id, "service_url": SERVICE, "application_id": APPLICATION, @@ -90,6 +107,7 @@ class DirectARRExecutor: "manifest_sha256": original_manifest}, replace=False) manifest_sha256 = original_manifest if review: + self.recover_data_review(request_id) completed = self.data_reviews.payload(request_id) if completed is None: return DownloadOutcome("needs_data_review") diff --git a/arr_web/arr_data_review.py b/arr_web/arr_data_review.py index 0aa52f5..2788963 100644 --- a/arr_web/arr_data_review.py +++ b/arr_web/arr_data_review.py @@ -31,6 +31,7 @@ VERSION = "arr-source-field-review/v1" REANALYSIS_POLICY = "oracle-optional-association/v1" REANALYSIS_VERSION = "arr-source-reanalysis/v1" STATUS_POLICY = "arr-exclude-cancelled/v1" +EMPTY_REVIEW_ACTOR = "system:arr-empty-source-review" LIMIT = 100 * 1024 * 1024 OPTIONAL = frozenset({"BLOCK_CODE", "RES_COMMENT", "PRODUCTS", "ROOM_CATEGORY_LABEL"}) LABELS = { @@ -381,9 +382,11 @@ class DataFieldReviews: require(fingerprint(actual) == fingerprint(meta), "data_review_prepare_conflict") review = self._public(directory, actual, original, state) # A complete ordinary source needs no human step or derived source. - # A reanalysis remains an explicit review until finalized, even if - # its validated source interpretation resolves every listed gap. - return review if review["total_count"] or "source_reanalysis" in state or "reservation_status_evidence" in state else None + # Validated overlays still need a frozen derived source. The + # executor may finish a genuinely empty review under a system actor. + return review if (review["total_count"] or state["decisions"] or state["status"] == "finalized" + or "source_reanalysis" in state or "reservation_status_evidence" in state + or (directory / "finalize-intent.json").exists()) else None def get(self, request_id): directory = self._directory(request_id) @@ -583,33 +586,57 @@ class DataFieldReviews: return json_bytes(data), binding def finalize(self, request_id, revision, actor): + directory = self._directory(request_id) + with self._lock(directory): + meta, original, state = self._read(directory) + return self._finalize_locked(directory, meta, original, state, revision, actor) + + def recover_ready(self, request_id): + """Freeze validated zero-item saved reviews, or replay a verified prior freeze. + + This internal recovery operation has no HTTP route. A completed manual + list is deliberately not empty: its explicit human confirmation remains + required. Read/derive/guard/freeze share the same revision lock. + """ directory = self._directory(request_id) with self._lock(directory): meta, original, state = self._read(directory) if state["status"] == "finalized": self._verify_frozen(directory, meta, original, state) - return self._public(directory, meta, original, state) - self._revision(state, revision) - actor = self._actor(actor) - if self._issues(directory, self._derive(directory, original, state)): - raise _error("INCOMPLETE", "请先完善并确认所有异常字段", 409) - intent = directory / "finalize-intent.json" - prior_sha256 = fingerprint(state) - if intent.exists(): - frozen = strict_json(protected_read(intent, LIMIT)) - require(frozen["prior_sha256"] == prior_sha256, "data_review_finalize_context_changed") - state = frozen["state"] - else: - state["revision"] += 1 - state["status"] = "finalized" - state["events"].append({"revision": state["revision"], "action": "finalize", "actor": actor, - "at": datetime.now(timezone.utc).isoformat()}) - atomic_json(intent, {"prior_sha256": prior_sha256, "state": state}, replace=False) - payload, binding = self._frozen(directory, meta, original, state) - _write_once(directory / "reviewed-source.json", payload) - state["derived_sha256"], state["binding_manifest_sha256"] = _hash(payload), binding - self._publish(directory, state) + return True + if original.get("status") not in {"collected", "collected_with_gaps"}: + return False + review = self._public(directory, meta, original, state) + if review["total_count"] != 0 or review["pending_count"] != 0 or not review["can_finalize"]: + return False + self._finalize_locked(directory, meta, original, state, state["revision"], EMPTY_REVIEW_ACTOR) + return True + + def _finalize_locked(self, directory, meta, original, state, revision, actor): + if state["status"] == "finalized": + self._verify_frozen(directory, meta, original, state) return self._public(directory, meta, original, state) + self._revision(state, revision) + actor = self._actor(actor) + if self._issues(directory, self._derive(directory, original, state)): + raise _error("INCOMPLETE", "请先完善并确认所有异常字段", 409) + intent = directory / "finalize-intent.json" + prior_sha256 = fingerprint(state) + if intent.exists(): + frozen = strict_json(protected_read(intent, LIMIT)) + require(frozen["prior_sha256"] == prior_sha256, "data_review_finalize_context_changed") + state = frozen["state"] + else: + state["revision"] += 1 + state["status"] = "finalized" + state["events"].append({"revision": state["revision"], "action": "finalize", "actor": actor, + "at": datetime.now(timezone.utc).isoformat()}) + atomic_json(intent, {"prior_sha256": prior_sha256, "state": state}, replace=False) + payload, binding = self._frozen(directory, meta, original, state) + _write_once(directory / "reviewed-source.json", payload) + state["derived_sha256"], state["binding_manifest_sha256"] = _hash(payload), binding + self._publish(directory, state) + return self._public(directory, meta, original, state) def _verify_frozen(self, directory, meta, original, state): payload, binding = self._frozen(directory, meta, original, state) diff --git a/arr_web/arr_downloads.py b/arr_web/arr_downloads.py index 47af62c..9f0941c 100644 --- a/arr_web/arr_downloads.py +++ b/arr_web/arr_downloads.py @@ -354,6 +354,7 @@ class PersistentARRDownloads: def _work(self) -> None: try: + self._recover_source_reviews() while True: with self._condition: if self._stopping: @@ -384,6 +385,35 @@ class PersistentARRDownloads: finally: os.close(self._lease) + def _recover_source_reviews(self): + """Recover saved review continuations at startup, without browser reads. + + Individual corrupt/conflicting reviews stay pending. A verified freeze + survives a crash before this queue update and is reused on next startup. + """ + recover = getattr(self._executor, "recover_data_review", None) + if not callable(recover): + return + with self._connect() as db: + rows = db.execute("""SELECT request_id FROM downloads WHERE status='needs_data_review' + ORDER BY created_at, request_id""").fetchall() + for row in rows: + with self._condition: + if self._stopping: + return + try: + if recover(row["request_id"]) is not True: + continue + with self._condition: + with self._connect() as db: + db.execute("""UPDATE downloads SET status='queued',retryable=0,updated_at=? + WHERE request_id=? AND status='needs_data_review'""", (_now(), row["request_id"])) + self._condition.notify_all() + except Exception: + # Recovery is local, optional maintenance, not evidence of a + # failed acquisition or permission to skip an unresolved review. + continue + def close(self, *, wait: bool = False) -> None: """Stop dispatch; owners must wait before closing executor dependencies.""" with self._condition: diff --git a/arr_web/local_ohip.py b/arr_web/local_ohip.py index bc2e83a..a475993 100644 --- a/arr_web/local_ohip.py +++ b/arr_web/local_ohip.py @@ -132,6 +132,10 @@ class AuthorizedExecutor: def finalize_data_review(self, *args): return self.executor.finalize_data_review(*args) + def recover_data_review(self, request_id): + recover = getattr(self.executor, "recover_data_review", None) + return recover(request_id) if callable(recover) else False + class LocalOHIPPortal(LocalReplayPortal): allow_xml_upload = True diff --git a/arr_web/static/app.js b/arr_web/static/app.js index 673cd2d..951d6e3 100644 --- a/arr_web/static/app.js +++ b/arr_web/static/app.js @@ -542,9 +542,15 @@ function scheduleARRDownloadPoll() { clearARRDownloadPoll(); if (document.hidden || state.arrDownloadBusy) return; + const review = state.arrDataReview; + const emptySourceReview = arrDownloadNeedsDataReview() && review?.status === "editing" + && state.arrDataReviewRequestId === state.arrDownloadTask.request_id && review.request_id === state.arrDataReviewRequestId + && review.report_date === state.arrDownloadTask.report_date && review.total_count === 0 && review.pending_count === 0 + && Array.isArray(review.items) && review.items.length === 0 && review.can_finalize === true + && !state.arrDataReviewLoading && !state.arrDataReviewDisconnected && !Object.keys(state.arrDataReviewDrafts).length; if (!state.arrDownloadReady || state.arrDownloadConfigDisconnected) { state.arrDownloadPollTimer = window.setTimeout(initARRDownload, 10000); - } else if (state.arrDownloadIntent && (arrDownloadActive() || state.arrDownloadDisconnected || !state.arrDownloadTask || state.arrDownloadTask.status === "needs_review" || state.arrDataReviewFinalizing)) { + } else if (state.arrDownloadIntent && (arrDownloadActive() || state.arrDownloadDisconnected || !state.arrDownloadTask || state.arrDownloadTask.status === "needs_review" || state.arrDataReviewFinalizing || emptySourceReview)) { state.arrDownloadPollTimer = window.setTimeout(loadARRDownloadTask, 4000); } } diff --git a/deploy/OHIP_RELEASE_HANDOVER.md b/deploy/OHIP_RELEASE_HANDOVER.md index 9f4a957..e2be383 100644 --- a/deploy/OHIP_RELEASE_HANDOVER.md +++ b/deploy/OHIP_RELEASE_HANDOVER.md @@ -16,6 +16,7 @@ - 直接处理接口数据,不生成中间XML;原人工XML上传入口保留。 - 日期选择、运行中预选、任务状态、重复提交保护和原任务恢复已实现;取数进度按预订笔数显示,最终完整性核对通过才显示100%。 - “需处理日期”保留不同日期的字段/价格待处理任务;统一在“人工核对”中完成。取数已完成不等于报表已生成。 +- 旧任务在已验证的来源解释/排除规则下清空整张字段清单时自动继续,重启可恢复原请求;清单仍有人工核对项或缺价时保留相应确认步骤,历史决定完整留痕。 - 缺少定价资料进入原有价格复核;影响处理的必需字段缺失或存在无法确定的候选时提示异常。 - 本机模拟已覆盖日报、复核、入库、月报、报表下载及断线恢复;用户确认日期交互验收通过。 diff --git a/tests/javascript/arr_data_review.cjs b/tests/javascript/arr_data_review.cjs index 90be084..639d31a 100644 --- a/tests/javascript/arr_data_review.cjs +++ b/tests/javascript/arr_data_review.cjs @@ -9,6 +9,7 @@ const task={request_id:requestId,report_date:'2026-10-07',status:'needs_data_rev const item=(extra={})=>({item_id:'1:BLOCK_CODE',source_sequence:1,confirmation_no:'SYNTHETIC-001',room_no:'SYNTHETIC-ROOM',company_name:'Synthetic Company',rate_code:'SYNTHETIC',field:'BLOCK_CODE',field_label:'团队代码',reason_code:'MISSING',can_be_empty:true,value:null,confirmed:false,...extra}); const review=(items,revision=1)=>({request_id:requestId,report_date:task.report_date,status:'editing',revision,items,pending_count:items.filter(x=>!x.confirmed).length,total_count:items.length,can_finalize:items.every(x=>x.confirmed)}); const plain=value=>JSON.parse(JSON.stringify(value)); +const activeTaskPoll=h=>h.timers.filter(timer=>!timer.cancelled&&timer.callback.name==='loadARRDownloadTask').at(-1); const priceTask={request_id:'c'.repeat(32),report_date:'2026-09-16',status:'needs_review', job_id:'fixture-price-september',can_retry:false}; const priceReview=(price=null)=>({case_id:'fixture-price-case',case_status:'open',business_date:priceTask.report_date,revision:1, @@ -27,6 +28,77 @@ function pendingAPI(url,currentPrice=priceTask) { throw new Error('unexpected request '+url); } +test('an empty validated field review polls the same task and follows processing then price review by GET only',async()=>{ + let nextStatus='processing'; + const jobId='fixture-recovered-empty-review'; + const h=harness(url=>{ + if(url.endsWith('/data-review')) return response(200,review([])); + if(url===`/api/arr-downloads/${requestId}`) return response(200,{...task,status:nextStatus,job_id:nextStatus==='processing'?null:jobId}); + if(url.startsWith('/api/jobs?')) return response(200,[]); + if(url.endsWith('/trace')) return response(200,{job:{job_id:jobId,business_date:task.report_date,status:nextStatus,active:true},logs:[]}); + if(url.includes('/review?')) return response(200,{...priceReview(),business_date:task.report_date}); + throw new Error('unexpected request '+url); + }); + await h.acceptARRDownloadTask(task,{sync:false}); + const emptyPoll=activeTaskPoll(h); + assert.equal(emptyPoll?.delay,4000); + await emptyPoll.callback(); + assert.equal(h.state.arrDownloadTask.status,'processing'); + assert.equal(h.state.arrDataReview,null); + assert.equal(h.element('#arr-data-review-panel').hidden,true); + nextStatus='needs_review'; + await activeTaskPoll(h).callback(); + assert.equal(h.state.arrDownloadTask.status,'needs_review'); + assert.equal(h.element('#arr-download-review').hidden,false); + assert.equal(h.element('#arr-data-review-panel').hidden,true); + assert.equal(h.state.arrDownloadPendingReviews[0].status,'needs_review'); + await h.openDailyPriceReview(jobId,false); + assert.equal(h.element('#daily-price-review-panel').hidden,false); + assert.equal(h.element('#manual-review-workspace').hidden,false); + assert.equal(h.calls.every(call=>call.method==='GET'),true); + assert.equal(h.calls.some(call=>call.url.includes('/finalize')),false); +}); + +for(const status of ['succeeded','failed']) { + test(`empty-review recovery follows the real ${status} task status without submitting a decision`,async()=>{ + const h=harness(url=>{ + if(url.endsWith('/data-review')) return response(200,review([])); + assert.equal(url,`/api/arr-downloads/${requestId}`); + return response(200,{...task,status,error_code:status==='failed'?'ARR_DATA_EXECUTION_FAILED':null}); + }); + await h.acceptARRDownloadTask(task,{sync:false}); + await activeTaskPoll(h).callback(); + assert.equal(h.state.arrDownloadTask.status,status); + assert.equal(h.element('#arr-data-review-panel').hidden,true); + assert.equal(h.element('#manual-review-workspace').hidden,true); + assert.equal(h.element('#arr-download-pending').hidden,true); + assert.match(h.element('#arr-download-status').textContent,new RegExp(`arr_download\\.${status}`)); + assert.equal(activeTaskPoll(h),undefined); + assert.equal(h.calls.every(call=>call.method==='GET'),true); + }); +} + +test('filled real review items, inconsistent zero counts and unsaved drafts never trigger zero-item recovery polling',async()=>{ + for(const snapshot of [review([item({confirmed:true,value:'SAVED'})]), + {...review([]),total_count:1}, {...review([]),pending_count:1}, {...review([]),can_finalize:false}, + {...review([]),items:[item()]}, {...review([]),request_id:'f'.repeat(32)}, + {...review([]),report_date:'2026-09-17'}]) { + const h=harness(()=>response(200,snapshot)); + await h.acceptARRDownloadTask(task,{sync:false}); + assert.equal(activeTaskPoll(h),undefined); + assert.equal(h.calls.length,1); + assert.equal(h.calls.every(call=>call.method==='GET'),true); + } + const h=harness(()=>response(200,review([item({confirmed:true,value:'SAVED'})]))); + await h.acceptARRDownloadTask(task,{sync:false}); + h.trackARRDataReviewDraft(h.row('1:BLOCK_CODE','UNSAVED').input); + h.state.arrDataReview=review([]); + await h.acceptARRDownloadTask(task,{sync:false}); + assert.equal(activeTaskPoll(h),undefined); + assert.equal(h.state.arrDataReviewDrafts['1:BLOCK_CODE'],'UNSAVED'); + assert.equal(h.calls.some(call=>call.url.includes('/finalize')),false); +}); + test('cross-month field and price dates restore, stay clickable after polling, and open their own pages by GET',async()=>{ const h=harness(url=>url==='/api/arr-downloads' ? response(200,{context_id:'production',ready:true,default_date:priceTask.report_date, diff --git a/tests/javascript/helpers/arr_ui_harness.cjs b/tests/javascript/helpers/arr_ui_harness.cjs index 17dd25c..852fe81 100644 --- a/tests/javascript/helpers/arr_ui_harness.cjs +++ b/tests/javascript/helpers/arr_ui_harness.cjs @@ -9,7 +9,7 @@ const boot = ' boot();\n})();'; assert(source.endsWith(boot+'\n') || source.endsWith(boot)); function harness(respond) { - const elements=new Map(), storage=new Map(), storageReads=[], calls=[]; + const elements=new Map(), storage=new Map(), storageReads=[], calls=[], timers=[]; const element=selector => { if (!elements.has(selector)) { const attributes=new Map(),classes=new Set(); @@ -27,7 +27,8 @@ function harness(respond) { const context=vm.createContext({Headers, console, Date, Intl, Uint8Array, document:{hidden:false,querySelector:element}, window:{ARRI18n:{t:(key,values)=>key==='manual_review.context'?`${values.date} · ${values.type}`:key,text:value=>value,errorMessage:code=>code,formatInteger:value=>String(value),formatMonth:value=>value},crypto:webcrypto, - setTimeout(){return 1;},clearTimeout(){},location:{replace(){}}}, + setTimeout(callback,delay){const timer={id:timers.length+1,callback,delay,cancelled:false};timers.push(timer);return timer.id;}, + clearTimeout(id){const timer=timers.find(entry=>entry.id===id);if(timer) timer.cancelled=true;},location:{replace(){}}}, localStorage:{setItem:(k,v)=>storage.set(k,v),getItem:k=>{storageReads.push(k);return storage.get(k);},removeItem:k=>storage.delete(k)}, fetch:async (url, options) => {calls.push({url,method:options.method||'GET',body:options.body,csrf:options.headers.get('X-ARR-CSRF')});return respond(url,options,calls);}, }); @@ -41,7 +42,7 @@ function harness(respond) { const node={dataset:{arrDataReviewItemId:itemId},querySelector:()=>input}; return {input,button:{closest:()=>node}}; }; - return {...subject,calls,storage,storageReads,element,row,i18n:context.window.ARRI18n,submit:()=>subject.submitARRDownload({preventDefault(){}})}; + return {...subject,calls,storage,storageReads,timers,element,row,i18n:context.window.ARRI18n,submit:()=>subject.submitARRDownload({preventDefault(){}})}; } const response=(status,data,code) => ({status,ok:status<400,json:async()=>code?{ok:false,error:{code}}:{ok:true,data}}); const scopedKey=(context='production',username='operator')=>`arr:last-download:v2:${encodeURIComponent(context)}:${encodeURIComponent(username)}`; diff --git a/tests/test_arr_empty_review_recovery.py b/tests/test_arr_empty_review_recovery.py new file mode 100644 index 0000000..eff8b65 --- /dev/null +++ b/tests/test_arr_empty_review_recovery.py @@ -0,0 +1,379 @@ +"""Saved zero-item source reviews continue locally without a human confirmation.""" +import hashlib +import json +from pathlib import Path +import sqlite3 +import tempfile +import threading +import unittest +from unittest.mock import Mock, patch + +from arr_processing.policy import load_processor_policy +from arr_web import arr_data_review as review_module +from arr_web.arr_data_review import DataFieldReviews, EMPTY_REVIEW_ACTOR +from arr_web.arr_downloads import DownloadOutcome, PersistentARRDownloads +from integrations.ohip.collect_arr_source import CollectionError +from tests import test_arr_data_review as review_fixtures +from tests.test_arr_data_review import ( + ACTOR, CONTEXT, DAY, PROJECT, REANALYSIS_MANIFEST, REANALYSIS_POLICY, + REQUEST, SOURCE_MANIFEST, gap, optional_omission_source, + raw, source_document, +) + + +def omission_source(): + document = source_document() + for field, reason in (("BLOCK_CODE", "missing_reservation_block"), + ("PRODUCTS", "missing_reservation_packages")): + gap(document, field) + document["records"][0]["fields"][field]["reason"] = reason + return document + + +class EmptyReviewRecoveryTests(unittest.TestCase): + @classmethod + def setUpClass(cls): + cls.policy = load_processor_policy(PROJECT) + + def setUp(self): + temporary = tempfile.TemporaryDirectory() + self.addCleanup(temporary.cleanup) + self.root = Path(temporary.name) / "reviews" + self.service = DataFieldReviews(root=self.root, policy=self.policy, context=CONTEXT) + self.document = omission_source() + + def prepare(self): + return self.service.prepare(REQUEST, raw(self.document), SOURCE_MANIFEST, DAY) + + def reanalyze(self, review): + return self.service.apply_source_reanalysis(REQUEST, raw(optional_omission_source(self.document)), + REANALYSIS_MANIFEST, REANALYSIS_POLICY, expected_revision=review["revision"]) + + def test_zero_item_reanalysis_freezes_under_system_actor_without_original_or_decision_changes(self): + review = self.reanalyze(self.prepare()) + original = self.service.original(REQUEST) + self.assertEqual((review["total_count"], review["pending_count"]), (0, 0)) + self.assertTrue(self.service.recover_ready(REQUEST)) + payload = self.service.payload(REQUEST) + audit = json.loads(payload[0])["manual_data_review"]["manifest"] + self.assertEqual(audit["changes"], []) + self.assertEqual([event["action"] for event in audit["events"]], ["source_reanalysis", "finalize"]) + self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR) + self.assertTrue(self.service.recover_ready(REQUEST)) + self.assertEqual(self.service.payload(REQUEST), payload) + self.assertEqual(self.service.original(REQUEST), original) + + def test_fully_filled_manual_items_still_require_explicit_confirmation(self): + review = self.prepare() + saved = self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "STAFF-VERIFIED", ACTOR) + review = self.reanalyze(saved) + self.assertEqual((review["total_count"], review["pending_count"]), (1, 0)) + self.assertFalse(self.service.recover_ready(REQUEST)) + self.assertIsNone(self.service.payload(REQUEST)) + self.assertEqual(self.service.get(REQUEST), review) + self.service.finalize(REQUEST, review["revision"], ACTOR) + self.assertTrue(self.service.recover_ready(REQUEST), "an explicit saved freeze may resume") + audit = json.loads(self.service.payload(REQUEST)[0])["manual_data_review"]["manifest"] + self.assertEqual(audit["events"][-1]["actor"], ACTOR) + + def test_pending_failed_or_ambiguous_observations_do_not_auto_continue(self): + for state in ("missing", "failed", "ambiguous"): + with self.subTest(state=state): + document = omission_source() + gap(document, "DISP_ROOM_NO", state=state) + key = {"missing": "b", "failed": "c", "ambiguous": "d"}[state] * 32 + review = self.service.prepare(key, raw(document), SOURCE_MANIFEST, DAY) + review = self.service.apply_source_reanalysis(key, raw(optional_omission_source(document)), + REANALYSIS_MANIFEST, REANALYSIS_POLICY, expected_revision=review["revision"]) + self.assertGreater(review["pending_count"], 0) + self.assertFalse(self.service.recover_ready(key)) + self.assertIsNone(self.service.payload(key)) + + def test_complete_ordinary_prepare_still_needs_no_review_or_derived_source(self): + self.document = source_document() + self.assertIsNone(self.prepare()) + self.assertIsNone(self.service.payload(REQUEST)) + + def test_cancelled_and_pm_hidden_history_decisions_freeze_without_overlay_and_remain_audited(self): + for kind in ("cancelled", "pm"): + with self.subTest(kind=kind): + document = source_document(2) + gap(document, "BLOCK_CODE") + if kind == "cancelled": + document["records"][0]["reservation_status"] = "Cancelled" + else: + document["records"][0]["fields"]["ROOM_CATEGORY_LABEL"] = {"state": "available", "value": "PM"} + key = ("b" if kind == "cancelled" else "c") * 32 + old_rule = "is_cancelled_record" if kind == "cancelled" else "is_pm_record" + # Represent an audited legacy decision saved before the approved + # exclusion; the original source and rule identity stay pinned. + with patch.object(self.service.rules, old_rule, return_value=False): + review = self.service.prepare(key, raw(document), SOURCE_MANIFEST, DAY) + self.service.update(key, "1:BLOCK_CODE", review["revision"], "PRESERVED-" + kind, ACTOR) + original = self.service.original(key) + review = self.service.prepare(key, raw(document), SOURCE_MANIFEST, DAY) + self.assertEqual((review["total_count"], review["pending_count"]), (0, 0)) + self.assertTrue(self.service.recover_ready(key)) + audit = json.loads(self.service.payload(key)[0])["manual_data_review"]["manifest"] + self.assertEqual(audit["changes"][0]["value"], "PRESERVED-" + kind) + self.assertEqual(audit["changes"][0]["actor"], ACTOR) + self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR) + self.assertNotIn("source_reanalysis", audit) + self.assertNotIn("reservation_status_evidence", audit) + self.assertEqual(self.service.original(key), original) + + def test_failed_original_source_is_not_automatically_frozen(self): + self.document["status"] = "failed" + review = self.reanalyze(self.prepare()) + self.assertEqual(review["total_count"], 0) + self.assertFalse(self.service.recover_ready(REQUEST)) + self.assertIsNone(self.service.payload(REQUEST)) + + def test_invalid_receipt_blocks_freeze_without_state_change(self): + review = self.reanalyze(self.prepare()) + pointer = (self.root / REQUEST / "state.json").read_bytes() + reference = next((self.root / REQUEST).glob("source-reanalysis-*.json")) + reference.write_bytes(reference.read_bytes() + b" ") + with self.assertRaises(CollectionError): + self.service.recover_ready(REQUEST) + self.assertEqual((self.root / REQUEST / "state.json").read_bytes(), pointer) + self.assertFalse((self.root / REQUEST / "finalize-intent.json").exists()) + + def test_original_hash_or_context_conflict_never_auto_freezes(self): + self.reanalyze(self.prepare()) + path = self.root / REQUEST / "original.json" + original = path.read_bytes() + path.write_bytes(original + b" ") + with self.assertRaises(CollectionError): + self.service.recover_ready(REQUEST) + path.write_bytes(original) + changed = DataFieldReviews(root=self.root, policy=self.policy, context={**CONTEXT, "hotel_id": "OTHER"}) + with self.assertRaises(CollectionError): + changed.recover_ready(REQUEST) + self.assertFalse((self.root / REQUEST / "finalize-intent.json").exists()) + + def test_interrupted_freeze_replays_same_bytes_and_one_system_event(self): + review = self.reanalyze(self.prepare()) + publish = review_module.atomic_json + def interrupt(path, document, **kwargs): + if path.name == "state.json": + raise OSError("synthetic freeze pointer interruption") + return publish(path, document, **kwargs) + with patch.object(review_module, "atomic_json", side_effect=interrupt), self.assertRaises(OSError): + self.service.recover_ready(REQUEST) + frozen = (self.root / REQUEST / "reviewed-source.json").read_bytes() + self.assertEqual(self.service.get(REQUEST), review) + restored = DataFieldReviews(root=self.root, policy=self.policy, context=CONTEXT) + self.assertTrue(restored.recover_ready(REQUEST)) + self.assertEqual(restored.payload(REQUEST)[0], frozen) + audit = json.loads(frozen)["manual_data_review"]["manifest"] + self.assertEqual(sum(event["action"] == "finalize" for event in audit["events"]), 1) + self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR) + + def test_saved_status_zero_items_preserves_prior_decisions_in_audit(self): + self.document = source_document(2) + gap(self.document, "BLOCK_CODE") + review = self.prepare() + review = self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "PRESERVED", ACTOR) + evidence = {"version": "arr-reservation-status-evidence/v1", "policy_id": review_module.STATUS_POLICY, + "original_sha256": hashlib.sha256(raw(self.document)).hexdigest(), "source_manifest_sha256": SOURCE_MANIFEST, + "hotel_id": CONTEXT["hotel_id"], "report_date": DAY, + "records": [{"source_sequence": row["source_sequence"], "reservation_id": row["reservation_id"], + "reservation_status": status, "sources": [{"file": "synthetic.response.bin", "sha256": "f" * 64}]} + for row, status in zip(self.document["records"], ["Cancelled", "InHouse"])]} + with patch("integrations.ohip.reservation_status.saved_status_evidence", return_value=evidence): + review = self.service.apply_saved_status_evidence(REQUEST, Path("synthetic-unused"), + expected_revision=review["revision"]) + self.assertEqual((review["total_count"], review["pending_count"]), (0, 0)) + original = self.service.original(REQUEST) + self.assertTrue(self.service.recover_ready(REQUEST)) + audit = json.loads(self.service.payload(REQUEST)[0])["manual_data_review"]["manifest"] + self.assertEqual(audit["changes"][0]["value"], "PRESERVED") + self.assertEqual(audit["changes"][0]["actor"], ACTOR) + self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR) + self.assertEqual(self.service.original(REQUEST), original) + + +class EmptyReviewExecutorTests(unittest.TestCase): + @classmethod + def setUpClass(cls): + cls.policy = load_processor_policy(PROJECT) + + setUp = review_fixtures.DirectFieldReviewTests.setUp + configure_executor = review_fixtures.DirectFieldReviewTests.configure_executor + execute = review_fixtures.DirectFieldReviewTests.execute + wait_task = review_fixtures.DirectFieldReviewTests.wait_task + + def fixture(self): + document = omission_source() + payload = raw(document) + path = self.files / "synthetic-collected.json" + path.write_bytes(payload) + path.chmod(0o600) + fetched = {"collection_complete": True, "data_path": str(path), + "data_sha256": hashlib.sha256(payload).hexdigest(), "manifest_sha256": SOURCE_MANIFEST} + return document, fetched + + def reanalyze(self, document): + review = self.executor.get_data_review(REQUEST) + return self.executor.data_reviews.apply_source_reanalysis(REQUEST, raw(optional_omission_source(document)), + REANALYSIS_MANIFEST, REANALYSIS_POLICY, expected_revision=review["revision"]) + + def test_executor_zero_item_saved_review_continues_same_source_with_one_finance_version(self): + document, fetched = self.fixture() + with patch.object(self.source, "fetch", return_value=fetched) as fetch: + self.assertEqual(self.execute().status, "needs_data_review") + self.reanalyze(document) + original = self.executor.data_reviews.original(REQUEST) + self.assertEqual(self.execute().status, "succeeded") + self.assertEqual(self.execute().status, "succeeded") + self.assertEqual(fetch.call_count, 1) + self.assertEqual(len(self.repository._versions), 1) + self.assertEqual(self.executor.data_reviews.original(REQUEST), original) + + def test_missing_or_changed_saved_checkpoint_never_auto_freezes_or_refetches(self): + document, fetched = self.fixture() + with patch.object(self.source, "fetch", return_value=fetched): + self.assertEqual(self.execute().status, "needs_data_review") + self.reanalyze(document) + path = self.files / "executor" / "requests" / REQUEST / "collected.json" + saved = path.read_bytes() + with patch.object(self.source, "fetch", side_effect=AssertionError("must not refetch")) as fetch: + path.unlink() + with self.assertRaises(FileNotFoundError): + self.executor.recover_data_review(REQUEST) + path.write_bytes(saved) + path.chmod(0o600) + changed = json.loads(saved) + changed["data_sha256"] = "0" * 64 + path.write_bytes(raw(changed)) + with self.assertRaises(CollectionError): + self.executor.recover_data_review(REQUEST) + fetch.assert_not_called() + self.assertEqual(self.executor.get_data_review(REQUEST)["status"], "editing") + self.assertIsNone(self.executor.data_reviews.payload(REQUEST)) + self.assertEqual(len(self.repository._versions), 0) + + def test_startup_recovers_without_browser_get_or_another_fetch_and_restart_is_idempotent(self): + document, fetched = self.fixture() + with patch.object(self.source, "fetch", return_value=fetched): + queue = PersistentARRDownloads(self.files / "queue", self.executor) + try: + queue.create(DAY, REQUEST) + self.assertEqual(self.wait_task(queue)["status"], "needs_data_review") + finally: + queue.close(wait=True) + self.reanalyze(document) + completed = threading.Event() + execute = self.executor.execute + def run(**kwargs): + try: + return execute(**kwargs) + finally: + completed.set() + with patch.object(self.source, "fetch", side_effect=AssertionError("must replay saved checkpoint")), \ + patch.object(self.executor, "execute", side_effect=run): + queue = PersistentARRDownloads(self.files / "queue", self.executor) + try: + self.assertTrue(completed.wait(30), "startup must recover without any browser read") + self.assertEqual(self.wait_task(queue)["status"], "succeeded") + finally: + queue.close(wait=True) + queue = PersistentARRDownloads(self.files / "queue", self.executor) + try: + self.assertEqual(queue.get(REQUEST)["status"], "succeeded") + finally: + queue.close(wait=True) + self.assertEqual(len(self.repository._versions), 1) + + def test_successful_freeze_before_lost_queue_update_recovers_on_next_startup(self): + document, fetched = self.fixture() + with patch.object(self.source, "fetch", return_value=fetched): + queue = PersistentARRDownloads(self.files / "queue", self.executor) + try: + queue.create(DAY, REQUEST) + self.assertEqual(self.wait_task(queue)["status"], "needs_data_review") + finally: + queue.close(wait=True) + self.reanalyze(document) + recover = self.executor.recover_data_review + recovered = threading.Event() + def lost_update(request_id): + recover(request_id) + recovered.set() + raise OSError("synthetic loss before queue update") + with patch.object(self.executor, "recover_data_review", side_effect=lost_update): + queue = PersistentARRDownloads(self.files / "queue", self.executor) + try: + self.assertTrue(recovered.wait(10)) + with sqlite3.connect(self.files / "queue" / "arr-downloads.sqlite3") as db: + self.assertEqual(db.execute("SELECT status FROM downloads WHERE request_id=?", (REQUEST,)).fetchone()[0], + "needs_data_review") + self.assertTrue(queue.ready) + finally: + queue.close(wait=True) + frozen = self.executor.data_reviews.payload(REQUEST) + with patch.object(self.source, "fetch", side_effect=AssertionError("must replay saved checkpoint")): + queue = PersistentARRDownloads(self.files / "queue", self.executor) + try: + self.assertEqual(self.wait_task(queue)["status"], "succeeded") + finally: + queue.close(wait=True) + self.assertEqual(self.executor.data_reviews.payload(REQUEST), frozen) + self.assertEqual(len(self.repository._versions), 1) + + def test_authorized_wrapper_recovery_is_local_without_access_refresh(self): + from arr_web.local_ohip import AuthorizedExecutor + executor, access = Mock(), Mock() + executor.recover_data_review.return_value = True + self.assertTrue(AuthorizedExecutor(executor, access).recover_data_review(REQUEST)) + executor.recover_data_review.assert_called_once_with(REQUEST) + access.require_ready.assert_not_called() + + +class EmptyReviewQueueRecoveryTests(unittest.TestCase): + def test_one_bad_review_cannot_kill_worker_and_get_never_runs_new_recovery(self): + temporary = tempfile.TemporaryDirectory() + self.addCleanup(temporary.cleanup) + root = Path(temporary.name) + entered = threading.Event() + executor = Mock() + executor.source_identity.return_value = {"kind": "synthetic-empty-review-recovery"} + executor.get_acquisition_progress.return_value = None + executor.get_data_review.return_value = {"status": "editing"} + executor.recover_data_review.return_value = False + queue = PersistentARRDownloads(root, executor) + queue.close(wait=True) + now = "2026-10-08T00:00:00+00:00" + with sqlite3.connect(root / "arr-downloads.sqlite3") as db: + db.executemany("INSERT INTO downloads(request_id,report_date,status,created_at,updated_at) VALUES(?,?,?,?,?)", + [("a" * 32, DAY, "needs_data_review", now, now), + ("b" * 32, "2026-09-16", "needs_data_review", now, now), + ("c" * 32, "2026-09-17", "needs_data_review", now, now)]) + def recover(request_id): + if request_id == "a" * 32: + raise CollectionError("synthetic conflicting evidence") + return request_id == "c" * 32 + def execute(**kwargs): + entered.set() + return DownloadOutcome("succeeded", "arrjob-synthetic-recovered") + executor.recover_data_review.side_effect = recover + executor.execute.side_effect = execute + queue = PersistentARRDownloads(root, executor) + try: + self.assertTrue(entered.wait(5), "startup resumes the good task without a browser read") + queue.close(wait=True) + self.assertEqual(queue.get("c" * 32)["status"], "succeeded") + before = executor.recover_data_review.call_count + self.assertEqual(queue.get("a" * 32)["status"], "needs_data_review") + self.assertEqual(queue.get("b" * 32)["status"], "needs_data_review") + queue.latest() + queue.pending_reviews() + self.assertEqual(executor.recover_data_review.call_count, before) + self.assertEqual(executor.execute.call_count, 1) + finally: + queue.close(wait=True) + + +if __name__ == "__main__": + unittest.main()