telosieve 0.2.0-rc.4

Read-only infrastructure instruction evaluation that refuses when trusted evidence cannot agree
Documentation
#!/usr/bin/env python3
"""Qualify Telosieve against orchestrated real HTTP endpoints."""
from __future__ import annotations
import http.client, json, os, pathlib, resource, shutil, ssl, stat, subprocess, tempfile, threading, time
from concurrent.futures import ThreadPoolExecutor
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer

ROOT = pathlib.Path(__file__).resolve().parents[1]
ADAPTER = (ROOT / "scripts/http-json-integration-adapter.py").resolve()
PRODUCER = (ROOT / "scripts/http-json-observation-producer.py").resolve()
PKI_CHECK = (ROOT / "scripts/http-json-pki-check.py").resolve()
PKI_MONITOR = (ROOT / "scripts/http-json-pki-monitor.py").resolve()
PATH = "/v1/telosieve/snapshot"
LOAD_EVALUATIONS, LOAD_CONCURRENCY = 8, 4
TOKENS = {"adapter":"A"*32,"producer-a":"B"*32,"producer-b":"C"*32}
BASE_SNAPSHOT = {"schema_version":"telosieve.http-json-snapshot/v1","revision":"revision-81","complete":True,
                 "desired":{"cluster/epoch":"7","user/message":"new"},
                 "replicas":{name:{"cluster/epoch":"7","user/message":"old"} for name in ("replica-a","replica-b","replica-c")}}

class State:
    def __init__(self): self.lock=threading.Lock(); self.modes={token:"valid" for token in TOKENS.values()}; self.gets=0; self.mutations=0

class QualificationHTTPServer(ThreadingHTTPServer):
    def handle_error(self, request, client_address): pass

class Handler(BaseHTTPRequestHandler):
    protocol_version = "HTTP/1.1"
    def log_message(self, *_): pass
    def _mutation(self):
        with self.server.state.lock: self.server.state.mutations += 1
        self.send_response(405); self.send_header("Content-Length","0"); self.end_headers()
    do_POST=_mutation; do_PUT=_mutation; do_PATCH=_mutation; do_DELETE=_mutation
    def do_GET(self):
        token=self.headers.get("Authorization","").removeprefix("Bearer ")
        with self.server.state.lock:
            mode=self.server.state.modes.get(token,"unauthorized"); self.server.state.gets += 1
        if self.path != PATH: mode="not-found"
        if mode=="timeout": time.sleep(2)
        if mode=="redirect":
            self.send_response(302); self.send_header("Location",PATH); self.send_header("Content-Length","0"); self.end_headers(); return
        if mode in {"unauthorized","not-found"}:
            self.send_response(401 if mode=="unauthorized" else 404); self.send_header("Content-Length","0"); self.end_headers(); return
        if mode=="malformed": body=b"{"
        elif mode=="oversized":
            self.send_response(200); self.send_header("Content-Type","application/json"); self.send_header("Content-Length",str(2*1024*1024+1)); self.end_headers(); return
        else:
            value=json.loads(json.dumps(BASE_SNAPSHOT))
            if mode=="incomplete": value["complete"]=False
            if mode=="disagreement": value["revision"]="revision-82"
            body=json.dumps(value,separators=(",",":")).encode()
        try:
            self.send_response(200); self.send_header("Content-Type","application/json"); self.send_header("Content-Length",str(len(body)))
            self.send_header("Cache-Control","no-store"); self.end_headers(); self.wfile.write(body)
        except (BrokenPipeError, ConnectionResetError): pass

def run(args, *, success=True, timeout=10):
    result=subprocess.run(args,cwd=ROOT,capture_output=True,check=False,timeout=timeout)
    if (result.returncode==0)!=success: raise RuntimeError(f"http-json-e2e: unexpected exit {result.returncode}: {result.stderr.decode(errors='replace')[-4096:]}")
    return result

