Wire VedAstro orchestrator into strict workflows

This commit is contained in:
732642856
2026-06-29 10:26:54 +08:00
parent 3452beee1b
commit 0f9b06d9d4
18 changed files with 904 additions and 162 deletions
+3 -86
View File
@@ -14,7 +14,6 @@ import io
import json, sys, os, math
import importlib.util
import re
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime, timedelta
from http.server import HTTPServer, BaseHTTPRequestHandler
from urllib.parse import urlparse
@@ -57,7 +56,7 @@ def _attach_vedastro_main_entry_overview(chart_result, birth_payload):
return chart_result
try:
adapter = _load_local_module('vedastro_service_adapter')
orchestrator = _load_local_module('vedastro_evidence_orchestrator')
except Exception:
return chart_result
@@ -66,9 +65,7 @@ def _attach_vedastro_main_entry_overview(chart_result, birth_payload):
or birth_payload.get('today')
or datetime.utcnow().strftime('%Y-%m-%d')
)[:10]
start_date = datetime.strptime(reference_date, '%Y-%m-%d').date()
end_date = start_date
case = {
modules['vedastro_range_scan_result'] = orchestrator.orchestrate_vedastro_evidence({
'year': birth_payload.get('year'),
'month': birth_payload.get('month'),
'day': birth_payload.get('day'),
@@ -80,87 +77,7 @@ def _attach_vedastro_main_entry_overview(chart_result, birth_payload):
'tz': birth_payload.get('tz'),
'ayanamsa_policy': birth_payload.get('ayanamsa') or 'lahiri',
'node_policy': birth_payload.get('node_mode') or birth_payload.get('nodeMode') or 'mean',
}
def _scan_domain(domain: str):
return domain, adapter.run_range_scan_for_case(
case,
domain=domain,
start_date=start_date.isoformat(),
end_date=end_date.isoformat(),
case_id=f"api_chart_{domain}",
)
domain_reports = {}
combined_events = []
domain_statuses = {}
top_events = {}
failure_reason = None
availability = True
with ThreadPoolExecutor(max_workers=3) as executor:
for domain, domain_report in executor.map(_scan_domain, ('career', 'marriage', 'wealth')):
domain_reports[domain] = domain_report
for domain in ('career', 'marriage', 'wealth'):
domain_report = domain_reports[domain]
domain_statuses[domain] = domain_report.get('status')
availability = availability and bool(domain_report.get('available', False))
if domain_report.get('status') != 'ok' and failure_reason is None:
failure_reason = domain_report.get('reason')
for event in domain_report.get('evidence_ledger') or []:
if isinstance(event, dict):
combined_events.append(event)
top_event = domain_report.get('top_event')
if isinstance(top_event, dict):
top_events[domain] = top_event
primary_status = next(
(
domain_reports[domain].get('status')
for domain in ('career', 'marriage', 'wealth')
if domain_reports.get(domain, {}).get('status') == 'ok'
),
domain_reports.get('marriage', {}).get('status') or 'blocked',
)
source_metadata = {
'ingestion_profile': 'main_entry_overview',
'search_scope': 'single_day_overview',
'reference_date': reference_date,
'scan_window': {'start': start_date.isoformat(), 'end': end_date.isoformat()},
'domain_statuses': domain_statuses,
'domain_event_counts': {
domain: int((domain_reports.get(domain, {}) or {}).get('event_count', 0) or 0)
for domain in ('career', 'marriage', 'wealth')
},
}
for domain in ('career', 'marriage', 'wealth'):
metadata = domain_reports.get(domain, {}).get('source_metadata')
if isinstance(metadata, dict):
for key in (
'endpoint',
'endpoint_host',
'transport',
'provenance_mode',
'timeout_seconds',
'retry_policy',
):
if key in metadata and key not in source_metadata:
source_metadata[key] = metadata[key]
modules['vedastro_range_scan_result'] = {
'backend': 'vedastro_service_adapter_candidate',
'available': availability,
'status': primary_status,
'operation': 'range_scan',
'domain': 'overview',
'event_count': len(combined_events),
'top_event': top_events.get('marriage') or next(iter(top_events.values()), None),
'top_events_by_domain': top_events,
'evidence_ledger': combined_events,
'source_metadata': source_metadata,
'reason': failure_reason,
'domain_reports': domain_reports,
}
}, route='overview', reference_date=reference_date, case_id='api_chart')
return chart_result
+4
View File
@@ -13,6 +13,10 @@ import tempfile
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT / "scripts"))
from local_env import load_local_env
load_local_env(ROOT)
APP = ROOT / "jyotish-app"
PYTHON = sys.executable
+133
View File
@@ -0,0 +1,133 @@
#!/usr/bin/env python3
"""Shared VedAstro evidence orchestration layer.
The orchestrator does not expose hundreds of VedAstro nodes to users. It picks
the smallest useful domain scan set for the current route and delegates the
actual external boundary to ``vedastro_service_adapter``.
"""
from __future__ import annotations
from datetime import datetime, timedelta
from typing import Any
try:
from scripts.vedastro_service_adapter import (
VEDASTRO_CALCULATION_COVERAGE,
run_range_scan_for_case,
)
except ModuleNotFoundError: # pragma: no cover - script execution path
from vedastro_service_adapter import VEDASTRO_CALCULATION_COVERAGE, run_range_scan_for_case
ROUTE_DOMAIN_MAP = {
"relationship": ["marriage"],
"marriage": ["marriage"],
"career": ["career"],
"finance": ["wealth"],
"wealth": ["wealth"],
"rectification": ["marriage", "career", "wealth"],
"timing": ["career", "marriage", "wealth"],
"general": ["career", "marriage", "wealth"],
"overview": ["career", "marriage", "wealth"],
}
def _default_window(reference_date: str | None, days: int = 180) -> tuple[str, str]:
raw = str(reference_date or datetime.utcnow().strftime("%Y-%m-%d"))[:10]
try:
start = datetime.strptime(raw, "%Y-%m-%d").date()
except ValueError:
return raw, raw
return start.isoformat(), (start + timedelta(days=days)).isoformat()
def _normalize_case(birth_payload: dict[str, Any]) -> dict[str, Any]:
return {
"year": birth_payload.get("year"),
"month": birth_payload.get("month"),
"day": birth_payload.get("day"),
"hour": birth_payload.get("hour"),
"minute": birth_payload.get("minute"),
"second": birth_payload.get("second", 0),
"lat": birth_payload.get("lat"),
"lon": birth_payload.get("lon"),
"tz": birth_payload.get("tz"),
"ayanamsa_policy": birth_payload.get("ayanamsa_policy")
or birth_payload.get("ayanamsa")
or "lahiri",
"node_policy": birth_payload.get("node_policy")
or birth_payload.get("node_mode")
or birth_payload.get("nodeMode")
or "mean",
}
def orchestrate_vedastro_evidence(
birth_payload: dict[str, Any],
*,
route: str = "general",
reference_date: str | None = None,
start_date: str | None = None,
end_date: str | None = None,
case_id: str = "vedastro_orchestrator",
) -> dict[str, Any]:
domains = ROUTE_DOMAIN_MAP.get(route, ROUTE_DOMAIN_MAP["general"])
window_start, window_end = (start_date, end_date) if start_date and end_date else _default_window(reference_date)
case = _normalize_case(birth_payload)
domain_reports: dict[str, Any] = {}
evidence_ledger: list[dict[str, Any]] = []
top_events_by_domain: dict[str, Any] = {}
domain_statuses: dict[str, Any] = {}
domain_event_counts: dict[str, int] = {}
available = False
first_reason = None
for domain in domains:
report = run_range_scan_for_case(
case,
domain,
str(window_start),
str(window_end),
case_id=f"{case_id}_{domain}",
)
domain_reports[domain] = report
domain_statuses[domain] = report.get("status")
domain_event_counts[domain] = int(report.get("event_count", 0) or 0)
available = available or bool(report.get("available"))
first_reason = first_reason or report.get("reason")
if isinstance(report.get("top_event"), dict):
top_events_by_domain[domain] = report["top_event"]
for event in report.get("evidence_ledger") or []:
if isinstance(event, dict):
evidence_ledger.append(event)
status = "ok" if any(item == "ok" for item in domain_statuses.values()) else next(iter(domain_statuses.values()), "blocked")
return {
"backend": "vedastro_service_adapter_candidate",
"available": available,
"status": status,
"operation": "range_scan",
"domain": domains[0] if len(domains) == 1 else "overview",
"route": route,
"event_count": len(evidence_ledger),
"top_event": next(iter(top_events_by_domain.values()), None),
"top_events_by_domain": top_events_by_domain,
"evidence_ledger": evidence_ledger,
"reason": None if status == "ok" else first_reason,
"domain_reports": domain_reports,
"source_metadata": {
"auto_ingested_by": "VedAstroEvidenceOrchestrator",
"strategy": "minimal_route_scoped_orchestration",
"node_coverage": {
"strategy": "domain_scoped_range_scan",
"official_calculation_coverage": VEDASTRO_CALCULATION_COVERAGE,
"selected_domains": domains,
"not_user_exposed": True,
},
"route": route,
"scan_window": {"start_date": str(window_start), "end_date": str(window_end)},
"domain_statuses": domain_statuses,
"domain_event_counts": domain_event_counts,
},
}
+189 -45
View File
@@ -9,6 +9,7 @@ workspace can evolve from research notes to an executable adapter contract.
from __future__ import annotations
import argparse
import http.client
import hashlib
import json
import os
@@ -19,8 +20,14 @@ from typing import Any
from urllib import request, error
from urllib.parse import urlparse
try:
from scripts.local_env import load_local_env
except ModuleNotFoundError: # pragma: no cover - script execution path
from local_env import load_local_env
ROOT = Path(__file__).resolve().parents[1]
load_local_env(ROOT)
PARITY_CASES = {
@@ -493,9 +500,11 @@ def _range_scan_preview(case: dict[str, Any], domain: str, start_date: str, end_
"start_date": start_date,
"end_date": end_date,
"event_model": "vedastro_events_at_range_candidate",
"search_mode": "single_point",
**case,
}
preview["official_request_profile"] = _build_official_search_events_profile(preview)
preview["live_sampling_request_profile"] = _build_live_sampling_search_events_profile(preview)
return preview
@@ -559,6 +568,28 @@ def _build_official_search_events_profile(request_preview: dict[str, Any]) -> di
}
def _build_live_sampling_search_events_profile(request_preview: dict[str, Any]) -> dict[str, Any]:
case = dict(request_preview)
case["tz"] = _normalize_tz(case)
body = {
"BirthTime": _time_json_from_case(case, f"{case['year']:04d}-{case['month']:02d}-{case['day']:02d}"),
"Ayanamsa": str(case.get("ayanamsa_policy") or "lahiri"),
"EventTagList": OFFICIAL_RANGE_SCAN_EVENT_TAGS.get(str(request_preview.get("domain") or ""), ["General"]),
"AtTime": _time_json_from_case(case, str(request_preview["start_date"])),
}
headers: dict[str, str] = {"Content-Type": "application/json"}
api_key = os.environ.get("VEDASTRO_API_KEY", "").strip()
if api_key:
headers["x-api-key"] = api_key
return {
"profile_version": f"{OFFICIAL_SEARCH_EVENTS_PROFILE_VERSION}_live_sampling",
"endpoint_path": OFFICIAL_SEARCH_EVENTS_ENDPOINT_PATH,
"method": OFFICIAL_SEARCH_EVENTS_METHOD,
"headers": headers,
"body": body,
}
def _external_technique_preview(
case: dict[str, Any],
domain: str,
@@ -685,9 +716,15 @@ def _normalize_range_scan_success(
attempt_count: int = 1,
retry_error_codes: list[int] | None = None,
) -> dict[str, Any]:
# Handle actual VedAstro response format: {"Status": "Pass", "Payload": [...]}
# Handle actual VedAstro response formats:
# {"Status": "Pass", "Payload": {"SearchEvents": [...]}}
# {"Status": "Pass", "Payload": [...]}
if payload.get("Status") == "Pass":
events = payload.get("Payload", [])
payload_body = payload.get("Payload", [])
if isinstance(payload_body, dict):
events = payload_body.get("SearchEvents", [])
else:
events = payload_body
else:
# Fallback to local stub format if not VedAstro format
events = payload.get("events", [])
@@ -716,10 +753,12 @@ def _normalize_range_scan_success(
for index, event in enumerate(events, start=1):
if not isinstance(event, dict):
continue
# VedAstro uses "Name" for event id, and "EventTags" for tags
# VedAstro uses "Name" for event id, and may expose tags as EventTags, tags, or Tag.
event_id = event.get("Name") or event.get("id") or event.get("name") or f"event_{index}"
tags = event.get("EventTags") or event.get("tags") or []
if not isinstance(tags, list):
tags = event.get("EventTags") or event.get("tags") or event.get("Tag") or []
if isinstance(tags, str):
tags = [part.strip() for part in tags.split(",") if part.strip()]
elif not isinstance(tags, list):
tags = []
tag_set = {str(tag) for tag in tags}
matched_by = "rejected"
@@ -865,7 +904,12 @@ def _source_metadata(endpoint: str) -> dict[str, Any]:
def _post_json(endpoint: str, request_preview: dict[str, Any]) -> dict[str, Any] | str:
official_request_profile = request_preview.get("official_request_profile") if isinstance(request_preview, dict) else None
official_request_profile = None
if isinstance(request_preview, dict):
official_request_profile = (
request_preview.get("live_sampling_request_profile")
or request_preview.get("official_request_profile")
)
request_url = endpoint
headers = {"Content-Type": "application/json"}
vedastro_payload = request_preview
@@ -913,6 +957,97 @@ def _post_json_with_retry(endpoint: str, request_preview: dict[str, Any]) -> tup
return {}, max_attempts, retry_error_codes
def _iter_sample_dates(start_date: str, end_date: str) -> list[str]:
from datetime import datetime, timedelta
start = datetime.strptime(start_date, "%Y-%m-%d").date()
end = datetime.strptime(end_date, "%Y-%m-%d").date()
if end <= start:
return [start_date]
span_days = max((end - start).days, 1)
step_days = max(1, span_days // 11)
dates: list[str] = []
current = start
while current <= end:
dates.append(current.isoformat())
current += timedelta(days=step_days)
if dates[-1] != end.isoformat():
dates.append(end.isoformat())
return dates
def _merge_range_scan_reports(
reports: list[dict[str, Any]],
endpoint: str,
base_preview: dict[str, Any],
) -> dict[str, Any]:
if not reports:
return _normalize_range_scan_success({"Status": "Pass", "Payload": []}, endpoint, base_preview)
deduped: dict[tuple[str, str, str], dict[str, Any]] = {}
for report in reports:
for item in report.get("evidence_ledger", []):
key = (
str(item.get("event_id") or ""),
str(item.get("start") or ""),
str(item.get("end") or ""),
)
existing = deduped.get(key)
existing_score = existing.get("signal_lift", 0) if existing else -1
score = item.get("signal_lift", 0)
if existing is None or score > existing_score:
deduped[key] = item
evidence_ledger = list(deduped.values())
top_event = None
if evidence_ledger:
top = max(
evidence_ledger,
key=lambda item: (
item.get("signal_lift") if isinstance(item.get("signal_lift"), (int, float)) else 0,
item.get("confidence") == "high",
),
)
top_event = {
"event_id": top.get("event_id"),
"signal_key": top.get("signal_key"),
"signal_label": top.get("signal_label"),
"signal_family": top.get("signal_family"),
"score": top.get("score"),
"start": top.get("start"),
"end": top.get("end"),
"tags": top.get("tags") or [],
}
metadata = dict(reports[-1].get("source_metadata") or {})
metadata["sampling_mode"] = "at_time_sweep"
metadata["sample_dates"] = [report.get("request_preview", {}).get("start_date") for report in reports]
metadata["sample_count"] = len(reports)
metadata["raw_event_count"] = sum(int((report.get("source_metadata") or {}).get("raw_event_count", 0)) for report in reports)
metadata["filtered_event_count"] = len(evidence_ledger)
metadata["allowlist_event_count"] = len(evidence_ledger)
metadata["attempt_count"] = sum(int((report.get("source_metadata") or {}).get("attempt_count", 1)) for report in reports)
retry_codes: list[int] = []
for report in reports:
retry_codes.extend(list((report.get("source_metadata") or {}).get("retry_error_codes") or []))
metadata["retry_error_codes"] = retry_codes
result = {
"backend": "vedastro_service_adapter_candidate",
"available": True,
"status": "ok",
"operation": "range_scan",
"domain": base_preview.get("domain"),
"request_preview": base_preview,
"event_count": len(evidence_ledger),
"top_event": top_event,
"evidence_ledger": evidence_ledger,
"source_metadata": metadata,
}
result["source_metadata"]["artifact_path"] = _write_artifact(result)
return result
def run_case(case_id: str) -> dict[str, Any]:
if case_id not in PARITY_CASES:
return {
@@ -1007,46 +1142,55 @@ def _run_range_scan_case(case: dict[str, Any], domain: str, start_date: str, end
"source_metadata": _source_metadata(endpoint),
}
try:
payload, attempt_count, retry_error_codes = _post_json_with_retry(endpoint, request_preview)
except error.HTTPError as exc:
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "http_error",
"reason": f"VedAstro range scan HTTP error: {exc.code}",
"request_preview": request_preview,
"source_metadata": _source_metadata(endpoint),
}
except error.URLError as exc:
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "network_error",
"reason": f"VedAstro range scan network error: {exc.reason}",
"request_preview": request_preview,
"source_metadata": _source_metadata(endpoint),
}
except (TimeoutError, socket.timeout):
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "timeout",
"reason": "VedAstro range scan timed out",
"request_preview": request_preview,
"source_metadata": _source_metadata(endpoint),
}
except json.JSONDecodeError:
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "invalid_json",
"reason": "VedAstro range scan received non-JSON response",
"request_preview": request_preview,
"source_metadata": _source_metadata(endpoint),
}
sample_dates = _iter_sample_dates(start_date, end_date)
reports: list[dict[str, Any]] = []
for sample_date in sample_dates:
sample_preview = dict(request_preview)
sample_preview["start_date"] = sample_date
sample_preview["end_date"] = sample_date
sample_preview["official_request_profile"] = _build_official_search_events_profile(sample_preview)
sample_preview["live_sampling_request_profile"] = _build_live_sampling_search_events_profile(sample_preview)
try:
payload, attempt_count, retry_error_codes = _post_json_with_retry(endpoint, sample_preview)
except error.HTTPError as exc:
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "http_error",
"reason": f"VedAstro range scan HTTP error: {exc.code}",
"request_preview": sample_preview,
"source_metadata": _source_metadata(endpoint),
}
except (error.URLError, http.client.RemoteDisconnected) as exc:
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "network_error",
"reason": f"VedAstro range scan network error: {getattr(exc, 'reason', str(exc))}",
"request_preview": sample_preview,
"source_metadata": _source_metadata(endpoint),
}
except (TimeoutError, socket.timeout):
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "timeout",
"reason": "VedAstro range scan timed out",
"request_preview": sample_preview,
"source_metadata": _source_metadata(endpoint),
}
except json.JSONDecodeError:
return {
"backend": "vedastro_service_adapter_candidate",
"available": False,
"status": "invalid_json",
"reason": "VedAstro range scan received non-JSON response",
"request_preview": sample_preview,
"source_metadata": _source_metadata(endpoint),
}
reports.append(_normalize_range_scan_success(payload, endpoint, sample_preview, attempt_count, retry_error_codes))
return _normalize_range_scan_success(payload, endpoint, request_preview, attempt_count, retry_error_codes)
return _merge_range_scan_reports(reports, endpoint, request_preview)
def run_range_scan(case_id: str, domain: str, start_date: str, end_date: str) -> dict[str, Any]: