785 lines
32 KiB
Python
Executable File
785 lines
32 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""Plan-18 Chat turn entrypoint.
|
|
|
|
Thin harness-owned router for Control Panel: classify once, then delegate to the
|
|
mode primitive. NestJS calls this file only; it does not own governance verdicts.
|
|
"""
|
|
import argparse
|
|
import hashlib
|
|
import json
|
|
import os
|
|
import subprocess
|
|
import sys
|
|
import tempfile
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
|
|
|
|
def project_root() -> str:
|
|
d = os.path.abspath(os.path.dirname(__file__))
|
|
p = d
|
|
while p != os.path.dirname(p):
|
|
if os.path.isdir(os.path.join(p, ".specify")) or os.path.isdir(os.path.join(p, "packages/casan-harness")):
|
|
return p
|
|
p = os.path.dirname(p)
|
|
return os.path.abspath(os.path.join(d, "..", "..", ".."))
|
|
|
|
|
|
ROOT = project_root()
|
|
BIN = os.path.join(ROOT, "packages", "casan-harness", "scripts", "bash")
|
|
ROUTER = os.path.join(BIN, "prompt-mode-router.py")
|
|
AGENT_RESOLVER = os.path.join(BIN, "chat-agent-resolver.py")
|
|
MODEL_RESOLVER = os.path.join(BIN, "chat-model-resolver.py")
|
|
READONLY = os.path.join(BIN, "chat-readonly.py")
|
|
OPERATOR = os.path.join(BIN, "chat-operator.py")
|
|
APPROVAL_INBOX = os.path.join(BIN, "approval-inbox.py")
|
|
LOOP_RUN = os.path.join(BIN, "loop-run.sh")
|
|
LOOP_TRACE = os.path.join(BIN, "loop-trace.py")
|
|
PREFLIGHT = os.path.join(BIN, "harness-preflight.sh")
|
|
CONTEXT_SCAN = os.path.join(BIN, "context-assemble-scan.sh")
|
|
TOOL_OUTPUT_SCAN = os.path.join(BIN, "tool-output-scan.sh")
|
|
ARTIFACT_SCAN = os.path.join(BIN, "artifact-scan.sh")
|
|
TENANT_STORE = os.path.join(BIN, "tenant-store.sh")
|
|
TENANT_CRYPT = os.path.join(BIN, "tenant-crypt.sh")
|
|
KILL_SWITCH = os.path.join(BIN, "kill-switch.sh")
|
|
COST_SPIKE = os.path.join(BIN, "cost-spike-detect.sh")
|
|
MODEL_ROUTER = os.path.join(BIN, "model-router.sh")
|
|
MODEL_PROVIDERS = os.path.join(ROOT, "packages", "casan-harness", "config", "model-providers.yaml")
|
|
|
|
|
|
def state_root() -> str:
|
|
return os.environ.get("CASAN_STATE_ROOT") or os.path.join(ROOT, ".specify")
|
|
|
|
|
|
def tenant_path(logical: str, fallback: str) -> str:
|
|
if os.environ.get("CASAN_TENANT_ID"):
|
|
r = subprocess.run(["bash", TENANT_STORE, "resolve", logical], cwd=ROOT, capture_output=True, text=True)
|
|
if r.returncode != 0:
|
|
raise SystemExit((r.stderr or r.stdout or "TENANT_DENIED").strip())
|
|
return r.stdout.strip()
|
|
return os.path.join(state_root(), fallback)
|
|
|
|
|
|
def guarded_override(path: str) -> str:
|
|
if path and os.environ.get("CASAN_TENANT_ID"):
|
|
r = subprocess.run(["bash", TENANT_STORE, "guard", path], cwd=ROOT, capture_output=True, text=True)
|
|
if r.returncode != 0:
|
|
raise SystemExit((r.stderr or r.stdout or "TENANT_DENIED").strip())
|
|
return path
|
|
|
|
|
|
def now_iso() -> str:
|
|
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
|
|
|
|
def audit_path() -> str:
|
|
return guarded_override(os.environ["CASAN_CHAT_AUDIT_LOG"]) if os.environ.get("CASAN_CHAT_AUDIT_LOG") else tenant_path("chat/chat-turns.jsonl", "logs/chat/chat-turns.jsonl")
|
|
|
|
|
|
def head_path() -> str:
|
|
return guarded_override(os.environ["CASAN_CHAT_AUDIT_HEAD"]) if os.environ.get("CASAN_CHAT_AUDIT_HEAD") else tenant_path("chat/chat-head.txt", "logs/chat/chat-head.txt")
|
|
|
|
|
|
def sha(text: str) -> str:
|
|
return hashlib.sha256(text.encode("utf-8")).hexdigest()
|
|
|
|
|
|
def sha_file(path: str) -> str:
|
|
h = hashlib.sha256()
|
|
with open(path, "rb") as fh:
|
|
for chunk in iter(lambda: fh.read(65536), b""):
|
|
h.update(chunk)
|
|
return h.hexdigest()
|
|
|
|
|
|
def classify(message: str):
|
|
r = subprocess.run(["python3", ROUTER, "classify", "--message", message], cwd=ROOT, capture_output=True, text=True)
|
|
try:
|
|
return json.loads(r.stdout)
|
|
except Exception:
|
|
return {"mode": "BLOCK", "risk": "high", "reason": "router_invalid_json", "matched_rules": [r.stderr.strip()]}
|
|
|
|
|
|
def scan_tool_output(text: str, label: str):
|
|
with tempfile.TemporaryDirectory() as td:
|
|
path = os.path.join(td, "tool-output.txt")
|
|
with open(path, "w", encoding="utf-8") as fh:
|
|
fh.write(text)
|
|
r = subprocess.run(["bash", TOOL_OUTPUT_SCAN, path, label], cwd=ROOT, capture_output=True, text=True)
|
|
return r.returncode, (r.stdout + r.stderr).strip()
|
|
|
|
|
|
def run_mode(args, binding, loop_run=None, scan_output_label=""):
|
|
env = os.environ.copy()
|
|
if loop_run is not None:
|
|
env["CASAN_CHAT_LOOP_RUN_JSON"] = json.dumps(loop_run, ensure_ascii=False, sort_keys=True)
|
|
if binding and binding.get("model_binding", {}).get("decision") == "BOUND":
|
|
env["CASAN_CHAT_MODEL_PROVIDER"] = binding["model_binding"].get("provider", "")
|
|
r = subprocess.run(args, cwd=ROOT, capture_output=True, text=True, env=env)
|
|
if scan_output_label:
|
|
scan_rc, scan_msg = scan_tool_output(r.stdout + "\n" + r.stderr, scan_output_label)
|
|
if scan_rc != 0:
|
|
print(json.dumps({
|
|
"success": False,
|
|
"mode": binding.get("mode") if binding else "OPERATOR",
|
|
"risk": "high",
|
|
"decision": "DENIED",
|
|
"answer": "Denied by H4 tool-output scan.",
|
|
"sources": [],
|
|
"certified": False,
|
|
"audit": {},
|
|
"router": {"reason": "tool_output_scan_denied", "matched_rules": [scan_msg]},
|
|
"agent_binding": binding,
|
|
"loop_run": loop_run,
|
|
}, ensure_ascii=False))
|
|
return 2
|
|
try:
|
|
payload = json.loads(r.stdout)
|
|
payload["agent_binding"] = binding
|
|
if loop_run is not None:
|
|
payload["loop_run"] = loop_run
|
|
print(json.dumps(payload, ensure_ascii=False))
|
|
except Exception:
|
|
if r.stdout:
|
|
sys.stdout.write(r.stdout)
|
|
if r.stderr:
|
|
sys.stderr.write(r.stderr)
|
|
return r.returncode
|
|
|
|
|
|
def run_and_passthrough(args, env=None):
|
|
r = subprocess.run(args, cwd=ROOT, text=True, env=env)
|
|
return r.returncode
|
|
|
|
|
|
def default_agent_for_mode(mode: str) -> str:
|
|
if mode == "OPERATOR":
|
|
return "ops-operator"
|
|
if mode == "CODEGEN":
|
|
return "codegen-draft"
|
|
return "evidence-reader"
|
|
|
|
|
|
def write_json(path: str, payload):
|
|
os.makedirs(os.path.dirname(path), exist_ok=True)
|
|
with open(path, "w", encoding="utf-8") as fh:
|
|
json.dump(payload, fh, ensure_ascii=False, indent=2, sort_keys=True)
|
|
|
|
|
|
def write_text(path: str, text: str):
|
|
os.makedirs(os.path.dirname(path), exist_ok=True)
|
|
with open(path, "w", encoding="utf-8") as fh:
|
|
fh.write(text)
|
|
|
|
|
|
def loop_dir(run_id: str) -> str:
|
|
return tenant_path(f"chat/loop-runs/{run_id}", f"logs/chat/loop-runs/{run_id}")
|
|
|
|
|
|
def codegen_dir(run_id: str) -> str:
|
|
return tenant_path(f"chat/codegen-artifacts/{run_id}", f"logs/chat/codegen-artifacts/{run_id}")
|
|
|
|
|
|
def load_chat_head() -> str:
|
|
try:
|
|
return open(head_path(), encoding="utf-8").read().strip() or ("0" * 64)
|
|
except OSError:
|
|
return "0" * 64
|
|
|
|
|
|
def append_chat_turn(base):
|
|
path = audit_path()
|
|
os.makedirs(os.path.dirname(path), exist_ok=True)
|
|
seq = 1
|
|
if os.path.isfile(path):
|
|
with open(path, encoding="utf-8") as fh:
|
|
seq = sum(1 for line in fh if line.strip()) + 1
|
|
prev = load_chat_head()
|
|
core = {"seq": seq, **base, "prev_hash": prev}
|
|
record_hash = sha(json.dumps(core, sort_keys=True, ensure_ascii=False))
|
|
rec = {**core, "record_hash": record_hash}
|
|
with open(path, "a", encoding="utf-8") as fh:
|
|
fh.write(json.dumps(rec, ensure_ascii=False) + "\n")
|
|
os.makedirs(os.path.dirname(head_path()), exist_ok=True)
|
|
with open(head_path(), "w", encoding="utf-8") as fh:
|
|
fh.write(record_hash + "\n")
|
|
encrypt_chat_audit_snapshot(path)
|
|
return rec
|
|
|
|
|
|
def encrypt_chat_audit_snapshot(path: str):
|
|
if not os.environ.get("CASAN_TENANT_ID"):
|
|
return
|
|
if not os.path.isfile(path):
|
|
return
|
|
enc = path + ".enc"
|
|
subprocess.run(["bash", TENANT_CRYPT, "encrypt", path, enc], cwd=ROOT, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
|
|
|
|
|
|
def tenant_runtime_guard(args):
|
|
if args.tenant and args.tenant != "default":
|
|
os.environ["CASAN_TENANT_ID"] = args.tenant
|
|
if os.environ.get("CASAN_TENANT_ID"):
|
|
ks = subprocess.run(["bash", KILL_SWITCH, "check", "tenant", os.environ["CASAN_TENANT_ID"]], cwd=ROOT, capture_output=True, text=True)
|
|
if ks.returncode != 0:
|
|
print(json.dumps({
|
|
"success": False,
|
|
"mode": "BLOCK",
|
|
"risk": "high",
|
|
"decision": "HALTED",
|
|
"answer": f"Tenant kill-switch active for {os.environ['CASAN_TENANT_ID']}.",
|
|
"sources": [],
|
|
"certified": False,
|
|
"audit": {},
|
|
"router": {"reason": "tenant_kill_switch_active", "matched_rules": ["tenant_kill_switch"]},
|
|
}, ensure_ascii=False))
|
|
return 3
|
|
if os.environ.get("CASAN_COST_CUMULATIVE_BUDGET_TOKENS"):
|
|
quota = subprocess.run(["bash", COST_SPIKE], cwd=ROOT, capture_output=True, text=True)
|
|
if quota.returncode == 2:
|
|
print(json.dumps({
|
|
"success": False,
|
|
"mode": "BLOCK",
|
|
"risk": "high",
|
|
"decision": "HALTED",
|
|
"answer": f"Tenant quota exceeded for {os.environ['CASAN_TENANT_ID']}.",
|
|
"sources": [],
|
|
"certified": False,
|
|
"audit": {},
|
|
"router": {"reason": "tenant_quota_exceeded", "matched_rules": ["tenant_quota"]},
|
|
}, ensure_ascii=False))
|
|
return 3
|
|
return 0
|
|
|
|
|
|
def loop_env_for(args):
|
|
env = {
|
|
**os.environ,
|
|
"CASAN_LOOP_STATE_ROOT": os.environ.get("CASAN_LOOP_STATE_ROOT") or tenant_path("chat/loop-state", "logs/chat/loop-state"),
|
|
}
|
|
if os.environ.get("CASAN_TENANT_ID"):
|
|
env["CASAN_TENANT_ID"] = os.environ["CASAN_TENANT_ID"]
|
|
return env
|
|
|
|
|
|
def run_preflight_and_context(run_id: str, draft_path: str):
|
|
preflight_out = os.path.join(loop_dir(run_id), "preflight.json")
|
|
r = subprocess.run(["bash", PREFLIGHT, draft_path, preflight_out, "--model", "local:chat-turn"], cwd=ROOT, capture_output=True, text=True)
|
|
if r.returncode != 0:
|
|
return False, "harness_preflight_failed", (r.stdout + r.stderr).strip()
|
|
r = subprocess.run(["bash", CONTEXT_SCAN, draft_path], cwd=ROOT, capture_output=True, text=True)
|
|
if r.returncode != 0:
|
|
return False, "context_assemble_scan_failed", (r.stdout + r.stderr).strip()
|
|
return True, "preflight_pass", (r.stdout + r.stderr).strip()
|
|
|
|
|
|
def run_artifact_scan(path: str, label: str):
|
|
r = subprocess.run(["bash", ARTIFACT_SCAN, path, label], cwd=ROOT, capture_output=True, text=True)
|
|
return r.returncode == 0, (r.stdout + r.stderr).strip()
|
|
|
|
|
|
def certify_operator_draft(args, router, binding):
|
|
run_id = f"chat-{sha('|'.join([args.chat_id or 'chat-default', args.turn_id or '', args.message]))[:16]}"
|
|
d = loop_dir(run_id)
|
|
draft_path = os.path.join(d, "draft.txt")
|
|
criteria_path = os.path.join(d, "success-criteria.json")
|
|
draft = "\n".join([
|
|
"CHAT_TURN",
|
|
"ACTION_HELD",
|
|
f"chat_id={args.chat_id or 'chat-default'}",
|
|
f"turn_id={args.turn_id or 'auto'}",
|
|
f"mode={router.get('mode')}",
|
|
f"agent={binding.get('agent_selected')}",
|
|
f"skill={binding.get('skill_selected')}",
|
|
f"delegation_level={binding.get('delegation_level')}",
|
|
f"user_msg_ref={sha(args.message)}",
|
|
f"user_message_preview={args.message[:160]}",
|
|
"",
|
|
])
|
|
write_text(draft_path, draft)
|
|
write_json(criteria_path, {
|
|
"must_contain": ["CHAT_TURN", "ACTION_HELD", f"agent={binding.get('agent_selected')}"],
|
|
"must_not_contain": ["BYPASS_LOOP_GATE"],
|
|
})
|
|
|
|
ok, preflight_reason, preflight_detail = run_preflight_and_context(run_id, draft_path)
|
|
if not ok:
|
|
return {
|
|
"success": False,
|
|
"decision": "DENIED",
|
|
"reason": preflight_reason,
|
|
"detail": preflight_detail,
|
|
"run_id": run_id,
|
|
"draft_ref": sha(draft),
|
|
"draft_certified": False,
|
|
"side_effect_released": False,
|
|
"artifact": draft_path,
|
|
"success_criteria": criteria_path,
|
|
}, 2
|
|
|
|
loop_env = loop_env_for(args)
|
|
dlevel = f"L{binding.get('delegation_level', 0)}"
|
|
r = subprocess.run([
|
|
"bash", LOOP_RUN,
|
|
"--run-id", run_id,
|
|
"--artifact", draft_path,
|
|
"--success-criteria", criteria_path,
|
|
"--profile", os.environ.get("CASAN_PROFILE", "dev"),
|
|
"--delegation-level", dlevel,
|
|
"--project", args.project,
|
|
"--max-steps", "2",
|
|
"--tokens-per-step", "128",
|
|
"--cost-per-step", "0",
|
|
"--context", draft_path,
|
|
], cwd=ROOT, capture_output=True, text=True, env=loop_env)
|
|
trace_rc = subprocess.run(["python3", LOOP_TRACE, "verify-chain", "--run-id", run_id], cwd=ROOT, capture_output=True, text=True, env=loop_env)
|
|
replay_rc = subprocess.run(["python3", LOOP_TRACE, "replay", "--run-id", run_id, "--profile", os.environ.get("CASAN_PROFILE", "dev")], cwd=ROOT, capture_output=True, text=True, env=loop_env)
|
|
certified = r.returncode == 0 and "final=DONE" in r.stdout and trace_rc.returncode == 0 and replay_rc.returncode == 0
|
|
return {
|
|
"success": certified,
|
|
"decision": "CERTIFIED" if certified else "HALTED",
|
|
"reason": "loop_run_pass" if certified else "loop_run_failed",
|
|
"run_id": run_id,
|
|
"draft_ref": sha(draft),
|
|
"draft_certified": certified,
|
|
"side_effect_released": certified,
|
|
"artifact": draft_path,
|
|
"success_criteria": criteria_path,
|
|
"output": (r.stdout + r.stderr).strip(),
|
|
"trace_verify": {"ok": trace_rc.returncode == 0, "output": (trace_rc.stdout + trace_rc.stderr).strip()},
|
|
"replay": {"ok": replay_rc.returncode == 0, "output": (replay_rc.stdout + replay_rc.stderr).strip()},
|
|
}, (0 if certified else 3)
|
|
|
|
|
|
def render_codegen_draft(args, binding, body=None, meta=None) -> str:
|
|
fn = "generated_chat_draft"
|
|
scaffold = "\n".join([
|
|
"# CODEGEN_DRAFT",
|
|
"# GENERATED_BY_CASAN_CHAT",
|
|
f"# agent={binding.get('agent_selected')}",
|
|
f"# skill={binding.get('skill_selected')}",
|
|
f"# user_message_preview={args.message[:200]}",
|
|
"",
|
|
f"def {fn}():",
|
|
" \"\"\"Deterministic draft generated by governed chat; review before use.\"\"\"",
|
|
" return {",
|
|
f" \"request_hash\": \"{sha(args.message)}\",",
|
|
" \"status\": \"draft_only\",",
|
|
" }",
|
|
"",
|
|
])
|
|
if body:
|
|
meta = meta or {}
|
|
scaffold += "\n".join([
|
|
"# === MODEL_DRAFT BEGIN (review-only; never auto-applied) ===",
|
|
f"# provider={meta.get('provider')} model={meta.get('model')}",
|
|
body,
|
|
"# === MODEL_DRAFT END ===",
|
|
"",
|
|
])
|
|
return scaffold
|
|
|
|
|
|
def _model_codegen_body(args):
|
|
"""Item 4: full model-router CODEGEN path. Offline-first — model synthesis is
|
|
gated behind CASAN_CHAT_MODEL_MODE=model; any failure falls SAFE back to the
|
|
deterministic scaffold. The generated code stays draft-only and is still run
|
|
through artifact-scan + Plan-17 loop certification by the caller.
|
|
"""
|
|
mode = os.environ.get("CASAN_CHAT_MODEL_MODE", "off").strip().lower()
|
|
if mode != "model":
|
|
return None, {"mode": "deterministic", "reason": "model_mode_off"}
|
|
try:
|
|
cfg = json.load(open(os.environ.get("CASAN_MODEL_PROVIDERS_FILE") or MODEL_PROVIDERS, encoding="utf-8"))
|
|
except Exception:
|
|
return None, {"mode": "deterministic", "reason": "providers_unreadable"}
|
|
providers = cfg.get("providers", {})
|
|
bindings = cfg.get("role_bindings", {})
|
|
provider_id = os.environ.get("CASAN_CHAT_MODEL_PROVIDER") or bindings.get("codegen", "") or bindings.get("read_only", "")
|
|
provider = providers.get(provider_id, {})
|
|
model_spec = os.environ.get("CASAN_CHAT_SELECTED_MODEL") or provider.get("model")
|
|
if not model_spec:
|
|
return None, {"mode": "deterministic", "reason": "provider_unresolved"}
|
|
pclass = provider.get("class", "local")
|
|
if provider.get("requires_key"):
|
|
key_env = provider.get("key_env", "")
|
|
if key_env and not os.environ.get(key_env):
|
|
return None, {"mode": "deterministic", "reason": "provider_key_unset"}
|
|
router = os.environ.get("CASAN_CHAT_MODEL_ROUTER") or MODEL_ROUTER
|
|
env = os.environ.copy()
|
|
if pclass == "cloud" or provider.get("requires_preflight"):
|
|
env["CASAN_PREFLIGHT"] = "1"
|
|
prompt = "\n".join([
|
|
"You are CASAN's governed codegen assistant. Produce a SMALL Python draft",
|
|
"fulfilling the request. Output code only. No shell, no network, no file I/O.",
|
|
f"REQUEST: {args.message[:400]}",
|
|
])
|
|
with tempfile.TemporaryDirectory() as td:
|
|
pf = os.path.join(td, "p.txt")
|
|
oj = os.path.join(td, "o.json")
|
|
write_text(pf, prompt)
|
|
r = subprocess.run(["bash", router, pf, oj, "--role", "generate", "--model", model_spec],
|
|
cwd=ROOT, capture_output=True, text=True, env=env)
|
|
fallback_from = ""
|
|
if (r.returncode != 0 or not os.path.isfile(oj)) and provider_id != "local":
|
|
fallback = providers.get(os.environ.get("CASAN_CHAT_LOCAL_FALLBACK_PROVIDER", "local"), {})
|
|
fallback_model = fallback.get("model")
|
|
if fallback_model and fallback.get("class", "local") == "local":
|
|
fallback_from = provider_id
|
|
r = subprocess.run(["bash", router, pf, oj, "--role", "generate", "--model", fallback_model],
|
|
cwd=ROOT, capture_output=True, text=True, env=os.environ.copy())
|
|
if r.returncode == 0 and os.path.isfile(oj):
|
|
provider_id, model_spec = "local", fallback_model
|
|
if r.returncode != 0 or not os.path.isfile(oj):
|
|
return None, {"mode": "deterministic", "reason": "model_unavailable"}
|
|
try:
|
|
out = json.load(open(oj, encoding="utf-8"))
|
|
except Exception:
|
|
return None, {"mode": "deterministic", "reason": "model_output_unreadable"}
|
|
text = (out.get("text") or "").strip()
|
|
if not text:
|
|
return None, {"mode": "deterministic", "reason": "model_empty"}
|
|
return text, {
|
|
"mode": "model",
|
|
"provider": provider_id,
|
|
"model": model_spec,
|
|
"input_tokens": int(out.get("input_tokens") or 0),
|
|
"output_tokens": int(out.get("output_tokens") or 0),
|
|
"fallback_from": fallback_from or None,
|
|
}
|
|
|
|
|
|
def certify_codegen_draft(args, router, binding):
|
|
run_id = f"chat-codegen-{sha('|'.join([args.chat_id or 'chat-default', args.turn_id or '', args.message]))[:16]}"
|
|
d = codegen_dir(run_id)
|
|
artifact_path = os.path.join(d, "draft.py")
|
|
criteria_path = os.path.join(d, "success-criteria.json")
|
|
body, synth_meta = _model_codegen_body(args)
|
|
draft = render_codegen_draft(args, binding, body, synth_meta)
|
|
write_text(artifact_path, draft)
|
|
write_json(criteria_path, {
|
|
"must_contain": ["CODEGEN_DRAFT", "GENERATED_BY_CASAN_CHAT", "def generated_chat_draft"],
|
|
"must_not_contain": ["BYPASS_LOOP_GATE"],
|
|
})
|
|
|
|
scan_ok, scan_out = run_artifact_scan(artifact_path, "chat-codegen-draft")
|
|
if not scan_ok:
|
|
return {
|
|
"success": False,
|
|
"decision": "DENIED",
|
|
"reason": "artifact_scan_failed",
|
|
"run_id": run_id,
|
|
"draft_certified": False,
|
|
"side_effect_released": False,
|
|
"artifact": artifact_path,
|
|
"success_criteria": criteria_path,
|
|
"artifact_scan": {"ok": False, "output": scan_out},
|
|
"synthesis": synth_meta,
|
|
}, 2
|
|
|
|
loop_env = loop_env_for(args)
|
|
dlevel = f"L{binding.get('delegation_level', 0)}"
|
|
r = subprocess.run([
|
|
"bash", LOOP_RUN,
|
|
"--run-id", run_id,
|
|
"--artifact", artifact_path,
|
|
"--success-criteria", criteria_path,
|
|
"--profile", os.environ.get("CASAN_PROFILE", "dev"),
|
|
"--delegation-level", dlevel,
|
|
"--project", args.project,
|
|
"--max-steps", "2",
|
|
"--tokens-per-step", "128",
|
|
"--cost-per-step", "0",
|
|
"--context", artifact_path,
|
|
], cwd=ROOT, capture_output=True, text=True, env=loop_env)
|
|
trace_rc = subprocess.run(["python3", LOOP_TRACE, "verify-chain", "--run-id", run_id], cwd=ROOT, capture_output=True, text=True, env=loop_env)
|
|
replay_rc = subprocess.run(["python3", LOOP_TRACE, "replay", "--run-id", run_id, "--profile", os.environ.get("CASAN_PROFILE", "dev")], cwd=ROOT, capture_output=True, text=True, env=loop_env)
|
|
output_scan_rc, output_scan_msg = scan_tool_output(draft, "chat-codegen-output")
|
|
certified = (
|
|
r.returncode == 0
|
|
and "final=DONE" in r.stdout
|
|
and trace_rc.returncode == 0
|
|
and replay_rc.returncode == 0
|
|
and output_scan_rc == 0
|
|
)
|
|
return {
|
|
"success": certified,
|
|
"decision": "CERTIFIED" if certified else "HALTED",
|
|
"reason": "codegen_loop_pass" if certified else "codegen_loop_failed",
|
|
"run_id": run_id,
|
|
"draft_ref": sha(draft),
|
|
"draft_certified": certified,
|
|
"side_effect_released": False,
|
|
"artifact": artifact_path,
|
|
"success_criteria": criteria_path,
|
|
"artifact_scan": {"ok": True, "output": scan_out},
|
|
"tool_output_scan": {"ok": output_scan_rc == 0, "output": output_scan_msg},
|
|
"synthesis": synth_meta,
|
|
"output": (r.stdout + r.stderr).strip(),
|
|
"trace_verify": {"ok": trace_rc.returncode == 0, "output": (trace_rc.stdout + trace_rc.stderr).strip()},
|
|
"replay": {"ok": replay_rc.returncode == 0, "output": (replay_rc.stdout + replay_rc.stderr).strip()},
|
|
}, (0 if certified else 3)
|
|
|
|
|
|
def finish_codegen(args, router, binding, loop_run, rc: int):
|
|
artifact = loop_run.get("artifact", "")
|
|
source = []
|
|
if artifact and os.path.isfile(artifact):
|
|
source.append({
|
|
"path": os.path.relpath(artifact, ROOT) if artifact.startswith(ROOT) else artifact,
|
|
"line": 1,
|
|
"excerpt": "CODEGEN_DRAFT generated as review-only artifact.",
|
|
"score": 10 if rc == 0 else 1,
|
|
"hash": sha_file(artifact),
|
|
})
|
|
decision = "ANSWERED" if rc == 0 else ("DENIED" if loop_run.get("decision") == "DENIED" else "HALTED")
|
|
rec = append_chat_turn({
|
|
"timestamp": now_iso(),
|
|
"trace_id": loop_run.get("run_id"),
|
|
"chat_id": args.chat_id or "chat-default",
|
|
"turn_id": args.turn_id or str(uuid.uuid4()),
|
|
"tenant_id": args.tenant,
|
|
"actor": args.actor,
|
|
"mode": "CODEGEN",
|
|
"risk": router.get("risk", "high"),
|
|
"decision": decision,
|
|
"answer": f"Codegen draft certified: {os.path.relpath(artifact, ROOT) if artifact else 'n/a'}" if rc == 0 else f"Codegen draft held: {loop_run.get('reason')}",
|
|
"sources": source,
|
|
"router": router,
|
|
"agent_binding": binding,
|
|
"loop_run": loop_run,
|
|
})
|
|
print(json.dumps({
|
|
"success": rc == 0,
|
|
"mode": "CODEGEN",
|
|
"risk": router.get("risk", "high"),
|
|
"decision": decision,
|
|
"answer": f"Codegen draft certified: {os.path.relpath(artifact, ROOT) if artifact else 'n/a'}" if rc == 0 else f"Codegen draft held: {loop_run.get('reason')}",
|
|
"sources": source,
|
|
"certified": rc == 0,
|
|
"audit": {"seq": rec["seq"], "record_hash": rec["record_hash"], "head": rec["record_hash"], "path": audit_path()},
|
|
"router": router,
|
|
"agent_binding": binding,
|
|
"loop_run": loop_run,
|
|
"codegen": {
|
|
"artifact": os.path.relpath(artifact, ROOT) if artifact and artifact.startswith(ROOT) else artifact,
|
|
"artifact_scan": loop_run.get("artifact_scan"),
|
|
"tool_output_scan": loop_run.get("tool_output_scan"),
|
|
"synthesis": loop_run.get("synthesis"),
|
|
},
|
|
}, ensure_ascii=False))
|
|
return rc
|
|
|
|
|
|
def denied_from_loop(router, binding, loop_run):
|
|
print(json.dumps({
|
|
"success": False,
|
|
"mode": router.get("mode", "OPERATOR"),
|
|
"risk": router.get("risk", "medium"),
|
|
"decision": "DENIED" if loop_run.get("decision") == "DENIED" else "HALTED",
|
|
"answer": f"Operator side-effect held: {loop_run.get('reason')}",
|
|
"sources": [{
|
|
"path": os.path.relpath(loop_run.get("artifact", ""), ROOT) if loop_run.get("artifact") else "",
|
|
"line": 1,
|
|
"excerpt": "UNCERTIFIED draft held before side-effect.",
|
|
"score": 1,
|
|
}],
|
|
"certified": False,
|
|
"audit": {},
|
|
"router": router,
|
|
"agent_binding": binding,
|
|
"loop_run": loop_run,
|
|
}, ensure_ascii=False))
|
|
|
|
|
|
def bind_agent(args, router):
|
|
agent = args.agent or default_agent_for_mode(router.get("mode", "READ_ONLY"))
|
|
tools = []
|
|
if router.get("mode") == "OPERATOR":
|
|
# The operator primitive resolves the exact registered action later; the
|
|
# agent bind still proves this turn is allowed to use the operator class.
|
|
tools.append("run-chat-tests")
|
|
if router.get("mode") == "CODEGEN":
|
|
tools.extend(["artifact-scan", "tool-output-scan"])
|
|
cmd = [
|
|
"python3", AGENT_RESOLVER, "bind",
|
|
"--agent", agent,
|
|
"--actor", args.actor,
|
|
"--role", args.role,
|
|
"--project", args.project,
|
|
"--tenant", args.tenant,
|
|
"--delegation-level", str(args.delegation_level),
|
|
]
|
|
if args.skill:
|
|
cmd += ["--skill", args.skill]
|
|
for tool in tools:
|
|
cmd += ["--tool", tool]
|
|
r = subprocess.run(cmd, cwd=ROOT, capture_output=True, text=True)
|
|
try:
|
|
return r.returncode, json.loads(r.stdout)
|
|
except Exception:
|
|
if r.returncode != 0 and (r.stdout or r.stderr):
|
|
sys.stdout.write((r.stdout or r.stderr).strip() + "\n")
|
|
return r.returncode, None
|
|
print(json.dumps({"success": False, "decision": "DENIED", "reason": "agent_resolver_invalid_json"}, ensure_ascii=False))
|
|
return 2, None
|
|
|
|
|
|
def bind_model(args, binding):
|
|
r = subprocess.run([
|
|
"python3", MODEL_RESOLVER, "bind",
|
|
"--provider", args.model_provider,
|
|
"--actor", args.actor,
|
|
"--role", args.role,
|
|
"--model-role", binding.get("model_role") or "read_only",
|
|
], cwd=ROOT, capture_output=True, text=True)
|
|
try:
|
|
return r.returncode, json.loads(r.stdout)
|
|
except Exception:
|
|
return 2, {"success": False, "decision": "DENIED", "reason": "model_resolver_invalid_json"}
|
|
|
|
|
|
def submit_escalation(args, router, binding):
|
|
turn_id = args.turn_id or str(uuid.uuid4())
|
|
payload = {
|
|
"chat_id": args.chat_id or "chat-default",
|
|
"turn_id": turn_id,
|
|
"tenant_id": args.tenant,
|
|
"actor": args.actor,
|
|
"role": args.role,
|
|
"message_ref": sha(args.message),
|
|
"message_preview": args.message[:240],
|
|
"router": router,
|
|
"agent_binding": binding,
|
|
}
|
|
reason = f"chat turn requires approval: {binding.get('reason', 'requires_approval')}"
|
|
r = subprocess.run([
|
|
"python3", APPROVAL_INBOX, "submit",
|
|
"--project", args.project,
|
|
"--action", "chat.escalate",
|
|
"--target", f"chat:{args.chat_id or 'chat-default'}:{turn_id}",
|
|
"--risk", router.get("risk", "high"),
|
|
"--sensitive",
|
|
"--proposer", args.actor,
|
|
"--reason", reason,
|
|
"--payload", json.dumps(payload, ensure_ascii=False, sort_keys=True),
|
|
], cwd=ROOT, capture_output=True, text=True)
|
|
try:
|
|
proposal = json.loads(r.stdout)
|
|
except Exception:
|
|
proposal = {"status": "failed", "error": (r.stdout + r.stderr).strip()}
|
|
rec = append_chat_turn({
|
|
"timestamp": proposal.get("created_at") or "",
|
|
"trace_id": "chat-escalation-" + sha(f"{args.chat_id}|{turn_id}|{args.message}")[:12],
|
|
"chat_id": args.chat_id or "chat-default",
|
|
"turn_id": turn_id,
|
|
"tenant_id": args.tenant,
|
|
"actor": args.actor,
|
|
"mode": router.get("mode", "READ_ONLY"),
|
|
"risk": router.get("risk", "high"),
|
|
"decision": "ESCALATED" if r.returncode == 0 else "DENIED",
|
|
"answer": "Chat turn escalated to approval inbox." if r.returncode == 0 else "Chat escalation failed closed.",
|
|
"sources": [],
|
|
"router": router,
|
|
"agent_binding": binding,
|
|
"approval": {"proposal_id": proposal.get("id"), "status": proposal.get("status"), "target": proposal.get("target")},
|
|
})
|
|
print(json.dumps({
|
|
"success": False,
|
|
"mode": router.get("mode", "READ_ONLY"),
|
|
"risk": router.get("risk", "high"),
|
|
"decision": "ESCALATED" if r.returncode == 0 else "DENIED",
|
|
"answer": "Chat turn escalated to approval inbox." if r.returncode == 0 else "Chat escalation failed closed.",
|
|
"sources": [],
|
|
"certified": False,
|
|
"audit": {"seq": rec["seq"], "record_hash": rec["record_hash"], "head": rec["record_hash"], "path": audit_path()},
|
|
"router": router,
|
|
"agent_binding": binding,
|
|
"approval": {"proposal": proposal, "submit_rc": r.returncode, "output": (r.stdout + r.stderr).strip()[:400]},
|
|
}, ensure_ascii=False))
|
|
return 3 if r.returncode == 0 else 2
|
|
|
|
|
|
def ask(args) -> int:
|
|
guard_rc = tenant_runtime_guard(args)
|
|
if guard_rc != 0:
|
|
return guard_rc
|
|
router = classify(args.message)
|
|
bind_rc, binding = bind_agent(args, router)
|
|
if bind_rc != 0:
|
|
if binding and binding.get("decision") == "REQUIRES_APPROVAL":
|
|
return submit_escalation(args, router, binding)
|
|
if binding:
|
|
print(json.dumps(binding, ensure_ascii=False))
|
|
return bind_rc
|
|
model_rc, model_binding = bind_model(args, binding)
|
|
if model_rc != 0:
|
|
print(json.dumps({"success": False, "mode": router.get("mode"), "decision": "DENIED", "reason": model_binding.get("reason"), "agent_binding": binding, "model_binding": model_binding}, ensure_ascii=False))
|
|
return model_rc
|
|
binding["model_binding"] = model_binding
|
|
os.environ["CASAN_CHAT_MODEL_PROVIDER"] = model_binding.get("provider", "")
|
|
common = [
|
|
"--message", args.message,
|
|
"--actor", args.actor,
|
|
"--chat-id", args.chat_id,
|
|
"--turn-id", args.turn_id,
|
|
"--tenant", args.tenant,
|
|
]
|
|
if router.get("mode") == "OPERATOR":
|
|
loop_run, loop_rc = certify_operator_draft(args, router, binding)
|
|
if loop_rc != 0:
|
|
denied_from_loop(router, binding, loop_run)
|
|
return loop_rc
|
|
return run_mode(["python3", OPERATOR, "run", *common], binding, loop_run, "chat-operator")
|
|
if router.get("mode") == "CODEGEN":
|
|
loop_run, loop_rc = certify_codegen_draft(args, router, binding)
|
|
return finish_codegen(args, router, binding, loop_run, loop_rc)
|
|
if getattr(args, "stream", False) and router.get("mode") in ("READ_ONLY", "ANALYSIS"):
|
|
env = os.environ.copy()
|
|
env["CASAN_CHAT_MODEL_PROVIDER"] = model_binding.get("provider", "")
|
|
return run_and_passthrough(["python3", READONLY, "ask", "--stream", *common], env)
|
|
return run_mode(["python3", READONLY, "ask", *common], binding)
|
|
|
|
|
|
def verify_audit() -> int:
|
|
return run_and_passthrough(["python3", READONLY, "verify-audit"])
|
|
|
|
|
|
def history(args) -> int:
|
|
command = ["python3", READONLY, "history", "--actor", args.actor, "--tenant", args.tenant, "--limit", str(args.limit)]
|
|
if args.chat_id:
|
|
command += ["--chat-id", args.chat_id]
|
|
return run_and_passthrough(command)
|
|
|
|
|
|
def main() -> int:
|
|
ap = argparse.ArgumentParser()
|
|
sub = ap.add_subparsers(dest="cmd", required=True)
|
|
askp = sub.add_parser("ask")
|
|
askp.add_argument("--message", required=True)
|
|
askp.add_argument("--actor", default="anonymous")
|
|
askp.add_argument("--role", default="viewer")
|
|
askp.add_argument("--project", default="default")
|
|
askp.add_argument("--chat-id", default="")
|
|
askp.add_argument("--turn-id", default="")
|
|
askp.add_argument("--tenant", default="default")
|
|
askp.add_argument("--agent", default="")
|
|
askp.add_argument("--skill", default="")
|
|
askp.add_argument("--model-provider", default="")
|
|
askp.add_argument("--delegation-level", type=int, default=0)
|
|
askp.add_argument("--stream", action="store_true")
|
|
askp.set_defaults(func=ask)
|
|
sub.add_parser("verify-audit").set_defaults(func=lambda _args: verify_audit())
|
|
hp = sub.add_parser("history")
|
|
hp.add_argument("--actor", required=True)
|
|
hp.add_argument("--chat-id", default="")
|
|
hp.add_argument("--tenant", default="default")
|
|
hp.add_argument("--limit", type=int, default=50)
|
|
hp.set_defaults(func=history)
|
|
args = ap.parse_args()
|
|
return args.func(args)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|