def credential(path, token):
    path.write_text(json.dumps({"schema_version":"telosieve.http-json-credentials/v1","bearer_token":token},separators=(",",":"))); path.chmod(0o600)

def producer_material(work,binary,port,credentials,transport="http"):
    keys=[]; sources=[]; configs=[]
    for index,name in enumerate(("producer-a","producer-b")):
        seed=("41" if index==0 else "42")*32; key=work/f"{name}.key"; key.write_text(seed); key.chmod(0o600)
        public=run([str(binary),"observation-public-key",str(key)]).stdout.decode().strip()
        domain=f"http-reader-{name[-1]}"
        keys.append({"producer":name,"key_id":f"key-{name[-1]}","fault_domain":domain,"public_key":public,
                     "not_before":1750000000,"not_after":1750001000})
        config=work/f"{name}.json"; config.write_text(json.dumps({"host":"127.0.0.1","port":port,"path":PATH,"transport":transport,
            "credentials":str(credentials[index]),"integration_id":"http-json","resource_kind":"replicated-key-value",
            "target_id":"http/qualification","subject":"kv/research","evaluation_time":1750000000,
            "telosieve":str(binary),"key":str(key),"producer":name,"key_id":f"key-{name[-1]}","domain":domain,
            "issued":1750000000,"expires":1750000300},separators=(",",":")))
        configs.append(config); sources.append({"executable_path":str(PRODUCER),"arguments":["--config",str(config)]})
    trust=work/"trust.json"; trust.write_text(json.dumps({"schema_version":"telosieve.observation-trust/v1",
        "evaluation_time":1750000100,"required_distinct_domains":2,"keys":keys},separators=(",",":")))
    return trust,sources,configs

def make_config(work,binary,port,credentials,stem,transport="http"):
    trust,sources,producer_configs=producer_material(work,binary,port,credentials[1:],transport)
    value={"schema_version":"telosieve.evaluation-config/v7","mode":"external-read-only",
        "scenario_path":str((ROOT/"scenarios/benign.json").resolve()),"certificate_path":str(work/f"{stem}-certificate.json"),
        "ledger_path":str(work/f"{stem}-ledger.jsonl"),"observation_trust_path":str(trust),"observation_sources":sources,
        "adapter":{"executable_path":str(ADAPTER),"arguments":["--host","127.0.0.1","--port",str(port),"--path",PATH,"--transport",transport,
          "--credentials",str(credentials[0])],"integration_id":"http-json","resource_kind":"replicated-key-value",
          "target_id":"http/qualification"}}
    path=work/f"{stem}-evaluation.json"; path.write_text(json.dumps(value,separators=(",",":")))
    return path,value,producer_configs

