276 lines
13 KiB
Python
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()
|