|
@@ -16,7 +16,7 @@ from urllib.parse import urljoin
|
|
|
import httpx
|
|
import httpx
|
|
|
|
|
|
|
|
from script_build_host.domain.errors import ProtocolViolation
|
|
from script_build_host.domain.errors import ProtocolViolation
|
|
|
-from script_build_host.infrastructure.canonical_json import canonical_sha256
|
|
|
|
|
|
|
+from script_build_host.infrastructure.canonical_json import canonical_json_bytes, canonical_sha256
|
|
|
from script_build_host.infrastructure.outbound import OutboundPolicy
|
|
from script_build_host.infrastructure.outbound import OutboundPolicy
|
|
|
from script_build_host.infrastructure.redaction import redact, redact_text
|
|
from script_build_host.infrastructure.redaction import redact, redact_text
|
|
|
from script_build_host.tools.gateway import RetrievalResult
|
|
from script_build_host.tools.gateway import RetrievalResult
|
|
@@ -105,12 +105,14 @@ class PatternRetrievalAdapter:
|
|
|
self,
|
|
self,
|
|
|
endpoint: str,
|
|
endpoint: str,
|
|
|
http: SafeHttpClient,
|
|
http: SafeHttpClient,
|
|
|
|
|
+ raw_store: RawArtifactStore,
|
|
|
*,
|
|
*,
|
|
|
poll_interval_seconds: float = 1.0,
|
|
poll_interval_seconds: float = 1.0,
|
|
|
poll_timeout_seconds: float = 600.0,
|
|
poll_timeout_seconds: float = 600.0,
|
|
|
) -> None:
|
|
) -> None:
|
|
|
self.endpoint = endpoint
|
|
self.endpoint = endpoint
|
|
|
self.http = http
|
|
self.http = http
|
|
|
|
|
+ self.raw_store = raw_store
|
|
|
self.poll_interval_seconds = poll_interval_seconds
|
|
self.poll_interval_seconds = poll_interval_seconds
|
|
|
self.poll_timeout_seconds = poll_timeout_seconds
|
|
self.poll_timeout_seconds = poll_timeout_seconds
|
|
|
|
|
|
|
@@ -129,12 +131,18 @@ class PatternRetrievalAdapter:
|
|
|
"queued",
|
|
"queued",
|
|
|
}:
|
|
}:
|
|
|
if time.monotonic() >= deadline:
|
|
if time.monotonic() >= deadline:
|
|
|
|
|
+ safe_payload, raw_ref, response_digest = await _freeze_json_response(
|
|
|
|
|
+ self.raw_store,
|
|
|
|
|
+ payload,
|
|
|
|
|
+ )
|
|
|
return _limited_result(
|
|
return _limited_result(
|
|
|
"pattern",
|
|
"pattern",
|
|
|
"pattern retrieval timed out",
|
|
"pattern retrieval timed out",
|
|
|
"timeout",
|
|
"timeout",
|
|
|
task_id=task_id,
|
|
task_id=task_id,
|
|
|
- payload=payload,
|
|
|
|
|
|
|
+ payload=safe_payload,
|
|
|
|
|
+ raw_artifact_ref=raw_ref,
|
|
|
|
|
+ response_digest=response_digest,
|
|
|
)
|
|
)
|
|
|
poll_url = payload.get("poll_url")
|
|
poll_url = payload.get("poll_url")
|
|
|
if not isinstance(poll_url, str) or not poll_url:
|
|
if not isinstance(poll_url, str) or not poll_url:
|
|
@@ -144,60 +152,80 @@ class PatternRetrievalAdapter:
|
|
|
await asyncio.sleep(self.poll_interval_seconds)
|
|
await asyncio.sleep(self.poll_interval_seconds)
|
|
|
payload = await self.http.json("GET", poll_url)
|
|
payload = await self.http.json("GET", poll_url)
|
|
|
if isinstance(payload, dict) and payload.get("status") in {"failed", "error"}:
|
|
if isinstance(payload, dict) and payload.get("status") in {"failed", "error"}:
|
|
|
|
|
+ safe_payload, raw_ref, response_digest = await _freeze_json_response(
|
|
|
|
|
+ self.raw_store,
|
|
|
|
|
+ payload,
|
|
|
|
|
+ )
|
|
|
return _limited_result(
|
|
return _limited_result(
|
|
|
"pattern",
|
|
"pattern",
|
|
|
"pattern retrieval failed",
|
|
"pattern retrieval failed",
|
|
|
"upstream failure",
|
|
"upstream failure",
|
|
|
task_id=task_id,
|
|
task_id=task_id,
|
|
|
- payload=payload,
|
|
|
|
|
|
|
+ payload=safe_payload,
|
|
|
|
|
+ raw_artifact_ref=raw_ref,
|
|
|
|
|
+ response_digest=response_digest,
|
|
|
)
|
|
)
|
|
|
- compact_payload = _compact_pattern_payload(payload)
|
|
|
|
|
|
|
+ safe_payload, raw_ref, response_digest = await _freeze_json_response(
|
|
|
|
|
+ self.raw_store,
|
|
|
|
|
+ payload,
|
|
|
|
|
+ )
|
|
|
|
|
+ compact_payload = _compact_pattern_payload(safe_payload)
|
|
|
return _result(
|
|
return _result(
|
|
|
"pattern",
|
|
"pattern",
|
|
|
compact_payload,
|
|
compact_payload,
|
|
|
metadata={
|
|
metadata={
|
|
|
"task_id": _safe_source_id(task_id, rank=0) if task_id else None,
|
|
"task_id": _safe_source_id(task_id, rank=0) if task_id else None,
|
|
|
- "response_sha256": canonical_sha256(payload).wire,
|
|
|
|
|
|
|
+ "response_sha256": response_digest,
|
|
|
},
|
|
},
|
|
|
|
|
+ raw_artifact_ref=raw_ref,
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
class ExternalRetrievalAdapter:
|
|
class ExternalRetrievalAdapter:
|
|
|
"""Pure HTTP client; it contains no legacy-log repository dependency."""
|
|
"""Pure HTTP client; it contains no legacy-log repository dependency."""
|
|
|
|
|
|
|
|
- def __init__(self, endpoint: str, http: SafeHttpClient) -> None:
|
|
|
|
|
|
|
+ def __init__(self, endpoint: str, http: SafeHttpClient, raw_store: RawArtifactStore) -> None:
|
|
|
self.endpoint = endpoint
|
|
self.endpoint = endpoint
|
|
|
self.http = http
|
|
self.http = http
|
|
|
|
|
+ self.raw_store = raw_store
|
|
|
|
|
|
|
|
async def retrieve(self, *, query: Mapping[str, Any], snapshot: Any) -> RetrievalResult:
|
|
async def retrieve(self, *, query: Mapping[str, Any], snapshot: Any) -> RetrievalResult:
|
|
|
del snapshot
|
|
del snapshot
|
|
|
payload = await self.http.json("POST", self.endpoint, json_body=query)
|
|
payload = await self.http.json("POST", self.endpoint, json_body=query)
|
|
|
|
|
+ safe_payload, raw_ref, response_digest = await _freeze_json_response(
|
|
|
|
|
+ self.raw_store,
|
|
|
|
|
+ payload,
|
|
|
|
|
+ )
|
|
|
return _result(
|
|
return _result(
|
|
|
"external",
|
|
"external",
|
|
|
- payload,
|
|
|
|
|
- metadata={"response_sha256": canonical_sha256(payload).wire},
|
|
|
|
|
|
|
+ safe_payload,
|
|
|
|
|
+ metadata={"response_sha256": response_digest},
|
|
|
|
|
+ raw_artifact_ref=raw_ref,
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
class HttpDecodeRetrievalAdapter:
|
|
class HttpDecodeRetrievalAdapter:
|
|
|
"""Call the configured Decode service with the legacy query contract."""
|
|
"""Call the configured Decode service with the legacy query contract."""
|
|
|
|
|
|
|
|
- def __init__(self, endpoint: str, http: SafeHttpClient) -> None:
|
|
|
|
|
|
|
+ def __init__(self, endpoint: str, http: SafeHttpClient, raw_store: RawArtifactStore) -> None:
|
|
|
self.endpoint = endpoint
|
|
self.endpoint = endpoint
|
|
|
self.http = http
|
|
self.http = http
|
|
|
|
|
+ self.raw_store = raw_store
|
|
|
|
|
|
|
|
async def retrieve(self, *, query: Mapping[str, Any], snapshot: Any) -> RetrievalResult:
|
|
async def retrieve(self, *, query: Mapping[str, Any], snapshot: Any) -> RetrievalResult:
|
|
|
payload = await self.http.json("POST", self.endpoint, json_body=query)
|
|
payload = await self.http.json("POST", self.endpoint, json_body=query)
|
|
|
|
|
+ safe_payload, raw_ref, response_digest = await _freeze_json_response(
|
|
|
|
|
+ self.raw_store,
|
|
|
|
|
+ payload,
|
|
|
|
|
+ )
|
|
|
rows = (
|
|
rows = (
|
|
|
- payload.get("data", payload.get("results", payload.get("items", [])))
|
|
|
|
|
- if isinstance(payload, dict)
|
|
|
|
|
- else payload
|
|
|
|
|
|
|
+ safe_payload.get("data", safe_payload.get("results", safe_payload.get("items", [])))
|
|
|
|
|
+ if isinstance(safe_payload, dict)
|
|
|
|
|
+ else safe_payload
|
|
|
)
|
|
)
|
|
|
if not isinstance(rows, list):
|
|
if not isinstance(rows, list):
|
|
|
raise ProtocolViolation("decode response items must be an array")
|
|
raise ProtocolViolation("decode response items must be an array")
|
|
|
- safe_rows = redact(rows)
|
|
|
|
|
- if not isinstance(safe_rows, list):
|
|
|
|
|
- raise ProtocolViolation("redacted decode rows must remain an array")
|
|
|
|
|
|
|
+ safe_rows = rows
|
|
|
refs: list[str] = []
|
|
refs: list[str] = []
|
|
|
scores: list[float] = []
|
|
scores: list[float] = []
|
|
|
for row in safe_rows:
|
|
for row in safe_rows:
|
|
@@ -217,18 +245,19 @@ class HttpDecodeRetrievalAdapter:
|
|
|
model_manifest = getattr(snapshot, "model_manifest", {})
|
|
model_manifest = getattr(snapshot, "model_manifest", {})
|
|
|
return RetrievalResult(
|
|
return RetrievalResult(
|
|
|
source_refs=tuple(refs),
|
|
source_refs=tuple(refs),
|
|
|
- summary=json.dumps(safe_rows, ensure_ascii=False, sort_keys=True),
|
|
|
|
|
|
|
+ summary=_bounded_json_summary(safe_rows),
|
|
|
supports=tuple(str(query["return_field"]) for _ in refs),
|
|
supports=tuple(str(query["return_field"]) for _ in refs),
|
|
|
confidence="ranked",
|
|
confidence="ranked",
|
|
|
limitations=(),
|
|
limitations=(),
|
|
|
metadata={
|
|
metadata={
|
|
|
- "response_sha256": canonical_sha256(payload).wire,
|
|
|
|
|
|
|
+ "response_sha256": response_digest,
|
|
|
"decode_index": datasource_manifest.get("decode_index"),
|
|
"decode_index": datasource_manifest.get("decode_index"),
|
|
|
"embedding_model": model_manifest.get("embedding_model"),
|
|
"embedding_model": model_manifest.get("embedding_model"),
|
|
|
"hit_ids": refs,
|
|
"hit_ids": refs,
|
|
|
"ranks": list(range(1, len(refs) + 1)),
|
|
"ranks": list(range(1, len(refs) + 1)),
|
|
|
"scores": scores,
|
|
"scores": scores,
|
|
|
},
|
|
},
|
|
|
|
|
+ raw_artifact_ref=raw_ref,
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
@@ -407,6 +436,7 @@ def _result(
|
|
|
payload: Any,
|
|
payload: Any,
|
|
|
*,
|
|
*,
|
|
|
metadata: Mapping[str, Any],
|
|
metadata: Mapping[str, Any],
|
|
|
|
|
+ raw_artifact_ref: str | None = None,
|
|
|
) -> RetrievalResult:
|
|
) -> RetrievalResult:
|
|
|
rows = payload.get("data", payload.get("results", [])) if isinstance(payload, dict) else payload
|
|
rows = payload.get("data", payload.get("results", [])) if isinstance(payload, dict) else payload
|
|
|
if not isinstance(rows, list):
|
|
if not isinstance(rows, list):
|
|
@@ -422,18 +452,58 @@ def _result(
|
|
|
)
|
|
)
|
|
|
return RetrievalResult(
|
|
return RetrievalResult(
|
|
|
source_refs=refs,
|
|
source_refs=refs,
|
|
|
- summary=json.dumps(safe_rows, ensure_ascii=False, sort_keys=True),
|
|
|
|
|
|
|
+ summary=_bounded_json_summary(safe_rows),
|
|
|
supports=refs,
|
|
supports=refs,
|
|
|
confidence="upstream",
|
|
confidence="upstream",
|
|
|
limitations=(),
|
|
limitations=(),
|
|
|
|
|
+ raw_artifact_ref=raw_artifact_ref,
|
|
|
metadata=metadata,
|
|
metadata=metadata,
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
+async def _freeze_json_response(
|
|
|
|
|
+ raw_store: RawArtifactStore,
|
|
|
|
|
+ payload: Any,
|
|
|
|
|
+) -> tuple[Any, str, str]:
|
|
|
|
|
+ """Redact once, then freeze the exact canonical JSON represented by its digest."""
|
|
|
|
|
+
|
|
|
|
|
+ safe_payload = redact(payload)
|
|
|
|
|
+ content = canonical_json_bytes(safe_payload)
|
|
|
|
|
+ raw_ref = await raw_store.freeze_bytes(content, media_type="application/json")
|
|
|
|
|
+ return safe_payload, raw_ref, "sha256:" + sha256(content).hexdigest()
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+def _bounded_json_summary(payload: Any, *, max_chars: int = 20_000) -> str:
|
|
|
|
|
+ rendered = json.dumps(payload, ensure_ascii=False, sort_keys=True)
|
|
|
|
|
+ if len(rendered) <= max_chars:
|
|
|
|
|
+ return rendered
|
|
|
|
|
+ rows = payload if isinstance(payload, list) else [payload]
|
|
|
|
|
+ kept: list[Any] = []
|
|
|
|
|
+ for row in rows:
|
|
|
|
|
+ candidate = [*kept, row]
|
|
|
|
|
+ if len(json.dumps(candidate, ensure_ascii=False, sort_keys=True)) > max_chars:
|
|
|
|
|
+ break
|
|
|
|
|
+ kept.append(row)
|
|
|
|
|
+ while True:
|
|
|
|
|
+ bounded = json.dumps(
|
|
|
|
|
+ {
|
|
|
|
|
+ "items": kept,
|
|
|
|
|
+ "items_returned": len(kept),
|
|
|
|
|
+ "items_total": len(rows),
|
|
|
|
|
+ "truncated_for_context": True,
|
|
|
|
|
+ },
|
|
|
|
|
+ ensure_ascii=False,
|
|
|
|
|
+ sort_keys=True,
|
|
|
|
|
+ )
|
|
|
|
|
+ if len(bounded) <= max_chars or not kept:
|
|
|
|
|
+ return bounded
|
|
|
|
|
+ kept.pop()
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
def _compact_pattern_payload(
|
|
def _compact_pattern_payload(
|
|
|
payload: Any,
|
|
payload: Any,
|
|
|
*,
|
|
*,
|
|
|
- max_candidates: int = 3,
|
|
|
|
|
|
|
+ max_candidates: int = 5,
|
|
|
max_chars: int = 20_000,
|
|
max_chars: int = 20_000,
|
|
|
) -> Any:
|
|
) -> Any:
|
|
|
"""Keep a bounded Pattern evidence view while preserving the raw digest upstream."""
|
|
"""Keep a bounded Pattern evidence view while preserving the raw digest upstream."""
|
|
@@ -525,6 +595,8 @@ def _limited_result(
|
|
|
*,
|
|
*,
|
|
|
task_id: str,
|
|
task_id: str,
|
|
|
payload: Any,
|
|
payload: Any,
|
|
|
|
|
+ raw_artifact_ref: str | None = None,
|
|
|
|
|
+ response_digest: str | None = None,
|
|
|
) -> RetrievalResult:
|
|
) -> RetrievalResult:
|
|
|
return RetrievalResult(
|
|
return RetrievalResult(
|
|
|
source_refs=(),
|
|
source_refs=(),
|
|
@@ -532,9 +604,10 @@ def _limited_result(
|
|
|
supports=(),
|
|
supports=(),
|
|
|
confidence="unavailable",
|
|
confidence="unavailable",
|
|
|
limitations=(limitation,),
|
|
limitations=(limitation,),
|
|
|
|
|
+ raw_artifact_ref=raw_artifact_ref,
|
|
|
metadata={
|
|
metadata={
|
|
|
"task_id": _safe_source_id(task_id, rank=0) if task_id else None,
|
|
"task_id": _safe_source_id(task_id, rank=0) if task_id else None,
|
|
|
- "response_sha256": canonical_sha256(payload).wire,
|
|
|
|
|
|
|
+ "response_sha256": response_digest or canonical_sha256(payload).wire,
|
|
|
},
|
|
},
|
|
|
)
|
|
)
|
|
|
|
|
|