def certificates(work):
    openssl="openssl"
    def command(*args): run([openssl,*args],timeout=10)
    ca_key=work/"ca.key"; ca=work/"ca.pem"; rogue_key=work/"rogue-ca.key"; rogue_ca=work/"rogue-ca.pem"
    command("req","-x509","-newkey","rsa:2048","-nodes","-subj","/CN=Telosieve Qualification CA","-days","1","-keyout",str(ca_key),"-out",str(ca))
    command("req","-x509","-newkey","rsa:2048","-nodes","-subj","/CN=Rogue Qualification CA","-days","1","-keyout",str(rogue_key),"-out",str(rogue_ca))
    (work/"newcerts").mkdir(); (work/"index.txt").write_text(""); (work/"serial").write_text("1000\n"); (work/"crlnumber").write_text("1000\n")
    config=work/"ca.cnf"; config.write_text(f"""[ca]\ndefault_ca=main\n[main]\ndir={work}\ndatabase=$dir/index.txt\nnew_certs_dir=$dir/newcerts\ncertificate=$dir/ca.pem\nprivate_key=$dir/ca.key\nserial=$dir/serial\ncrlnumber=$dir/crlnumber\ndefault_md=sha256\ndefault_days=1\ndefault_crl_days=1\npolicy=policy\n[policy]\ncommonName=supplied\n[server]\nbasicConstraints=CA:FALSE\nkeyUsage=digitalSignature,keyEncipherment\nextendedKeyUsage=serverAuth\nsubjectAltName=IP:127.0.0.1\n[client]\nbasicConstraints=CA:FALSE\nkeyUsage=digitalSignature,keyEncipherment\nextendedKeyUsage=clientAuth\n""")
    def signed(name,usage):
        key=work/f"{name}.key"; csr=work/f"{name}.csr"; certificate=work/f"{name}.pem"
        command("req","-new","-newkey","rsa:2048","-nodes","-subj",f"/CN={name}","-keyout",str(key),"-out",str(csr))
        command("ca","-batch","-config",str(config),"-extensions",usage,"-in",str(csr),"-out",str(certificate))
        key.chmod(0o600); return certificate,key
    server_pair=signed("server","server")
    clients={name:signed(f"mtls-client-{name}","client") for name in TOKENS}
    clients["rotated-adapter"]=signed("mtls-client-rotated-adapter","client")
    rogue_keypair=work/"mtls-client-rogue.key",work/"mtls-client-rogue.csr",work/"mtls-client-rogue.pem"
    command("req","-new","-newkey","rsa:2048","-nodes","-subj","/CN=rogue","-keyout",str(rogue_keypair[0]),"-out",str(rogue_keypair[1]))
    command("x509","-req","-in",str(rogue_keypair[1]),"-CA",str(rogue_ca),"-CAkey",str(rogue_key),"-CAcreateserial","-days","1","-out",str(rogue_keypair[2])); rogue_keypair[0].chmod(0o600)
    clients["rogue"]=(rogue_keypair[2],rogue_keypair[0])
    crl=work/"ca.crl.pem"; command("ca","-gencrl","-config",str(config),"-out",str(crl))
    ca_key.chmod(0o600); rogue_key.chmod(0o600)
    return ca,crl,rogue_ca,server_pair,clients,config

def revoke(certificate,config,crl):
    run(["openssl","ca","-batch","-config",str(config),"-revoke",str(certificate)],timeout=10)
    run(["openssl","ca","-gencrl","-config",str(config),"-out",str(crl)],timeout=10)

def tls_credential(path,token,ca,crl,certificate,key):
    path.write_text(json.dumps({"schema_version":"telosieve.http-json-mtls-credentials/v2","bearer_token":token,
        "ca_certificate":str(ca),"certificate_revocation_list":str(crl),"client_certificate":str(certificate),"client_key":str(key)},separators=(",",":"))); path.chmod(0o600)

