From db42a5e087d2203927099b41a2805f8f0d700229 Mon Sep 17 00:00:00 2001 From: Wyndham ARR Date: Thu, 8 Oct 2026 23:42:31 +0800 Subject: [PATCH] Show durable reservation acquisition progress in download tasks --- .../tasks/20261008-production-review-9e7b.md | 12 ++ arr_web/arr_data_executor.py | 3 + arr_web/arr_downloads.py | 30 ++- arr_web/local_ohip.py | 4 + arr_web/static/app.js | 50 +++++ arr_web/static/i18n.js | 16 ++ arr_web/static/index.html | 10 + arr_web/static/styles.css | 19 ++ integrations/ohip/acquisition_progress.py | 47 ++++ integrations/ohip/arr_data.py | 122 +++++++++-- tests/javascript/arr_download_progress.cjs | 146 +++++++++++++ tests/javascript/helpers/arr_ui_harness.cjs | 15 +- tests/test_arr_acquisition_progress.py | 200 ++++++++++++++++++ tests/test_arr_data_review.py | 3 + tests/test_arr_web_daily_visual_ui.py | 13 ++ tests/test_arr_web_downloads.py | 90 ++++++++ 16 files changed, 757 insertions(+), 23 deletions(-) create mode 100644 integrations/ohip/acquisition_progress.py create mode 100644 tests/javascript/arr_download_progress.cjs create mode 100644 tests/test_arr_acquisition_progress.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 c222ade..32257f2 100644 --- a/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md +++ b/.project-docs/30-worklog/tasks/20261008-production-review-9e7b.md @@ -299,3 +299,15 @@ Read: memory-index, project-positioning, current-state latest September sections - Actual local UI: reloaded only the independent verification tab, then selected9/17 fields→history9/16 prices→same cached9/17 fields→9/15 fields→9/16 prices. At each step the date/type matched and only one review body was visible.9/15 retains its1 missing-room item,9/16 retains its one blank price key0/1,9/17 retains0 active field items; no save, cancellation or generation was triggered. All3 complete review objects match the pre-layout authenticated baseline exactly. Private before/after JSON evidence and screenshots are in `production-validation-20261007/unified-manual-check-{before,after}-20261008.json` and `unified-manual-check-layout-20261008.png`; raw material stays outside Git. - Responsive result: observed CSS widths375/767/1024/1440 verified. Initial767 check found7px page overflow from the absolute screen-reader label inside the price table; setting the new parent to position:relative bounds that label, and final375/767 checks now have scrollWidth equal clientWidth. Wide tables retain their existing contained scrolling. Temporary viewport overrides were reset. The preserved visual palette/typography was preferred over the design search's unrelated liquid-glass/marketing recommendations. Existing user tabs were not reloaded; result tab retained. Static changes are served by the running8875 instance, without a service restart, new acquisition, hotel write, real review mutation, report generation or remote push. - Promotion candidate: integration-owned product/data-flow docs should describe independent “需处理日期” navigation followed by a single “人工核对” workspace, keeping source completeness and price decisions as distinct business stages. No conflict with ADR006/007 or additional business-policy confirmation. Remaining action for users is refreshing an existing page to load the new layout, then continuing original pending tasks. + +## Same-task Follow-up: Acquisition Progress Display + +- User requests a progress bar, acquisition counts and fetching/completed status. Concurrent Task Gate Passed for same task/feature/codex/owned checkout/branch/base, no peers and clean initial checkout. Project Context Loaded: retained required memory/positioning/decisions004/006/007/architecture/domain/evidence/reflection/commitments/stale context, refreshed task record and planning entry/gate. Goal: truthful reservation-level acquisition visibility without changing completion, source-field or price rules. Modules: OHIP source progress, executor/queue public snapshots, static download UI and focused tests. Constraints: preserve immutable capture/checkpoints, user decisions and completed reports; no new Oracle acquisition for testing and no fabricated progress/time estimates. Existing design skill/system applies to the same download-card layout. Planning Gate Passed. +- Plan: persist privacy-minimized progress outside immutable capture inventories as records complete; expose it in existing task GET snapshots. Show completed/total reservations and a determinate bar once the initial search establishes the total, with an indeterminate state before that. Keep below100% until final completeness/recheck succeeds, distinguish failed/interrupted collection, and label successful acquisition “已完成” even when subsequent human checking is still pending. Recover counts from previously completed local capture metadata where validated, without refetching. Implement collection/API and UI in parallel, verify synthetic in-progress/failure/retry/restart/completion states, safely activate the local service after checking for active jobs, and confirm the actual existing completed acquisitions in the page. + +- Backend outcome: reservation-level observations now persist in a private request-level `progress.json`, outside immutable attempt inventories. The existing task snapshots expose only phase/completed/total/percent/updated_at through source/executor/authorized wrapper/queue; reading progress requires no OHIP access refresh, HTTP, execution lock or guest payload read. Initial search establishes total, completed increments after each complete record, and final source recheck/checkpoint publication gates100%. Acquisition failure stays below100%; retry resets the observation, restart exposes interruption, and later report failure does not erase proven acquisition completion. Completed legacy sources recover counts from bounded validated request/ready/result/capture metadata. Source/version/hash/processor identity, shared business rules and original payload remain unchanged; observation write failure cannot fail acquisition. +- Frontend outcome: the download card now shows a compact progress bar, count and “取数中/已完成” status, with three-language copy and accessible progress values. Unknown totals show an indeterminate indicator; completeness verification displays at most99%. Completed acquisition remains separate from manual field/price review and overall report status. Failed/interrupted collection is explicit, unrelated selected dates hide the old task's progress, and legacy completion without counts does not invent a total. Existing dates/manual workspace and all decisions remain. +- Verification:118 backend cases passed with `runtime/bin/python -m unittest tests.test_arr_acquisition_progress tests.test_arr_web_downloads tests.test_ohip_arr_data tests.test_arr_data_review.DirectFieldReviewTests tests.test_arr_download_runtime`;84 JavaScript cases passed with `node --test tests/javascript/*.cjs`;13 static UI cases passed with `runtime/bin/python -m unittest tests.test_arr_web_daily_price_review_ui tests.test_arr_web_daily_visual_ui`. Focused tests cover true row counts rather than request counts, final recheck/checkpoint failure at6/6 still99%, optional field gaps with acquisition completed, observation writer failure isolation, unlocked/no-HTTP getters, read-only legacy metadata, corrupt/context-mismatched observations, queued/retry/restart/later-report-failure precedence, API/wrapper forwarding and date/context UI recovery. A preliminary read-only queue inspection used the wrong table name and was corrected to `downloads`; it did not mutate storage and is not counted as a test failure or successful verification. +- Actual UI verification: a temporary loopback-only synthetic fixture verified unknown total →58/85 at68% →85/85 verification at99% →100% while report processing continues →failed68%, without Oracle, real database or report mutations. Fixture/tab were stopped and closed. Before activation the real queue had no active tasks. Gracefully restarted the same8875 LaunchAgent once; all required processing/acquisition/monthly/download health checks are ready. In a new independent tab,9/15 now shows85/85,100%,取数状态已完成, while its original room item remains0/1 and generation disabled.9/17 API reports106/106 completed. The replayed9/16 request lacks an original source folder, so it truthfully shows completed with historical counts unavailable rather than borrowing counts from another request.375px CSS-width check confirms scrollWidth=clientWidth and visible counts; temporary viewport override reset. User tabs were not reloaded. +- Preservation/evidence: exact authenticated before/after comparison confirms all3 complete manual-review objects, all15 original capture metadata hashes, and all6 queue statuses/job IDs/attempt counters unchanged. No new Oracle request, hotel write, price/source decision, finalization, report/monthly generation, schema/config/credential change, primary checkout edit or remote push occurred. Private evidence remains in `production-validation-20261007/acquisition-progress-{before,after}-20261008.json`, synthetic fixture source, and `acquisition-progress-local-20261008.png`; actual result tab retained. Node syntax, Git whitespace and project-document gates apply before commit. +- Promotion candidate: integration-owned product/data-flow/current-state documentation should describe reservation-level acquisition progress as a distinct observed stage, with unknown totals before paging completes and100% only after completeness/checkpoint success. Acquisition completion does not establish report completion or remove manual review. No new business-policy decision or ADR conflict. Remaining limitation: some replayed historical requests have no own acquisition counts; current new acquisitions persist them. Existing unresolved local-status fetch timeout resilience and same-page price redraw draft behavior remain separate follow-ups, not claimed fixed here. diff --git a/arr_web/arr_data_executor.py b/arr_web/arr_data_executor.py index 95a68ad..cf98d24 100644 --- a/arr_web/arr_data_executor.py +++ b/arr_web/arr_data_executor.py @@ -25,6 +25,9 @@ class DirectARRExecutor: def get_data_review(self, request_id): return self.data_reviews.get(request_id) + def get_acquisition_progress(self, *, request_id, report_date): + return self.source.get_acquisition_progress(request_id=request_id, report_date=report_date) + def update_data_review_item(self, request_id, item_id, revision, value, actor): return self.data_reviews.update(request_id, item_id, revision, value, actor) diff --git a/arr_web/arr_downloads.py b/arr_web/arr_downloads.py index 68cebd9..47af62c 100644 --- a/arr_web/arr_downloads.py +++ b/arr_web/arr_downloads.py @@ -22,6 +22,7 @@ from typing import Callable, Protocol from zoneinfo import ZoneInfo from arr_web.contracts import PortalError, validate_job_id +from integrations.ohip import acquisition_progress ACTIVE = frozenset({"queued", "downloading", "processing"}) @@ -188,8 +189,32 @@ class PersistentARRDownloads: def ready(self) -> bool: return not self._stopping and self._thread.is_alive() - @staticmethod - def _public(row: sqlite3.Row) -> dict: + def _acquisition_progress(self, row): + observation = None + getter = getattr(self._executor, "get_acquisition_progress", None) + if callable(getter): + try: + value = getter(request_id=row["request_id"], report_date=row["report_date"]) + if value is not None: + observation = acquisition_progress.validate(value) + except Exception: + pass + # Acquisition completion remains true when later report processing fails. + # Other observations cannot override the queue's restart/retry state. + if observation and observation["phase"] == "completed": + return observation + status = row["status"] + if status == "queued": + return acquisition_progress.snapshot("queued", updated_at=row["updated_at"]) + if status in {"failed", "interrupted"}: + return acquisition_progress.snapshot(status, observation["completed"] if observation else 0, + observation["total"] if observation else None, updated_at=row["updated_at"]) + if status == "downloading" and (observation is None + or datetime.fromisoformat(observation["updated_at"]) < datetime.fromisoformat(row["updated_at"])): + return acquisition_progress.snapshot("searching", updated_at=row["updated_at"]) + return observation + + def _public(self, row: sqlite3.Row) -> dict: return { "request_id": row["request_id"], "report_date": row["report_date"], "from_date": row["report_date"], "to_date": row["report_date"], @@ -199,6 +224,7 @@ class PersistentARRDownloads: if not row["job_id"] else None, "can_retry": bool(row["retryable"]) and row["attempts"] < MAX_ATTEMPTS, "created_at": row["created_at"], "updated_at": row["updated_at"], + "acquisition_progress": self._acquisition_progress(row), } def latest(self) -> dict | None: diff --git a/arr_web/local_ohip.py b/arr_web/local_ohip.py index 90164b5..bc2e83a 100644 --- a/arr_web/local_ohip.py +++ b/arr_web/local_ohip.py @@ -122,6 +122,10 @@ class AuthorizedExecutor: def source_identity(self): return self.executor.source_identity() + def get_acquisition_progress(self, **kwargs): + getter = getattr(self.executor, "get_acquisition_progress", None) + return getter(**kwargs) if callable(getter) else None + def get_data_review(self, request_id): return self.executor.get_data_review(request_id) def update_data_review_item(self, *args): return self.executor.update_data_review_item(*args) diff --git a/arr_web/static/app.js b/arr_web/static/app.js index 428f043..673cd2d 100644 --- a/arr_web/static/app.js +++ b/arr_web/static/app.js @@ -437,6 +437,55 @@ scheduleARRDownloadPoll(); } + function renderARRAcquisitionProgress(task = state.arrDownloadTask) { + const panel = $("#arr-acquisition-progress"); + const raw = task?.acquisition_progress; + const phases = ["queued", "searching", "fetching", "verifying", "completed", "failed", "interrupted"]; + const progress = raw && phases.includes(raw.phase) ? raw : null; + const completedStatuses = ["processing", "needs_data_review", "needs_review", "succeeded"]; + const phase = progress?.phase || (completedStatuses.includes(task?.status) ? "completed" + : task?.status === "downloading" ? "fetching" : ["queued", "failed", "interrupted"].includes(task?.status) ? task.status : ""); + panel.hidden = !phase || !task || task.report_date !== $("#arr-download-date").value; + if (panel.hidden) return; + const complete = phase === "completed"; + const failed = ["failed", "interrupted"].includes(phase); + const completed = Number.isInteger(progress?.completed) && progress.completed >= 0 ? progress.completed : null; + const total = completed !== null && Number.isInteger(progress?.total) && progress.total >= completed ? progress.total : null; + // Only server-reported percentages with a known total are determinate. + // All records fetched still means at most 99% until completeness is verified. + const percent = total !== null && Number.isInteger(progress?.percent) && progress.percent >= 0 && progress.percent <= 100 + ? Math.min(progress.percent, complete ? 100 : 99) : null; + const statusKey = ["queued", "completed", "failed", "interrupted"].includes(phase) ? phase : "active"; + const status = I18N.t(`arr_acquisition.${statusKey}`); + const counts = total !== null ? I18N.t("arr_acquisition.counts", { completed: formatInteger(completed), total: formatInteger(total) }) : ""; + let detail; + if (phase === "queued") detail = I18N.t("arr_acquisition.waiting"); + else if (phase === "verifying") detail = total !== null + ? I18N.t("arr_acquisition.verifying_counts", { completed: formatInteger(completed), total: formatInteger(total) }) + : I18N.t("arr_acquisition.verifying"); + else if (complete) detail = counts || I18N.t("arr_acquisition.legacy_completed"); + else if (failed) detail = [counts, I18N.t(`arr_acquisition.${phase}_detail`)].filter(Boolean).join(" · "); + else if (counts) detail = counts; + else if (phase === "searching") detail = I18N.t("arr_acquisition.searching"); + else detail = completed !== null ? I18N.t("arr_acquisition.count_unknown", { completed: formatInteger(completed) }) + : I18N.t("arr_acquisition.fetching"); + panel.classList.toggle("is-completed", complete); + panel.classList.toggle("is-error", failed); + $("#arr-acquisition-label").textContent = I18N.t("arr_acquisition.label"); + $("#arr-acquisition-status").textContent = status; + $("#arr-acquisition-detail").textContent = detail; + $("#arr-acquisition-percent").hidden = percent === null; + $("#arr-acquisition-percent").textContent = percent === null ? "" : `${percent}%`; + const track = $("#arr-acquisition-track"); + const indeterminate = percent === null && !complete && !failed; + track.hidden = percent === null && !indeterminate; + track.classList.toggle("is-indeterminate", indeterminate); + track.setAttribute("aria-valuetext", `${status} · ${detail}${percent === null ? "" : ` · ${percent}%`}`); + if (percent === null) track.removeAttribute("aria-valuenow"); + else track.setAttribute("aria-valuenow", String(percent)); + $("#arr-acquisition-bar").style.width = percent === null ? "" : `${percent}%`; + } + function renderARRDownload() { const task = state.arrDownloadTask; const pending = Boolean(state.arrDownloadIntent && !task); @@ -481,6 +530,7 @@ $("#arr-download-review").hidden = task?.status !== "needs_review" || !task.job_id; $("#arr-download-data-review").hidden = !arrDownloadNeedsDataReview(); $("#arr-download-data-review").disabled = locked; + renderARRAcquisitionProgress(task); renderARRPendingReviews(); } diff --git a/arr_web/static/i18n.js b/arr_web/static/i18n.js index d2ffa07..54a44e1 100644 --- a/arr_web/static/i18n.js +++ b/arr_web/static/i18n.js @@ -113,6 +113,22 @@ "arr_download.date_resume": ["可先选择日期;请先确认或继续 {date} 的原任务。", "You can choose a date. First check or resume the original task for {date}.", "เลือกวันที่ได้ แต่โปรดตรวจสอบหรือดำเนินงานเดิมของวันที่ {date} ต่อก่อน"], "arr_download.help": ["选择一天,自动获取到店数据并生成报表。", "Choose one day to fetch arrivals and generate reports.", "เลือกหนึ่งวันเพื่อดาวน์โหลดและประมวลผลรายงานผู้เข้าพัก"], "arr_download.start": ["下载并处理", "Download & process", "ดาวน์โหลดและประมวลผล"], + "arr_acquisition.label": ["取数状态", "Data retrieval", "สถานะการดึงข้อมูล"], + "arr_acquisition.queued": ["等待取数", "Queued", "รอดึงข้อมูล"], + "arr_acquisition.active": ["取数中", "Retrieving", "กำลังดึงข้อมูล"], + "arr_acquisition.completed": ["已完成", "Completed", "เสร็จแล้ว"], + "arr_acquisition.failed": ["取数失败", "Retrieval failed", "ดึงข้อมูลไม่สำเร็จ"], + "arr_acquisition.interrupted": ["取数中断", "Retrieval interrupted", "การดึงข้อมูลขัดจังหวะ"], + "arr_acquisition.waiting": ["正在等待取数…", "Waiting to retrieve data…", "กำลังรอดึงข้อมูล…"], + "arr_acquisition.searching": ["正在确认预订总数…", "Determining the reservation total…", "กำลังตรวจสอบจำนวนการจองทั้งหมด…"], + "arr_acquisition.fetching": ["正在获取预订数据…", "Retrieving reservation data…", "กำลังดึงข้อมูลการจอง…"], + "arr_acquisition.counts": ["已获取 {completed} / {total} 笔", "Retrieved {completed} / {total} reservations", "ดึงข้อมูลแล้ว {completed} / {total} การจอง"], + "arr_acquisition.count_unknown": ["已获取 {completed} 笔 · 总数尚未确定", "Retrieved {completed} reservations · total not yet known", "ดึงข้อมูลแล้ว {completed} การจอง · ยังไม่ทราบจำนวนทั้งหมด"], + "arr_acquisition.verifying": ["正在校验完整性…", "Verifying completeness…", "กำลังตรวจสอบความครบถ้วน…"], + "arr_acquisition.verifying_counts": ["已获取 {completed} / {total} 笔 · 正在校验完整性", "Retrieved {completed} / {total} reservations · verifying completeness", "ดึงข้อมูลแล้ว {completed} / {total} การจอง · กำลังตรวจสอบความครบถ้วน"], + "arr_acquisition.legacy_completed": ["该任务未记录取数笔数。", "Retrieval counts were not recorded for this task.", "งานนี้ไม่ได้บันทึกจำนวนการจองที่ดึงข้อมูล"], + "arr_acquisition.failed_detail": ["取数未完成,请查看任务提示。", "Data retrieval did not complete. Check the task message.", "การดึงข้อมูลยังไม่เสร็จ โปรดดูข้อความของงาน"], + "arr_acquisition.interrupted_detail": ["取数已中断,请继续原任务。", "Data retrieval was interrupted. Resume the original task.", "การดึงข้อมูลขัดจังหวะ โปรดดำเนินงานเดิมต่อ"], "arr_download.connecting": ["正在连接下载服务…", "Connecting to download service…", "กำลังเชื่อมต่อบริการดาวน์โหลด…"], "arr_download.service_reconnecting": ["下载服务连接中断,正在重新连接…", "Download service disconnected. Reconnecting…", "การเชื่อมต่อบริการดาวน์โหลดขัดข้อง กำลังเชื่อมต่อใหม่…"], "arr_download.unavailable": ["自动下载服务暂未就绪", "Automatic download is not yet available", "บริการดาวน์โหลดอัตโนมัติยังไม่พร้อมใช้งาน"], diff --git a/arr_web/static/index.html b/arr_web/static/index.html index 0ea404e..4b23709 100644 --- a/arr_web/static/index.html +++ b/arr_web/static/index.html @@ -74,6 +74,16 @@ 正在连接下载服务… +
diff --git a/arr_web/static/styles.css b/arr_web/static/styles.css index 7f5a4db..eb37165 100644 --- a/arr_web/static/styles.css +++ b/arr_web/static/styles.css @@ -197,6 +197,25 @@ button { color: inherit; } .arr-download-feedback.is-success { color: #087a55; } .arr-download-feedback.is-review { color: var(--amber); } .arr-download-feedback.is-error { color: #b42318; } +.arr-acquisition-progress { min-width: 0; margin: 0 0 12px; } +.arr-acquisition-progress[hidden], .arr-acquisition-track[hidden] { display: none; } +.arr-acquisition-meta { display: flex; align-items: baseline; justify-content: space-between; flex-wrap: wrap; gap: 4px 10px; color: var(--muted); font-size: 10px; line-height: 1.5; } +.arr-acquisition-meta > span { display: flex; flex-wrap: wrap; gap: 7px; } +#arr-acquisition-status { color: var(--blue-dark); font-weight: 650; } +#arr-acquisition-percent { color: var(--blue-dark); font-variant-numeric: tabular-nums; } +.arr-acquisition-track { height: 5px; margin-top: 6px; overflow: hidden; border-radius: 999px; background: #e8eef9; } +.arr-acquisition-track > span { display: block; height: 100%; border-radius: inherit; background: var(--blue); transition: width .25s ease; } +.arr-acquisition-track.is-indeterminate > span { width: 32%; animation: arr-acquisition-pulse 1.7s ease-in-out infinite; } +.arr-acquisition-progress.is-completed #arr-acquisition-status { color: #087a55; } +.arr-acquisition-progress.is-completed .arr-acquisition-track > span { background: #087a55; } +.arr-acquisition-progress.is-error #arr-acquisition-status { color: #b42318; } +.arr-acquisition-progress.is-error .arr-acquisition-track > span { background: #b42318; } +#arr-acquisition-detail { margin: 6px 0 0; color: var(--muted); font-size: 10px; line-height: 1.5; overflow-wrap: anywhere; font-variant-numeric: tabular-nums; } +@keyframes arr-acquisition-pulse { from { transform: translateX(-100%); } to { transform: translateX(350%); } } +@media (prefers-reduced-motion: reduce) { + .arr-acquisition-track > span { transition: none; } + .arr-acquisition-track.is-indeterminate > span { animation: none; } +} .arr-download-actions { display: flex; align-items: center; justify-content: flex-end; flex-wrap: wrap; gap: 8px; } .arr-download-actions .primary-button { min-height: 34px; margin: 0; padding: 0 12px; border-radius: 8px; font-size: 11px; } .arr-download-review { padding: 4px 0; border: 0; color: var(--blue-dark); background: transparent; font-size: 11px; font-weight: 650; cursor: pointer; text-decoration: underline; text-underline-offset: 3px; } diff --git a/integrations/ohip/acquisition_progress.py b/integrations/ohip/acquisition_progress.py new file mode 100644 index 0000000..38c802e --- /dev/null +++ b/integrations/ohip/acquisition_progress.py @@ -0,0 +1,47 @@ +"""Private, best-effort acquisition observations, separate from source evidence.""" +from datetime import datetime, timezone + +from .capture_job import atomic_json, fingerprint + + +VERSION = "arr-acquisition-progress/v1" +PHASES = frozenset({"queued", "searching", "fetching", "verifying", "completed", "failed", "interrupted"}) +PUBLIC_FIELDS = frozenset({"phase", "completed", "total", "percent", "updated_at"}) + + +def snapshot(phase, completed=0, total=None, *, updated_at=None): + percent = (100 if phase == "completed" else 0 if phase == "queued" else + min(99, completed * 100 // total) if total else 99 if total == 0 else None) + return validate({"phase": phase, "completed": completed, "total": total, "percent": percent, + "updated_at": updated_at or datetime.now(timezone.utc).isoformat()}) + + +def validate(value): + """Return only the fixed privacy-safe public contract; reject malformed data.""" + if type(value) is not dict or set(value) != PUBLIC_FIELDS or value["phase"] not in PHASES: + raise ValueError("invalid acquisition progress") + completed, total, percent = value["completed"], value["total"], value["percent"] + if type(completed) is not int or not 0 <= completed <= 10000000: + raise ValueError("invalid acquisition count") + if total is not None and (type(total) is not int or not completed <= total <= 10000000): + raise ValueError("invalid acquisition total") + if percent is not None and (type(percent) is not int or not 0 <= percent <= 100): + raise ValueError("invalid acquisition percent") + if value["phase"] == "completed": + if total is None or completed != total or percent != 100: + raise ValueError("incomplete acquisition") + elif percent == 100: + raise ValueError("premature acquisition completion") + stamp = value["updated_at"] + if type(stamp) is not str or len(stamp) > 40 or datetime.fromisoformat(stamp).tzinfo is None: + raise ValueError("invalid acquisition timestamp") + return dict(value) + + +def publish(folder, identity, phase, completed=0, total=None): + """An observation failure must never change the source collection outcome.""" + try: + atomic_json(folder / "progress.json", {"version": VERSION, "identity_sha256": fingerprint(identity), + "progress": snapshot(phase, completed, total)}, replace=True) + except Exception: + pass diff --git a/integrations/ohip/arr_data.py b/integrations/ohip/arr_data.py index aad1048..fe3d024 100644 --- a/integrations/ohip/arr_data.py +++ b/integrations/ohip/arr_data.py @@ -19,10 +19,11 @@ import stat import time from . import collect_arr_source as base +from . import acquisition_progress as progress from . import profile_summary as profiles from . import rate_info, room_calendar_evidence as calendar, source_fields as fields from .audit_arr_capture import protected_read -from .capture_job import atomic_json, job_lock, private_directory, sync_directory +from .capture_job import atomic_json, fingerprint, job_lock, private_directory, sync_directory from .data_client import DataReader, document, json_bytes, typed_id @@ -379,7 +380,7 @@ def _record(search, detail, reader, sequence): "fields": values, "related": related} -def collect(reader): +def collect(reader, report_progress=None): """Collect all source rows; return no Finance/business success claims.""" options = reader.options rows, fatal = [], None @@ -387,8 +388,18 @@ def collect(reader): "service_url": base.SERVICE, "application_id": base.APPLICATION, "source_kind": reader.source_kind, "max_requests": reader.max_requests}) complete = False + def report(phase, completed=0, total=None): + if report_progress is not None: + try: + report_progress(phase, completed, total) + except Exception: + pass + total = None + report("searching") try: searches = base.search_day(reader, options, server_sort=False) + total = len(searches) + report("fetching", 0, total) search_refs = list(reader.used) for index, search in enumerate(searches, 1): before = len(reader.used) @@ -415,6 +426,8 @@ def collect(reader): row = _record(search, detail, reader, index) row["sources"] = list(dict.fromkeys(search_refs + reader.used[before:])) rows.append(row) + report("fetching", len(rows), total) + report("verifying", len(rows), total) pending = [r for r in rows if r["fields"]["DISP_ROOM_NO"]["state"] not in {"available", "empty"}] if pending: before = len(reader.used) @@ -443,6 +456,8 @@ def collect(reader): f["state"] == "available" or (key in OPTIONAL_FIELDS and f["state"] == "empty") for row in rows for key, f in row["fields"].items()) failed = fatal is not None or states["failed"] > 0 + if failed: + report("failed", len(rows), total) status = "failed" if failed else "collected" if ready else "collected_with_gaps" payload = {"version": VERSION, "hotel_id": options.hotel_id, "report_date": options.arrival_date, "source_kind": reader.source_kind, @@ -481,13 +496,74 @@ class ARRDataSource: self.page_size, self.max_pages, self.max_records = page_size, max_pages, max_records self.max_requests, self.sleep = max_requests, sleep - def fetch(self, report_date: str, request_id: str) -> dict: + def _identity(self, report_date, request_id): options = base.Options(report_date, self.hotel_id, self.page_size, self.max_pages, self.max_records) options.validate() require(type(request_id) is str and re.fullmatch(r"[0-9a-f]{32}", request_id), "invalid_data_request_id") identity = {"version": VERSION, "options": vars(options), "max_requests": self.max_requests, "request_id": request_id, "service_url": base.SERVICE, "application_id": base.APPLICATION, "source_kind": "ohip_platform" if self.transport_factory is None else "test_transport"} + return options, identity + + def get_acquisition_progress(self, *, report_date: str, request_id: str): + """Read local observations/legacy completion metadata without an execution lock or HTTP.""" + try: + _, identity = self._identity(report_date, request_id) + folder = self.root / request_id + self._safe_directory(self.root) + self._safe_directory(folder) + stored = document(protected_read(folder / "request.json", 65536)) + require(fingerprint(stored) == fingerprint(identity), "data_progress_context_mismatch") + pointer = folder / "ready.json" + if pointer.exists(): + return self._completed_progress(folder, document(protected_read(pointer, 65536)), identity) + stored_progress = document(protected_read(folder / "progress.json", 4096)) + require(set(stored_progress) == {"version", "identity_sha256", "progress"} + and stored_progress["version"] == progress.VERSION + and stored_progress["identity_sha256"] == fingerprint(identity), "data_progress_context_mismatch") + observation = progress.validate(stored_progress["progress"]) + # Only a published completion checkpoint establishes completion. + return observation if observation["phase"] != "completed" else None + except Exception: + return None + + @staticmethod + def _safe_directory(path): + info = path.lstat() + require(stat.S_ISDIR(info.st_mode) and info.st_uid == os.getuid() + and stat.S_IMODE(info.st_mode) == 0o700, "unsafe_data_directory") + + def _completed_progress(self, folder, pointer, identity): + require(set(pointer) == {"attempt", "manifest_sha256", "data_sha256"} + and type(pointer["attempt"]) is str and re.fullmatch(r"attempt-[0-9]{4}", pointer["attempt"]), + "invalid_data_checkpoint") + directory = folder / pointer["attempt"] + self._safe_directory(directory) + raw = protected_read(directory / "result.json", base.MAX_MANIFEST_BYTES) + require(hashlib.sha256(raw).hexdigest() == pointer["manifest_sha256"], "data_manifest_changed") + result = document(raw) + require(result.get("version") == VERSION and result.get("status") in {"collected", "collected_with_gaps"} + and result.get("collection_complete") is True + and result.get("source_kind") == identity["source_kind"] + and result.get("data_sha256") == pointer["data_sha256"], "invalid_completed_data") + capture_raw = protected_read(directory / "capture.json", 65536) + capture = document(capture_raw) + require(set(capture) == {"version", "options", "service_url", "application_id", "max_requests", "source_kind"} + and json_bytes(capture) == json_bytes({k: identity[k] for k in capture}), "data_capture_context_mismatch") + files = result.get("files") + require(type(files) is list and sum(item.get("name") == "capture.json" for item in files if type(item) is dict) == 1, + "invalid_data_inventory") + entry = next(item for item in files if type(item) is dict and item.get("name") == "capture.json") + require(entry.get("bytes") == len(capture_raw) and entry.get("sha256") == hashlib.sha256(capture_raw).hexdigest(), + "data_capture_changed") + count = result.get("records") + require(type(count) is int and 0 <= count <= self.max_records, "invalid_data_count") + stamp = directory.joinpath("result.json").stat().st_mtime + from datetime import datetime, timezone + return progress.snapshot("completed", count, count, updated_at=datetime.fromtimestamp(stamp, timezone.utc).isoformat()) + + def fetch(self, report_date: str, request_id: str) -> dict: + options, identity = self._identity(report_date, request_id) private_directory(self.root) folder = self.root / request_id private_directory(folder) @@ -504,21 +580,31 @@ class ARRDataSource: attempts = [p for p in folder.iterdir() if re.fullmatch(r"attempt-[0-9]{4}", p.name)] attempt = max([int(p.name[-4:]) for p in attempts], default=0) + 1 require(attempt <= 9999, "attempt_limit_exceeded") - key = "" - if self.transport_factory is None: - key = base.load_key(Path(self.credential_file)) - transport = base.HTTPTransport(key) - else: - transport = self.transport_factory() - archive = base.Archive(folder / f"attempt-{attempt:04d}") - reader = DataReader(archive, options, transport, key=key, sleep=self.sleep, max_requests=self.max_requests, - source_kind=identity["source_kind"]) - summary = collect(reader) - summary.update(data_path=str(archive.path / "arr-data.json"), attempt=attempt) - if summary["status"] != "failed": - atomic_json(pointer, {"attempt": archive.path.name, "manifest_sha256": summary["manifest_sha256"], - "data_sha256": summary["data_sha256"]}, replace=False) - return summary + progress.publish(folder, identity, "searching") + latest_counts = [0, None] + def report_progress(phase, completed, total): + latest_counts[:] = [completed, total] + progress.publish(folder, identity, phase, completed, total) + try: + key = "" + if self.transport_factory is None: + key = base.load_key(Path(self.credential_file)) + transport = base.HTTPTransport(key) + else: + transport = self.transport_factory() + archive = base.Archive(folder / f"attempt-{attempt:04d}") + reader = DataReader(archive, options, transport, key=key, sleep=self.sleep, max_requests=self.max_requests, + source_kind=identity["source_kind"]) + summary = collect(reader, report_progress) + summary.update(data_path=str(archive.path / "arr-data.json"), attempt=attempt) + if summary["status"] != "failed": + atomic_json(pointer, {"attempt": archive.path.name, "manifest_sha256": summary["manifest_sha256"], + "data_sha256": summary["data_sha256"]}, replace=False) + progress.publish(folder, identity, "completed", summary["records"], summary["records"]) + return summary + except Exception: + progress.publish(folder, identity, "failed", *latest_counts) + raise @staticmethod def _replay(folder, pointer, identity): diff --git a/tests/javascript/arr_download_progress.cjs b/tests/javascript/arr_download_progress.cjs new file mode 100644 index 0000000..a10f45f --- /dev/null +++ b/tests/javascript/arr_download_progress.cjs @@ -0,0 +1,146 @@ +// Actual progress rendering against synthetic task snapshots; no Oracle requests. +const test = require('node:test'); +const assert = require('node:assert/strict'); +const fs = require('node:fs'); +const path = require('node:path'); +const vm = require('node:vm'); +const {harness,response,scopedKey} = require('./helpers/arr_ui_harness.cjs'); +const localeContext=vm.createContext({window:{},document:{readyState:'loading',addEventListener(){}}}); +vm.runInContext(fs.readFileSync(path.resolve(__dirname,'../../arr_web/static/i18n.js'),'utf8'),localeContext); +const catalog=localeContext.window.ARRI18n.CATALOG; +const date='2026-09-16'; +const task={request_id:'d'.repeat(32),report_date:date,status:'downloading',job_id:null,can_retry:false}; +const progress=(phase,completed,total,percent)=>({phase,completed,total,percent,updated_at:'2026-10-08T05:00:00Z'}); +function setup(respond=()=>{throw new Error('rendering must not send a request');}) { + const h=harness(respond); + h.i18n.t=(key,values={})=>(catalog[key]?.[0]||key).replace(/\{(\w+)\}/g,(_,name)=>String(values[name]??'')); + return h; +} +async function render(h,snapshot) { + await h.acceptARRDownloadTask(snapshot,{sync:false}); +} + +test('unknown reservation totals remain indeterminate and never invent zero totals or percentages',async()=>{ + const h=setup(); + for(const phase of ['queued','searching','fetching']) { + await render(h,{...task,status:phase==='queued'?'queued':'downloading',acquisition_progress:progress(phase,0,null,null)}); + assert.equal(h.element('#arr-acquisition-progress').hidden,false); + assert.equal(h.element('#arr-acquisition-track').getAttribute('aria-valuenow'),null); + assert.equal(h.element('#arr-acquisition-track').classList.contains('is-indeterminate'),true); + assert.equal(h.element('#arr-acquisition-percent').hidden,true); + assert.doesNotMatch(h.element('#arr-acquisition-detail').textContent,/0\s*\/\s*0|%/); + } + assert.equal(h.calls.length,0); +}); + +test('a real 58 of 85 snapshot renders server progress and accessible counts',async()=>{ + const h=setup(); + await render(h,{...task,acquisition_progress:progress('fetching',58,85,68)}); + assert.equal(h.element('#arr-acquisition-status').textContent,'取数中'); + assert.equal(h.element('#arr-acquisition-detail').textContent,'已获取 58 / 85 笔'); + assert.equal(h.element('#arr-acquisition-percent').textContent,'68%'); + assert.equal(h.element('#arr-acquisition-track').getAttribute('aria-valuenow'),'68'); + assert.match(h.element('#arr-acquisition-track').getAttribute('aria-valuetext'),/58 \/ 85.*68%/); + assert.equal(h.element('#arr-acquisition-track').classList.contains('is-indeterminate'),false); + assert.equal(h.element('#arr-acquisition-bar').style.width,'68%'); +}); + +test('all records still show 99 percent while completeness is being verified',async()=>{ + const h=setup(); + for(const percent of [99,100]) { + await render(h,{...task,acquisition_progress:progress('verifying',85,85,percent)}); + assert.equal(h.element('#arr-acquisition-status').textContent,'取数中'); + assert.equal(h.element('#arr-acquisition-percent').textContent,'99%'); + assert.match(h.element('#arr-acquisition-detail').textContent,/85 \/ 85.*正在校验完整性/); + assert.equal(h.element('#arr-acquisition-track').getAttribute('aria-valuenow'),'99'); + } +}); + +test('a verified empty result is a real completed zero total',async()=>{ + const h=setup(); + await render(h,{...task,status:'succeeded',acquisition_progress:progress('completed',0,0,100)}); + assert.equal(h.element('#arr-acquisition-status').textContent,'已完成'); + assert.equal(h.element('#arr-acquisition-detail').textContent,'已获取 0 / 0 笔'); + assert.equal(h.element('#arr-acquisition-percent').textContent,'100%'); +}); + +test('completed retrieval remains distinct from pending field and price review',async()=>{ + for(const status of ['needs_data_review','needs_review','processing','succeeded']) { + const h=setup(url=>{ + assert.equal(url,`/api/arr-downloads/${task.request_id}/data-review`); + return response(200,{request_id:task.request_id,report_date:date,status:'editing',items:[],pending_count:0,total_count:0}); + }); + await render(h,{...task,status,job_id:status==='needs_data_review'?null:'fixture-job',acquisition_progress:progress('completed',85,85,100)}); + assert.equal(h.element('#arr-acquisition-status').textContent,'已完成'); + assert.equal(h.element('#arr-acquisition-percent').textContent,'100%'); + assert.equal(h.element('#arr-acquisition-detail').textContent,'已获取 85 / 85 笔'); + if(status==='needs_data_review') assert.equal(h.element('#arr-download-data-review').hidden,false); + if(status==='needs_review') assert.equal(h.element('#arr-download-review').hidden,false); + assert.doesNotMatch(h.element('#arr-acquisition-detail').textContent,/生成|报表已完成/); + } +}); + +test('failed and interrupted retrieval never appear as 100 percent complete',async()=>{ + const h=setup(); + for(const phase of ['failed','interrupted']) { + await render(h,{...task,status:phase,acquisition_progress:progress(phase,85,85,100)}); + assert.equal(h.element('#arr-acquisition-status').textContent,phase==='failed'?'取数失败':'取数中断'); + assert.equal(h.element('#arr-acquisition-progress').classList.contains('is-completed'),false); + assert.equal(h.element('#arr-acquisition-progress').classList.contains('is-error'),true); + assert.equal(h.element('#arr-acquisition-percent').textContent,'99%'); + assert.equal(h.element('#arr-acquisition-track').classList.contains('is-indeterminate'),false); + } +}); + +test('legacy metadata absence shows truthful status without fabricated counts',async()=>{ + const h=setup(); + await render(h,task); + assert.equal(h.element('#arr-acquisition-status').textContent,'取数中'); + assert.equal(h.element('#arr-acquisition-track').classList.contains('is-indeterminate'),true); + assert.equal(h.element('#arr-acquisition-detail').textContent,'正在获取预订数据…'); + await render(h,{...task,status:'needs_review',job_id:'fixture-job'}); + assert.equal(h.element('#arr-acquisition-status').textContent,'已完成'); + assert.equal(h.element('#arr-acquisition-percent').hidden,true); + assert.equal(h.element('#arr-acquisition-track').hidden,true); + assert.equal(h.element('#arr-acquisition-detail').textContent,'该任务未记录取数笔数。'); + for(const status of ['failed','interrupted']) { + await render(h,{...task,status}); + assert.equal(h.element('#arr-acquisition-percent').hidden,true); + assert.equal(h.element('#arr-acquisition-track').hidden,true); + assert.notEqual(h.element('#arr-acquisition-status').textContent,'已完成'); + } +}); + +test('progress hides for another selected date and returns for the original running task',async()=>{ + const h=setup(); + await render(h,{...task,acquisition_progress:progress('fetching',58,85,68)}); + h.element('#arr-download-date').value='2026-09-17'; + h.handleARRDownloadDateChange(); + assert.equal(h.element('#arr-acquisition-progress').hidden,true); + assert.equal(h.state.arrDownloadTask.report_date,date); + assert.match(h.element('#arr-download-status').textContent,/2026-09-16/); + h.element('#arr-download-date').value=date; + h.handleARRDownloadDateChange(); + assert.equal(h.element('#arr-acquisition-progress').hidden,false); + assert.equal(h.element('#arr-acquisition-percent').textContent,'68%'); + assert.equal(h.calls.length,0); +}); + +test('context changes hide the former task progress and empty selection has none',async()=>{ + const h=setup(()=>response(200,{context_id:'new-context',ready:true,default_date:date,pending_reviews:[],latest_task:null})); + h.state.arrDownloadContextId='old-context'; + h.state.arrDownloadStorageKey=scopedKey('old-context'); + await render(h,{...task,acquisition_progress:progress('fetching',58,85,68)}); + await h.initARRDownload(); + assert.equal(h.state.arrDownloadTask,null); + assert.equal(h.element('#arr-acquisition-progress').hidden,true); +}); + +test('acquisition labels, counts and completeness guidance are available in all three languages',()=>{ + for(const [key,values] of Object.entries(catalog).filter(([key])=>key.startsWith('arr_acquisition.'))) { + assert.equal(values.length,3,key); + assert.equal(values.every(value=>typeof value==='string'&&value.length>0),true,key); + assert.match(values[2],/[\u0e00-\u0e7f]/,key); + } + assert.match(catalog['arr_acquisition.verifying_counts'][0],/校验完整性/); +}); diff --git a/tests/javascript/helpers/arr_ui_harness.cjs b/tests/javascript/helpers/arr_ui_harness.cjs index a1e2577..17dd25c 100644 --- a/tests/javascript/helpers/arr_ui_harness.cjs +++ b/tests/javascript/helpers/arr_ui_harness.cjs @@ -11,8 +11,17 @@ assert(source.endsWith(boot+'\n') || source.endsWith(boot)); function harness(respond) { const elements=new Map(), storage=new Map(), storageReads=[], calls=[]; const element=selector => { - if (!elements.has(selector)) elements.set(selector,{value:'', disabled:false, hidden:false, textContent:'', innerHTML:'',style:{}, - classList:{add(){},remove(){},toggle(){},contains(){return false;}},setAttribute(){},removeAttribute(){},querySelectorAll(){return [];},scrollIntoView(){},focus(){}}); + if (!elements.has(selector)) { + const attributes=new Map(),classes=new Set(); + elements.set(selector,{value:'', disabled:false, hidden:false, textContent:'', innerHTML:'',style:{}, + classList:{add:value=>classes.add(value),remove:value=>classes.delete(value),toggle(value,force){ + const enabled=force===undefined?!classes.has(value):force; + if(enabled) classes.add(value); else classes.delete(value); + return enabled; + },contains:value=>classes.has(value)},setAttribute:(key,value)=>attributes.set(key,String(value)), + removeAttribute:key=>attributes.delete(key),getAttribute:key=>attributes.get(key)??null, + querySelectorAll(){return [];},scrollIntoView(){},focus(){}}); + } return elements.get(selector); }; const context=vm.createContext({Headers, console, Date, Intl, Uint8Array, @@ -32,7 +41,7 @@ function harness(respond) { const node={dataset:{arrDataReviewItemId:itemId},querySelector:()=>input}; return {input,button:{closest:()=>node}}; }; - return {...subject,calls,storage,storageReads,element,row,submit:()=>subject.submitARRDownload({preventDefault(){}})}; + return {...subject,calls,storage,storageReads,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_acquisition_progress.py b/tests/test_arr_acquisition_progress.py new file mode 100644 index 0000000..80e9b01 --- /dev/null +++ b/tests/test_arr_acquisition_progress.py @@ -0,0 +1,200 @@ +"""Offline progress observations; no credentials, Oracle, or business writes.""" +import hashlib +import json +from pathlib import Path +import tempfile +import unittest +from unittest.mock import Mock, patch + +from integrations.ohip import acquisition_progress as progress +from integrations.ohip import arr_data as data +from integrations.ohip import collect_arr_source as base +from integrations.ohip.capture_job import job_lock +from tests.test_ohip_arr_data import DAY, HOTEL, REQUEST, SimulatedOHIP + + +class AcquisitionProgressTests(unittest.TestCase): + def setUp(self): + temporary = tempfile.TemporaryDirectory() + self.addCleanup(temporary.cleanup) + self.root = Path(temporary.name) / "source" + self.transport = SimulatedOHIP(6) + self.source = data.ARRDataSource(self.root, HOTEL, transport_factory=lambda: self.transport, + sleep=lambda _: None, page_size=2) + + def get(self): + return self.source.get_acquisition_progress(report_date=DAY, request_id=REQUEST) + + def observe(self): + observations = [] + original = progress.publish + def publish(folder, identity, phase, completed=0, total=None): + original(folder, identity, phase, completed, total) + observations.append((phase, completed, total, self.get())) + return observations, publish + + def test_counts_rows_only_and_verifies_before_completion_without_changing_evidence(self): + observations, publish = self.observe() + with patch.object(progress, "publish", side_effect=publish): + result = self.source.fetch(DAY, REQUEST) + self.assertTrue(result["collection_complete"]) + fetching = [(done, total) for phase, done, total, _ in observations if phase == "fetching"] + self.assertEqual(fetching, [(number, 6) for number in range(7)]) + verifying = next(value for phase, _, _, value in observations if phase == "verifying") + self.assertEqual(verifying["completed"], 6) + self.assertEqual(verifying["total"], 6) + self.assertEqual(verifying["percent"], 99) + completed = self.get() + self.assertEqual({key: completed[key] for key in ("phase", "completed", "total", "percent")}, + {"phase": "completed", "completed": 6, "total": 6, "percent": 100}) + self.assertEqual(set(completed), progress.PUBLIC_FIELDS) + attempt = Path(result["data_path"]).parent + manifest = json.loads((attempt / "result.json").read_bytes()) + self.assertNotIn("progress.json", [item["name"] for item in manifest["files"]]) + self.assertNotIn("acquisition_progress", json.loads(Path(result["data_path"]).read_bytes())) + self.assertEqual(manifest["version"], data.VERSION) + + def test_all_rows_collected_but_final_recheck_failure_never_reaches_100(self): + def edit(operation, envelope, request): + if operation == base.SEARCH and sum(call["operation"] == base.SEARCH for call in self.transport.calls) > 3: + envelope["data"]["reservations"]["reservationInfo"][0]["lastModifyDateTime"] = "2026-09-15T11:00:00+07:00" + return envelope + self.transport.edit = edit + result = self.source.fetch(DAY, REQUEST) + self.assertEqual(result["status"], "failed") + self.assertEqual(result["records"], 6) + observed = self.get() + self.assertEqual(observed["phase"], "failed") + self.assertEqual((observed["completed"], observed["total"], observed["percent"]), (6, 6, 99)) + self.assertFalse((self.root / REQUEST / "ready.json").exists()) + + def test_nonbusiness_observation_write_failure_does_not_fail_collection(self): + with patch.object(progress, "atomic_json", side_effect=OSError("synthetic observation disk failure")): + result = self.source.fetch(DAY, REQUEST) + self.assertEqual(result["status"], "collected") + self.assertTrue(result["collection_complete"]) + self.assertFalse((self.root / REQUEST / "progress.json").exists()) + self.assertEqual(self.get()["percent"], 100) + + def test_getter_is_unlocked_and_no_http_while_collection_is_running(self): + observations = [] + def status(operation, count): + before = len(self.transport.calls) + observations.append(self.get()) + self.assertEqual(len(self.transport.calls), before) + return 200 + self.transport.status = status + self.source.fetch(DAY, REQUEST) + self.assertEqual(observations[0]["phase"], "searching") + self.assertIn("fetching", [item["phase"] for item in observations]) + self.assertIn("verifying", [item["phase"] for item in observations]) + with job_lock(self.root / REQUEST): + self.assertEqual(self.get()["phase"], "completed") + + def test_legacy_complete_metadata_recovers_exact_counts_read_only_without_guest_reads(self): + result = self.source.fetch(DAY, REQUEST) + (self.root / REQUEST / "progress.json").unlink() + metadata = [self.root / REQUEST / name for name in ("request.json", "ready.json")] + metadata += [Path(result["data_path"]).parent / name for name in ("result.json", "capture.json")] + before = {str(path): hashlib.sha256(path.read_bytes()).hexdigest() for path in metadata} + calls = len(self.transport.calls) + read = data.protected_read + def read_metadata(path, maximum): + self.assertIn(Path(path).name, {"request.json", "ready.json", "result.json", "capture.json"}) + return read(path, maximum) + with patch.object(data, "protected_read", side_effect=read_metadata), patch.object(data, "job_lock", side_effect=AssertionError("getter locked")): + self.assertEqual(self.get()["completed"], 6) + self.assertEqual(self.get()["total"], 6) + self.assertEqual(before, {str(path): hashlib.sha256(path.read_bytes()).hexdigest() for path in metadata}) + self.assertEqual(len(self.transport.calls), calls) + self.assertFalse((self.root / REQUEST / "progress.json").exists()) + + def test_completed_with_field_gaps_is_acquisition_complete(self): + self.transport.rows[0]["roomStay"]["currentRoomInfo"].pop("roomId") + self.transport.calendar_rooms = [] + result = self.source.fetch(DAY, REQUEST) + self.assertEqual(result["status"], "collected_with_gaps") + self.assertEqual(self.get()["phase"], "completed") + self.assertEqual(self.get()["percent"], 100) + + def test_failed_attempt_retry_resets_counts_and_replay_uses_no_http(self): + self.transport.status = lambda operation, _: 503 if operation == "getProfiles" else 200 + failed = self.source.fetch(DAY, REQUEST) + self.assertEqual(failed["status"], "failed") + self.assertEqual(self.get()["phase"], "failed") + self.transport.status = lambda *_: 200 + observations, publish = self.observe() + with patch.object(progress, "publish", side_effect=publish): + complete = self.source.fetch(DAY, REQUEST) + self.assertEqual(complete["attempt"], 2) + self.assertEqual(observations[0][0:3], ("searching", 0, None)) + calls = len(self.transport.calls) + self.assertEqual(self.source.fetch(DAY, REQUEST), complete) + self.assertEqual(len(self.transport.calls), calls) + self.assertEqual(self.get()["completed"], 6) + + def test_context_changed_corrupt_large_and_symlink_progress_stays_unknown(self): + self.transport.status = lambda *_: 503 + self.source.fetch(DAY, REQUEST) + other = data.ARRDataSource(self.root, HOTEL, transport_factory=lambda: self.transport, page_size=1) + self.assertIsNone(other.get_acquisition_progress(report_date=DAY, request_id=REQUEST)) + self.assertIsNone(self.source.get_acquisition_progress(report_date="2026-09-16", request_id=REQUEST)) + path = self.root / REQUEST / "progress.json" + original = path.read_bytes() + for raw in (b"{", b" " * 4097): + path.write_bytes(raw) + self.assertIsNone(self.get()) + path.unlink() + target = self.root / "unrelated.json" + target.write_bytes(original) + target.chmod(0o600) + path.symlink_to(target) + self.assertIsNone(self.get()) + + def test_unpinned_completed_observation_is_not_completion(self): + self.transport.status = lambda *_: 503 + self.source.fetch(DAY, REQUEST) + _, identity = self.source._identity(DAY, REQUEST) + progress.publish(self.root / REQUEST, identity, "completed", 6, 6) + self.assertIsNone(self.get()) + + def test_changed_legacy_manifest_or_capture_does_not_claim_completion(self): + result = self.source.fetch(DAY, REQUEST) + attempt = Path(result["data_path"]).parent + for name in ("result.json", "capture.json"): + path = attempt / name + original = path.read_bytes() + path.write_bytes(original + b" ") + self.assertIsNone(self.get()) + path.write_bytes(original) + + def test_missing_old_source_remains_unknown_without_creating_files_or_transport(self): + self.assertIsNone(self.get()) + self.assertFalse(self.root.exists()) + self.assertEqual(self.transport.calls, []) + + def test_completion_checkpoint_publication_failure_preserves_collected_count_but_not_100(self): + write = data.atomic_json + def publish(path, document, *, replace): + if Path(path).name == "ready.json": + raise OSError("synthetic completion publication interruption") + return write(path, document, replace=replace) + with patch.object(data, "atomic_json", side_effect=publish), self.assertRaises(OSError): + self.source.fetch(DAY, REQUEST) + observed = self.get() + self.assertEqual((observed["phase"], observed["completed"], observed["total"], observed["percent"]), + ("failed", 6, 6, 99)) + + def test_local_authorized_wrapper_keeps_metadata_read_local_without_access_refresh(self): + from arr_web.local_ohip import AuthorizedExecutor + access, executor = Mock(), Mock() + expected = progress.snapshot("fetching", 3, 6) + executor.get_acquisition_progress.return_value = expected + wrapped = AuthorizedExecutor(executor, access) + self.assertEqual(wrapped.get_acquisition_progress(request_id=REQUEST, report_date=DAY), expected) + executor.get_acquisition_progress.assert_called_once_with(request_id=REQUEST, report_date=DAY) + access.require_ready.assert_not_called() + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_arr_data_review.py b/tests/test_arr_data_review.py index dc94086..c3a6332 100644 --- a/tests/test_arr_data_review.py +++ b/tests/test_arr_data_review.py @@ -694,6 +694,9 @@ class DirectFieldReviewTests(unittest.TestCase): outcome = self.execute() process.assert_not_called() self.assertEqual(outcome.status, "needs_data_review") + acquisition = self.executor.get_acquisition_progress(request_id=REQUEST, report_date=DAY) + self.assertEqual((acquisition["phase"], acquisition["completed"], acquisition["total"], acquisition["percent"]), + ("completed", 1, 1, 100)) self.assertIsNone(outcome.job_id) self.assertIsNone(self.repository.current_version_id(date.fromisoformat(DAY))) self.assertFalse(list((self.files / "executor" / "handoffs").rglob("delivery.json"))) diff --git a/tests/test_arr_web_daily_visual_ui.py b/tests/test_arr_web_daily_visual_ui.py index 27d5b05..1829804 100644 --- a/tests/test_arr_web_daily_visual_ui.py +++ b/tests/test_arr_web_daily_visual_ui.py @@ -73,6 +73,19 @@ class DailyVisualStaticContractTests(unittest.TestCase): self.assertIn("min-width: 118px", self.styles) self.assertNotIn("width: calc(100% - 48px)", self.styles) + def test_download_card_exposes_accessible_retrieval_progress_with_reduced_motion(self) -> None: + start = self.html.index('id="arr-acquisition-progress"') + end = self.html.index('
', start) + progress = self.html[start:end] + self.assertLess(self.html.index('id="arr-download-status"'), start) + self.assertIn('role="progressbar"', progress) + self.assertIn('aria-labelledby="arr-acquisition-label"', progress) + self.assertIn('aria-describedby="arr-acquisition-detail"', progress) + self.assertIn('aria-valuemin="0" aria-valuemax="100"', progress) + self.assertNotIn('aria-valuenow=', progress) + self.assertIn('@media (prefers-reduced-motion: reduce)', self.styles) + self.assertIn('.arr-acquisition-track.is-indeterminate > span { animation: none; }', self.styles) + def test_review_dates_are_independent_and_both_review_bodies_share_one_workspace(self) -> None: class Regions(HTMLParser): def __init__(self) -> None: diff --git a/tests/test_arr_web_downloads.py b/tests/test_arr_web_downloads.py index 1d20919..be5f04d 100644 --- a/tests/test_arr_web_downloads.py +++ b/tests/test_arr_web_downloads.py @@ -17,6 +17,7 @@ from zoneinfo import ZoneInfo from arr_web.app import PortalApplication, RuntimeHealth from arr_web.arr_downloads import DownloadOutcome, PersistentARRDownloads, default_report_date, validate_report_date from arr_web.contracts import PortalError +from integrations.ohip.acquisition_progress import snapshot as acquisition_snapshot from tests.test_arr_web import TEST_CREDENTIALS, FakeRepository, login @@ -27,6 +28,10 @@ class Executor: self.entered = threading.Event() self.outcome = outcome or DownloadOutcome("succeeded", "arrjob-download-test") self.raise_once = False + self.progress = None + + def get_acquisition_progress(self, *, request_id, report_date): + return self.progress def execute(self, **kwargs): self.calls.append(kwargs) @@ -200,6 +205,78 @@ class DownloadTests(unittest.TestCase): self.assertTrue(self.service.get("a" * 32)["can_retry"]) self.assertEqual(len(self.executor.calls), 1) + def test_acquisition_snapshot_survives_get_latest_pending_and_restart(self): + self.executor.progress = acquisition_snapshot("completed", 6, 6) + self.executor.outcome = DownloadOutcome("needs_data_review") + self.executor.release.set() + self.service.create("2026-09-15", "a" * 32) + task = self.finished() + self.assertEqual(task["acquisition_progress"], self.executor.progress) + self.assertEqual(self.service.latest()["acquisition_progress"], self.executor.progress) + self.assertEqual(self.service.pending_data_reviews()[0]["acquisition_progress"], self.executor.progress) + self.assertEqual(self.service.pending_reviews()[0]["acquisition_progress"], self.executor.progress) + self.service.close(wait=True) + self.service = PersistentARRDownloads(self.root, self.executor) + self.assertEqual(self.service.get("a" * 32)["acquisition_progress"], self.executor.progress) + self.assertEqual(len(self.executor.calls), 1) + + def test_restart_interrupted_state_overrides_unfinished_source_observation(self): + self.executor.progress = acquisition_snapshot("fetching", 5, 6) + self.executor.release.set() + self.service.create("2026-09-15", "a" * 32) + self.finished() + self.service.close(wait=True) + with sqlite3.connect(self.root / "arr-downloads.sqlite3") as db: + db.execute("UPDATE downloads SET status='downloading'") + self.service = PersistentARRDownloads(self.root, self.executor) + observed = self.service.get("a" * 32)["acquisition_progress"] + self.assertEqual((observed["phase"], observed["completed"], observed["total"]), ("interrupted", 5, 6)) + self.assertLess(observed["percent"], 100) + self.assertEqual(len(self.executor.calls), 1) + + def test_failed_acquisition_never_turns_full_row_count_into_completion(self): + self.executor.progress = acquisition_snapshot("verifying", 6, 6) + self.executor.outcome = DownloadOutcome("failed", retryable=True) + self.executor.release.set() + self.service.create("2026-09-15", "a" * 32) + task = self.finished() + self.assertEqual(task["acquisition_progress"]["phase"], "failed") + self.assertEqual(task["acquisition_progress"]["percent"], 99) + + def test_later_report_failure_preserves_authoritative_acquisition_completion(self): + self.executor.progress = acquisition_snapshot("completed", 6, 6) + self.executor.outcome = DownloadOutcome("failed", "arrjob-report-failed", retryable=True) + self.executor.release.set() + self.service.create("2026-09-15", "a" * 32) + self.assertEqual(self.finished()["acquisition_progress"]["phase"], "completed") + + def test_retry_hides_stale_failed_attempt_until_new_source_observation(self): + self.executor.progress = acquisition_snapshot("failed", 5, 6) + self.executor.raise_once = True + self.executor.release.set() + self.service.create("2026-09-15", "a" * 32) + self.assertEqual(self.finished()["status"], "interrupted") + self.executor.release.clear() + self.executor.entered.clear() + self.service.retry("a" * 32) + self.assertTrue(self.executor.entered.wait(1)) + observed = self.service.get("a" * 32)["acquisition_progress"] + self.assertEqual(observed["phase"], "searching") + self.assertEqual(observed["completed"], 0) + self.assertIsNone(observed["total"]) + + def test_bad_or_unavailable_observations_never_break_task_or_expose_extra_fields(self): + self.executor.release.set() + self.service.create("2026-09-15", "a" * 32) + self.finished() + for value in ({"guest": "private guest"}, {**acquisition_snapshot("completed", 6, 6), "guest": "private guest"}): + self.executor.progress = value + task = self.service.get("a" * 32) + self.assertIsNone(task["acquisition_progress"]) + self.assertNotIn("private guest", json.dumps(task)) + with patch.object(self.executor, "get_acquisition_progress", side_effect=OSError("private file")): + self.assertEqual(self.service.get("a" * 32)["status"], "succeeded") + def test_second_owner_cannot_start_same_queue(self): with self.assertRaises(BlockingIOError): PersistentARRDownloads(self.root, self.executor) @@ -275,13 +352,26 @@ class DownloadRoutesTests(unittest.TestCase): self.assertEqual(response.status, 202) result = json.loads(response.body)["data"] self.assertEqual(result["from_date"], result["to_date"]) + self.assertIn("acquisition_progress", result) get = self.app.handle("GET", "/api/arr-downloads/" + "a" * 32, self.headers) self.assertEqual(get.status, 200) + self.assertIn("acquisition_progress", json.loads(get.body)["data"]) retry_path = "/api/arr-downloads/" + "a" * 32 + "/retry" self.assertEqual(self.app.handle("POST", retry_path, {}, b"{}").status, 401) self.assertEqual(self.app.handle("POST", retry_path, self.headers, b'{"report_date":"2026-09-16"}').status, 400) self.assertEqual(self.app.handle("POST", retry_path, self.headers, b"{}").status, 202) + def test_progress_is_retained_by_public_get_and_configuration_without_new_execution(self): + self.post({"report_date": "2026-09-15", "request_id": "a" * 32}) + self.assertTrue(self.executor.entered.wait(1)) + self.executor.progress = acquisition_snapshot("fetching", 3, 6) + task = json.loads(self.app.handle("GET", "/api/arr-downloads/" + "a" * 32, self.headers).body)["data"] + config = json.loads(self.app.handle("GET", "/api/arr-downloads", self.headers).body)["data"] + self.assertEqual(task["acquisition_progress"], self.executor.progress) + self.assertEqual(config["latest_task"]["acquisition_progress"], self.executor.progress) + self.assertEqual(set(task["acquisition_progress"]), {"phase", "completed", "total", "percent", "updated_at"}) + self.assertEqual(len(self.executor.calls), 1) + def test_default_runtime_is_explicitly_unavailable(self): app = PortalApplication(credentials=TEST_CREDENTIALS) _, headers = login(app)