Files

276 lines
13 KiB
Python

"""Local storage and UI with an explicitly configured OHIP sandbox source.
Startup never calls Oracle. check-access reads platform application/key metadata only.
New business reads require a user download action and complete application grants.
"""
from __future__ import annotations
import argparse
from contextlib import contextmanager
from datetime import datetime, timezone
from html import escape
import json
import os
from pathlib import Path
import secrets
import signal
import subprocess
import tempfile
from arr_processing.policy import load_processor_policy
from arr_web.arr_data_executor import DirectARRExecutor
from arr_web.contracts import PortalError
from arr_web.local_replay import ReplayRuntime
from arr_web.local_replay_database import ReplayDatabase
from arr_web.local_xml_replay import LocalReplayPortal, PROJECT, document
from arr_web.server import serve
from integrations.ohip.arr_data import ARRDataSource
from integrations.ohip.capture_job import atomic_json, job_lock, private_directory
from integrations.ohip.collect_arr_source import APPLICATION, SERVICE, CollectionError, load_key
HOTEL = "OHIPSB02"
OWNER = "admin_DUMJ3Obu5Z_O627b"
PRINCIPAL = "auto_219f3af6dc2123dc29ee3138"
REQUIRED = frozenset({"reservations.read", "profiles.read", "blocks.read"})
def identity(root, credential_file, policy):
return {"version": "arr-local-ohip-sandbox/v1", "root": str(root),
"source_kind": "ohip_platform", "service_url": SERVICE, "application_id": APPLICATION,
"hotel_id": HOTEL, "credential_file": str(credential_file),
"processor_version": policy.processor_version, "rule_set_sha256": policy.rule_set_sha256}
def create(parent, credential_file, login_file=None):
parent = parent.expanduser().resolve()
if parent.is_relative_to(PROJECT):
raise ValueError("local_ohip_state_must_be_outside_repository")
credential_file = credential_file.expanduser().absolute()
load_key(credential_file) # Local validation only; no value is logged or copied.
login = document(login_file) if login_file else {"username": "arr-data", "password": secrets.token_urlsafe(24)}
from arr_web.auth import LoginCredentials
LoginCredentials(**login)
private_directory(parent)
root = Path(tempfile.mkdtemp(prefix="ohip-sandbox-", dir=parent)).resolve()
atomic_json(root / "instance.json", identity(root, credential_file, load_processor_policy(PROJECT)), replace=False)
atomic_json(root / "login.json", login, replace=False)
return root
class AccessState:
def __init__(self, root):
self.root = root
def status(self):
try:
data = document(self.root / "access.json")
if (data.get("version") != "arr-ohip-access/v1" or data.get("application_id") != APPLICATION
or data.get("service_url") != SERVICE or data.get("hotel_id") != HOTEL
or data.get("owner_user_id") != OWNER or data.get("automation_principal_id") != PRINCIPAL):
raise ValueError("access_identity_mismatch")
missing = sorted(REQUIRED - set(data["capability_groups"]))
ready = data.get("enabled") is True and data.get("key_active") is True and not missing
return {"ready": ready, "missing_capabilities": missing, "checked_at": data.get("checked_at")}
except (OSError, ValueError, KeyError, TypeError, CollectionError):
return {"ready": False, "missing_capabilities": sorted(REQUIRED), "checked_at": None}
def require_ready(self):
if not self.status()["ready"]:
raise PortalError("OHIP_ACCESS_NOT_READY", "平台读取权限尚未补齐,请联系平台负责人。", 503)
class PermissionDownloads:
def __init__(self, queue, access):
self.queue, self.access = queue, access
@property
def ready(self):
return self.queue.ready and self.access.status()["ready"]
def latest(self): return self.queue.latest()
def pending_data_reviews(self): return self.queue.pending_data_reviews()
def pending_reviews(self): return self.queue.pending_reviews()
def get(self, request_id): return self.queue.get(request_id)
@property
def context_id(self): return self.queue.context_id
def get_data_review(self, request_id): return self.queue.get_data_review(request_id)
def update_data_review_item(self, *args): return self.queue.update_data_review_item(*args)
def finalize_data_review(self, *args):
self.access.require_ready()
return self.queue.finalize_data_review(*args)
def create(self, report_date, request_id):
self.access.require_ready()
return self.queue.create(report_date, request_id)
def retry(self, request_id):
self.access.require_ready()
return self.queue.retry(request_id)
class AuthorizedExecutor:
def __init__(self, executor, access):
self.executor, self.access = executor, access
def execute(self, **kwargs):
self.access.require_ready()
return self.executor.execute(**kwargs)
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)
def finalize_data_review(self, *args): return self.executor.finalize_data_review(*args)
def recover_data_review(self, request_id):
recover = getattr(self.executor, "recover_data_review", None)
return recover(request_id) if callable(recover) else False
class LocalOHIPPortal(LocalReplayPortal):
allow_xml_upload = True
cookie_name = "arr_local_ohip_session"
environment = "local-ohip-sandbox"
source_kind = "ohip_platform"
filename_prefix = "OHIP-SANDBOX-"
page_title = "OHIP 沙箱"
def _health_ready(self):
# The process can be healthy while the download button awaits grants.
return self.app._health.database_ready
def _download_context(self):
return {"environment": self.environment, "source_kind": self.source_kind, "hotel_id": HOTEL,
"oracle_connection_verified": False, "access": self.snapshot.status()}
def _banner_html(self):
status = self.snapshot.status()
if status["ready"]:
message = "平台连接已配置。选日期后点击下载才会查询订单,报表保存在本机。"
elif status["missing_capabilities"] == ["blocks.read"]:
message = "连接已配置,等待平台补齐团队资料读取权限;下载按钮暂不可用。"
else:
message = "连接已配置,等待核对平台访问权限;下载按钮暂不可用。"
return (f'<aside class="local-replay-banner" role="note"><strong>本机运行 · OHIP 沙箱 {escape(HOTEL)}</strong>'
f' · {message}尚未验证实际取数。</aside>')
def check_access(root, automation_file, cli):
"""Only platform control-plane commands are permitted here, never business requests."""
config = document(root / "instance.json")
if config != identity(root, Path(config["credential_file"]), load_processor_policy(PROJECT)):
raise ValueError("local_ohip_identity_changed")
# Mark unavailable before checking so a failed refresh cannot leave stale readiness.
atomic_json(root / "access.json", {"version": "arr-ohip-access-pending/v1"}, replace=True)
env = dict(os.environ)
env.pop("OHIP_EDGE_AUTOMATION_TOKEN", None)
env.update(OHIP_EDGE_URL=SERVICE, OHIP_EDGE_AUTOMATION_FILE=str(automation_file), OHIP_EDGE_TIMEOUT="20s")
request_ids = []
def read(*arguments):
result = subprocess.run([str(cli), *arguments], env=env, capture_output=True, text=True, timeout=30)
reply = json.loads(result.stdout)
if result.returncode or reply.get("ok") is not True:
raise RuntimeError("platform_access_check_failed")
request_ids.append(reply.get("request_id"))
return reply["data"]
principal = read("automation", "whoami")["principal"]
if principal.get("id") != PRINCIPAL or principal.get("owner_user_id") != OWNER:
raise ValueError("platform_principal_mismatch")
apps = read("integration", "list")["items"]
selected = [x for x in apps if x.get("id") == APPLICATION]
if len(selected) != 1:
raise ValueError("platform_application_missing")
app = selected[0]
if app.get("owner_user_id") != OWNER or app.get("automation_principal_id") != PRINCIPAL:
raise ValueError("platform_application_owner_mismatch")
credential_file = Path(config["credential_file"])
load_key(credential_file)
issuance = document(credential_file).get("issuance_id")
keys = read("integration", "key", "list", "--application-id", APPLICATION)["items"]
now = datetime.now(timezone.utc)
matches = [k for k in keys if issuance and k.get("issuance_id") == issuance and not k.get("revoked_at")
and (not k.get("expires_at") or datetime.fromisoformat(k["expires_at"].replace("Z", "+00:00")) > now)]
access = {"version": "arr-ohip-access/v1", "service_url": SERVICE, "application_id": APPLICATION,
"owner_user_id": OWNER, "automation_principal_id": PRINCIPAL, "hotel_id": HOTEL,
"enabled": app.get("enabled") is True, "capability_groups": app.get("capability_groups", []),
"key_active": len(matches) == 1, "key_id": matches[0]["id"] if len(matches) == 1 else None,
"checked_at": now.isoformat(), "request_ids": request_ids}
atomic_json(root / "access.json", access, replace=True)
return AccessState(root).status()
@contextmanager
def open_portal(root, port):
if root.is_symlink() or root.resolve().is_relative_to(PROJECT) or not 1024 <= port <= 65535:
raise ValueError("invalid_local_ohip_instance")
root = root.resolve(); private_directory(root)
config = document(root / "instance.json")
policy = load_processor_policy(PROJECT)
credential_file = Path(config["credential_file"])
if config != identity(root, credential_file, policy):
raise ValueError("local_ohip_identity_or_rules_changed")
load_key(credential_file)
access = AccessState(root)
with job_lock(root):
database = ReplayDatabase(root, schema_version=21)
try:
database.start()
def factory(*, root, snapshot, **dependencies):
dependencies["repository"].assert_data_source_schema()
source = ARRDataSource(root / "source", HOTEL, credential_file=credential_file)
executor = DirectARRExecutor(root=root / "processing", source=source, **dependencies)
return AuthorizedExecutor(executor, access)
runtime = ReplayRuntime(root, database, access, policy, port, executor_factory=factory,
portal_type=LocalOHIPPortal, enable_upload=True, downloads_wrapper=PermissionDownloads)
try:
yield runtime
finally:
runtime.close()
finally:
database.close()
def main():
parser = argparse.ArgumentParser(description=__doc__)
actions = parser.add_subparsers(dest="action", required=True)
init = actions.add_parser("init")
init.add_argument("--parent", type=Path, required=True)
init.add_argument("--credential-file", type=Path, required=True)
init.add_argument("--login-file", type=Path)
check = actions.add_parser("check-access")
check.add_argument("--root", type=Path, required=True)
check.add_argument("--automation-file", type=Path, required=True)
check.add_argument("--cli", type=Path, default=Path.home() / ".local/bin/ohipctl")
run = actions.add_parser("serve")
run.add_argument("--root", type=Path, required=True)
run.add_argument("--port", type=int, default=8875)
args = parser.parse_args(); os.umask(0o077)
if args.action == "init":
print(create(args.parent, args.credential_file, args.login_file)); return
if args.action == "check-access":
print(json.dumps(check_access(args.root.resolve(), args.automation_file, args.cli))); return
def stop(*_): raise KeyboardInterrupt
signal.signal(signal.SIGTERM, stop)
try:
with open_portal(args.root, args.port) as runtime:
runtime.start_monthly_worker()
print(f"Local OHIP sandbox: http://127.0.0.1:{args.port}/", flush=True)
serve(runtime.app, "127.0.0.1", args.port)
except KeyboardInterrupt:
pass
if __name__ == "__main__": main()