def main():
    binary=(ROOT/"target/debug/telosieve").resolve()
    if not binary.is_file(): raise SystemExit("http-json-e2e: build target/debug/telosieve first")
    started=time.monotonic(); state=State(); server=QualificationHTTPServer(("127.0.0.1",0),Handler); server.state=state
    server.daemon_threads=True; thread=threading.Thread(target=server.serve_forever); thread.start()
    failures=0
    try:
      with tempfile.TemporaryDirectory(prefix="telosieve-http-json-e2e-") as raw:
        work=pathlib.Path(raw); credentials=[]
        for name in ("adapter","producer-a","producer-b"):
            path=work/f"{name}-credential.json"; credential(path,TOKENS[name]); credentials.append(path)
        config_path,base,producer_configs=make_config(work,binary,server.server_port,credentials,"valid")
        report=json.loads(run([str(binary),"evaluate",str(config_path)]).stdout)
        certificate=json.loads((work/"valid-certificate.json").read_bytes())
        if report["target_mutated"] or certificate["integration"]["target_revision"]!="revision-81": raise RuntimeError("http-json-e2e: invalid success")
        def evaluate_case(index):
            case=work/f"load-{index}"; case.mkdir(); value=json.loads(json.dumps(base)); value["certificate_path"]=str(case/"certificate.json"); value["ledger_path"]=str(case/"ledger.jsonl")
            path=case/"evaluation.json"; path.write_text(json.dumps(value)); run([str(binary),"evaluate",str(path)])
        load_started=time.monotonic()
        with ThreadPoolExecutor(max_workers=LOAD_CONCURRENCY) as pool:
            for future in [pool.submit(evaluate_case,i) for i in range(LOAD_EVALUATIONS)]: future.result(timeout=20)
        load_elapsed=time.monotonic()-load_started
        denied=0
        for method in ("POST","PUT","PATCH","DELETE"):
            connection=http.client.HTTPConnection("127.0.0.1",server.server_port,timeout=1)
            connection.request(method,PATH,body=b"{}",headers={"Authorization":f"Bearer {TOKENS['adapter']}"})
            if connection.getresponse().status!=405: raise RuntimeError("http-json-e2e: mutation method was not denied")
            connection.close(); denied+=1
        def fault(stem,token=None,mode=None,config=None):
            nonlocal failures
            if token:
                with state.lock: state.modes[token]=mode
            try:
                value=json.loads(json.dumps(config or base)); value["certificate_path"]=str(work/f"{stem}-certificate.json"); value["ledger_path"]=str(work/f"{stem}-ledger.jsonl")
                path=work/f"{stem}-evaluation.json"; path.write_text(json.dumps(value)); run([str(binary),"evaluate",str(path)],success=False)
                if (work/f"{stem}-certificate.json").exists() or (work/f"{stem}-ledger.jsonl").exists(): raise RuntimeError(f"http-json-e2e: {stem} emitted evidence")
            finally:
                if token:
                    with state.lock: state.modes[token]="valid"
            failures+=1
        for stem,mode in (("redirect","redirect"),("malformed","malformed"),("oversized","oversized"),("incomplete","incomplete"),("timeout","timeout")):
            fault(stem,TOKENS["adapter"],mode)
        fault("unauthorized",TOKENS["adapter"],"unauthorized")
        fault("producer-disagreement",TOKENS["producer-a"],"disagreement")
        server.shutdown(); server.server_close(); thread.join(timeout=3)
        fault("outage")
        ca,crl,rogue_ca,(server_cert,server_key),clients,ca_config=certificates(work)
        state=State(); server=QualificationHTTPServer(("127.0.0.1",0),Handler); server.state=state; server.daemon_threads=True
        server_context=ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER); server_context.minimum_version=ssl.TLSVersion.TLSv1_3
        server_context.maximum_version=ssl.TLSVersion.TLSv1_3; server_context.verify_mode=ssl.CERT_REQUIRED
        server_context.load_cert_chain(server_cert,server_key); server_context.load_verify_locations(cafile=ca); server_context.load_verify_locations(cafile=crl)
        server_context.verify_flags |= ssl.VERIFY_CRL_CHECK_LEAF
        server.socket=server_context.wrap_socket(server.socket,server_side=True)
        thread=threading.Thread(target=server.serve_forever); thread.start()
        tls_credentials=[]
        for name in ("adapter","producer-a","producer-b"):
            path=work/f"mtls-{name}-credential.json"; tls_credential(path,TOKENS[name],ca,crl,*clients[name]); tls_credentials.append(path)
        openssl=pathlib.Path(shutil.which("openssl") or "").resolve()
        for path in tls_credentials: run([str(PKI_CHECK),"--openssl",str(openssl),"--credentials",str(path),"--renew-before-seconds","3600"])
        pki_refusals=0
        run([str(PKI_CHECK),"--openssl",str(openssl),"--credentials",str(tls_credentials[0]),"--renew-before-seconds","172800"],success=False); pki_refusals+=1
        pki_original=tls_credentials[0].read_text(); pki_bad=json.loads(pki_original); pki_bad["client_key"]=str(clients["producer-a"][1]); tls_credentials[0].write_text(json.dumps(pki_bad))
        try: run([str(PKI_CHECK),"--openssl",str(openssl),"--credentials",str(tls_credentials[0]),"--renew-before-seconds","3600"],success=False); pki_refusals+=1
        finally: tls_credentials[0].write_text(pki_original)
        monitor_config=work/"pki-monitor.json"; monitor_status=work/"pki-monitor-status.json"
        monitor_config.write_text(json.dumps({"schema_version":"telosieve.http-json-pki-monitor/v1","checker":str(PKI_CHECK),
            "openssl":str(openssl),"credentials":[str(path) for path in tls_credentials],"renew_before_seconds":3600},separators=(",",":")))
        monitor_status.write_text(json.dumps({"checked_at":1,"status":"stale"},separators=(",",":"))); monitor_status.chmod(0o600)
        monitored=json.loads(run([str(PKI_MONITOR),"--config",str(monitor_config),"--status",str(monitor_status)]).stdout)
        if monitored["status"]!="ready" or monitored["checked_at"]<=1 or stat.S_IMODE(monitor_status.stat().st_mode)!=0o600: raise RuntimeError("http-json-e2e: PKI monitor ready publication is invalid")
        tls_path,tls_base,tls_producer_configs=make_config(work,binary,server.server_port,tls_credentials,"mtls-valid","https")
        tls_report=json.loads(run([str(binary),"evaluate",str(tls_path)]).stdout)
        if tls_report["target_mutated"]: raise RuntimeError("http-json-e2e: mTLS success reported mutation")
        client_context=ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT); client_context.minimum_version=ssl.TLSVersion.TLSv1_3
        client_context.maximum_version=ssl.TLSVersion.TLSv1_3; client_context.load_verify_locations(cafile=ca); client_context.load_verify_locations(cafile=crl)
        client_context.verify_flags |= ssl.VERIFY_CRL_CHECK_LEAF
        client_context.load_cert_chain(*clients["adapter"])
        tls_denied=0
        for method in ("POST","PUT","PATCH","DELETE"):
            connection=http.client.HTTPSConnection("127.0.0.1",server.server_port,timeout=1,context=client_context)
            connection.request(method,PATH,body=b"{}",headers={"Authorization":f"Bearer {TOKENS['adapter']}"})
            if connection.getresponse().status!=405: raise RuntimeError("http-json-e2e: mTLS mutation method was not denied")
            connection.close(); tls_denied+=1
        original=tls_credentials[0].read_text()
        bad=json.loads(original); bad["ca_certificate"]=str(rogue_ca); tls_credentials[0].write_text(json.dumps(bad))
        try: fault("mtls-wrong-ca",config=tls_base)
        finally: tls_credentials[0].write_text(original)
        bad=json.loads(original); bad["client_key"]=str(work/"missing.key"); tls_credentials[0].write_text(json.dumps(bad))
        try: fault("mtls-missing-key",config=tls_base)
        finally: tls_credentials[0].write_text(original)
        clients["adapter"][1].chmod(0o644)
        try: fault("mtls-unsafe-key-permissions",config=tls_base)
        finally: clients["adapter"][1].chmod(0o600)
        bad=json.loads(original); bad["client_certificate"]=str(clients["rogue"][0]); bad["client_key"]=str(clients["rogue"][1]); tls_credentials[0].write_text(json.dumps(bad))
        try: fault("mtls-untrusted-client",config=tls_base)
        finally: tls_credentials[0].write_text(original)
        downgrade=json.loads(json.dumps(tls_base)); arguments=downgrade["adapter"]["arguments"]
        arguments[arguments.index("https")]="http"; fault("mtls-plaintext-downgrade",config=downgrade)
        fault("mtls-producer-disagreement",TOKENS["producer-a"],"disagreement",tls_base)
        fault("mtls-timeout",TOKENS["adapter"],"timeout",tls_base)
        bad=json.loads(original); bad["certificate_revocation_list"]=str(work/"missing.crl"); tls_credentials[0].write_text(json.dumps(bad))
        try: fault("mtls-missing-crl",config=tls_base)
        finally: tls_credentials[0].write_text(original)
        malformed_crl=work/"malformed.crl"; malformed_crl.write_text("not-a-crl\n"); bad=json.loads(original); bad["certificate_revocation_list"]=str(malformed_crl); tls_credentials[0].write_text(json.dumps(bad))
        try: fault("mtls-malformed-crl",config=tls_base)
        finally: tls_credentials[0].write_text(original)
        rotated=json.loads(original); rotated["client_certificate"]=str(clients["rotated-adapter"][0]); rotated["client_key"]=str(clients["rotated-adapter"][1]); tls_credentials[0].write_text(json.dumps(rotated))
        rotation=json.loads(json.dumps(tls_base)); rotation["certificate_path"]=str(work/"mtls-rotation-certificate.json"); rotation["ledger_path"]=str(work/"mtls-rotation-ledger.jsonl")
        rotation_path=work/"mtls-rotation-evaluation.json"; rotation_path.write_text(json.dumps(rotation)); run([str(binary),"evaluate",str(rotation_path)])
        tls_credentials[0].write_text(original); revoke(server_cert,ca_config,crl); fault("mtls-revoked-server",config=tls_base)
        revoke(clients["adapter"][0],ca_config,crl)
        run([str(PKI_CHECK),"--openssl",str(openssl),"--credentials",str(tls_credentials[0]),"--renew-before-seconds","3600"],success=False); pki_refusals+=1
        run([str(PKI_MONITOR),"--config",str(monitor_config),"--status",str(monitor_status)],success=False)
        monitored_bytes=monitor_status.read_bytes(); monitored=json.loads(monitored_bytes)
        if monitored["status"]!="not-ready" or monitored["not_ready"]!=1 or any(str(path).encode() in monitored_bytes for path in tls_credentials):
            raise RuntimeError("http-json-e2e: PKI monitor refusal publication is invalid")
        server.shutdown(); server.server_close(); thread.join(timeout=3); fault("mtls-outage",config=tls_base)
        elapsed=time.monotonic()-started; peak=resource.getrusage(resource.RUSAGE_CHILDREN).ru_maxrss
        if os.uname().sysname=="Linux": peak*=1024
        if elapsed>60 or load_elapsed>20 or thread.is_alive(): raise RuntimeError("http-json-e2e: resource or cleanup bound exceeded")
        print(json.dumps({"schema_version":"telosieve.http-json-integration-qualification/v5","transports":["http/1.1","https-tls1.3-mtls-crl"],
          "orchestrated_endpoints":True,"external_endpoints":False,"loopback_only":True,"bearer_identities":3,
          "mtls_client_identities":3,"mutation_methods_refused":denied+tls_denied,"observation_producers":2,"separate_control_planes":False,
          "successful_evaluations":3,"client_rotations":1,"revocations_refused":1,"pki_preflights":3,"pki_preflight_refusals":pki_refusals,"pki_monitor_runs":2,"pki_monitor_transitions":1,"pki_monitor_stale_replacements":1,"load_evaluations":LOAD_EVALUATIONS,"load_concurrency":LOAD_CONCURRENCY,
          "load_elapsed_seconds":round(load_elapsed,3),"fail_closed_evaluations":failures,"target_mutated":False,
          "elapsed_seconds":round(elapsed,3),"peak_child_rss_bytes":peak,"independent_evidence":False,"status":"passed"},separators=(",",":")))
    finally:
      try: server.shutdown(); server.server_close()
      except Exception: pass
      thread.join(timeout=3)

if __name__=="__main__": main()