152 lines
7.6 KiB
Python
152 lines
7.6 KiB
Python
from __future__ import annotations
|
|
|
|
import argparse
|
|
import json
|
|
import sys
|
|
from dataclasses import asdict
|
|
|
|
from .config import settings
|
|
from .database import SessionLocal
|
|
from .importer import import_records
|
|
from .community_importer import stage_observations, update_waterbody_detail, update_waterbody_details, upsert_waterbody_catalog
|
|
from .retention import RetentionPolicy, apply_retention
|
|
from .storage import delete_screenshot
|
|
from .catalog_audit import audit_catalog, audit_waterbody_catalog
|
|
from .community_scheduler import run_source, configured_sources
|
|
|
|
# Static registry for argparse choices — no DB required for --help
|
|
STATIC_SOURCE_CHOICES = [
|
|
"rf4db",
|
|
"rf4stat-fishing",
|
|
"rf4stat-post",
|
|
"rf4map",
|
|
"rf4posts-spot",
|
|
]
|
|
|
|
|
|
def main() -> int:
|
|
parser = argparse.ArgumentParser(prog="python -m app.cli")
|
|
sub = parser.add_subparsers(dest="command", required=True)
|
|
command = sub.add_parser("import-records")
|
|
command.add_argument("--url", default="https://rf4game.de/records/region/RU/")
|
|
command.add_argument("--region", default="RU")
|
|
command.add_argument("--category", default="records")
|
|
community = sub.add_parser("stage-community-json")
|
|
community.add_argument("--input", default="-", help="JSON array path or - for stdin")
|
|
community.add_argument("--limit", type=int, default=500)
|
|
waterbodies = sub.add_parser("import-waterbody-catalog")
|
|
waterbodies.add_argument("--input", required=True, help="JSON snapshot path or - for stdin")
|
|
waterbodies.add_argument("--limit", type=int, default=100)
|
|
detail = sub.add_parser("import-waterbody-detail")
|
|
detail.add_argument("--input", required=True, help="JSON detail snapshot path")
|
|
details = sub.add_parser("import-waterbody-details")
|
|
details.add_argument("--input", required=True, help="JSON array of detail snapshots")
|
|
fetch_community = sub.add_parser("fetch-community")
|
|
fetch_community.add_argument("source", choices=STATIC_SOURCE_CHOICES)
|
|
cleanup = sub.add_parser("cleanup-retention")
|
|
cleanup.add_argument("--apply", action="store_true", help="apply changes; default is dry-run")
|
|
sub.add_parser("audit-catalog")
|
|
waterbody_audit = sub.add_parser("audit-waterbody-catalog")
|
|
waterbody_audit.add_argument("--input", required=True, help="JSON snapshot path")
|
|
args = parser.parse_args()
|
|
with SessionLocal() as session:
|
|
if args.command == "import-records":
|
|
run = import_records(session, url=args.url, region=args.region, category=args.category)
|
|
print(f"import {run.status.value}: seen={run.rows_seen} created={run.rows_created} updated={run.rows_updated}")
|
|
elif args.command == "stage-community-json":
|
|
if not 1 <= args.limit <= 5000:
|
|
parser.error("--limit must be between 1 and 5000")
|
|
stream = sys.stdin if args.input == "-" else open(args.input, encoding="utf-8")
|
|
try:
|
|
payload = json.load(stream)
|
|
finally:
|
|
if stream is not sys.stdin:
|
|
stream.close()
|
|
if not isinstance(payload, list):
|
|
parser.error("input must be a JSON array")
|
|
created, updated = stage_observations(session, payload[:args.limit])
|
|
print(f"staged: created={created} updated={updated}")
|
|
elif args.command == "import-waterbody-catalog":
|
|
if not 1 <= args.limit <= 500:
|
|
parser.error("--limit must be between 1 and 500")
|
|
stream = sys.stdin if args.input == "-" else open(args.input, encoding="utf-8")
|
|
try:
|
|
snapshot = json.load(stream)
|
|
finally:
|
|
if stream is not sys.stdin:
|
|
stream.close()
|
|
if isinstance(snapshot, dict):
|
|
payload = snapshot.get("items")
|
|
source_system = snapshot.get("source_system")
|
|
if isinstance(payload, list) and isinstance(source_system, str):
|
|
payload = [
|
|
{"source_system": source_system, **item}
|
|
for item in payload if isinstance(item, dict)
|
|
]
|
|
else:
|
|
payload = snapshot
|
|
if not isinstance(payload, list):
|
|
parser.error("input must be a JSON array or an object with an items array")
|
|
created, updated = upsert_waterbody_catalog(session, payload[:args.limit])
|
|
print(f"waterbodies: created={created} updated={updated}")
|
|
elif args.command == "import-waterbody-detail":
|
|
with open(args.input, encoding="utf-8") as stream:
|
|
payload = json.load(stream)
|
|
if not isinstance(payload, dict):
|
|
parser.error("input must be a JSON object")
|
|
update_waterbody_detail(session, payload)
|
|
print(f"waterbody detail: updated={payload.get('source_external_id', 'unknown')}")
|
|
elif args.command == "import-waterbody-details":
|
|
with open(args.input, encoding="utf-8") as stream:
|
|
payload = json.load(stream)
|
|
if not isinstance(payload, list):
|
|
parser.error("input must be a JSON array")
|
|
created, updated = update_waterbody_details(session, payload)
|
|
print(f"waterbody details: created={created} updated={updated}")
|
|
elif args.command == "fetch-community":
|
|
# A09: Verify source is enabled at runtime (not just in static choices)
|
|
enabled = configured_sources()
|
|
if args.source not in enabled:
|
|
print(f"source {args.source!r} is disabled or not configured", file=sys.stderr)
|
|
return 1
|
|
started = run_source(args.source)
|
|
print("community fetch started" if started else "community fetch skipped: locked or cooling down")
|
|
elif args.command == "cleanup-retention":
|
|
policy = RetentionPolicy(
|
|
submission_days=settings.retention_submission_days,
|
|
unreviewed_days=settings.retention_unreviewed_days,
|
|
approved_personal_days=settings.retention_approved_personal_days,
|
|
staging_days=settings.retention_staging_days,
|
|
audit_days=settings.retention_audit_days,
|
|
published_payload_days=settings.retention_published_payload_days,
|
|
)
|
|
counts = apply_retention(session, policy=policy, dry_run=not args.apply, delete_object=delete_screenshot)
|
|
print(json.dumps({"mode": "apply" if args.apply else "dry-run", "policy": asdict(policy), "counts": counts}, ensure_ascii=False))
|
|
elif args.command == "audit-waterbody-catalog":
|
|
stream = sys.stdin if args.input == "-" else open(args.input, encoding="utf-8")
|
|
try:
|
|
snapshot = json.load(stream)
|
|
finally:
|
|
if stream is not sys.stdin:
|
|
stream.close()
|
|
items = snapshot.get("items") if isinstance(snapshot, dict) else snapshot
|
|
if not isinstance(items, list):
|
|
parser.error("input must be a JSON array or an object with an items array")
|
|
expected_ids = {
|
|
str(item["source_external_id"])
|
|
for item in items
|
|
if isinstance(item, dict) and item.get("source_external_id")
|
|
}
|
|
result = audit_waterbody_catalog(session, expected_ids)
|
|
print(json.dumps(result, ensure_ascii=False))
|
|
return 1 if result["failures"] else 0
|
|
else:
|
|
result = audit_catalog(session)
|
|
print(json.dumps(result, ensure_ascii=False))
|
|
return 1 if result["failures"] else 0
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|