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)