Files
ARR-2.0-0918/arr_web/arr_data_review.py
T

617 lines
35 KiB
Python

"""Revisioned source-field completion before deterministic OHIP processing.
Original acquisition bytes remain immutable. Staff may resolve only fields that
block the unchanged processor's candidate validation. Finalization freezes a
reproducible derived source with the original observations and decision audit.
"""
from __future__ import annotations
import copy
from contextlib import contextmanager
from datetime import date, datetime, timezone
from decimal import Decimal, InvalidOperation
import hashlib
import importlib.util
import os
from pathlib import Path
import re
import sys
import tempfile
import threading
import uuid
from arr_web.arr_downloads import validate_request_id
from arr_web.contracts import PortalError
from integrations.ohip.audit_arr_capture import protected_read
from integrations.ohip.capture_job import atomic_json, fingerprint, job_lock, private_directory, sync_directory
from integrations.ohip.collect_arr_source import json_bytes, require, strict_json
from integrations.ohip.collect_arr_source import CollectionError
VERSION = "arr-source-field-review/v1"
REANALYSIS_POLICY = "oracle-optional-association/v1"
REANALYSIS_VERSION = "arr-source-reanalysis/v1"
STATUS_POLICY = "arr-exclude-cancelled/v1"
LIMIT = 100 * 1024 * 1024
OPTIONAL = frozenset({"BLOCK_CODE", "RES_COMMENT", "PRODUCTS", "ROOM_CATEGORY_LABEL"})
LABELS = {
"BLOCK_CODE": "团队代码", "ADULTS": "成人人数", "CHILDREN": "儿童人数",
"COMPANY_NAME": "公司名称", "CONFIRMATION_NO": "确认号", "DISP_ROOM_NO": "房号",
"EFFECTIVE_RATE_AMOUNT": "Opera 房价", "FULL_NAME": "客人姓名", "RES_COMMENT": "预订备注",
"NO_OF_ROOMS": "房间数", "PRODUCTS": "包价项目", "RATE_CODE": "费率代码",
"ROOM_CATEGORY_LABEL": "房型", "ARRIVAL": "到店日期", "DEPARTURE": "离店日期",
}
FIELD_ALIASES = {"CF_CHILDREN": "CHILDREN", "TRUNC_BEGIN": "ARRIVAL", "TRUNC_END": "DEPARTURE"}
def _error(code="INVALID", message="请填写有效的字段值", status=400):
return PortalError("ARR_DATA_REVIEW_" + code, message, status)
def _hash(raw):
return hashlib.sha256(raw).hexdigest()
def _write_once(path, raw):
if path.exists():
require(protected_read(path, LIMIT) == raw, "data_review_immutable_file_changed")
return
fd, temporary = tempfile.mkstemp(prefix=".review-", dir=path.parent)
temporary = Path(temporary)
try:
with os.fdopen(fd, "wb") as stream:
stream.write(raw)
stream.flush()
os.fsync(stream.fileno())
os.link(temporary, path) # Publish complete bytes atomically and exclusively.
sync_directory(path.parent)
finally:
temporary.unlink(missing_ok=True)
def _rules(policy):
name = "_arr_source_review_" + uuid.uuid4().hex
spec = importlib.util.spec_from_file_location(name, policy.skill_root / "scripts" / "process_daily.py")
module = importlib.util.module_from_spec(spec)
sys.modules[name] = module
try:
spec.loader.exec_module(module)
finally:
sys.modules.pop(name, None)
require(module.PROCESSOR_VERSION == policy.processor_version
and module.rule_set_sha256() == policy.rule_set_sha256, "data_review_rules_changed")
return module
class DataFieldReviews:
def __init__(self, *, root, policy, context):
self.root = Path(root).absolute()
require(not self.root.resolve().is_relative_to(Path(__file__).resolve().parents[1]),
"data_review_store_must_be_outside_repository")
self.policy, self.context = policy, copy.deepcopy(context)
self.rules = _rules(policy)
self._mutex = threading.RLock()
@contextmanager
def _lock(self, directory):
with self._mutex:
try:
with job_lock(directory):
yield
except CollectionError as error:
if str(error) == "batch_busy":
raise _error("CONFLICT", "数据正在更新,请重新查看后再保存", 409) from None
raise
def _directory(self, request_id, *, create=False):
validate_request_id(request_id)
directory = self.root / request_id
if create:
private_directory(self.root)
private_directory(directory)
sync_directory(self.root)
if not directory.exists():
raise _error("NOT_FOUND", "该任务没有待完善的数据", 404)
private_directory(directory)
return directory
def _read(self, directory):
meta = strict_json(protected_read(directory / "identity.json", 65536))
require(meta["version"] == VERSION and fingerprint(meta["context"]) == fingerprint(self.context)
and meta["request_id"] == directory.name, "data_review_identity_changed")
raw = protected_read(directory / "original.json", LIMIT)
require(_hash(raw) == meta["original_sha256"], "data_review_original_changed")
original = strict_json(raw)
require(original["report_date"] == meta["report_date"] and original.get("collection_complete") is True,
"data_review_source_context_changed")
pointer = strict_json(protected_read(directory / "state.json", 65536))
filename = pointer.get("file", "")
require(bool(re.fullmatch(r"revision-[0-9]{6}-[0-9a-f]{64}\.json", filename)), "data_review_pointer_invalid")
state_raw = protected_read(directory / filename, LIMIT)
require(_hash(state_raw) == pointer["sha256"], "data_review_state_changed")
state = strict_json(state_raw)
require(state["identity_sha256"] == fingerprint(meta) and type(state["revision"]) is int
and state["status"] in {"editing", "finalized"}, "data_review_state_context_changed")
if "source_reanalysis" in state:
reference = state["source_reanalysis"]
require(type(reference) is dict and reference.get("parent_original_sha256") == meta["original_sha256"]
and reference.get("source_manifest_sha256") == meta["source_manifest_sha256"],
"data_review_reanalysis_context_changed")
self._reanalyzed_source(directory, original, state)
self._status_evidence(directory, meta, original, state)
return meta, original, state
def _publish(self, directory, state):
raw = json_bytes(state)
digest = _hash(raw)
filename = f"revision-{state['revision']:06d}-{digest}.json"
_write_once(directory / filename, raw)
atomic_json(directory / "state.json", {"file": filename, "sha256": digest}, replace=True)
@staticmethod
def _observation(field, value):
empty = value == "" or value == []
if empty and field in OPTIONAL:
return {"state": "empty", "value": [] if field in {"PRODUCTS", "RES_COMMENT"} else ""}
if field == "PRODUCTS":
value = [{"package": {"packageCode": code}} for code in value]
elif field == "RES_COMMENT":
value = [value]
return {"state": "available", "value": value}
@staticmethod
def _reanalysis_projection(original, reanalyzed, policy_id):
require(policy_id == REANALYSIS_POLICY and type(reanalyzed) is dict,
"data_review_reanalysis_policy_invalid")
expected = copy.deepcopy(original)
require(type(reanalyzed.get("records")) is list
and len(reanalyzed["records"]) == len(original["records"]),
"data_review_reanalysis_records_changed")
changes = []
allowed = {"BLOCK_CODE": ("missing_reservation_block", ""),
"PRODUCTS": ("missing_reservation_packages", [])}
for index, (old, new) in enumerate(zip(original["records"], reanalyzed["records"])):
require(type(new) is dict and type(new.get("fields")) is dict,
"data_review_reanalysis_record_invalid")
for field, (reason, empty_value) in allowed.items():
before, after = old["fields"][field], new["fields"].get(field)
if fingerprint(before) == fingerprint(after):
continue
empty = {"state": "empty", "value": empty_value}
require(before.get("state") == "missing" and before.get("reason") == reason
and fingerprint(after) == fingerprint(empty),
"data_review_reanalysis_observation_invalid")
expected["records"][index]["fields"][field] = copy.deepcopy(empty)
changes.append({"item_id": f"{old['source_sequence']}:{field}",
"source_sequence": old["source_sequence"], "field": field,
"original_observation": copy.deepcopy(before), "reanalyzed_observation": copy.deepcopy(empty)})
require(bool(changes), "data_review_reanalysis_no_changes")
unresolved = any(observation.get("state") != "available"
and not (field in OPTIONAL and observation.get("state") == "empty")
for row in expected["records"] for field, observation in row["fields"].items())
expected["input_complete"] = not unresolved
expected["status"] = "collected_with_gaps" if unresolved else "collected"
expected["optional_omission_policy"] = REANALYSIS_POLICY
require(fingerprint(expected) == fingerprint(reanalyzed), "data_review_reanalysis_source_changed")
return changes
def _reanalyzed_source(self, directory, original, state):
reference = state.get("source_reanalysis")
if reference is None:
return original, None
require(type(reference) is dict and reference.get("policy_id") == REANALYSIS_POLICY,
"data_review_reanalysis_reference_invalid")
digest, receipt_digest = reference.get("sha256", ""), reference.get("receipt_sha256", "")
manifest_digest = reference.get("manifest_sha256", "")
require(type(digest) is str and bool(re.fullmatch(r"[0-9a-f]{64}", digest))
and type(receipt_digest) is str and bool(re.fullmatch(r"[0-9a-f]{64}", receipt_digest))
and type(manifest_digest) is str and bool(re.fullmatch(r"[0-9a-f]{64}", manifest_digest))
and reference.get("file") == f"reanalyzed-source-{digest}.json"
and reference.get("receipt_file") == f"source-reanalysis-{receipt_digest}.json",
"data_review_reanalysis_reference_invalid")
payload = protected_read(directory / reference["file"], LIMIT)
receipt_raw = protected_read(directory / reference["receipt_file"], LIMIT)
require(_hash(payload) == digest and _hash(receipt_raw) == receipt_digest,
"data_review_reanalysis_changed")
reanalyzed, receipt = strict_json(payload), strict_json(receipt_raw)
changes = self._reanalysis_projection(original, reanalyzed, REANALYSIS_POLICY)
expected_receipt = {"version": REANALYSIS_VERSION, "policy_id": REANALYSIS_POLICY,
"request_id": directory.name, "report_date": original["report_date"], "context": self.context,
"parent_original_sha256": reference["parent_original_sha256"],
"source_manifest_sha256": reference["source_manifest_sha256"],
"reanalyzed_sha256": digest, "reanalysis_manifest_sha256": reference["manifest_sha256"],
"changes": changes}
require(fingerprint(receipt) == fingerprint(expected_receipt), "data_review_reanalysis_receipt_changed")
return reanalyzed, receipt
def _status_evidence(self, directory, meta, original, state):
reference = state.get("reservation_status_evidence")
if reference is None:
return None
require(type(reference) is dict and reference.get("policy_id") == STATUS_POLICY,
"data_review_status_reference_invalid")
digest = reference.get("sha256", "")
require(type(digest) is str and bool(re.fullmatch(r"[0-9a-f]{64}", digest))
and reference.get("file") == f"reservation-status-{digest}.json",
"data_review_status_reference_invalid")
raw = protected_read(directory / reference["file"], LIMIT)
require(_hash(raw) == digest, "data_review_status_evidence_changed")
evidence = strict_json(raw)
require(evidence.get("version") == "arr-reservation-status-evidence/v1"
and evidence.get("policy_id") == STATUS_POLICY
and evidence.get("original_sha256") == meta["original_sha256"]
and evidence.get("source_manifest_sha256") == meta["source_manifest_sha256"]
and evidence.get("report_date") == meta["report_date"]
and evidence.get("hotel_id") == original["hotel_id"]
and type(evidence.get("records")) is list
and len(evidence["records"]) == len(original["records"]),
"data_review_status_context_changed")
for row, status in zip(original["records"], evidence["records"]):
require(type(status) is dict and status.get("source_sequence") == row["source_sequence"]
and status.get("reservation_id") == row["reservation_id"]
and type(status.get("reservation_status")) is str
and bool(status["reservation_status"].strip())
and ("reservation_status" not in row
or row["reservation_status"] == status["reservation_status"])
and type(status.get("sources")) is list and bool(status["sources"]),
"data_review_status_records_changed")
return evidence
def _derive(self, directory, original, state):
source, _receipt = self._reanalyzed_source(directory, original, state)
data = copy.deepcopy(source)
meta = strict_json(protected_read(directory / "identity.json", 65536))
evidence = self._status_evidence(directory, meta, original, state)
if evidence is not None:
for row, status in zip(data["records"], evidence["records"]):
row["reservation_status"] = status["reservation_status"]
for item_id, decision in state["decisions"].items():
sequence, field = item_id.split(":")
row = data["records"][int(sequence) - 1]
require(row["source_sequence"] == int(sequence) and field in LABELS, "data_review_decision_invalid")
row["fields"][field] = {**self._observation(field, decision["value"]),
"origin": "manual_review", "review_item_id": item_id}
unresolved = any(observation.get("state") != "available"
and not (field in OPTIONAL and observation.get("state") == "empty")
for row in data["records"] for field, observation in row["fields"].items())
data["input_complete"] = not unresolved
data["status"] = "collected_with_gaps" if unresolved else "collected"
return data
def _issues(self, directory, data):
# Use the exact frozen processor's normalizer and field validation; the
# acquisition layer never maintains a second whitelist or numeric rule.
with tempfile.NamedTemporaryFile(suffix=".json", dir=directory) as stream:
stream.write(json_bytes(data))
stream.flush()
business_date, records = self.rules.read_data_source(Path(stream.name))
issues = {}
for row, record in zip(data["records"], records):
if self.rules.is_cancelled_record(record):
continue
rate = record["_NORMALIZED_RATE_CODE"]
fields = {}
if not rate or "RATE_CODE" in record["_DATA_GAPS"]:
fields["RATE_CODE"] = "rate_code_unresolved"
elif rate in self.rules.RATE_WHITELIST:
for field in record["_DATA_GAPS"]:
fields[field] = row["fields"][field].get("reason", "field_unresolved")
for error in self.rules.validate_whitelisted_record(record, business_date):
field = str(error.source_location or "").rsplit("/", 1)[-1]
if error.code == "DATA_NEGATIVE_NIGHTS":
field = "DEPARTURE"
field = FIELD_ALIASES.get(field, field)
if field in LABELS:
fields.setdefault(field, error.code)
for field, reason in fields.items():
issues[f"{row['source_sequence']}:{field}"] = str(reason)
return issues
@staticmethod
def _display(row, field):
observation = row["fields"][field]
value = observation.get("value")
return value if observation.get("state") == "available" and isinstance(value, (str, int)) else ""
def _public(self, directory, meta, original, state):
data = self._derive(directory, original, state)
issues = self._issues(directory, data)
cancelled_sequences = {row["source_sequence"] for row in data["records"]
if self.rules.is_cancelled_record({"_RESERVATION_STATUS": row.get("reservation_status", "")})}
keys = set(issues) | {key for key in state["decisions"]
if int(key.split(":")[0]) not in cancelled_sequences}
items = []
for item_id in sorted(keys, key=lambda key: (int(key.split(":")[0]), key.split(":")[1])):
sequence, field = item_id.split(":")
row = data["records"][int(sequence) - 1]
decision = state["decisions"].get(item_id)
items.append({"item_id": item_id, "source_sequence": int(sequence), "field": field,
"field_label": LABELS[field], "confirmation_no": self._display(row, "CONFIRMATION_NO"),
"room_no": self._display(row, "DISP_ROOM_NO"), "company_name": self._display(row, "COMPANY_NAME"),
"rate_code": self._display(row, "RATE_CODE"), "reason_code": issues.get(item_id, "confirmed"),
"source_state": row["fields"][field].get("state"),
"can_be_empty": field in OPTIONAL, "value": decision["value"] if decision else None,
"confirmed": decision is not None and item_id not in issues})
return {"request_id": meta["request_id"], "report_date": meta["report_date"],
"revision": state["revision"], "status": state["status"], "items": items,
"pending_count": len(issues), "total_count": len(items),
"excluded_cancelled_count": len(cancelled_sequences),
"can_finalize": not issues and state["status"] == "editing"}
def prepare(self, request_id, payload, manifest_sha256, report_date):
directory = self._directory(request_id, create=True)
with self._lock(directory):
meta = {"version": VERSION, "context": self.context, "request_id": request_id,
"report_date": report_date, "original_sha256": _hash(payload),
"source_manifest_sha256": manifest_sha256}
if not (directory / "identity.json").exists():
require(bool(re.fullmatch(r"[0-9a-f]{64}", manifest_sha256)), "data_review_manifest_invalid")
original = strict_json(payload)
require(original.get("report_date") == report_date and original.get("collection_complete") is True,
"data_review_requires_complete_collection")
for key in ("hotel_id", "source_kind"):
if key in self.context:
require(original.get(key) == self.context[key], "data_review_source_context_mismatch")
_write_once(directory / "original.json", payload)
atomic_json(directory / "identity.json", meta, replace=False)
require(fingerprint(strict_json(protected_read(directory / "identity.json", 65536))) == fingerprint(meta),
"data_review_prepare_conflict")
if not (directory / "state.json").exists():
state = {"identity_sha256": fingerprint(meta), "revision": 0, "status": "editing",
"decisions": {}, "events": []}
initial_file = f"revision-000000-{_hash(json_bytes(state))}.json"
require({p.name for p in directory.glob("revision-*.json")} <= {initial_file}
and not (directory / "finalize-intent.json").exists(), "data_review_state_missing")
self._publish(directory, state)
actual, original, state = self._read(directory)
require(fingerprint(actual) == fingerprint(meta), "data_review_prepare_conflict")
review = self._public(directory, actual, original, state)
# A complete ordinary source needs no human step or derived source.
# A reanalysis remains an explicit review until finalized, even if
# its validated source interpretation resolves every listed gap.
return review if review["total_count"] or "source_reanalysis" in state or "reservation_status_evidence" in state else None
def get(self, request_id):
directory = self._directory(request_id)
with self._lock(directory):
return self._public(directory, *self._read(directory))
def apply_saved_status_evidence(self, request_id, capture_root, *, expected_revision):
"""Bind saved search/detail statuses; no HTTP route or new business read.
Evidence is reconstructed from the complete original capture, never from
staff-entered status values. Fields, order and original bytes stay intact.
"""
from integrations.ohip.reservation_status import saved_status_evidence
directory = self._directory(request_id)
with self._lock(directory):
meta, original, state = self._read(directory)
if state["status"] != "editing" or (directory / "finalize-intent.json").exists():
raise _error("FROZEN", "已确认生成,预订状态不能再补充", 409)
evidence = saved_status_evidence(capture_root, meta["source_manifest_sha256"],
protected_read(directory / "original.json", LIMIT))
raw = json_bytes(evidence)
digest = _hash(raw)
if type(expected_revision) is not int:
self._revision(state, expected_revision)
prior = state.get("reservation_status_evidence")
if prior is not None:
if prior.get("sha256") == digest:
return self._public(directory, meta, original, state)
raise _error("CONFLICT", "该来源已补充预订状态,请核对现有记录", 409)
self._revision(state, expected_revision)
reference = {"policy_id": STATUS_POLICY, "file": f"reservation-status-{digest}.json", "sha256": digest}
_write_once(directory / reference["file"], raw)
state["reservation_status_evidence"] = reference
state["revision"] += 1
state["events"].append({"revision": state["revision"], "action": "reservation_status_evidence",
"policy_id": STATUS_POLICY, "evidence_sha256": digest,
"at": datetime.now(timezone.utc).isoformat()})
result = self._public(directory, meta, original, state)
self._publish(directory, state)
return result
def apply_source_reanalysis(self, request_id, payload, manifest_sha256, policy_id, *, expected_revision):
"""Attach a validated offline reinterpretation, without changing acquisition or staff decisions.
This internal maintenance operation has no HTTP route. Its caller must
establish the policy's omission evidence from the saved complete capture.
The projection here prevents any other source or business-value edits.
"""
require(policy_id == REANALYSIS_POLICY, "data_review_reanalysis_policy_invalid")
require(type(payload) is bytes and len(payload) <= LIMIT,
"data_review_reanalysis_payload_invalid")
require(type(manifest_sha256) is str and bool(re.fullmatch(r"[0-9a-f]{64}", manifest_sha256)),
"data_review_reanalysis_manifest_invalid")
reanalyzed = strict_json(payload)
directory = self._directory(request_id)
with self._lock(directory):
meta, original, state = self._read(directory)
if state["status"] != "editing" or (directory / "finalize-intent.json").exists():
raise _error("FROZEN", "已确认生成,来源解释不能再修改", 409)
if type(expected_revision) is not int:
self._revision(state, expected_revision)
digest = _hash(payload)
prior = state.get("source_reanalysis")
if prior is not None:
if prior["sha256"] == digest and prior["manifest_sha256"] == manifest_sha256:
return self._public(directory, meta, original, state)
raise _error("CONFLICT", "该来源已完成重新解释,请核对现有记录", 409)
self._revision(state, expected_revision)
changes = self._reanalysis_projection(original, reanalyzed, policy_id)
receipt = {"version": REANALYSIS_VERSION, "policy_id": policy_id,
"request_id": request_id, "report_date": meta["report_date"], "context": self.context,
"parent_original_sha256": meta["original_sha256"],
"source_manifest_sha256": meta["source_manifest_sha256"], "reanalyzed_sha256": digest,
"reanalysis_manifest_sha256": manifest_sha256, "changes": changes}
receipt_raw = json_bytes(receipt)
receipt_digest = _hash(receipt_raw)
reference = {"policy_id": policy_id, "file": f"reanalyzed-source-{digest}.json", "sha256": digest,
"receipt_file": f"source-reanalysis-{receipt_digest}.json", "receipt_sha256": receipt_digest,
"parent_original_sha256": meta["original_sha256"],
"source_manifest_sha256": meta["source_manifest_sha256"], "manifest_sha256": manifest_sha256}
_write_once(directory / reference["file"], payload)
_write_once(directory / reference["receipt_file"], receipt_raw)
state["source_reanalysis"] = reference
state["revision"] += 1
state["events"].append({"revision": state["revision"], "action": "source_reanalysis",
"policy_id": policy_id, "parent_original_sha256": meta["original_sha256"],
"reanalyzed_sha256": digest, "reanalysis_manifest_sha256": manifest_sha256,
"receipt_sha256": receipt_digest, "at": datetime.now(timezone.utc).isoformat()})
result = self._public(directory, meta, original, state)
self._publish(directory, state)
return result
@staticmethod
def _revision(state, revision):
if type(revision) is not int or revision != state["revision"]:
raise _error("CONFLICT", "数据已更新,请重新查看后再保存", 409)
@staticmethod
def _actor(actor):
if not isinstance(actor, str) or not actor.strip() or len(actor) > 128 or re.search(r"[\x00-\x1f\x7f]", actor):
raise _error()
return actor
def _value(self, field, value, report_date, data, sequence):
if field == "PRODUCTS":
if (not isinstance(value, list) or len(value) > 100 or any(not isinstance(code, str)
or not code.strip() or len(code) > 100 or re.search(r"[\x00-\x1f\x7f]", code) for code in value)):
raise _error(message="请按顺序填写有效的包价代码,或确认无此项")
return [code.strip() for code in value]
if isinstance(value, bool) or not isinstance(value, (str, int)):
raise _error()
value = str(value).strip()
if len(value) > 4000 or re.search(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]", value):
raise _error()
if not value:
if field in OPTIONAL:
return ""
raise _error(message="此字段必须填写,不能确认为空")
if field in {"ADULTS", "CHILDREN", "NO_OF_ROOMS"}:
if not re.fullmatch(r"[0-9]{1,8}", value):
raise _error(message="人数和房间数必须是非负整数")
number = int(value)
if field == "NO_OF_ROOMS" and number == 0:
raise _error(message="房间数必须大于零")
return number
if field == "EFFECTIVE_RATE_AMOUNT":
if not re.fullmatch(r"[0-9]{1,16}(?:\.[0-9]{1,2})?", value):
raise _error(message="Opera 房价必须是非负金额,最多两位小数")
try:
amount = Decimal(value)
except InvalidOperation:
raise _error() from None
return format(amount, "f")
if field in {"ARRIVAL", "DEPARTURE"}:
try:
day = date.fromisoformat(value)
if day.isoformat() != value:
raise ValueError()
except ValueError:
raise _error(message="日期请填写 YYYY-MM-DD") from None
if (field == "ARRIVAL" and value != report_date) or (field == "DEPARTURE" and value < report_date):
raise _error(message="到店日期须与报表日期一致,离店日期不得早于到店日期")
return value.upper() if field == "RATE_CODE" else value
def update(self, request_id, item_id, revision, value, actor):
directory = self._directory(request_id)
with self._lock(directory):
meta, original, state = self._read(directory)
if state["status"] != "editing" or (directory / "finalize-intent.json").exists():
raise _error("FROZEN", "已确认生成,字段不能再修改", 409)
self._revision(state, revision)
actor = self._actor(actor)
data = self._derive(directory, original, state)
sequence = int(item_id.split(":")[0]) if isinstance(item_id, str) and re.fullmatch(r"[1-9][0-9]*:[A-Z_]+", item_id) else 0
if 1 <= sequence <= len(data["records"]) and self.rules.is_cancelled_record(
{"_RESERVATION_STATUS": data["records"][sequence - 1].get("reservation_status", "")}):
raise _error(message="已取消预订不参与日报,无需完善字段")
if item_id not in self._issues(directory, data) and item_id not in state["decisions"]:
raise _error(message="只能完善当前任务列出的异常字段")
sequence, field = item_id.split(":")
value = self._value(field, value, meta["report_date"], data, int(sequence))
previous = copy.deepcopy(state["decisions"].get(item_id))
state["revision"] += 1
state["decisions"][item_id] = {"value": value, "actor": actor, "revision": state["revision"]}
state["events"].append({"revision": state["revision"], "action": "set_field", "item_id": item_id,
"actor": actor, "value": value, "previous": previous, "at": datetime.now(timezone.utc).isoformat()})
result = self._public(directory, meta, original, state)
if item_id in {item["item_id"] for item in result["items"] if not item["confirmed"]}:
raise _error(message="该字段仍不满足报表规则,请核对填写内容")
self._publish(directory, state)
return result
def _frozen(self, directory, meta, original, state):
changes = []
for item_id, decision in sorted(state["decisions"].items()):
sequence, field = item_id.split(":")
changes.append({"item_id": item_id, "source_sequence": int(sequence), "field": field,
"original_observation": original["records"][int(sequence) - 1]["fields"][field],
**decision})
audit = {"version": VERSION, "request_id": meta["request_id"], "report_date": meta["report_date"],
"context": self.context, "original_sha256": meta["original_sha256"],
"source_manifest_sha256": meta["source_manifest_sha256"], "revision": state["revision"],
"changes": changes, "events": state["events"]}
_source, receipt = self._reanalyzed_source(directory, original, state)
if receipt is not None:
audit["source_reanalysis"] = {**receipt, "receipt_sha256": state["source_reanalysis"]["receipt_sha256"]}
evidence = self._status_evidence(directory, meta, original, state)
if evidence is not None:
audit["reservation_status_evidence"] = {**evidence,
"evidence_sha256": state["reservation_status_evidence"]["sha256"]}
data = self._derive(directory, original, state)
data["manual_data_review"] = {"manifest": audit, "manifest_sha256": fingerprint(audit)}
binding = fingerprint({"source_manifest_sha256": meta["source_manifest_sha256"],
"review_manifest_sha256": fingerprint(audit), "data_sha256": _hash(json_bytes(data))})
return json_bytes(data), binding
def finalize(self, request_id, revision, actor):
directory = self._directory(request_id)
with self._lock(directory):
meta, original, state = self._read(directory)
if state["status"] == "finalized":
self._verify_frozen(directory, meta, original, state)
return self._public(directory, meta, original, state)
self._revision(state, revision)
actor = self._actor(actor)
if self._issues(directory, self._derive(directory, original, state)):
raise _error("INCOMPLETE", "请先完善并确认所有异常字段", 409)
intent = directory / "finalize-intent.json"
prior_sha256 = fingerprint(state)
if intent.exists():
frozen = strict_json(protected_read(intent, LIMIT))
require(frozen["prior_sha256"] == prior_sha256, "data_review_finalize_context_changed")
state = frozen["state"]
else:
state["revision"] += 1
state["status"] = "finalized"
state["events"].append({"revision": state["revision"], "action": "finalize", "actor": actor,
"at": datetime.now(timezone.utc).isoformat()})
atomic_json(intent, {"prior_sha256": prior_sha256, "state": state}, replace=False)
payload, binding = self._frozen(directory, meta, original, state)
_write_once(directory / "reviewed-source.json", payload)
state["derived_sha256"], state["binding_manifest_sha256"] = _hash(payload), binding
self._publish(directory, state)
return self._public(directory, meta, original, state)
def _verify_frozen(self, directory, meta, original, state):
payload, binding = self._frozen(directory, meta, original, state)
require(_hash(payload) == state["derived_sha256"] and binding == state["binding_manifest_sha256"]
and protected_read(directory / "reviewed-source.json", LIMIT) == payload,
"data_review_frozen_source_changed")
return payload, binding
def payload(self, request_id):
directory = self._directory(request_id)
with self._lock(directory):
meta, original, state = self._read(directory)
if state["status"] != "finalized":
return None
return self._verify_frozen(directory, meta, original, state)
def original(self, request_id):
directory = self._directory(request_id)
with self._lock(directory):
meta, _original, _state = self._read(directory)
return protected_read(directory / "original.json", LIMIT), meta["source_manifest_sha256"]