from __future__ import annotations
|
|
from dataclasses import asdict
|
from datetime import datetime, timezone
|
import hashlib
|
import json
|
import os
|
from pathlib import Path
|
from typing import Any, Callable, Mapping, Sequence
|
|
from .models import (
|
ContractError, ErrorCode, EvidenceState, PackageState, TerminalReceipt,
|
TerminalStatus, TERMINAL_MAPPING,
|
)
|
|
|
TERMINAL_KEYS_V003 = (
|
"schema_version", "message_type", "contract_version", "project_id", "task_id", "handoff_id",
|
"source_role_instance_id", "source_thread_id", "target_role_instance_id", "target_thread_id",
|
"reply_thread_id", "requester", "review_owner", "run_id", "mode", "status",
|
"capability_status", "exit_code", "stop_code", "blocker", "observed_at_utc", "started_at_utc",
|
"ended_at_utc", "total_elapsed_ms", "work_deadline_reached", "close_deadline_reached",
|
"requested", "triggered", "succeeded", "failed", "duplicates", "gaps", "quota_confirmed",
|
"quota_uncertain", "quota_active", "quota_cumulative_consumed", "quota_safe_available",
|
"quota_ledger_path", "items", "manifest_path", "delivery_path", "timing_path",
|
"prohibited_action_attestation",
|
)
|
TERMINAL_SUFFIX_V009 = (
|
"terminal_target_path", "terminal_path", "terminal_presence", "terminal_present",
|
"terminal_persist_attempted", "terminal_persisted", "terminal_operation_completed",
|
"terminal_liveness_unknown",
|
)
|
TERMINAL_KEYS_V009 = TERMINAL_KEYS_V003 + TERMINAL_SUFFIX_V009
|
|
|
def utc_now() -> str:
|
return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
|
|
|
def build_terminal(base: Mapping[str, Any], receipt: TerminalReceipt, *, status: TerminalStatus,
|
stop_code: ErrorCode | None, blocker: str | None = None) -> dict[str, Any]:
|
receipt.validate()
|
capability, exit_code = TERMINAL_MAPPING[status]
|
output = {key: base.get(key) for key in TERMINAL_KEYS_V003}
|
output.update({
|
"schema_version": "HIBOR_FAST_REPORT_COLLECTION_TERMINAL_V009",
|
"message_type": "report_collection_terminal",
|
"contract_version": "REPORT-COLLECTION-CAPABILITY-V1",
|
"status": status.value,
|
"capability_status": capability,
|
"exit_code": exit_code,
|
"stop_code": None if stop_code is None or stop_code is ErrorCode.NONE else stop_code.value,
|
"blocker": blocker,
|
"ended_at_utc": base.get("ended_at_utc") or utc_now(),
|
"terminal_target_path": str(receipt.target_path),
|
"terminal_path": str(receipt.terminal_path) if receipt.terminal_path else None,
|
"terminal_presence": receipt.presence.value,
|
"terminal_present": receipt.presence is EvidenceState.V,
|
"terminal_persist_attempted": receipt.persist_attempted,
|
"terminal_persisted": receipt.persisted,
|
"terminal_operation_completed": receipt.operation_completed,
|
"terminal_liveness_unknown": receipt.liveness_unknown,
|
})
|
if tuple(output.keys()) != TERMINAL_KEYS_V009:
|
raise ContractError(ErrorCode.TASK_SPEC_INVALID, "terminal.keys", "51-key order mismatch")
|
_validate_terminal(output)
|
return output
|
|
|
def _validate_terminal(value: Mapping[str, Any]) -> None:
|
if tuple(value.keys()) != TERMINAL_KEYS_V009:
|
raise ContractError(ErrorCode.TASK_SPEC_INVALID, "terminal.keys", "mismatch")
|
if value["terminal_present"] != (value["terminal_presence"] == "V"):
|
raise ContractError(ErrorCode.TASK_SPEC_INVALID, "terminal_present", "cross-field")
|
if value["terminal_presence"] == "V" and value["terminal_path"] != value["terminal_target_path"]:
|
raise ContractError(ErrorCode.TASK_SPEC_INVALID, "terminal_path", "valid path mismatch")
|
if value["terminal_presence"] != "V" and value["terminal_path"] is not None:
|
raise ContractError(ErrorCode.TASK_SPEC_INVALID, "terminal_path", "non-valid must be null")
|
if value["status"] not in {item.value for item in TerminalStatus}:
|
raise ContractError(ErrorCode.TASK_SPEC_INVALID, "status", "enum")
|
|
|
def canonical_terminal_bytes(value: Mapping[str, Any]) -> bytes:
|
_validate_terminal(value)
|
return json.dumps(value, ensure_ascii=False, separators=(",", ":"), sort_keys=False).encode("utf-8")
|
|
|
def legacy_v003_projection(value: Mapping[str, Any]) -> dict[str, Any]:
|
_validate_terminal(value)
|
projected = {key: value[key] for key in TERMINAL_KEYS_V003}
|
projected["schema_version"] = "HIBOR_FAST_REPORT_COLLECTION_TERMINAL_V003"
|
return projected
|
|
|
class TerminalWriter:
|
def read_valid(self, target: Path, *, task_id: str, handoff_id: str,
|
run_id: str) -> dict[str, Any] | None:
|
if not target.exists() and not target.is_symlink():
|
return None
|
if not target.is_file() or target.is_symlink():
|
raise ContractError(ErrorCode.PERSIST_LATE, "terminal", "existing non-file")
|
try:
|
raw = target.read_bytes()
|
value = json.loads(raw.decode("utf-8"))
|
except (OSError, UnicodeError, json.JSONDecodeError) as exc:
|
raise ContractError(ErrorCode.PERSIST_LATE, "terminal", "existing invalid") from exc
|
if not isinstance(value, dict) or canonical_terminal_bytes(value) != raw:
|
raise ContractError(ErrorCode.PERSIST_LATE, "terminal", "existing noncanonical")
|
if (value["task_id"], value["handoff_id"], value["run_id"]) != (task_id, handoff_id, run_id):
|
raise ContractError(ErrorCode.QUOTA_REPLAY_CONFLICT, "terminal", "identity mismatch")
|
if not value["terminal_persisted"] or value["terminal_presence"] != "V":
|
raise ContractError(ErrorCode.PERSIST_LATE, "terminal", "existing not closed")
|
return value
|
|
def persist(self, target: Path, value: Mapping[str, Any], *, source_state: PackageState,
|
checkpoint: Callable[[], None] | None = None) -> TerminalReceipt:
|
data = canonical_terminal_bytes(value)
|
attempted = False
|
created = False
|
try:
|
# The pre-attempt deadline is part of T00B. Keep it inside the
|
# receipt boundary so a spent close budget returns an in-memory
|
# N/no-attempt receipt instead of escaping BudgetExpired.
|
if checkpoint:
|
checkpoint()
|
target.parent.mkdir(parents=True, exist_ok=True)
|
attempted = True
|
flags = os.O_CREAT | os.O_EXCL | os.O_WRONLY | getattr(os, "O_BINARY", 0)
|
fd = os.open(target, flags, 0o600)
|
created = True
|
try:
|
view = memoryview(data)
|
while view:
|
if checkpoint:
|
checkpoint()
|
count = os.write(fd, view)
|
if count <= 0:
|
raise OSError("short write")
|
view = view[count:]
|
os.fsync(fd)
|
finally:
|
os.close(fd)
|
reopened = target.read_bytes()
|
if checkpoint:
|
checkpoint()
|
if reopened != data:
|
digest = hashlib.sha256(reopened).hexdigest() if reopened else None
|
receipt = TerminalReceipt(target, None, True, False, EvidenceState.I, True, False,
|
True, len(reopened) if reopened else None, digest,
|
ErrorCode.PERSIST_LATE, source_state)
|
receipt.validate()
|
return receipt
|
receipt = TerminalReceipt(target, target, True, True, EvidenceState.V, True, False,
|
True, len(data), hashlib.sha256(data).hexdigest(), None, source_state)
|
receipt.validate()
|
return receipt
|
except FileExistsError:
|
try:
|
if target.is_file() and not target.is_symlink():
|
raw = target.read_bytes()
|
if raw == data:
|
receipt = TerminalReceipt(
|
target, target, True, True, EvidenceState.V, True, False,
|
True, len(raw), hashlib.sha256(raw).hexdigest(), None, source_state,
|
)
|
else:
|
receipt = TerminalReceipt(
|
target, None, True, False, EvidenceState.I, True, False,
|
True, len(raw), hashlib.sha256(raw).hexdigest(),
|
ErrorCode.PERSIST_LATE, source_state,
|
)
|
else:
|
receipt = TerminalReceipt(
|
target, None, True, False, EvidenceState.I, True, False,
|
True, None, None, ErrorCode.PERSIST_LATE, source_state,
|
)
|
except OSError:
|
receipt = TerminalReceipt(
|
target, None, True, None, EvidenceState.U, None, True,
|
True, None, None, ErrorCode.PROCESS_LIVENESS_UNKNOWN, source_state,
|
)
|
receipt.validate()
|
return receipt
|
except (OSError, ContractError):
|
if not attempted and not target.exists():
|
receipt = TerminalReceipt(target, None, False, False, EvidenceState.N, False, False,
|
True, None, None, None, source_state)
|
elif not created and not target.exists():
|
receipt = TerminalReceipt(target, None, True, False, EvidenceState.N, True, False,
|
True, None, None, ErrorCode.PERSIST_LATE, source_state)
|
else:
|
try:
|
if target.is_file() and not target.is_symlink():
|
raw = target.read_bytes()
|
receipt = TerminalReceipt(target, None, True, False, EvidenceState.I, True, False,
|
True, len(raw), hashlib.sha256(raw).hexdigest(),
|
ErrorCode.PERSIST_LATE, source_state)
|
else:
|
receipt = TerminalReceipt(target, None, True, None, EvidenceState.U, True, False,
|
True, None, None, ErrorCode.RECOVERY_UNKNOWN, source_state)
|
except OSError:
|
receipt = TerminalReceipt(target, None, True, None, EvidenceState.U, True, False,
|
True, None, None, ErrorCode.RECOVERY_UNKNOWN, source_state)
|
receipt.validate()
|
return receipt
|