"""Bounded, local JSON tool executor. Standard library only; no model or network.""" import argparse import csv from decimal import Decimal, InvalidOperation, localcontext import hashlib import io import json import os from pathlib import Path import re import sys import uuid PROTOCOL = "1.0" SOURCE_PATHS = ("data/transactions.csv", "data/groups.csv", "data/customers.csv", "contracts/metric-definitions.json") IDENTIFIER = re.compile(r"[A-Za-z0-9][A-Za-z0-9_-]{0,63}\Z") HASH = re.compile(r"[a-f0-9]{64}\Z") NUMBER = re.compile(r"-?\d{1,12}(?:\.\d{1,6})?\Z") DEFINITIONS = { "definition_id": "synthetic-v1", "units": "fictional_currency_units", "conversion": "sum(successes) / sum(eligible); zero denominator is unknown", "net_revenue": "quantity * unit_price * (1 - discount_rate); preserve signed returns", "missing_price": "unknown contribution; known subtotal only; blocks complete revenue release", "zero_price": "observed zero contribution; include in known coverage", "customer_join": "many-to-one; duplicate dimension keys block join; no deduplication", "release": "eligible only if revenue complete, customer join valid, and denominator positive; human approval remains separate", } SCOPE = { "native_bi": "not_run", "host_cancellation": "not_tested", "os_sandbox": "not_claimed", "real_model_quality": "not_tested", "human_review": "not_performed", } class ToolError(Exception): def __init__(self, code, message, retryable=False): super().__init__(message) self.code, self.message, self.retryable = code, message, retryable def fail(code, message, retryable=False): raise ToolError(code, message, retryable) def canonical(value): return (json.dumps(value, sort_keys=True, ensure_ascii=True, separators=(",", ":"), allow_nan=False) + "\n").encode("utf-8") def sha(content): return hashlib.sha256(content).hexdigest() def strict_json(text): def unique(pairs): result = {} for key, value in pairs: if key in result: raise ValueError("duplicate JSON key") result[key] = value return result def invalid_constant(value): raise ValueError("nonstandard JSON constant") return json.loads(text, object_pairs_hook=unique, parse_constant=invalid_constant) def identifier(value): return isinstance(value, str) and IDENTIFIER.fullmatch(value) is not None def validate_request(request): if not isinstance(request, dict) or set(request) != {"request_id", "tool", "arguments"}: fail("INVALID_REQUEST", "Require exactly request_id, tool, and arguments.") if not identifier(request["request_id"]): fail("INVALID_REQUEST", "request_id must be 1 to 64 ASCII letters, digits, underscores or hyphens, starting with a letter or digit.") tool, args = request["tool"], request["arguments"] if not isinstance(tool, str) or tool not in ("inspect_inputs", "analyze", "verify_result"): fail("UNKNOWN_TOOL", "The requested tool is not allowlisted.") if not isinstance(args, dict): fail("INVALID_ARGUMENTS", "arguments must be a JSON object.") if tool == "inspect_inputs" and args: fail("INVALID_ARGUMENTS", "inspect_inputs takes an empty arguments object.") if tool == "analyze": if not args: fail("MISSING_CONTEXT", "Inspect inputs first and supply expected_source_digest.") if set(args) != {"expected_source_digest"} or not isinstance(args["expected_source_digest"], str) or not HASH.fullmatch(args["expected_source_digest"]): fail("INVALID_ARGUMENTS", "analyze requires exactly one SHA256 expected_source_digest.") if tool == "verify_result" and (set(args) != {"run_id"} or not identifier(args["run_id"])): fail("INVALID_ARGUMENTS", "verify_result requires exactly one valid run_id.") def within(path, parent): return path == parent or parent in path.parents def checked_path(workspace, path): resolved = path.resolve() if not within(resolved, workspace): fail("PATH_BOUNDARY", "A selected path or existing link escapes the workspace.") return resolved def paths(workspace_arg, output_arg, cancel_arg): workspace = Path(workspace_arg).resolve() if not workspace.is_dir(): fail("MISSING_WORKSPACE", "The operator-selected workspace does not exist.") selected = Path(output_arg) output = checked_path(workspace, selected if selected.is_absolute() else workspace / selected) for name in ("data", "src", "contracts"): protected = checked_path(workspace, workspace / name) if within(output, protected) or within(protected, output): fail("PATH_BOUNDARY", "The output directory overlaps protected inputs, definitions, or source code.") if output.exists() and not output.is_dir(): fail("INVALID_OUTPUT", "The selected output location is not a directory.") cancel = None if cancel_arg: selected = Path(cancel_arg) cancel = checked_path(workspace, selected if selected.is_absolute() else workspace / selected) return workspace, output, cancel def cancelled(cancel): if cancel is not None and cancel.exists(): fail("CANCELLED", "The local cancellation marker was present before commit; no committed result was created.") def read_bounded(path, missing_code="MISSING_INPUT", max_bytes=1_048_576): if not path.is_file(): fail(missing_code, "A required file is missing or is not a regular file.") with path.open("rb") as handle: content = handle.read(max_bytes + 1) if len(content) > max_bytes: fail("INPUT_LIMIT", "A file exceeds its configured size limit.") return content def source_snapshot(workspace): contents = {name: read_bounded(checked_path(workspace, workspace / name)) for name in SOURCE_PATHS} hashes = {name: sha(content) for name, content in contents.items()} return contents, hashes, sha(canonical(hashes)) def rows(contents, name, fields): try: reader = csv.DictReader(io.StringIO(contents[name].decode("utf-8-sig")), strict=True) if reader.fieldnames != fields: fail("INVALID_DATA", "CSV headers must match the fixed analytical contract.") data = list(reader) except (UnicodeError, csv.Error): fail("INVALID_DATA", "Source data must be valid UTF-8 CSV.") if not data or len(data) > 10000 or any(set(row) != set(fields) or any(value is None for value in row.values()) for row in data): fail("INVALID_DATA", "CSV rows must be nonempty, rectangular, and at most 10000 rows per source.") return data def number(value, integer=False, nonnegative=False): if not NUMBER.fullmatch(value): fail("INVALID_DATA", "Numeric cells require finite plain decimal values within the documented bounds.") result = Decimal(value) if (integer and result != result.to_integral_value()) or (nonnegative and result < 0): fail("INVALID_DATA", "Counts must be integral and nonnegative where required.") return result def unique_keys(data, key): values = [row[key] for row in data] if any(not identifier(value) for value in values) or len(values) != len(set(values)): fail("INVALID_DATA", "Transaction and group identifiers must be valid and unique.") def load_data(contents): try: definitions = strict_json(contents["contracts/metric-definitions.json"].decode("utf-8-sig")) except (ValueError, UnicodeError): fail("INVALID_DEFINITIONS", "The metric definition file is malformed.") if definitions != DEFINITIONS: fail("INVALID_DEFINITIONS", "This engine supports only the complete synthetic-v1 definitions.") transactions = rows(contents, "data/transactions.csv", ["transaction_id", "customer_id", "quantity", "unit_price", "discount_rate"]) groups = rows(contents, "data/groups.csv", ["group", "successes", "eligible"]) customers = rows(contents, "data/customers.csv", ["customer_id", "segment"]) unique_keys(transactions, "transaction_id") unique_keys(groups, "group") for row in transactions: if not identifier(row["customer_id"]): fail("INVALID_DATA", "Customer identifiers must satisfy the identifier contract.") row["quantity"] = number(row["quantity"], integer=True) row["unit_price"] = None if row["unit_price"] == "" else number(row["unit_price"], nonnegative=True) row["discount_rate"] = number(row["discount_rate"], nonnegative=True) if row["discount_rate"] > 1: fail("INVALID_DATA", "Discount rates must be between zero and one.") for row in groups: row["successes"] = int(number(row["successes"], integer=True, nonnegative=True)) row["eligible"] = int(number(row["eligible"], integer=True, nonnegative=True)) if row["successes"] > row["eligible"]: fail("INVALID_DATA", "Successes cannot exceed eligible counts.") for row in customers: if not identifier(row["customer_id"]) or not identifier(row["segment"]): fail("INVALID_DATA", "Customer and segment identifiers must satisfy the identifier contract.") return transactions, groups, customers def decimal_text(value): text = format(value, "f") return text.rstrip("0").rstrip(".") if "." in text else text def compute(contents, digest): transactions, groups, customers = load_data(contents) successes, eligible = sum(row["successes"] for row in groups), sum(row["eligible"] for row in groups) contributions, subtotal, known = [], Decimal(0), 0 with localcontext() as context: context.prec = 80 rate = format((Decimal(successes) / Decimal(eligible)).quantize(Decimal("0.000000000001")), "f") if eligible else None for row in transactions: value = None if row["unit_price"] is None else row["quantity"] * row["unit_price"] * (1 - row["discount_rate"]) contributions.append({"transaction_id": row["transaction_id"], "net_revenue": decimal_text(value) if value is not None else None}) if value is not None: subtotal += value known += 1 counts = {} for row in customers: counts[row["customer_id"]] = counts.get(row["customer_id"], 0) + 1 duplicates = sorted(key for key, count in counts.items() if count > 1) unmatched = sorted({row["customer_id"] for row in transactions} - set(counts)) blockers = [] if known != len(transactions): blockers.append("INCOMPLETE_REVENUE") if duplicates: blockers.append("DUPLICATE_CUSTOMER_KEY") if unmatched: blockers.append("UNMATCHED_CUSTOMER_KEY") if not eligible: blockers.append("ZERO_ELIGIBLE") return { "definition_id": "synthetic-v1", "source_digest": digest, "synthetic": True, "conversion": {"successes": successes, "eligible": eligible, "rate": rate, "method": "pooled_denominator"}, "revenue": {"units": "fictional_currency_units", "known_subtotal": decimal_text(subtotal), "known_rows": known, "total_rows": len(transactions), "complete_total": decimal_text(subtotal) if known == len(transactions) else None, "contributions": contributions}, "customer_join": {"status": "blocked" if duplicates or unmatched else "eligible", "duplicate_keys": duplicates, "unmatched_keys": unmatched, "joined_output": "not_produced"}, "release": {"eligible": not blockers, "blockers": blockers, "human_approval": "not_granted"}, } def write_file(path, content): with path.open("xb") as handle: handle.write(content) handle.flush() os.fsync(handle.fileno()) def commit(workspace, output, cancel, run_id, result, hashes, digest): cancelled(cancel) checked_path(workspace, output) output.mkdir(parents=True, exist_ok=True) lock = checked_path(workspace, output / ".writer.lock") try: lock_handle = lock.open("x", encoding="utf-8") except FileExistsError: fail("OUTPUT_BUSY", "An output writer lock already exists; investigate its owner before retrying.", True) stage = None created = [] try: with lock_handle: lock_handle.write(run_id + "\n") destination = checked_path(workspace, output / run_id) if destination.exists(): fail("OUTPUT_COLLISION", "This run_id already exists; choose a new request_id.") cancelled(cancel) stage = output / (".pending-" + run_id + "-" + uuid.uuid4().hex) stage.mkdir() result_bytes = canonical(result) receipt = {"protocol_version": PROTOCOL, "run_id": run_id, "source_digest": digest, "sources": hashes, "result_sha256": sha(result_bytes), "engine_sha256": sha(Path(__file__).read_bytes()), "definition_id": "synthetic-v1"} receipt_bytes = canonical(receipt) files = {"result.json": result_bytes, "receipt.json": receipt_bytes, "checksums.json": canonical({"result.json": sha(result_bytes), "receipt.json": sha(receipt_bytes)})} for name, content in files.items(): cancelled(cancel) path = stage / name created.append(path) write_file(path, content) if source_snapshot(workspace)[2] != digest: fail("STALE_CONTEXT", "Inputs changed while computing; no committed result was created.") checked_path(workspace, output) checked_path(workspace, destination) cancelled(cancel) if destination.exists(): fail("OUTPUT_COLLISION", "This run_id already exists; choose a new request_id.") os.rename(stage, destination) stage = None return {"result": (destination / "result.json").relative_to(workspace).as_posix(), "receipt": (destination / "receipt.json").relative_to(workspace).as_posix(), "checksums": (destination / "checksums.json").relative_to(workspace).as_posix()} finally: # Only files created by this invocation are candidates for cleanup. if stage is not None: for path in reversed(created): if path.exists(): path.unlink() stage.rmdir() lock.unlink() def verify(workspace, output, run_id): location = checked_path(workspace, output / run_id) raw = {name: read_bounded(checked_path(workspace, location / name), "MISSING_RESULT", max_bytes=4_194_304) for name in ("result.json", "receipt.json", "checksums.json")} try: checksums = strict_json(raw["checksums.json"].decode("utf-8")) if not isinstance(checksums, dict) or set(checksums) != {"result.json", "receipt.json"}: raise ValueError("invalid checksum map") if checksums["receipt.json"] != sha(raw["receipt.json"]): fail("RECEIPT_TAMPERED", "The saved receipt checksum does not match.") if checksums["result.json"] != sha(raw["result.json"]): fail("RESULT_TAMPERED", "The saved result checksum does not match.") receipt, result = strict_json(raw["receipt.json"].decode("utf-8")), strict_json(raw["result.json"].decode("utf-8")) if not isinstance(receipt, dict) or set(receipt) != {"protocol_version", "run_id", "source_digest", "sources", "result_sha256", "engine_sha256", "definition_id"}: raise ValueError("invalid receipt") if receipt["run_id"] != run_id or receipt["protocol_version"] != PROTOCOL or receipt["definition_id"] != "synthetic-v1": fail("RECEIPT_TAMPERED", "The receipt identity does not match this run.") if receipt["result_sha256"] != sha(raw["result.json"]): fail("RESULT_TAMPERED", "The result does not match its receipt.") contents, hashes, digest = source_snapshot(workspace) if receipt["sources"] != hashes or receipt["source_digest"] != digest: fail("STALE_SOURCES", "Current source bytes or metric definitions differ from the saved run.") if receipt["engine_sha256"] != sha(Path(__file__).read_bytes()): fail("ENGINE_CHANGED", "Engine code differs from the implementation bound to the saved run.") if result != compute(contents, digest): fail("RESULT_TAMPERED", "The saved metrics do not match a fresh computation.") except (UnicodeError, ValueError, KeyError, TypeError): fail("INVALID_SAVED_RESULT", "A saved result, receipt or checksum file is malformed.") return {"run_id": run_id, "integrity": "verified", "source_digest": digest, "release_eligible": result["release"]["eligible"], "human_approval": "not_granted", "checks": ["source_hashes", "definition_hash", "engine_hash", "result_checksum", "receipt_checksum", "recomputed_metrics"]} class Parser(argparse.ArgumentParser): def error(self, message): fail("INVALID_OPTIONS", "Require --workspace and --output-dir, with optional --cancel-file.") def main(argv=None): request_id, tool = None, None response = {"protocol_version": PROTOCOL, "request_id": None, "ok": False, "status": "error", "tool": None, "result": None, "error": None, "artifacts": {}, "verification": dict(SCOPE)} exit_code = 1 try: parser = Parser(add_help=False) parser.add_argument("--workspace", required=True) parser.add_argument("--output-dir", required=True) parser.add_argument("--cancel-file") options = parser.parse_args(argv) text = sys.stdin.read(65_537) if len(text) > 65_536: fail("REQUEST_LIMIT", "The JSON request exceeds 65536 characters.") try: request = strict_json(text) except (ValueError, RecursionError): fail("INVALID_JSON", "The request is not one strict JSON document.") if isinstance(request, dict): request_id = request.get("request_id") if identifier(request.get("request_id")) else None tool = request.get("tool") if request.get("tool") in ("inspect_inputs", "analyze", "verify_result") else None validate_request(request) workspace, output, cancel = paths(options.workspace, options.output_dir, options.cancel_file) cancelled(cancel) if tool == "verify_result": result = verify(workspace, output, request["arguments"]["run_id"]) else: contents, hashes, digest = source_snapshot(workspace) if tool == "analyze" and request["arguments"]["expected_source_digest"] != digest: fail("STALE_CONTEXT", "Inputs changed since inspection; inspect again before analysis.") result = compute(contents, digest) if tool == "inspect_inputs": loaded = load_data(contents) result = {"source_digest": digest, "sources": hashes, "definition_id": "synthetic-v1", "row_counts": {"transactions": len(loaded[0]), "groups": len(loaded[1]), "customers": len(loaded[2])}, "release_blockers": result["release"]["blockers"]} else: response["artifacts"] = commit(workspace, output, cancel, request_id, result, hashes, digest) response.update(ok=True, status="success", result=result) exit_code = 0 except ToolError as error: response["error"] = {"code": error.code, "message": error.message, "retryable": error.retryable} response["status"] = "cancelled" if error.code == "CANCELLED" else "error" exit_code = 3 if error.code == "CANCELLED" else 1 except (OSError, InvalidOperation): response["error"] = {"code": "IO_OR_NUMERIC_ERROR", "message": "A local filesystem or numeric operation failed; no success is claimed.", "retryable": False} except Exception: response["error"] = {"code": "INTERNAL_ERROR", "message": "An unexpected executor failure occurred; no success is claimed.", "retryable": False} response.update(request_id=request_id, tool=tool) sys.stdout.write(canonical(response).decode("utf-8")) return exit_code if __name__ == "__main__": sys.exit(main())