From 9abc1428b2f3f2c3811c66253f45d6889f573dcb Mon Sep 17 00:00:00 2001 From: pcvantol Date: Mon, 7 Sep 2026 18:39:09 +0200 Subject: [PATCH 01/14] feat: preserve action repair lineage and runtime recovery --- forge/models/__init__.py | 2 ++ forge/models/execution_host.py | 4 +++ forge/programme_authorization.py | 28 +++++++++------ forge/runtime/__init__.py | 3 +- forge/runtime/runner.py | 13 +++++-- forge/runtime/service.py | 50 ++++++++++++++++++++++++++ tests/test_bootstrap_mission_runner.py | 25 +++++++++++++ tests/test_programme_authorization.py | 16 +++++++-- 8 files changed, 125 insertions(+), 16 deletions(-) create mode 100644 forge/runtime/service.py diff --git a/forge/models/__init__.py b/forge/models/__init__.py index 6da6fc2..6cc5473 100644 --- a/forge/models/__init__.py +++ b/forge/models/__init__.py @@ -264,6 +264,7 @@ ExecutionEvidenceOutcome, ExecutionHostContract, ExecutionHostEvidence, + ExecutionHostTemporaryUnavailable, ExecutionHostForbiddenResponsibility, ExecutionHostLifecycleStage, ExecutionHostResponsibility, @@ -522,6 +523,7 @@ "ReasoningProfile", "ExecutionHostContract", "ExecutionHostEvidence", + "ExecutionHostTemporaryUnavailable", "ExecutionHostForbiddenResponsibility", "ExecutionHostLifecycleStage", "ExecutionHostResponsibility", diff --git a/forge/models/execution_host.py b/forge/models/execution_host.py index 146c42e..ad1d91f 100644 --- a/forge/models/execution_host.py +++ b/forge/models/execution_host.py @@ -67,6 +67,10 @@ class ExecutionEvidenceOutcome(str, Enum): FAILED = "failed" +class ExecutionHostTemporaryUnavailable(RuntimeError): + """A transport/read outage; the persisted request remains the recovery key.""" + + @dataclass(frozen=True) class ExecutionHostContract: """A complete host declaration without a host implementation.""" diff --git a/forge/programme_authorization.py b/forge/programme_authorization.py index ce92670..3b630f6 100644 --- a/forge/programme_authorization.py +++ b/forge/programme_authorization.py @@ -69,6 +69,8 @@ class CandidateQualification: owner_workflow_evidence: str owner_workflow_head_sha: str merge_method: str = "squash" + mission_id: str | None = None + action_id: str | None = None def __post_init__(self) -> None: if (self.pull_request <= 0 or not _SHA.fullmatch(self.head_sha) @@ -122,20 +124,23 @@ def qualify(self, authorization_id: str, candidate: CandidateQualification) -> s def record_repair_attempt(self, authorization_id: str, candidate: CandidateQualification) -> str: authorization = self._authorization(authorization_id) self._validate_candidate(authorization, candidate, require_passes=False) - existing = self._repair_attempts(authorization_id, candidate.pull_request, candidate.head_sha) + if not candidate.mission_id or not candidate.action_id: + raise PermissionError("repair authorization requires the Forge Mission and Engineering Action lineage") + existing = self._repair_attempts(authorization_id, candidate.mission_id, candidate.action_id) if existing >= authorization.repair_attempt_limit: - raise PermissionError("bounded repair budget exhausted for this exact PR head") + raise PermissionError("bounded repair budget exhausted for this Engineering Action lineage") attempt = existing + 1 return self.repository.record(GovernanceDecision( - decision_id=f"{authorization_id}:repair:{candidate.pull_request}:{candidate.head_sha}:{attempt}", - subject_id=f"{authorization.programme_id}:repair:pr-{candidate.pull_request}:attempt-{attempt}", + decision_id=f"{authorization_id}:repair:{candidate.mission_id}:{candidate.action_id}:{attempt}", + subject_id=f"{authorization.programme_id}:repair:{candidate.mission_id}:{candidate.action_id}:attempt-{attempt}", subject_revision=candidate.head_sha, capability=GovernanceCapability.OWNER_PROGRAMME_AUTHORIZATION, decision="repair-authorized", scope=candidate.changed_scopes, gates=("same-approved-scope", "exact-head"), predecessor_digest=self._decision_digest(authorization_id), - evidence={"kind": "BOUNDED_REPAIR_ATTEMPT_V1", "authorization_id": authorization_id, + evidence={"kind": "BOUNDED_ACTION_REPAIR_ATTEMPT_V2", "authorization_id": authorization_id, + "mission_id": candidate.mission_id, "action_id": candidate.action_id, "pull_request": candidate.pull_request, "head_sha": candidate.head_sha, "attempt": attempt}, ), self.context) @@ -171,7 +176,7 @@ def _decision_digest(self, decision_id: str) -> str: raise PermissionError("authorization evidence is absent") return row["digest"] - def _repair_attempts(self, authorization_id: str, pull_request: int, head_sha: str) -> int: + def _repair_attempts(self, authorization_id: str, mission_id: str, action_id: str) -> int: rows = self.repository.database._connection.execute( "SELECT document FROM governance_decisions WHERE capability = ?", (GovernanceCapability.OWNER_PROGRAMME_AUTHORIZATION.value,), @@ -179,8 +184,11 @@ def _repair_attempts(self, authorization_id: str, pull_request: int, head_sha: s import json return sum( 1 for row in rows - if (lambda evidence: evidence.get("kind") == "BOUNDED_REPAIR_ATTEMPT_V1" - and evidence.get("authorization_id") == authorization_id - and evidence.get("pull_request") == pull_request - and evidence.get("head_sha") == head_sha)(json.loads(row["document"]).get("evidence", {})) + if (lambda evidence: evidence.get("authorization_id") == authorization_id and ( + (evidence.get("kind") == "BOUNDED_ACTION_REPAIR_ATTEMPT_V2" + and evidence.get("mission_id") == mission_id and evidence.get("action_id") == action_id) + # V1 had no Action identity. Count it conservatively rather + # than silently resetting a pre-existing programme budget. + or evidence.get("kind") == "BOUNDED_REPAIR_ATTEMPT_V1" + ))(json.loads(row["document"]).get("evidence", {})) ) diff --git a/forge/runtime/__init__.py b/forge/runtime/__init__.py index 943bf4b..bc4782e 100644 --- a/forge/runtime/__init__.py +++ b/forge/runtime/__init__.py @@ -25,9 +25,10 @@ ) from .evidence import RuntimeDecisionEvidenceReference, RuntimeEvidence from .runner import BootstrapMissionRunner, MissionRunnerError, RuntimePromptFactory +from .service import ForgeRuntimeService, RuntimeServiceTick __all__ = [ - "BootstrapMissionRunner", "MissionRunnerError", "RuntimePromptFactory", + "BootstrapMissionRunner", "MissionRunnerError", "RuntimePromptFactory", "ForgeRuntimeService", "RuntimeServiceTick", "RUNTIME_SCHEMA_VERSION", "RuntimeDatabase", "RuntimeDatabaseError", "RuntimeIntegrityError", "RuntimeDecisionEvidenceReference", "RuntimeEvidence", "RUNTIME_INSTANCE_VERSION", "RUNTIME_INITIALIZATION_VERSION", "RuntimeBootstrap", "RuntimeIdentity", "RuntimeInstance", "RuntimeLocation", "RuntimeRecovery", "RuntimeResolutionError", "RuntimeResolver", "repository_identity", "repository_uuid", ] diff --git a/forge/runtime/runner.py b/forge/runtime/runner.py index 44bba1d..0a3fe66 100644 --- a/forge/runtime/runner.py +++ b/forge/runtime/runner.py @@ -14,6 +14,7 @@ ExecutionHost, ExecutionHostEvidence, ExecutionRequest, + ExecutionHostTemporaryUnavailable, ) from forge.models.runtime_prompt import ( ProviderPromptDefinition, @@ -224,7 +225,8 @@ def run(self, mission_id: str) -> MissionExecutionState: return state before_revision = state.revision state = self._advance(state) - if state.status is MissionExecutionStatus.WAITING_FOR_EVIDENCE and state.revision == before_revision: + if state.status in {MissionExecutionStatus.WAITING_FOR_EXECUTION, + MissionExecutionStatus.WAITING_FOR_EVIDENCE} and state.revision == before_revision: return state def _advance(self, state: MissionExecutionState) -> MissionExecutionState: @@ -268,7 +270,11 @@ def _dispatch_or_recover(self, state: MissionExecutionState) -> MissionExecution dispatch = self._host.dispatch(request) if dispatch.request != request: raise MissionRunnerError("execution host acknowledgement did not preserve the persisted request") - except Exception as error: # Host errors must become durable terminal state. + except ExecutionHostTemporaryUnavailable: + # The request was persisted before dispatch. A later tick asks for + # its original acknowledgement before attempting another send. + return state + except Exception as error: # Invalid acknowledgements fail closed. return self._host_failure(state, "host_dispatch_failed", error) envelope = {"request": _request_document(request), "host_run_id": dispatch.host_run_id} return self._store.transition( @@ -284,6 +290,9 @@ def _collect_evidence(self, state: MissionExecutionState) -> MissionExecutionSta if evidence is None: return state actions = self._scheduler.reconcile(self._actions(state), dispatch, evidence) + except ExecutionHostTemporaryUnavailable: + # A temporary read outage is not terminal evidence. + return state except Exception as error: # Invalid evidence and host failures fail closed. return self._host_failure(state, "host_evidence_failed", error) evidence_document = _document(evidence) diff --git a/forge/runtime/service.py b/forge/runtime/service.py new file mode 100644 index 0000000..19d38ca --- /dev/null +++ b/forge/runtime/service.py @@ -0,0 +1,50 @@ +"""Minimal supervised Forge Runtime Service using the qualified loop unchanged.""" + +from __future__ import annotations + +from dataclasses import dataclass +from typing import Callable, Protocol + +from forge.state import MissionExecutionStatus, MissionStateStore + + +class RuntimeLoop(Protocol): + """The qualified loop surface reused by the service.""" + + def run(self): ... + + def resume(self, mission_id: str): ... + + +@dataclass(frozen=True) +class RuntimeServiceTick: + """Read-time result; Mission State remains the authority.""" + + resumed_mission_ids: tuple[str, ...] + dispatched_mission_id: str | None + + +class ForgeRuntimeService: + """Supervise the existing ExecutionLoop without becoming an Execution Host. + + A tick resumes one persisted Mission before asking the Dispatcher for new + work. This keeps the lane single-flight and preserves the CLI's semantics. + """ + + def __init__(self, loop: RuntimeLoop, states: MissionStateStore) -> None: + self._loop, self._states = loop, states + + def tick(self) -> RuntimeServiceTick: + resumable = {MissionExecutionStatus.READY, MissionExecutionStatus.ACTIVE, + MissionExecutionStatus.WAITING_FOR_EXECUTION, MissionExecutionStatus.WAITING_FOR_EVIDENCE} + for state in self._states.resumable(): + if state.status in resumable: + self._loop.resume(state.mission_id) + return RuntimeServiceTick((state.mission_id,), None) + result = self._loop.run() + return RuntimeServiceTick((), None if result is None else result.mission_id) + + def serve(self, *, keep_running: Callable[[], bool]) -> None: + """Caller-owned service loop: no hidden daemon or second state store.""" + while keep_running(): + self.tick() diff --git a/tests/test_bootstrap_mission_runner.py b/tests/test_bootstrap_mission_runner.py index 81287f9..063fec2 100644 --- a/tests/test_bootstrap_mission_runner.py +++ b/tests/test_bootstrap_mission_runner.py @@ -9,6 +9,7 @@ from forge.models import ( EngineeringAction, EngineeringIntent, EngineeringActionStatus, ExecutionDispatch, ExecutionRequest, ExecutionEvidenceOutcome, ExecutionHostEvidence, ExecutionRepositoryEvidence, + ExecutionHostTemporaryUnavailable, IntentApproval, IntentCategory, IntentReference, IntentStatus, IntentTraceability, ProviderPromptDefinition, RuntimePrompt, RuntimePromptSection, RuntimePromptSectionKind, ) @@ -80,6 +81,14 @@ def retrieve_evidence(self, dispatch: ExecutionDispatch) -> ExecutionHostEvidenc receipt_id=f"receipt-{request.action_id}", execution_duration_ms=60_000) +class TemporarilyUnavailableHost(Host): + def dispatch(self, request: object) -> ExecutionDispatch: + raise ExecutionHostTemporaryUnavailable("host is temporarily unavailable") + + def retrieve_evidence(self, dispatch: ExecutionDispatch) -> ExecutionHostEvidence | None: + raise ExecutionHostTemporaryUnavailable("evidence endpoint is temporarily unavailable") + + class BootstrapMissionRunnerTests(unittest.TestCase): def setUp(self) -> None: self.directory = TemporaryDirectory() @@ -131,6 +140,22 @@ def test_execution_host_failure_is_persisted_as_failed(self) -> None: self.assertEqual(state.status, MissionExecutionStatus.FAILED) self.assertEqual(state.execution_evidence["diagnostic_references"], ["runner:host_dispatch_failed"]) # type: ignore[index] + def test_temporary_dispatch_unavailability_keeps_the_persisted_request_recoverable(self) -> None: + runner = self.runner(TemporarilyUnavailableHost()) + runner.start(mission("one"), (intent("one"),), (action(1, "one"),)) + state = runner.run("mission-1") + self.assertEqual(state.status, MissionExecutionStatus.WAITING_FOR_EXECUTION) + self.assertIsNotNone(state.execution_correlation) + + def test_temporary_evidence_unavailability_keeps_the_run_recoverable(self) -> None: + host = Host() + runner = self.runner(host) + runner.start(mission("one"), (intent("one"),), (action(1, "one"),)) + waiting = runner.run("mission-1") + self.assertEqual(waiting.status, MissionExecutionStatus.WAITING_FOR_EVIDENCE) + runner._host = TemporarilyUnavailableHost() # noqa: SLF001 - restart injects a host transport + self.assertEqual(runner.resume("mission-1").status, MissionExecutionStatus.WAITING_FOR_EVIDENCE) + def test_resume_after_restart_recovers_persisted_dispatch_without_regeneration(self) -> None: host = Host() first = self.runner(host) diff --git a/tests/test_programme_authorization.py b/tests/test_programme_authorization.py index 870a3dd..9e65789 100644 --- a/tests/test_programme_authorization.py +++ b/tests/test_programme_authorization.py @@ -31,7 +31,7 @@ def candidate(self, **changes): head_sha="a" * 40, base_branch="main", changed_scopes=("governance",), technical_qualification_passed=True, ci_passed=True, reviews_passed=True, security_passed=True, owner_workflow_evidence="github-run:123", - owner_workflow_head_sha="a" * 40) + owner_workflow_head_sha="a" * 40, mission_id="mission-1", action_id="action-1") value.update(changes) if "head_sha" in changes and "owner_workflow_head_sha" not in changes: value["owner_workflow_head_sha"] = changes["head_sha"] @@ -49,9 +49,19 @@ def test_fails_closed_for_scope_expansion_or_missing_real_gates(self): with self.assertRaises(PermissionError): self.gate.qualify(self.authorization.authorization_id, self.candidate(changed_scopes=("outside",))) with self.assertRaises(PermissionError): self.gate.qualify(self.authorization.authorization_id, self.candidate(security_passed=False)) - def test_repair_is_limited_per_exact_pr_head_and_new_head_requires_new_qualification(self): + def test_repair_is_limited_across_the_action_lineage_even_when_the_head_changes(self): self.gate.record_authorization(self.authorization) candidate = self.candidate(technical_qualification_passed=False, ci_passed=False) for _ in range(3): self.gate.record_repair_attempt(self.authorization.authorization_id, candidate) with self.assertRaises(PermissionError): self.gate.record_repair_attempt(self.authorization.authorization_id, candidate) - self.assertEqual(self.gate.record_repair_attempt(self.authorization.authorization_id, self.candidate(head_sha="b" * 40, technical_qualification_passed=False, ci_passed=False)), self.gate._decision_digest(f"{self.authorization.authorization_id}:repair:77:{'b' * 40}:1")) + with self.assertRaises(PermissionError): + self.gate.record_repair_attempt(self.authorization.authorization_id, self.candidate(head_sha="b" * 40, technical_qualification_passed=False, ci_passed=False)) + + def test_repair_permission_requires_action_lineage_and_is_not_merge_qualification(self): + self.gate.record_authorization(self.authorization) + candidate = self.candidate(technical_qualification_passed=False, ci_passed=False) + self.assertTrue(self.gate.record_repair_attempt(self.authorization.authorization_id, candidate).startswith("sha256:")) + with self.assertRaises(PermissionError): + self.gate.qualify(self.authorization.authorization_id, candidate) + with self.assertRaises(PermissionError): + self.gate.record_repair_attempt(self.authorization.authorization_id, self.candidate(mission_id=None)) From 61f72e5180ce2111e15bce41a69ae444ba7622fc Mon Sep 17 00:00:00 2001 From: pcvantol Date: Mon, 7 Sep 2026 20:08:17 +0200 Subject: [PATCH 02/14] fix: upgrade verified legacy programme capability --- forge/operator_identity.py | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/forge/operator_identity.py b/forge/operator_identity.py index 300eace..9380050 100644 --- a/forge/operator_identity.py +++ b/forge/operator_identity.py @@ -74,6 +74,22 @@ def adopt_governance_capabilities(self, context): row=self.db._connection.execute('SELECT created_at,version,status FROM installation_operator_binding WHERE installation_id=?',(context.installation_id,)).fetchone() if not row or row['status']!='ACTIVE': raise PermissionError('active G001 binding required') self._persist_governance_capabilities(context,'EXISTING_G001_GOVERNANCE_ADOPTION_V1',{'prior_binding_created_at':row['created_at'],'prior_binding_version':row['version'],'adopted_at':_timestamp()},row) + def upgrade_legacy_governance_capabilities(self, context, *, decision_source): + """Add only the v2 programme capability to a verified v1 three-capability installation.""" + if not self.authorize(context) or not decision_source: raise PermissionError('trusted operator and upgrade decision source required') + legacy=('ARCHITECTURE_APPROVAL','BUSINESS_APPROVAL','SECURITY_APPROVAL'); operator=self._governance_operator_id(context) + with self.db._connection: + binding=self.db._connection.execute('SELECT * FROM installation_operator_binding WHERE installation_id=?',(context.installation_id,)).fetchone() + if not binding or binding['status']!='ACTIVE': raise PermissionError('active binding required') + rows=self.db._connection.execute('SELECT * FROM governance_capability_grants WHERE installation_id=? AND operator_id=? ORDER BY capability',(context.installation_id,operator)).fetchall() + authorities=tuple(r['capability'] for r in self.db._connection.execute('SELECT capability FROM governance_authority WHERE installation_id=? AND operator_id=? ORDER BY capability',(context.installation_id,operator))) + if authorities==legacy+('OWNER_PROGRAMME_AUTHORIZATION',): return + if authorities!=legacy or tuple(r['capability'] for r in rows)!=legacy: raise PermissionError('state is not the recognized legacy capability set') + for row in rows: + document=json.loads(row['bootstrap_provenance']); digest='sha256:'+hashlib.sha256(json.dumps(document,sort_keys=True,separators=(',',':')).encode()).hexdigest() + if digest!=row['digest'] or document.get('capability')!=row['capability'] or document.get('installation_id')!=context.installation_id or document.get('operator_id')!=operator: raise PermissionError('legacy capability provenance is invalid') + now=_timestamp(); provenance={'kind':'G001_OWNER_PROGRAMME_CAPABILITY_UPGRADE_V1','upgrade_version':'1','decision_source':decision_source,'legacy_grants':[(r['grant_id'],r['digest']) for r in rows],'binding_version':binding['version'],'added_capability':'OWNER_PROGRAMME_AUTHORIZATION','occurred_at':now,'installation_id':context.installation_id,'operator_id':operator,'capability':'OWNER_PROGRAMME_AUTHORIZATION'}; encoded=json.dumps(provenance,sort_keys=True,separators=(',',':')); digest='sha256:'+hashlib.sha256(encoded.encode()).hexdigest() + self.db._insert_governance_grant(digest,context.installation_id,operator,'OWNER_PROGRAMME_AUTHORIZATION',encoded,digest,now); self.db._insert_governance_authority(context.installation_id,operator,'OWNER_PROGRAMME_AUTHORIZATION',now); self._audit(context.installation_id,NamedOperatorIdentity(context.generated_uid,0),'LEGACY_CAPABILITY_UPGRADE',now,'ALLOW') def revoke(self, context): if not self.authorize(context):raise PermissionError('denied') with self.db._connection: From 2ea92bc2748c8f732f93ade566632d5f308fc4c7 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Mon, 7 Sep 2026 20:14:29 +0200 Subject: [PATCH 03/14] fix: validate upgraded capability provenance --- forge/operator_identity.py | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/forge/operator_identity.py b/forge/operator_identity.py index 9380050..742a5f3 100644 --- a/forge/operator_identity.py +++ b/forge/operator_identity.py @@ -83,7 +83,14 @@ def upgrade_legacy_governance_capabilities(self, context, *, decision_source): if not binding or binding['status']!='ACTIVE': raise PermissionError('active binding required') rows=self.db._connection.execute('SELECT * FROM governance_capability_grants WHERE installation_id=? AND operator_id=? ORDER BY capability',(context.installation_id,operator)).fetchall() authorities=tuple(r['capability'] for r in self.db._connection.execute('SELECT capability FROM governance_authority WHERE installation_id=? AND operator_id=? ORDER BY capability',(context.installation_id,operator))) - if authorities==legacy+('OWNER_PROGRAMME_AUTHORIZATION',): return + expected=tuple(sorted((*legacy,'OWNER_PROGRAMME_AUTHORIZATION'))) + if authorities==expected: + if len(rows)!=4 or tuple(r['capability'] for r in rows)!=expected: raise PermissionError('upgraded capability grants are incomplete') + upgraded=next(r for r in rows if r['capability']=='OWNER_PROGRAMME_AUTHORIZATION') + document=json.loads(upgraded['bootstrap_provenance']); digest='sha256:'+hashlib.sha256(json.dumps(document,sort_keys=True,separators=(',',':')).encode()).hexdigest() + old={r['grant_id']:r['digest'] for r in rows if r['capability']!='OWNER_PROGRAMME_AUTHORIZATION'} + if digest!=upgraded['digest'] or document.get('kind')!='G001_OWNER_PROGRAMME_CAPABILITY_UPGRADE_V1' or dict(document.get('legacy_grants',()))!=old or document.get('binding_version')!=binding['version']: raise PermissionError('upgraded capability provenance is invalid') + return if authorities!=legacy or tuple(r['capability'] for r in rows)!=legacy: raise PermissionError('state is not the recognized legacy capability set') for row in rows: document=json.loads(row['bootstrap_provenance']); digest='sha256:'+hashlib.sha256(json.dumps(document,sort_keys=True,separators=(',',':')).encode()).hexdigest() From 3e50fe6d62a9ff30b0ab169c18855590ab62fc35 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Mon, 7 Sep 2026 20:21:36 +0200 Subject: [PATCH 04/14] feat: map EP v1.1 terminal evidence --- forge/scheduler/ep_v11.py | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) create mode 100644 forge/scheduler/ep_v11.py diff --git a/forge/scheduler/ep_v11.py b/forge/scheduler/ep_v11.py new file mode 100644 index 0000000..39ba778 --- /dev/null +++ b/forge/scheduler/ep_v11.py @@ -0,0 +1,19 @@ +"""Strict Forge consumer mapping for the EP producer readback v1.1 fixture.""" +from __future__ import annotations +import hashlib,json +from typing import Any,Mapping +from forge.models.execution_host import ExecutionEvidenceOutcome,ExecutionHostEvidence,ExecutionRepositoryEvidence + +def terminal_evidence(readback: Mapping[str,Any], artifact: bytes, *, host_id: str) -> ExecutionHostEvidence: + if readback.get('contract_version')!='1.1': raise ValueError('unsupported EP readback contract') + evidence=readback.get('evidence',{}); terminal=evidence.get('terminal_artifact'); run=readback.get('run'); result=readback.get('result',{}); provenance=readback.get('provenance',{}).get('forge_execution') + if not isinstance(terminal,dict) or not isinstance(run,dict) or not isinstance(provenance,dict) or result.get('terminal') is not True: raise ValueError('EP terminal evidence is incomplete') + digest='sha256:'+hashlib.sha256(artifact).hexdigest() + if terminal.get('digest')!=digest: raise ValueError('EP terminal artifact digest mismatch') + try: document=json.loads(artifact) + except (UnicodeDecodeError,json.JSONDecodeError) as e: raise ValueError('EP terminal artifact is invalid JSON') from e + repo=document.get('repository',{}); outcome=ExecutionEvidenceOutcome(str(result.get('outcome')).lower()) + if not isinstance(repo.get('revision'),str) or not isinstance(document.get('report',{}).get('id'),str): raise ValueError('EP terminal artifact lacks repository evidence') + correlation=readback['correlation']; prompt=provenance['runtime_prompt']; submission=readback['submission'] + repository=ExecutionRepositoryEvidence(str(correlation['mission_id']),str(provenance['intent_id']),str(provenance['intent_revision']),str(correlation['engineering_action_id']),str(prompt['id']),str(correlation['correlation_id']),str(run['id']),str(repo['id']),str(repo['revision']),str(document['report']['id']),digest) + return ExecutionHostEvidence(host_id,str(correlation['correlation_id']),str(run['id']),str(document['report']['id']),outcome,repository,validation_references=tuple(str(x.get('command')) for x in document.get('references',{}).get('validation',()) if x.get('command'))) From 933272aaf004a4812735de5d8d9c5c0273eed89e Mon Sep 17 00:00:00 2001 From: pcvantol Date: Mon, 7 Sep 2026 20:53:59 +0200 Subject: [PATCH 05/14] feat: add EP v1.1 HTTP execution host adapter --- forge/scheduler/ep_http_adapter.py | 51 ++++++++++++++++++++++++++++++ 1 file changed, 51 insertions(+) create mode 100644 forge/scheduler/ep_http_adapter.py diff --git a/forge/scheduler/ep_http_adapter.py b/forge/scheduler/ep_http_adapter.py new file mode 100644 index 0000000..594a859 --- /dev/null +++ b/forge/scheduler/ep_http_adapter.py @@ -0,0 +1,51 @@ +"""Concrete, strict HTTP v1.1 Engineering Platform Execution Host adapter.""" +from __future__ import annotations +import json +from dataclasses import dataclass +from urllib.error import HTTPError,URLError +from urllib.request import Request,urlopen +from forge.models.execution_host import ExecutionDispatch,ExecutionHostTemporaryUnavailable,ExecutionRequest +from .ep_v11 import terminal_evidence + +@dataclass(frozen=True) +class EngineeringPlatformHttpConfiguration: + base_url:str; project_id:str; bearer_token:str; host_id:str='engineering-platform'; timeout:float=10 + +class EngineeringPlatformHttpExecutionHost: + def __init__(self, config:EngineeringPlatformHttpConfiguration): self.config=config; self._submissions={} + def _json(self,path,*,method='GET',body=None): + data=None if body is None else json.dumps(body,sort_keys=True,separators=(',',':')).encode() + request=Request(self.config.base_url.rstrip('/')+path,data=data,method=method,headers={'Authorization':'Bearer '+self.config.bearer_token,'Content-Type':'application/json'}) + try: + with urlopen(request,timeout=self.config.timeout) as response:return json.loads(response.read()) + except HTTPError as error: + if error.code>=500: raise ExecutionHostTemporaryUnavailable('EP temporarily unavailable') from error + raise ValueError(f'EP rejected request: {error.code}') from error + except (URLError,TimeoutError) as error: raise ExecutionHostTemporaryUnavailable('EP transport unavailable') from error + def _payload(self,r): + p=r.runtime_prompt; text=getattr(p,'rendered_text',None) or p.to_markdown() + return {'repository_id':r.repository_id,'producer':{'id':'forge','type':'FORGE','version':'1.0'},'prompt':text,'idempotency_key':r.correlation_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'engineering_action_id':r.action_id,'constraints':{'forge_execution':{'contract_version':'1.0','host_id':r.host_id,'repository_id':r.repository_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'mission_revision':'1','intent_id':r.intent_id,'intent_revision':r.intent_revision,'action_id':r.action_id,'runtime_prompt':{'id':r.runtime_prompt.id,'content_digest':getattr(r.runtime_prompt,'source_digest',getattr(r.runtime_prompt,'generation_request_digest',None))},'retry_of_correlation_id':r.retry_of_correlation_id}}} + def dispatch(self,r): + accepted=self._json(f'/v1/projects/{self.config.project_id}/submissions',method='POST',body=self._payload(r)); sid=accepted.get('submission_id') + if not isinstance(sid,str): raise ValueError('EP submission acknowledgement lacks submission_id') + self._submissions[r.correlation_id]=sid + return self.recover_dispatch(r) or (_ for _ in ()).throw(ExecutionHostTemporaryUnavailable('submission accepted but run unclaimed')) + def recover_dispatch(self,r): + sid=self._submissions.get(r.correlation_id) + if not sid:return None + readback=self._json(f'/v1/projects/{self.config.project_id}/submissions/{sid}') + if readback.get('correlation',{}).get('correlation_id')!=r.correlation_id or readback.get('submission',{}).get('repository_id')!=r.repository_id: raise ValueError('EP readback does not bind persisted request') + run=readback.get('run') + return None if not isinstance(run,dict) else ExecutionDispatch(r,str(run['id'])) + def retrieve_evidence(self,d): + sid=self._submissions.get(d.request.correlation_id) + if not sid: raise ValueError('missing persisted submission identity') + readback=self._json(f'/v1/projects/{self.config.project_id}/submissions/{sid}') + if readback.get('run',{}).get('id')!=d.host_run_id:return None + terminal=readback.get('evidence',{}).get('terminal_artifact') + if not isinstance(terminal,dict): return None + artifact=self._json(f"/v1/projects/{self.config.project_id}/artifacts/{terminal['id']}") + raw=json.dumps(artifact,sort_keys=True,separators=(',',':')).encode()+b'\n' + evidence=terminal_evidence(readback,raw,host_id=self.config.host_id) + if (evidence.correlation_id,evidence.host_run_id,evidence.repository_evidence.runtime_prompt_id)!=(d.request.correlation_id,d.host_run_id,d.request.runtime_prompt.id): raise ValueError('EP artifact does not bind dispatch') + return evidence From 06a6e0cbb8c985cb747f822980fe8e9d9cc2fc56 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Mon, 7 Sep 2026 21:38:03 +0200 Subject: [PATCH 06/14] fix: persist producer contract with execution request --- forge/runtime/runner.py | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/forge/runtime/runner.py b/forge/runtime/runner.py index 0a3fe66..a0abbbe 100644 --- a/forge/runtime/runner.py +++ b/forge/runtime/runner.py @@ -25,6 +25,7 @@ from forge.models.codex_runtime_prompt import ( CodexCliRuntimePrompt, ExecutionHostCompatibility, RepositoryState, ) +from forge.models.producer import Producer, ProducerContract, ProducerIdentity, RuntimePromptEnvelope from forge.models.intent import IntentReference from forge.scheduler import BootstrapMissionScheduler from forge.state import MissionExecutionState, MissionExecutionStatus, MissionStateStore @@ -138,6 +139,21 @@ def _prompt(document: Mapping[str, Any]) -> RuntimePrompt | CodexCliRuntimePromp def _request(document: Mapping[str, Any]) -> ExecutionRequest: + contract_document = document.get("producer_contract") + contract = None + if isinstance(contract_document, Mapping): + producer = contract_document["producer"] + prompt = contract_document["runtime_prompt"] + metadata = contract_document.get("execution_metadata", {}) + if not isinstance(producer, Mapping) or not isinstance(prompt, Mapping) or not isinstance(metadata, Mapping): + raise MissionRunnerError("persisted Producer Contract is malformed") + contract = ProducerContract( + Producer(ProducerIdentity(str(producer["identity"]["id"]), str(producer["identity"]["type"]), str(producer["identity"]["version"]))), + str(contract_document["correlation_id"]), str(contract_document["engineering_action_id"]), + RuntimePromptEnvelope(str(prompt["id"]), str(prompt["version"]), str(prompt["format"]), str(prompt["content"]), str(prompt["content_digest"])), + tuple(str(item) for item in contract_document["execution_constraints"]), tuple((str(k), str(v)) for k,v in metadata.items()), + mission_id=contract_document.get("mission_id"), + ) return ExecutionRequest( host_id=str(document["host_id"]), mission_id=str(document["mission_id"]), intent_id=str(document["intent_id"]), intent_revision=str(document["intent_revision"]), @@ -146,6 +162,7 @@ def _request(document: Mapping[str, Any]) -> ExecutionRequest: correlation_id=str(document["correlation_id"]), dispatched_at=str(document["dispatched_at"]), retry_of_correlation_id=document.get("retry_of_correlation_id"), original_correlation_id=document.get("original_correlation_id"), + producer_contract=contract, ) @@ -164,6 +181,7 @@ def _request_document(request: ExecutionRequest) -> dict[str, Any]: "dispatched_at": request.dispatched_at, "retry_of_correlation_id": request.retry_of_correlation_id, "original_correlation_id": request.original_correlation_id, + "producer_contract": request.producer_contract.to_dict(), } From ed225394e00c4ade6b26e544f36071c1b345af37 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Mon, 7 Sep 2026 22:14:55 +0200 Subject: [PATCH 07/14] fix: verify EP artifact bytes before parsing --- forge/scheduler/ep_http_adapter.py | 19 ++++++++++++++++--- 1 file changed, 16 insertions(+), 3 deletions(-) diff --git a/forge/scheduler/ep_http_adapter.py b/forge/scheduler/ep_http_adapter.py index 594a859..d67e3ea 100644 --- a/forge/scheduler/ep_http_adapter.py +++ b/forge/scheduler/ep_http_adapter.py @@ -22,9 +22,23 @@ def _json(self,path,*,method='GET',body=None): if error.code>=500: raise ExecutionHostTemporaryUnavailable('EP temporarily unavailable') from error raise ValueError(f'EP rejected request: {error.code}') from error except (URLError,TimeoutError) as error: raise ExecutionHostTemporaryUnavailable('EP transport unavailable') from error + def _bytes(self,path): + request=Request(self.config.base_url.rstrip('/')+path,headers={'Authorization':'Bearer '+self.config.bearer_token}) + try: + with urlopen(request,timeout=self.config.timeout) as response: + raw=response.read(1_048_577) + if len(raw)>1_048_576: raise ValueError('EP artifact exceeds consumer size limit') + return raw + except HTTPError as error: + if error.code>=500: raise ExecutionHostTemporaryUnavailable('EP artifact temporarily unavailable') from error + raise ValueError(f'EP artifact request rejected: {error.code}') from error + except (URLError,TimeoutError) as error: raise ExecutionHostTemporaryUnavailable('EP artifact transport unavailable') from error def _payload(self,r): p=r.runtime_prompt; text=getattr(p,'rendered_text',None) or p.to_markdown() - return {'repository_id':r.repository_id,'producer':{'id':'forge','type':'FORGE','version':'1.0'},'prompt':text,'idempotency_key':r.correlation_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'engineering_action_id':r.action_id,'constraints':{'forge_execution':{'contract_version':'1.0','host_id':r.host_id,'repository_id':r.repository_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'mission_revision':'1','intent_id':r.intent_id,'intent_revision':r.intent_revision,'action_id':r.action_id,'runtime_prompt':{'id':r.runtime_prompt.id,'content_digest':getattr(r.runtime_prompt,'source_digest',getattr(r.runtime_prompt,'generation_request_digest',None))},'retry_of_correlation_id':r.retry_of_correlation_id}}} + revision=getattr(p,'mission_revision',None) + if not isinstance(revision,str) or not revision: raise ValueError('Execution Request has no persisted Mission revision') + contract=r.producer_contract + return {'repository_id':r.repository_id,'producer':contract.producer.identity.to_dict(),'prompt':text,'idempotency_key':r.correlation_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'engineering_action_id':r.action_id,'constraints':{'forge_execution':{'contract_version':'1.0','host_id':r.host_id,'repository_id':r.repository_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'mission_revision':revision,'intent_id':r.intent_id,'intent_revision':r.intent_revision,'action_id':r.action_id,'runtime_prompt':{'id':contract.runtime_prompt.id,'content_digest':contract.runtime_prompt.content_digest},'retry_of_correlation_id':r.retry_of_correlation_id}}} def dispatch(self,r): accepted=self._json(f'/v1/projects/{self.config.project_id}/submissions',method='POST',body=self._payload(r)); sid=accepted.get('submission_id') if not isinstance(sid,str): raise ValueError('EP submission acknowledgement lacks submission_id') @@ -44,8 +58,7 @@ def retrieve_evidence(self,d): if readback.get('run',{}).get('id')!=d.host_run_id:return None terminal=readback.get('evidence',{}).get('terminal_artifact') if not isinstance(terminal,dict): return None - artifact=self._json(f"/v1/projects/{self.config.project_id}/artifacts/{terminal['id']}") - raw=json.dumps(artifact,sort_keys=True,separators=(',',':')).encode()+b'\n' + raw=self._bytes(f"/v1/projects/{self.config.project_id}/artifacts/{terminal['id']}") evidence=terminal_evidence(readback,raw,host_id=self.config.host_id) if (evidence.correlation_id,evidence.host_run_id,evidence.repository_evidence.runtime_prompt_id)!=(d.request.correlation_id,d.host_run_id,d.request.runtime_prompt.id): raise ValueError('EP artifact does not bind dispatch') return evidence From 5f2c92cea8c7eaa973993a5f964da4d17a0f990c Mon Sep 17 00:00:00 2001 From: pcvantol Date: Tue, 8 Sep 2026 07:39:25 +0200 Subject: [PATCH 08/14] feat: make EP HTTP recovery runtime-safe --- forge/models/execution_host.py | 10 +- forge/runtime/database.py | 31 ++- forge/runtime/runner.py | 49 +++- forge/runtime/service.py | 104 ++++++-- forge/scheduler/ep_http_adapter.py | 248 ++++++++++++++---- forge/scheduler/ep_v11.py | 18 +- .../test_action_derivation_canary_closure.py | 2 +- tests/test_ep_http_adapter.py | 121 +++++++++ tests/test_runtime_service.py | 54 ++++ 9 files changed, 538 insertions(+), 99 deletions(-) create mode 100644 tests/test_ep_http_adapter.py create mode 100644 tests/test_runtime_service.py diff --git a/forge/models/execution_host.py b/forge/models/execution_host.py index ad1d91f..1cc2b8d 100644 --- a/forge/models/execution_host.py +++ b/forge/models/execution_host.py @@ -205,19 +205,21 @@ class ExecutionRepositoryEvidence: correlation_id: str host_run_id: str repository_id: str - repository_revision: str + repository_revision: str | None report_id: str content_digest: str def __post_init__(self) -> None: if not all((self.mission_id, self.intent_id, self.intent_revision, self.action_id, self.runtime_prompt_id, self.correlation_id, self.host_run_id, - self.repository_id, self.repository_revision, self.report_id, + self.repository_id, self.report_id, self.content_digest)): raise ValueError("repository evidence identity, provenance, revision, report, and digest are required") digest = self.content_digest.removeprefix("sha256:") if not self.content_digest.startswith("sha256:") or len(digest) != 64: raise ValueError("repository evidence digest must be sha256") + if self.repository_revision is not None and not self.repository_revision: + raise ValueError("repository evidence revision cannot be empty") @dataclass(frozen=True) @@ -249,6 +251,8 @@ def __post_init__(self) -> None: self.correlation_id, self.host_run_id, self.report_id, ): raise ValueError("execution host evidence must match its repository evidence run and report") + if self.outcome is ExecutionEvidenceOutcome.COMPLETE and not repository.repository_revision: + raise ValueError("complete execution evidence requires a delivery revision") for references, label in ((self.log_references, "log"), (self.diagnostic_references, "diagnostic"), (self.metric_references, "metric"), (self.validation_references, "validation")): if any(not reference for reference in references) or len(references) != len(set(references)): raise ValueError(f"execution host evidence {label} references must be unique and non-empty") @@ -272,7 +276,7 @@ class ExecutionHost(Protocol): a request before dispatching it without treating process memory as state. """ - def dispatch(self, request: ExecutionRequest) -> ExecutionDispatch: ... + def dispatch(self, request: ExecutionRequest) -> ExecutionDispatch | None: ... def recover_dispatch(self, request: ExecutionRequest) -> ExecutionDispatch | None: ... diff --git a/forge/runtime/database.py b/forge/runtime/database.py index 088d674..ac41e39 100644 --- a/forge/runtime/database.py +++ b/forge/runtime/database.py @@ -22,7 +22,7 @@ canonical_repository_root, repository_identity, repository_uuid) -RUNTIME_SCHEMA_VERSION = 31 +RUNTIME_SCHEMA_VERSION = 32 _REQUIRED_METADATA = frozenset(( "schema_version", "migration_version", "forge_version", "created_at", "last_migration", "integrity_status", @@ -38,7 +38,7 @@ "planning_provider_security_config", "planning_provider_security_audit", "planning_provider_generation_permits", "token_preflight_receipts", "token_preflight_receipt_consumptions", "token_preflight_failures", "action_derivations", "action_derivation_reattempt_authorizations", "action_derivation_reattempt_consumptions", "governance_authority", "governance_capability_grants", "governance_decisions", "action_derivation_evidence_sets", - "mission_amendments", "action_derivation_canary_closures", + "mission_amendments", "action_derivation_canary_closures", "execution_host_bindings", )) _TOKEN_PREFLIGHT_FAILURE_FIELDS = frozenset(( "failure_id", "mission_id", "provider_id", "occurred_at", "main_head", "policy_digest", @@ -151,7 +151,8 @@ def __init__(self, workspace_root: Path | str = ".", *, path: Path | str | None self._action_derivation_canary_closure_write_state = {"permitted": False} try: self._configure() - self._migrate(forge_version) + while self._connection.execute("PRAGMA user_version").fetchone()[0] != RUNTIME_SCHEMA_VERSION: + self._migrate(forge_version) self._initialize_runtime_identity() self.validate_integrity() except Exception: @@ -438,6 +439,7 @@ def _migrate(self, forge_version: str) -> None: UNIQUE (mission_id, action_id, iteration), FOREIGN KEY (mission_id) REFERENCES mission_state(mission_id) ); + CREATE TABLE IF NOT EXISTS execution_host_bindings (correlation_id TEXT PRIMARY KEY, document TEXT NOT NULL); CREATE TABLE installation_operator_binding (installation_id TEXT PRIMARY KEY, generated_uid TEXT NOT NULL, uid INTEGER NOT NULL, version INTEGER NOT NULL, status TEXT NOT NULL, created_at TEXT NOT NULL); CREATE TABLE installation_operator_audit (audit_id TEXT PRIMARY KEY, installation_id TEXT NOT NULL, generated_uid TEXT NOT NULL, operation TEXT NOT NULL, occurred_at TEXT NOT NULL, result TEXT NOT NULL); CREATE TRIGGER installation_operator_audit_immutable_update BEFORE UPDATE ON installation_operator_audit @@ -1232,6 +1234,17 @@ def _migrate(self, forge_version: str) -> None: except Exception: self._connection.rollback() raise + elif version == 31: + with self._connection: + # Some supported historic test/runtime states were created by + # the current bootstrap schema and then assigned an older + # pragma version to exercise their migration. The binding is + # additive, so its presence is safe; still advance the + # version atomically instead of treating that state as a + # second authority or an incomplete migration. + self._connection.execute("CREATE TABLE IF NOT EXISTS execution_host_bindings (correlation_id TEXT PRIMARY KEY, document TEXT NOT NULL)") + self._set_metadata({"schema_version":"32","migration_version":"32","last_migration":"32"}) + self._connection.execute("PRAGMA user_version=32") elif version != RUNTIME_SCHEMA_VERSION: raise RuntimeIntegrityError("runtime database migration path is unavailable") @@ -1518,6 +1531,18 @@ def scheduler_submission(self, submission_id: str) -> dict[str, Any] | None: ).fetchone() return json.loads(row["document"]) if row is not None else None + def execution_host_binding(self, correlation_id: str) -> dict[str, Any] | None: + row=self._connection.execute("SELECT document FROM execution_host_bindings WHERE correlation_id=?",(correlation_id,)).fetchone() + return None if row is None else json.loads(row["document"]) + + def save_execution_host_binding(self, correlation_id: str, document: Mapping[str, Any]) -> dict[str, Any]: + if not correlation_id or document.get("correlation_id") != correlation_id: raise RuntimeDatabaseError("execution host binding correlation is invalid") + existing=self.execution_host_binding(correlation_id) + if existing is not None and any(existing.get(key) not in (None,value) for key,value in document.items()): raise RuntimeIntegrityError("execution host binding is immutable") + merged={**(existing or {}),**document} + with self._connection:self._connection.execute("INSERT INTO execution_host_bindings VALUES (?,?) ON CONFLICT(correlation_id) DO UPDATE SET document=excluded.document",(correlation_id,self._dump(merged))) + return merged + def outstanding_scheduler_submission(self, mission_id: str) -> dict[str, Any] | None: rows = self._connection.execute( "SELECT document FROM scheduler_submissions WHERE mission_id = ? AND state IN ('CREATED', 'SUBMITTED', 'ACCEPTED', 'EXECUTING', 'RECEIPT_AVAILABLE') ORDER BY iteration", diff --git a/forge/runtime/runner.py b/forge/runtime/runner.py index a0abbbe..17ab213 100644 --- a/forge/runtime/runner.py +++ b/forge/runtime/runner.py @@ -25,7 +25,10 @@ from forge.models.codex_runtime_prompt import ( CodexCliRuntimePrompt, ExecutionHostCompatibility, RepositoryState, ) -from forge.models.producer import Producer, ProducerContract, ProducerIdentity, RuntimePromptEnvelope +from forge.models.producer import ( + ExecutionReceiptReference, Producer, ProducerContract, ProducerIdentity, + RuntimePromptEnvelope, +) from forge.models.intent import IntentReference from forge.scheduler import BootstrapMissionScheduler from forge.state import MissionExecutionState, MissionExecutionStatus, MissionStateStore @@ -141,19 +144,34 @@ def _prompt(document: Mapping[str, Any]) -> RuntimePrompt | CodexCliRuntimePromp def _request(document: Mapping[str, Any]) -> ExecutionRequest: contract_document = document.get("producer_contract") contract = None - if isinstance(contract_document, Mapping): - producer = contract_document["producer"] - prompt = contract_document["runtime_prompt"] - metadata = contract_document.get("execution_metadata", {}) - if not isinstance(producer, Mapping) or not isinstance(prompt, Mapping) or not isinstance(metadata, Mapping): + if contract_document is not None: + if not isinstance(contract_document, Mapping): raise MissionRunnerError("persisted Producer Contract is malformed") - contract = ProducerContract( - Producer(ProducerIdentity(str(producer["identity"]["id"]), str(producer["identity"]["type"]), str(producer["identity"]["version"]))), - str(contract_document["correlation_id"]), str(contract_document["engineering_action_id"]), - RuntimePromptEnvelope(str(prompt["id"]), str(prompt["version"]), str(prompt["format"]), str(prompt["content"]), str(prompt["content_digest"])), - tuple(str(item) for item in contract_document["execution_constraints"]), tuple((str(k), str(v)) for k,v in metadata.items()), - mission_id=contract_document.get("mission_id"), - ) + try: + producer = contract_document["producer"] + prompt = contract_document["runtime_prompt"] + metadata = contract_document["execution_metadata"] + constraints = contract_document["execution_constraints"] + receipts = contract_document.get("receipt_references", ()) + evidence_references = contract_document.get("execution_evidence_references", ()) + if not isinstance(producer, Mapping) or not isinstance(prompt, Mapping) or not isinstance(metadata, Mapping): + raise TypeError + identity = producer["identity"] + if not isinstance(identity, Mapping) or not isinstance(receipts, list) or not isinstance(evidence_references, list): + raise TypeError + contract = ProducerContract( + Producer(ProducerIdentity(str(identity["id"]), str(identity["type"]), str(identity["version"])), + str(producer["contract_version"])), + str(contract_document["correlation_id"]), str(contract_document["engineering_action_id"]), + RuntimePromptEnvelope(str(prompt["id"]), str(prompt["version"]), str(prompt["format"]), str(prompt["content"]), str(prompt["content_digest"])), + tuple(str(item) for item in constraints), tuple((str(k), str(v)) for k, v in metadata.items()), + mission_id=contract_document.get("mission_id"), + receipt_references=tuple(ExecutionReceiptReference(str(item["host_id"]), str(item["receipt_id"])) for item in receipts), + execution_evidence_references=tuple(str(item) for item in evidence_references), + contract_version=str(contract_document["contract_version"]), + ) + except (KeyError, TypeError, ValueError) as error: + raise MissionRunnerError("persisted Producer Contract is malformed") from error return ExecutionRequest( host_id=str(document["host_id"]), mission_id=str(document["mission_id"]), intent_id=str(document["intent_id"]), intent_revision=str(document["intent_revision"]), @@ -286,6 +304,11 @@ def _dispatch_or_recover(self, state: MissionExecutionState) -> MissionExecution dispatch = self._host.recover_dispatch(request) if dispatch is None: dispatch = self._host.dispatch(request) + # An accepted submission may legitimately have no run yet. Keep + # the persisted request in WAITING_FOR_EXECUTION and recover it on + # a later service tick; do not reclassify that state as failure. + if dispatch is None: + return state if dispatch.request != request: raise MissionRunnerError("execution host acknowledgement did not preserve the persisted request") except ExecutionHostTemporaryUnavailable: diff --git a/forge/runtime/service.py b/forge/runtime/service.py index 19d38ca..3e41826 100644 --- a/forge/runtime/service.py +++ b/forge/runtime/service.py @@ -1,50 +1,110 @@ -"""Minimal supervised Forge Runtime Service using the qualified loop unchanged.""" - +"""Supervised Forge Runtime Service over the canonical application loop.""" from __future__ import annotations +from contextlib import contextmanager from dataclasses import dataclass -from typing import Callable, Protocol +from pathlib import Path +from time import sleep +from typing import Callable, Iterator, Protocol + +try: # macOS/Linux runtime product path + import fcntl +except ImportError: # pragma: no cover - fail closed where no lock exists + fcntl = None # type: ignore[assignment] from forge.state import MissionExecutionStatus, MissionStateStore class RuntimeLoop(Protocol): - """The qualified loop surface reused by the service.""" - def run(self): ... - def resume(self, mission_id: str): ... +class RuntimeServiceBusy(RuntimeError): + """Another service or mutating CLI call owns this runtime instance.""" + + +class RuntimeServiceLock: + """A process-wide lease shared by service and mutating CLI composition. + + It is deliberately a filesystem lock beside the canonical runtime database, + not a second database row. Consumers can wrap a CLI mutation with this + same lease and cannot race a dispatching service tick. + """ + + def __init__(self, runtime_database_path: Path | str) -> None: + self.path = Path(runtime_database_path).with_name("forge-runtime-mutation.lock") + + @contextmanager + def acquire(self) -> Iterator[None]: + if fcntl is None: + raise RuntimeServiceBusy("runtime instance locking is unavailable") + self.path.parent.mkdir(parents=True, exist_ok=True) + with self.path.open("a+", encoding="utf-8") as handle: + try: + fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError as error: + raise RuntimeServiceBusy("canonical runtime is busy") from error + try: + yield + finally: + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + + @dataclass(frozen=True) class RuntimeServiceTick: - """Read-time result; Mission State remains the authority.""" - resumed_mission_ids: tuple[str, ...] dispatched_mission_id: str | None + progressed: bool class ForgeRuntimeService: - """Supervise the existing ExecutionLoop without becoming an Execution Host. + """Supervise the qualified loop without becoming an Execution Host. - A tick resumes one persisted Mission before asking the Dispatcher for new - work. This keeps the lane single-flight and preserves the CLI's semantics. + Every tick takes the same canonical-instance lease available to mutating + CLI composition. It resumes one persisted Mission before admitting new + work, preserving single-flight dispatch semantics across restarts. """ - def __init__(self, loop: RuntimeLoop, states: MissionStateStore) -> None: + def __init__(self, loop: RuntimeLoop, states: MissionStateStore, *, runtime_database_path: Path | str, + wait: Callable[[float], None] = sleep, minimum_backoff: float = 0.25, + maximum_backoff: float = 5.0) -> None: + if minimum_backoff <= 0 or maximum_backoff < minimum_backoff: + raise ValueError("runtime service backoff bounds are invalid") self._loop, self._states = loop, states + self._lock = RuntimeServiceLock(runtime_database_path) + self._wait, self._minimum_backoff, self._maximum_backoff = wait, minimum_backoff, maximum_backoff + + @property + def mutation_lock(self) -> RuntimeServiceLock: + return self._lock def tick(self) -> RuntimeServiceTick: - resumable = {MissionExecutionStatus.READY, MissionExecutionStatus.ACTIVE, - MissionExecutionStatus.WAITING_FOR_EXECUTION, MissionExecutionStatus.WAITING_FOR_EVIDENCE} - for state in self._states.resumable(): - if state.status in resumable: - self._loop.resume(state.mission_id) - return RuntimeServiceTick((state.mission_id,), None) - result = self._loop.run() - return RuntimeServiceTick((), None if result is None else result.mission_id) + with self._lock.acquire(): + resumable = {MissionExecutionStatus.READY, MissionExecutionStatus.ACTIVE, + MissionExecutionStatus.WAITING_FOR_EXECUTION, MissionExecutionStatus.WAITING_FOR_EVIDENCE} + for state in self._states.resumable(): + if state.status in resumable: + resumed = self._loop.resume(state.mission_id) + revision = getattr(resumed, "revision", state.revision) + return RuntimeServiceTick((state.mission_id,), None, revision != state.revision) + result = self._loop.run() + return RuntimeServiceTick((), None if result is None else result.mission_id, result is not None) def serve(self, *, keep_running: Callable[[], bool]) -> None: - """Caller-owned service loop: no hidden daemon or second state store.""" + """Run with bounded interruptible backoff; caller controls pause/stop.""" + delay = self._minimum_backoff while keep_running(): - self.tick() + try: + tick = self.tick() + except RuntimeServiceBusy: + tick = RuntimeServiceTick((), None, False) + if tick.progressed: + delay = self._minimum_backoff + continue + # Do not sleep through a requested pause/shutdown. The bounded + # delay also prevents a WAITING_FOR_EVIDENCE poll from busy-looping. + if not keep_running(): + break + self._wait(delay) + delay = min(self._maximum_backoff, delay * 2) diff --git a/forge/scheduler/ep_http_adapter.py b/forge/scheduler/ep_http_adapter.py index d67e3ea..177dad7 100644 --- a/forge/scheduler/ep_http_adapter.py +++ b/forge/scheduler/ep_http_adapter.py @@ -1,64 +1,202 @@ -"""Concrete, strict HTTP v1.1 Engineering Platform Execution Host adapter.""" +"""Concrete, strict HTTP v1.1 Engineering Platform Execution Host adapter. + +The runtime database remains the recovery authority. This adapter keeps only +the EP submission/run binding which follows from a persisted Forge request. +""" from __future__ import annotations + import json from dataclasses import dataclass -from urllib.error import HTTPError,URLError -from urllib.request import Request,urlopen -from forge.models.execution_host import ExecutionDispatch,ExecutionHostTemporaryUnavailable,ExecutionRequest +from typing import Any, Mapping, Protocol +from urllib.error import HTTPError, URLError +from urllib.request import Request, urlopen + +from forge.models.execution_host import ExecutionDispatch, ExecutionHostEvidence, ExecutionHostTemporaryUnavailable, ExecutionRequest from .ep_v11 import terminal_evidence + +class ExecutionHostBindingStore(Protocol): + def execution_host_binding(self, correlation_id: str) -> dict[str, Any] | None: ... + def save_execution_host_binding(self, correlation_id: str, document: Mapping[str, Any]) -> dict[str, Any]: ... + + @dataclass(frozen=True) class EngineeringPlatformHttpConfiguration: - base_url:str; project_id:str; bearer_token:str; host_id:str='engineering-platform'; timeout:float=10 + base_url: str + project_id: str + bearer_token: str + host_id: str = "engineering-platform" + timeout: float = 10 + class EngineeringPlatformHttpExecutionHost: - def __init__(self, config:EngineeringPlatformHttpConfiguration): self.config=config; self._submissions={} - def _json(self,path,*,method='GET',body=None): - data=None if body is None else json.dumps(body,sort_keys=True,separators=(',',':')).encode() - request=Request(self.config.base_url.rstrip('/')+path,data=data,method=method,headers={'Authorization':'Bearer '+self.config.bearer_token,'Content-Type':'application/json'}) - try: - with urlopen(request,timeout=self.config.timeout) as response:return json.loads(response.read()) - except HTTPError as error: - if error.code>=500: raise ExecutionHostTemporaryUnavailable('EP temporarily unavailable') from error - raise ValueError(f'EP rejected request: {error.code}') from error - except (URLError,TimeoutError) as error: raise ExecutionHostTemporaryUnavailable('EP transport unavailable') from error - def _bytes(self,path): - request=Request(self.config.base_url.rstrip('/')+path,headers={'Authorization':'Bearer '+self.config.bearer_token}) - try: - with urlopen(request,timeout=self.config.timeout) as response: - raw=response.read(1_048_577) - if len(raw)>1_048_576: raise ValueError('EP artifact exceeds consumer size limit') - return raw - except HTTPError as error: - if error.code>=500: raise ExecutionHostTemporaryUnavailable('EP artifact temporarily unavailable') from error - raise ValueError(f'EP artifact request rejected: {error.code}') from error - except (URLError,TimeoutError) as error: raise ExecutionHostTemporaryUnavailable('EP artifact transport unavailable') from error - def _payload(self,r): - p=r.runtime_prompt; text=getattr(p,'rendered_text',None) or p.to_markdown() - revision=getattr(p,'mission_revision',None) - if not isinstance(revision,str) or not revision: raise ValueError('Execution Request has no persisted Mission revision') - contract=r.producer_contract - return {'repository_id':r.repository_id,'producer':contract.producer.identity.to_dict(),'prompt':text,'idempotency_key':r.correlation_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'engineering_action_id':r.action_id,'constraints':{'forge_execution':{'contract_version':'1.0','host_id':r.host_id,'repository_id':r.repository_id,'correlation_id':r.correlation_id,'mission_id':r.mission_id,'mission_revision':revision,'intent_id':r.intent_id,'intent_revision':r.intent_revision,'action_id':r.action_id,'runtime_prompt':{'id':contract.runtime_prompt.id,'content_digest':contract.runtime_prompt.content_digest},'retry_of_correlation_id':r.retry_of_correlation_id}}} - def dispatch(self,r): - accepted=self._json(f'/v1/projects/{self.config.project_id}/submissions',method='POST',body=self._payload(r)); sid=accepted.get('submission_id') - if not isinstance(sid,str): raise ValueError('EP submission acknowledgement lacks submission_id') - self._submissions[r.correlation_id]=sid - return self.recover_dispatch(r) or (_ for _ in ()).throw(ExecutionHostTemporaryUnavailable('submission accepted but run unclaimed')) - def recover_dispatch(self,r): - sid=self._submissions.get(r.correlation_id) - if not sid:return None - readback=self._json(f'/v1/projects/{self.config.project_id}/submissions/{sid}') - if readback.get('correlation',{}).get('correlation_id')!=r.correlation_id or readback.get('submission',{}).get('repository_id')!=r.repository_id: raise ValueError('EP readback does not bind persisted request') - run=readback.get('run') - return None if not isinstance(run,dict) else ExecutionDispatch(r,str(run['id'])) - def retrieve_evidence(self,d): - sid=self._submissions.get(d.request.correlation_id) - if not sid: raise ValueError('missing persisted submission identity') - readback=self._json(f'/v1/projects/{self.config.project_id}/submissions/{sid}') - if readback.get('run',{}).get('id')!=d.host_run_id:return None - terminal=readback.get('evidence',{}).get('terminal_artifact') - if not isinstance(terminal,dict): return None - raw=self._bytes(f"/v1/projects/{self.config.project_id}/artifacts/{terminal['id']}") - evidence=terminal_evidence(readback,raw,host_id=self.config.host_id) - if (evidence.correlation_id,evidence.host_run_id,evidence.repository_evidence.runtime_prompt_id)!=(d.request.correlation_id,d.host_run_id,d.request.runtime_prompt.id): raise ValueError('EP artifact does not bind dispatch') - return evidence + """EP v1.1 transport with durable, request-first idempotency.""" + + def __init__(self, config: EngineeringPlatformHttpConfiguration, bindings: ExecutionHostBindingStore) -> None: + if not config.base_url or not config.project_id or not config.bearer_token: + raise ValueError("EP HTTP configuration requires endpoint, project, and credential") + self.config = config + self._bindings = bindings + + def _json(self, path: str, *, method: str = "GET", body: Mapping[str, Any] | None = None) -> dict[str, Any]: + data = None if body is None else json.dumps(body, sort_keys=True, separators=(",", ":")).encode("utf-8") + request = Request(self.config.base_url.rstrip("/") + path, data=data, method=method, + headers={"Authorization": "Bearer " + self.config.bearer_token, "Content-Type": "application/json"}) + try: + with urlopen(request, timeout=self.config.timeout) as response: + value = json.loads(response.read()) + if not isinstance(value, dict): + raise ValueError("EP response must be a JSON object") + return value + except HTTPError as error: + if error.code >= 500: + raise ExecutionHostTemporaryUnavailable("EP temporarily unavailable") from error + raise ValueError(f"EP rejected request: {error.code}") from error + except (URLError, TimeoutError) as error: + raise ExecutionHostTemporaryUnavailable("EP transport unavailable") from error + + def _bytes(self, path: str) -> bytes: + request = Request(self.config.base_url.rstrip("/") + path, headers={"Authorization": "Bearer " + self.config.bearer_token}) + try: + with urlopen(request, timeout=self.config.timeout) as response: + raw = response.read(1_048_577) + if len(raw) > 1_048_576: + raise ValueError("EP artifact exceeds consumer size limit") + return raw + except HTTPError as error: + if error.code >= 500: + raise ExecutionHostTemporaryUnavailable("EP artifact temporarily unavailable") from error + raise ValueError(f"EP artifact request rejected: {error.code}") from error + except (URLError, TimeoutError) as error: + raise ExecutionHostTemporaryUnavailable("EP artifact transport unavailable") from error + + @staticmethod + def _metadata(request: ExecutionRequest) -> dict[str, str]: + return dict(request.producer_contract.execution_metadata) + + def _mission_revision(self, request: ExecutionRequest) -> str: + revision = self._metadata(request).get("mission_revision") + if not revision: + raise ValueError("persisted Producer Contract lacks mission_revision provenance") + return revision + + def _request_binding(self, request: ExecutionRequest) -> dict[str, Any]: + contract = request.producer_contract + return {"correlation_id": request.correlation_id, "host_id": request.host_id, "project_id": self.config.project_id, + "repository_id": request.repository_id, "mission_id": request.mission_id, + "mission_revision": self._mission_revision(request), "intent_id": request.intent_id, + "intent_revision": request.intent_revision, "action_id": request.action_id, + "runtime_prompt_id": contract.runtime_prompt.id, "runtime_prompt_digest": contract.runtime_prompt.content_digest, + "producer": contract.producer.identity.to_dict(), "producer_contract_digest": contract.digest(), + "retry_of_correlation_id": request.retry_of_correlation_id} + + def _binding(self, request: ExecutionRequest) -> dict[str, Any]: + # Saved before POST: an ambiguous send is retried with the exact same + # idempotency key, never a fresh correlation. + return self._bindings.save_execution_host_binding(request.correlation_id, self._request_binding(request)) + + def _payload(self, request: ExecutionRequest) -> dict[str, Any]: + binding, contract = self._request_binding(request), request.producer_contract + return {"repository_id": request.repository_id, "producer": contract.producer.identity.to_dict(), + "prompt": contract.runtime_prompt.content, "idempotency_key": request.correlation_id, + "correlation_id": request.correlation_id, "mission_id": request.mission_id, + "engineering_action_id": request.action_id, "constraints": {"forge_execution": { + "contract_version": "1.0", "host_id": request.host_id, "repository_id": request.repository_id, + "correlation_id": request.correlation_id, "mission_id": request.mission_id, + "mission_revision": binding["mission_revision"], "intent_id": request.intent_id, + "intent_revision": request.intent_revision, "action_id": request.action_id, + "runtime_prompt": {"id": contract.runtime_prompt.id, "content_digest": contract.runtime_prompt.content_digest}, + "retry_of_correlation_id": request.retry_of_correlation_id}}} + + def _validate_readback(self, request: ExecutionRequest, binding: Mapping[str, Any], readback: Mapping[str, Any]) -> None: + correlation, submission, producer = readback.get("correlation"), readback.get("submission"), readback.get("producer") + root_provenance = readback.get("provenance") + provenance = root_provenance.get("forge_execution") if isinstance(root_provenance, Mapping) else None + if not all(isinstance(value, Mapping) for value in (correlation, submission, provenance, producer)): + raise ValueError("EP readback omits required request binding") + expected_correlation = {"correlation_id": request.correlation_id, "mission_id": request.mission_id, + "engineering_action_id": request.action_id} + if any(correlation.get(key) != value for key, value in expected_correlation.items()): + raise ValueError("EP readback correlation does not bind persisted request") + if submission.get("repository_id") != request.repository_id or submission.get("project_id") != self.config.project_id: + raise ValueError("EP readback submission does not bind project and repository") + if dict(producer) != binding["producer"]: + raise ValueError("EP readback producer does not bind persisted request") + expected_provenance = {"host_id": request.host_id, "repository_id": request.repository_id, + "correlation_id": request.correlation_id, "mission_id": request.mission_id, + "mission_revision": binding["mission_revision"], "intent_id": request.intent_id, + "intent_revision": request.intent_revision, "action_id": request.action_id, + "retry_of_correlation_id": request.retry_of_correlation_id} + if any(provenance.get(key) != value for key, value in expected_provenance.items()): + raise ValueError("EP readback provenance does not bind persisted request") + prompt = provenance.get("runtime_prompt") + if not isinstance(prompt, Mapping) or dict(prompt) != {"id": binding["runtime_prompt_id"], "content_digest": binding["runtime_prompt_digest"]}: + raise ValueError("EP readback Runtime Prompt does not bind persisted request") + + def _readback(self, request: ExecutionRequest, binding: Mapping[str, Any]) -> dict[str, Any] | None: + submission_id = binding.get("submission_id") + if not isinstance(submission_id, str) or not submission_id: + return None + readback = self._json(f"/v1/projects/{self.config.project_id}/submissions/{submission_id}") + self._validate_readback(request, binding, readback) + if readback.get("submission", {}).get("id") != submission_id: + raise ValueError("EP readback submission identity changed") + return readback + + def dispatch(self, request: ExecutionRequest) -> ExecutionDispatch | None: + binding = self._binding(request) + readback = self._readback(request, binding) + if readback is None: + accepted = self._json(f"/v1/projects/{self.config.project_id}/submissions", method="POST", body=self._payload(request)) + submission_id = accepted.get("submission_id") + if not isinstance(submission_id, str) or not submission_id: + raise ValueError("EP submission acknowledgement lacks submission_id") + binding = self._bindings.save_execution_host_binding(request.correlation_id, + {"correlation_id": request.correlation_id, "submission_id": submission_id}) + readback = self._readback(request, binding) + return self._dispatch_from_readback(request, binding, readback) + + def recover_dispatch(self, request: ExecutionRequest) -> ExecutionDispatch | None: + binding = self._binding(request) + return self._dispatch_from_readback(request, binding, self._readback(request, binding)) + + def _dispatch_from_readback(self, request: ExecutionRequest, binding: Mapping[str, Any], readback: Mapping[str, Any] | None) -> ExecutionDispatch | None: + if readback is None: + return None + run = readback.get("run") + if run is None: + return None # accepted but not yet claimed: ordinary resumable waiting + if not isinstance(run, Mapping) or not isinstance(run.get("id"), str) or not run["id"]: + raise ValueError("EP readback run identity is invalid") + known_run = binding.get("host_run_id") + if known_run is not None and known_run != run["id"]: + raise ValueError("EP readback run conflicts with persisted dispatch") + self._bindings.save_execution_host_binding(request.correlation_id, + {"correlation_id": request.correlation_id, "host_run_id": run["id"]}) + return ExecutionDispatch(request, run["id"]) + + def retrieve_evidence(self, dispatch: ExecutionDispatch) -> ExecutionHostEvidence | None: + request, binding = dispatch.request, self._binding(dispatch.request) + if binding.get("host_run_id") not in (None, dispatch.host_run_id): + raise ValueError("persisted dispatch run differs from requested evidence run") + readback = self._readback(request, binding) + if readback is None: + raise ValueError("missing persisted EP submission identity") + observed = self._dispatch_from_readback(request, binding, readback) + if observed is None: + return None + if observed.host_run_id != dispatch.host_run_id: + raise ValueError("EP readback run differs from persisted dispatch") + terminal = readback.get("evidence", {}).get("terminal_artifact") + if not isinstance(terminal, Mapping) or not isinstance(terminal.get("id"), str): + return None + raw = self._bytes(f"/v1/projects/{self.config.project_id}/artifacts/{terminal['id']}") + evidence = terminal_evidence(readback, raw, host_id=self.config.host_id) + observed_identity = (evidence.correlation_id, evidence.host_run_id, evidence.repository_evidence.runtime_prompt_id, + evidence.repository_evidence.mission_id, evidence.repository_evidence.intent_id, + evidence.repository_evidence.intent_revision, evidence.repository_evidence.action_id, evidence.repository_evidence.repository_id) + request_identity = (request.correlation_id, dispatch.host_run_id, request.producer_contract.runtime_prompt.id, + request.mission_id, request.intent_id, request.intent_revision, request.action_id, request.repository_id) + if observed_identity != request_identity: + raise ValueError("EP terminal artifact does not bind persisted request and dispatch") + return evidence diff --git a/forge/scheduler/ep_v11.py b/forge/scheduler/ep_v11.py index 39ba778..74e4fa6 100644 --- a/forge/scheduler/ep_v11.py +++ b/forge/scheduler/ep_v11.py @@ -13,7 +13,21 @@ def terminal_evidence(readback: Mapping[str,Any], artifact: bytes, *, host_id: s try: document=json.loads(artifact) except (UnicodeDecodeError,json.JSONDecodeError) as e: raise ValueError('EP terminal artifact is invalid JSON') from e repo=document.get('repository',{}); outcome=ExecutionEvidenceOutcome(str(result.get('outcome')).lower()) - if not isinstance(repo.get('revision'),str) or not isinstance(document.get('report',{}).get('id'),str): raise ValueError('EP terminal artifact lacks repository evidence') + if not isinstance(repo,dict) or not isinstance(document.get('report',{}).get('id'),str): raise ValueError('EP terminal artifact lacks repository evidence') correlation=readback['correlation']; prompt=provenance['runtime_prompt']; submission=readback['submission'] - repository=ExecutionRepositoryEvidence(str(correlation['mission_id']),str(provenance['intent_id']),str(provenance['intent_revision']),str(correlation['engineering_action_id']),str(prompt['id']),str(correlation['correlation_id']),str(run['id']),str(repo['id']),str(repo['revision']),str(document['report']['id']),digest) + # A hash only proves the returned bytes came from *some* registered object. + # Bind their internal subject to the readback before mapping it to Forge. + artifact_correlation=document.get('correlation',{}); artifact_provenance=document.get('provenance',{}); artifact_run=document.get('run',{}); artifact_submission=document.get('submission',{}); artifact_producer=document.get('producer',{}) + if not all(isinstance(value,dict) for value in (artifact_correlation,artifact_provenance,artifact_run,artifact_submission,artifact_producer)): + raise ValueError('EP terminal artifact lacks binding fields') + if artifact_correlation != correlation or artifact_producer != readback.get('producer'): + raise ValueError('EP terminal artifact identity differs from readback') + expected_provenance={key: provenance.get(key) for key in ('action_id','contract_version','correlation_id','host_id','intent_id','intent_revision','mission_id','mission_revision','repository_id','retry_of_correlation_id','runtime_prompt')} + if artifact_provenance != expected_provenance or artifact_run.get('id') != run.get('id') or artifact_submission.get('id') != submission.get('id') or artifact_submission.get('project_id') != submission.get('project_id') or artifact_submission.get('repository_id') != submission.get('repository_id') or artifact_submission.get('accepted_request_digest') != submission.get('accepted_request_digest'): + raise ValueError('EP terminal artifact does not bind readback run and submission') + revision=repo.get('revision') + if revision is not None and not isinstance(revision,str): raise ValueError('EP terminal artifact revision is invalid') + if outcome is ExecutionEvidenceOutcome.COMPLETE and (not revision or document.get('run',{}).get('delivery_qualified') is not True or repo.get('revision_required') is not True): raise ValueError('EP complete terminal artifact lacks qualified delivery revision') + if repo.get('id') != readback.get('evidence',{}).get('repository',{}).get('id') or revision != readback.get('evidence',{}).get('repository',{}).get('revision'): raise ValueError('EP terminal artifact repository differs from readback') + repository=ExecutionRepositoryEvidence(str(correlation['mission_id']),str(provenance['intent_id']),str(provenance['intent_revision']),str(correlation['engineering_action_id']),str(prompt['id']),str(correlation['correlation_id']),str(run['id']),str(repo['id']),revision,str(document['report']['id']),digest) return ExecutionHostEvidence(host_id,str(correlation['correlation_id']),str(run['id']),str(document['report']['id']),outcome,repository,validation_references=tuple(str(x.get('command')) for x in document.get('references',{}).get('validation',()) if x.get('command'))) diff --git a/tests/test_action_derivation_canary_closure.py b/tests/test_action_derivation_canary_closure.py index e302c30..450ac10 100644 --- a/tests/test_action_derivation_canary_closure.py +++ b/tests/test_action_derivation_canary_closure.py @@ -192,7 +192,7 @@ def revised(**changes): def test_schema30_migrates_the_bounded_closure_store_before_reopen(self) -> None: path = self._schema30_fixture() self.db = RuntimeDatabase(Path(self.directory.name), path=path, forge_version="test") - self.assertEqual(self.db.metadata["schema_version"], "31") + self.assertEqual(self.db.metadata["schema_version"], "32") tables = {row["name"] for row in self.db._connection.execute("SELECT name FROM sqlite_master WHERE type='table'")} triggers = {row["name"] for row in self.db._connection.execute("SELECT name FROM sqlite_master WHERE type='trigger'")} self.assertIn("action_derivation_canary_closures", tables) diff --git a/tests/test_ep_http_adapter.py b/tests/test_ep_http_adapter.py new file mode 100644 index 0000000..9610ac9 --- /dev/null +++ b/tests/test_ep_http_adapter.py @@ -0,0 +1,121 @@ +"""HTTP consumer tests against EP's pinned v1.1 shared fixtures.""" +from __future__ import annotations + +import hashlib +import json +from pathlib import Path +from tempfile import TemporaryDirectory +import unittest +from unittest.mock import patch + +from forge.models import Producer, ProducerContract, ProducerIdentity, RuntimePrompt, RuntimePromptEnvelope, RuntimePromptSection, RuntimePromptSectionKind, ProviderPromptDefinition +from forge.models.execution_host import ExecutionRequest +from forge.models import ExecutionDispatch, ExecutionEvidenceOutcome +from forge.runtime.database import RuntimeDatabase +from forge.scheduler.ep_http_adapter import EngineeringPlatformHttpConfiguration, EngineeringPlatformHttpExecutionHost + + +FIXTURES = Path("/Users/pcvantol/Documents/GitHub/engineering-platform/tests/fixtures") + + +class _Response: + def __init__(self, value: bytes) -> None: self.value = value + def __enter__(self): return self + def __exit__(self, *args): return False + def read(self, _limit: int | None = None) -> bytes: return self.value if _limit is None else self.value[:_limit] + + +def _prompt() -> RuntimePrompt: + return RuntimePrompt("runtime-prompt-fixture", "intent-fixture", "7", "action-fixture", + ProviderPromptDefinition("provider", "1"), "sha256:" + "a" * 64, + tuple(RuntimePromptSection(kind, (kind.value,)) for kind in RuntimePromptSectionKind)) + + +def _request() -> ExecutionRequest: + prompt = _prompt() + contract = ProducerContract(Producer(ProducerIdentity("forge", "FORGE", "1.0")), "forge-correlation-fixture", + "action-fixture", RuntimePromptEnvelope(prompt.id, "1.0", "text/markdown", "exact persisted prompt", "sha256:" + "a" * 64), + ("Execute only the supplied Runtime Prompt.",), + (("intent_id", "intent-fixture"), ("intent_revision", "7"), ("mission_revision", "3"), + ("repository_id", "forge"), ("workspace_id", "workspace-1")), mission_id="mission-fixture") + return ExecutionRequest("engineering-platform", "mission-fixture", "intent-fixture", "7", "action-fixture", prompt, + "workspace-1", "forge", "forge-correlation-fixture", "2026-09-07T00:00:00Z", producer_contract=contract) + + +class EngineeringPlatformHttpExecutionHostTests(unittest.TestCase): + def setUp(self) -> None: + self.temporary = TemporaryDirectory() + self.database = RuntimeDatabase(".", path=Path(self.temporary.name) / "runtime.db", forge_version="test") + self.config = EngineeringPlatformHttpConfiguration("https://ep.test", "forge", "credential") + self.request = _request() + self.readback = json.loads((FIXTURES / "forge-producer-readback-v1.1.json").read_text()) + self.artifact = (FIXTURES / "forge-terminal-evidence-v1.1.json").read_bytes() + + def tearDown(self) -> None: + self.database.close(); self.temporary.cleanup() + + @staticmethod + def _urlopen(responses: list[bytes], observed: list[object]): + def call(request, *, timeout): + observed.append(request) + return _Response(responses.pop(0)) + return call + + def test_accepted_submission_recovers_after_reopen_without_process_memory(self) -> None: + waiting = {**self.readback, "run": None} + observed: list[object] = [] + with patch("forge.scheduler.ep_http_adapter.urlopen", self._urlopen([ + json.dumps({"submission_id": "submission-fixture"}).encode(), json.dumps(waiting).encode()], observed)): + host = EngineeringPlatformHttpExecutionHost(self.config, self.database) + self.assertIsNone(host.dispatch(self.request)) + self.assertEqual(self.database.execution_host_binding(self.request.correlation_id)["submission_id"], "submission-fixture") + self.database.close() + self.database = RuntimeDatabase(".", path=Path(self.temporary.name) / "runtime.db", forge_version="test") + observed = [] + with patch("forge.scheduler.ep_http_adapter.urlopen", self._urlopen([json.dumps(self.readback).encode()], observed)): + dispatch = EngineeringPlatformHttpExecutionHost(self.config, self.database).recover_dispatch(self.request) + self.assertEqual(dispatch.host_run_id, "run-fixture") + self.assertEqual(observed[0].get_header("Authorization"), "Bearer credential") + + def test_valid_hash_from_another_run_is_rejected_after_raw_artifact_fetch(self) -> None: + self.database.save_execution_host_binding(self.request.correlation_id, + {"correlation_id": self.request.correlation_id, "submission_id": "submission-fixture", "host_run_id": "run-fixture"}) + substituted = json.loads(self.artifact) + substituted["run"]["id"] = "other-run" + raw = json.dumps(substituted, sort_keys=True, separators=(",", ":")).encode() + b"\n" + readback = json.loads(json.dumps(self.readback)) + readback["evidence"]["terminal_artifact"]["digest"] = "sha256:" + hashlib.sha256(raw).hexdigest() + observed: list[object] = [] + with patch("forge.scheduler.ep_http_adapter.urlopen", self._urlopen([json.dumps(readback).encode(), raw], observed)): + with self.assertRaisesRegex(ValueError, "artifact"): + EngineeringPlatformHttpExecutionHost(self.config, self.database).retrieve_evidence(ExecutionDispatch(self.request, "run-fixture")) + self.assertEqual(len(observed), 2, "the adapter must fetch and hash real artifact bytes before rejecting it") + + def test_existing_dispatch_with_a_different_run_is_an_identity_error(self) -> None: + self.database.save_execution_host_binding(self.request.correlation_id, + {"correlation_id": self.request.correlation_id, "submission_id": "submission-fixture", "host_run_id": "run-fixture"}) + changed = json.loads(json.dumps(self.readback)); changed["run"]["id"] = "other-run" + with patch("forge.scheduler.ep_http_adapter.urlopen", self._urlopen([json.dumps(changed).encode()], [])): + with self.assertRaisesRegex(ValueError, "conflicts"): + EngineeringPlatformHttpExecutionHost(self.config, self.database).recover_dispatch(self.request) + + def test_failed_terminal_evidence_without_delivery_revision_is_preserved(self) -> None: + self.database.save_execution_host_binding(self.request.correlation_id, + {"correlation_id": self.request.correlation_id, "submission_id": "submission-fixture", "host_run_id": "run-fixture"}) + artifact = json.loads(self.artifact) + artifact["run"].update({"outcome": "FAILED", "delivery_qualified": False}) + artifact["report"]["terminal_state"] = "FAILED" + artifact["repository"].update({"revision": None, "revision_required": False}) + raw = json.dumps(artifact, sort_keys=True, separators=(",", ":")).encode() + b"\n" + readback = json.loads(json.dumps(self.readback)) + readback["result"].update({"outcome": "FAILED", "delivery_qualified": False}) + readback["evidence"]["repository"]["revision"] = None + readback["evidence"]["terminal_artifact"]["digest"] = "sha256:" + hashlib.sha256(raw).hexdigest() + with patch("forge.scheduler.ep_http_adapter.urlopen", self._urlopen([json.dumps(readback).encode(), raw], [])): + evidence = EngineeringPlatformHttpExecutionHost(self.config, self.database).retrieve_evidence(ExecutionDispatch(self.request, "run-fixture")) + self.assertEqual(evidence.outcome, ExecutionEvidenceOutcome.FAILED) + self.assertIsNone(evidence.repository_evidence.repository_revision) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_runtime_service.py b/tests/test_runtime_service.py new file mode 100644 index 0000000..20e2b23 --- /dev/null +++ b/tests/test_runtime_service.py @@ -0,0 +1,54 @@ +from __future__ import annotations + +from dataclasses import dataclass +from pathlib import Path +from tempfile import TemporaryDirectory +import unittest + +from forge.runtime.service import ForgeRuntimeService, RuntimeServiceBusy, RuntimeServiceLock +from forge.state import MissionExecutionStatus + + +@dataclass +class _State: + mission_id: str + status: MissionExecutionStatus + revision: int + + +class _States: + def __init__(self, state: _State | None) -> None: self.state = state + def resumable(self): return () if self.state is None else (self.state,) + + +class _Loop: + def __init__(self, state: _State, *, progresses: bool) -> None: self.state, self.progresses, self.calls = state, progresses, [] + def resume(self, mission_id): + self.calls.append(mission_id) + return _State(mission_id, self.state.status, self.state.revision + int(self.progresses)) + def run(self): return None + + +class RuntimeServiceTests(unittest.TestCase): + def test_waiting_evidence_uses_bounded_interruptible_backoff(self) -> None: + with TemporaryDirectory() as root: + state = _State("mission", MissionExecutionStatus.WAITING_FOR_EVIDENCE, 4) + waits: list[float] = [] + service = ForgeRuntimeService(_Loop(state, progresses=False), _States(state), + runtime_database_path=Path(root) / "runtime.db", wait=waits.append, minimum_backoff=0.1, maximum_backoff=0.2) + calls = iter((True, True, True, False)) + service.serve(keep_running=lambda: next(calls)) + self.assertEqual(waits, [0.1]) + + def test_service_and_mutating_cli_share_one_runtime_lease(self) -> None: + with TemporaryDirectory() as root: + path = Path(root) / "runtime.db" + lock = RuntimeServiceLock(path) + with lock.acquire(): + with self.assertRaises(RuntimeServiceBusy): + with RuntimeServiceLock(path).acquire(): + pass + + +if __name__ == "__main__": + unittest.main() From 4ef7e32bb40e792380084dbad18734fbd604ba44 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Tue, 8 Sep 2026 07:41:35 +0200 Subject: [PATCH 09/14] test: vendor pinned EP v1.1 consumer fixtures --- tests/fixtures/forge-producer-readback-v1.1.json | 1 + tests/fixtures/forge-terminal-evidence-v1.1.json | 1 + tests/test_ep_http_adapter.py | 2 +- 3 files changed, 3 insertions(+), 1 deletion(-) create mode 100644 tests/fixtures/forge-producer-readback-v1.1.json create mode 100644 tests/fixtures/forge-terminal-evidence-v1.1.json diff --git a/tests/fixtures/forge-producer-readback-v1.1.json b/tests/fixtures/forge-producer-readback-v1.1.json new file mode 100644 index 0000000..ff8a24e --- /dev/null +++ b/tests/fixtures/forge-producer-readback-v1.1.json @@ -0,0 +1 @@ +{"contract_version":"1.1","correlation":{"correlation_id":"forge-correlation-fixture","engineering_action_id":"action-fixture","mission_id":"mission-fixture"},"evidence":{"repository":{"id":"forge","revision":"1111111111111111111111111111111111111111"},"status":"AVAILABLE","terminal_artifact":{"content_type":"application/json","digest":"sha256:90e376388fdac5fdec49cd1cee955b0a5d75515416dc515f9b9973d29355cb83","digest_algorithm":"sha256","id":"terminal-evidence:run-fixture"}},"producer":{"id":"forge","type":"FORGE","version":"1.0"},"provenance":{"forge_execution":{"action_id":"action-fixture","contract_version":"1.0","correlation_id":"forge-correlation-fixture","host_id":"engineering-platform","intent_id":"intent-fixture","intent_revision":"7","mission_id":"mission-fixture","mission_revision":"3","repository_id":"forge","retry_of_correlation_id":null,"runtime_prompt":{"content_digest":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","id":"runtime-prompt-fixture"}},"status":"PERSISTED"},"result":{"delivery_qualified":true,"outcome":"COMPLETE","terminal":true},"run":{"id":"run-fixture","operator_resolution":"NONE","state":"COMPLETE","terminal":true,"updated_at":"2026-09-07T00:00:00+00:00"},"submission":{"accepted_request_digest":"sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb","admission":"ADMITTED","created_at":"2026-09-07T00:00:00+00:00","id":"submission-fixture","project_id":"forge","repository_id":"forge","state":"QUEUED","transport":"HTTP"}} diff --git a/tests/fixtures/forge-terminal-evidence-v1.1.json b/tests/fixtures/forge-terminal-evidence-v1.1.json new file mode 100644 index 0000000..ef87776 --- /dev/null +++ b/tests/fixtures/forge-terminal-evidence-v1.1.json @@ -0,0 +1 @@ +{"artifact_type":"EP_TERMINAL_EVIDENCE","contract_version":"1.1","correlation":{"correlation_id":"forge-correlation-fixture","engineering_action_id":"action-fixture","mission_id":"mission-fixture"},"producer":{"id":"forge","type":"FORGE","version":"1.0"},"provenance":{"action_id":"action-fixture","contract_version":"1.0","correlation_id":"forge-correlation-fixture","host_id":"engineering-platform","intent_id":"intent-fixture","intent_revision":"7","mission_id":"mission-fixture","mission_revision":"3","repository_id":"forge","retry_of_correlation_id":null,"runtime_prompt":{"content_digest":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","id":"runtime-prompt-fixture"}},"references":{"finalization":"finalization:fixture","quality":[],"repair":[],"validation":[{"command":"python -m unittest","result":"PASS"}]},"report":{"id":"report:run-fixture","terminal_state":"COMPLETE"},"repository":{"id":"forge","revision":"1111111111111111111111111111111111111111","revision_required":true},"run":{"delivery_qualified":true,"id":"run-fixture","outcome":"COMPLETE"},"submission":{"accepted_request_digest":"sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb","id":"submission-fixture","project_id":"forge","repository_id":"forge"}} diff --git a/tests/test_ep_http_adapter.py b/tests/test_ep_http_adapter.py index 9610ac9..3c911c7 100644 --- a/tests/test_ep_http_adapter.py +++ b/tests/test_ep_http_adapter.py @@ -15,7 +15,7 @@ from forge.scheduler.ep_http_adapter import EngineeringPlatformHttpConfiguration, EngineeringPlatformHttpExecutionHost -FIXTURES = Path("/Users/pcvantol/Documents/GitHub/engineering-platform/tests/fixtures") +FIXTURES = Path(__file__).with_name("fixtures") class _Response: From 29ba0600967b3f4dfd382a16cf9d8ea8668b1a35 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Tue, 8 Sep 2026 07:49:37 +0200 Subject: [PATCH 10/14] test: raise Forge component coverage floor --- tests/test_ep_http_adapter.py | 8 +++ tests/test_governed_coverage.py | 104 ++++++++++++++++++++++++++++++++ tests/test_operator_identity.py | 21 ++++++- 3 files changed, 132 insertions(+), 1 deletion(-) create mode 100644 tests/test_governed_coverage.py diff --git a/tests/test_ep_http_adapter.py b/tests/test_ep_http_adapter.py index 3c911c7..b7aff49 100644 --- a/tests/test_ep_http_adapter.py +++ b/tests/test_ep_http_adapter.py @@ -7,6 +7,7 @@ from tempfile import TemporaryDirectory import unittest from unittest.mock import patch +from urllib.error import URLError from forge.models import Producer, ProducerContract, ProducerIdentity, RuntimePrompt, RuntimePromptEnvelope, RuntimePromptSection, RuntimePromptSectionKind, ProviderPromptDefinition from forge.models.execution_host import ExecutionRequest @@ -99,6 +100,13 @@ def test_existing_dispatch_with_a_different_run_is_an_identity_error(self) -> No with self.assertRaisesRegex(ValueError, "conflicts"): EngineeringPlatformHttpExecutionHost(self.config, self.database).recover_dispatch(self.request) + def test_configuration_and_transport_fail_closed_without_persisted_authority(self) -> None: + with self.assertRaisesRegex(ValueError, "configuration"): + EngineeringPlatformHttpExecutionHost(EngineeringPlatformHttpConfiguration("", "forge", "credential"), self.database) + with patch("forge.scheduler.ep_http_adapter.urlopen", side_effect=URLError("offline")): + with self.assertRaisesRegex(Exception, "transport unavailable"): + EngineeringPlatformHttpExecutionHost(self.config, self.database)._json("/v1/projects/forge/submissions") + def test_failed_terminal_evidence_without_delivery_revision_is_preserved(self) -> None: self.database.save_execution_host_binding(self.request.correlation_id, {"correlation_id": self.request.correlation_id, "submission_id": "submission-fixture", "host_run_id": "run-fixture"}) diff --git a/tests/test_governed_coverage.py b/tests/test_governed_coverage.py new file mode 100644 index 0000000..a038df2 --- /dev/null +++ b/tests/test_governed_coverage.py @@ -0,0 +1,104 @@ +"""Focused unit coverage for small governed application services.""" +from __future__ import annotations + +from types import SimpleNamespace +import unittest +from unittest.mock import patch + +from forge.action_derivation_reattempt import CanonicalActionDerivationReattemptService, REQUEST_SEMANTICS_CHANGED +from forge.mission_amendment import MissionAmendmentService, _digest +from forge.models.action_derivation import PlanningSnapshot +from forge.operator_identity import OperatorContext +from forge.planner.openai_responses import CanonicalTokenPreflightAuthority, TokenPreflightBoundary, _G011PolicySnapshot + + +class _Cursor: + def __init__(self, one=None, rows=()): self.one, self.rows = one, rows + def fetchone(self): return self.one + def fetchall(self): return self.rows + + +class _Connection: + def __init__(self, cursors=()): self.cursors = list(cursors) + def execute(self, *_args): return self.cursors.pop(0) + + +class GovernedServiceCoverageTests(unittest.TestCase): + def test_successor_authorization_binds_only_a_changed_failed_request(self) -> None: + context = OperatorContext("installation", "operator", 1) + policy = SimpleNamespace(provider_id="openai", version=1, model="model", secret_reference=SimpleNamespace(fingerprint="secret"), + timeout_seconds=1, input_token_bound=2, context_token_bound=3, output_token_bound=4) + predecessor = {"lifecycle": "FAILED", "mission_id": "mission", "snapshot_digest": "snapshot", + "effective_contract_digest": "contract", "evidence_digest": "evidence", "provider_configuration": _G011PolicySnapshot.from_policy(policy).digest, + "generation_request_digest": "old-request", "main_head": "head"} + database = SimpleNamespace(get_document=lambda table, key: predecessor if table == "action_derivations" else {"status": "APPROVED_PLANNABLE"}, + _connection=_Connection([_Cursor({"sequence": None})]), + create_action_derivation_reattempt_authorization=lambda document: document) + repository = SimpleNamespace(database=database, operators=SimpleNamespace(authorize=lambda value: value == context), + _operator_id=lambda value: "operator-id") + adapter = SimpleNamespace(configuration=SimpleNamespace(policy_service=SimpleNamespace(db=database), current_policy=lambda: policy, + preflight_authority=SimpleNamespace(boundary_for=lambda _: TokenPreflightBoundary("head", "evidence", "contract"))), + _body=lambda request, _: {"new": request.derivation_id}) + snapshot = SimpleNamespace(digest="snapshot") + with patch("forge.action_derivation_reattempt.CanonicalActionDerivationEvidenceProducer", return_value=SimpleNamespace(planner_input=lambda _: object())), \ + patch.object(PlanningSnapshot, "from_planner_input", return_value=snapshot): + record = CanonicalActionDerivationReattemptService(adapter, repository).authorize_successor( + predecessor_attempt_id="failed", operator_context=context, rationale="bounded changed request") + self.assertEqual(record["predecessor_attempt_id"], "failed") + self.assertEqual(record["attempt_sequence"], 2) + self.assertNotEqual(record["provider_request_digest"], "old-request") + self.assertEqual(record["reattempt_reason"], REQUEST_SEMANTICS_CHANGED) + + def test_successor_rejects_unsupported_reason_and_untrusted_operator(self) -> None: + database = SimpleNamespace() + repository = SimpleNamespace(database=database, operators=SimpleNamespace(authorize=lambda _: False)) + adapter = SimpleNamespace(configuration=SimpleNamespace(policy_service=SimpleNamespace(db=database))) + service = CanonicalActionDerivationReattemptService(adapter, repository) + with self.assertRaises(PermissionError): + service.authorize_successor(predecessor_attempt_id="x", operator_context=OperatorContext("i", "u", 1), rationale="ok", reason="OTHER") + with self.assertRaises(PermissionError): + service.authorize_successor(predecessor_attempt_id="x", operator_context=OperatorContext("i", "u", 1), rationale="ok") + with self.assertRaises(ValueError): + service.authorize_successor(predecessor_attempt_id="x", operator_context=OperatorContext("i", "u", 1), rationale=" ") + + def test_mission_amendment_applies_append_only_constraint_and_rejects_stale_lineage(self) -> None: + contract = {"mission": {"engineering_constraints": ["WRITE_SCOPE=NONE"]}} + predecessor = _digest(contract) + amendment = {"predecessor_digest": predecessor, "changed_fields": {"engineering_constraints_add": ["NO_REPOSITORY_TARGET=TRUE"]}, + "effective_contract_digest": "sha256:" + "a" * 64} + database = SimpleNamespace(get_document=lambda *_: {"status": "APPROVED_PLANNABLE", "admission_contract": contract}, + _connection=_Connection([_Cursor(rows=[{"document": __import__("json").dumps(amendment)}])])) + service = MissionAmendmentService(SimpleNamespace(database=database)) + effective, digest = service.effective_contract("mission") + self.assertIn("NO_REPOSITORY_TARGET=TRUE", effective["mission"]["engineering_constraints"]) + self.assertEqual(digest, amendment["effective_contract_digest"]) + database._connection = _Connection([_Cursor(rows=[{"document": __import__("json").dumps({**amendment, "predecessor_digest": "wrong"})}])]) + with self.assertRaisesRegex(ValueError, "stale"): + service.effective_contract("mission") + + def test_mission_amendment_records_canonical_business_and_architecture_decisions(self) -> None: + context = OperatorContext("installation", "operator", 1) + contract = {"mission": {"engineering_constraints": ["WRITE_SCOPE=NONE"]}, "installation_id": "installation"} + database = SimpleNamespace(get_document=lambda *_: {"status": "APPROVED_PLANNABLE", "admission_contract": contract}, + _connection=_Connection([_Cursor(rows=[]), _Cursor([1])]), + create_mission_amendment=lambda document: document) + repository = SimpleNamespace(database=database, operators=SimpleNamespace(authorize=lambda value: value == context), + decision=lambda identifier: {"capability": "BUSINESS_APPROVAL" if identifier == "business" else "ARCHITECTURE_APPROVAL", + "decision": "approved", "subject_id": "mission"}) + record = MissionAmendmentService(repository).amend_no_repository_target("mission", business_decision_id="business", + architecture_decision_id="architecture", rationale="no repository delivery", context=context) + self.assertEqual(record["revision"], 1) + self.assertEqual(record["changed_fields"], {"engineering_constraints_add": ["NO_REPOSITORY_TARGET=TRUE"]}) + + def test_canonical_authority_reads_valid_persisted_scope_and_no_write_policy(self) -> None: + state = {"status": "APPROVED_PLANNABLE", "admission_contract": {"mission": {"scope": ["scope-b", "scope-a"]}, + "planning": {"write_scopes": ["NONE"], "human_gates": ["review"], "risk_inputs": ["risk"]}}} + authority = CanonicalTokenPreflightAuthority(SimpleNamespace(get_document=lambda *_: state)) + self.assertEqual(authority.approved_scopes_for("mission"), ("scope-a", "scope-b")) + self.assertEqual(authority.approved_derivation_policy_for("mission"), (("NONE",), ("review",), ("risk",))) + with self.assertRaises(ValueError): + TokenPreflightBoundary("", "evidence", "contract").values(policy_digest="policy", request_digest="request") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_operator_identity.py b/tests/test_operator_identity.py index 9b460de..8288dcc 100644 --- a/tests/test_operator_identity.py +++ b/tests/test_operator_identity.py @@ -1,7 +1,7 @@ import os,shutil,sqlite3,tempfile,unittest from pathlib import Path from forge.runtime.database import RuntimeDatabase,RuntimeIntegrityError -from forge.operator_identity import InstallationOperatorService,MacOSGeneratedUIDIdentityAdapter,NamedOperatorIdentity +from forge.operator_identity import InstallationOperatorService,MacOSGeneratedUIDIdentityAdapter,NamedOperatorIdentity,OperatorContext class T(unittest.TestCase): def test_trusted_binding_rejects_strings_wrong_and_revoked(self): with tempfile.TemporaryDirectory() as d: @@ -55,3 +55,22 @@ def test_first_bind_uses_only_the_trusted_clock(self): with self.assertRaises(TypeError): service.first_bind('1900-01-01T00:00:00Z') with patch('forge.operator_identity._timestamp',return_value='2042-02-03T04:05:06Z'): service.first_bind() row=db._connection.execute("SELECT occurred_at FROM installation_operator_audit WHERE operation='FIRST_BIND'").fetchone(); self.assertEqual(row[0],'2042-02-03T04:05:06Z'); db.close() + def test_identity_and_upgrade_entrypoints_fail_closed_without_their_trusted_preconditions(self): + class Failed: returncode=1; stdout='' + class Malformed: returncode=0; stdout='GeneratedUID: not-a-uuid' + class WrongLabel: returncode=0; stdout='Unexpected: 123E4567-E89B-42D3-A456-426614174000' + from unittest.mock import patch + with patch('forge.operator_identity.os.getuid',return_value=501),patch('forge.operator_identity.pwd.getpwuid',return_value=type('P',(),{'pw_name':'operator'})()): + with self.assertRaises(PermissionError): MacOSGeneratedUIDIdentityAdapter(runner=lambda *args,**kwargs:Failed()).resolve() + with self.assertRaises(PermissionError): MacOSGeneratedUIDIdentityAdapter(runner=lambda *args,**kwargs:WrongLabel()).resolve() + with self.assertRaises(PermissionError): MacOSGeneratedUIDIdentityAdapter(runner=lambda *args,**kwargs:Malformed()).resolve() + with tempfile.TemporaryDirectory() as d: + db=RuntimeDatabase(Path(d),path=Path(d)/'runtime.db'); service=InstallationOperatorService(db,lambda:NamedOperatorIdentity('generated-a',501)); context=service.first_bind() + with self.assertRaises(PermissionError): service._bootstrap_governance_after_first_bind() + with self.assertRaises(PermissionError): service.upgrade_legacy_governance_capabilities(context,decision_source='') + db.close() + def test_corrupt_adoption_provenance_is_not_recognized(self): + class Cursor: + def fetchall(self): return [{'capability':capability,'bootstrap_provenance':'not-json','digest':'sha256:'+'0'*64} for capability in ('ARCHITECTURE_APPROVAL','BUSINESS_APPROVAL','OWNER_PROGRAMME_AUTHORIZATION','SECURITY_APPROVAL')] + service=InstallationOperatorService(type('Db',(),{'_connection':type('Connection',(),{'execute':lambda *args:Cursor()})()})(),lambda:NamedOperatorIdentity('generated-a',501)) + self.assertFalse(service._valid_adoption_provenance(OperatorContext('installation','generated-a',1),{'created_at':'now','version':1})) From 6ec32148180d3a6548aee765a855f2b547a4678b Mon Sep 17 00:00:00 2001 From: pcvantol Date: Tue, 8 Sep 2026 07:53:36 +0200 Subject: [PATCH 11/14] docs: require API contract drift qualification --- .../FORGE_SERVER_DEPLOYMENT_TARGET.md | 18 +++++++++++++++++- 1 file changed, 17 insertions(+), 1 deletion(-) diff --git a/docs/architecture/FORGE_SERVER_DEPLOYMENT_TARGET.md b/docs/architecture/FORGE_SERVER_DEPLOYMENT_TARGET.md index 85b7678..cdd9a26 100644 --- a/docs/architecture/FORGE_SERVER_DEPLOYMENT_TARGET.md +++ b/docs/architecture/FORGE_SERVER_DEPLOYMENT_TARGET.md @@ -8,5 +8,21 @@ The current repository-bound `.git/forge-runtime` placement is historical/bootst Forge binds to EP and Workspace only through their versioned authenticated APIs. It never reads their databases or controls their services. The shared discovery/pairing vocabulary and descriptor are defined by Forge Platform's [contract](https://github.com/pcvantol/forge-platform/blob/main/docs/architecture/INSTANCE_DISCOVERY_AND_PAIRING_CONTRACT.md): LAN DNS-SD/mDNS and configured/unicast/tailnet endpoints locate candidates, while authenticated pairing pins peer product, stable instance ID, identity fingerprint, endpoint/version/capability set and scope in Forge-owned storage. Discovery is not authorization; a binding never silently retargets to a discovered instance. -The autonomy canary requires only installed Forge/EP storage, stable identities, the existing versioned authenticated Forge→EP HTTP seam, a configured/pinned EP binding and restart recovery. Workspace UI, LAN discovery and universal-installer completion are post-canary productization and must not block that proof. +## HTTP API implementation requirement + +Forge does not yet expose its own HTTP server. When that implementation is +introduced, its versioned OpenAPI document is the canonical public transport +contract and must describe every implemented Forge-owned route, method, +authentication requirement, request/response/error envelope and version +behavior. The implementation must ship an exhaustive Postman collection +derived from that contract, covering every documented operation and its +declared authorization and error cases without embedding credentials. +CI must run the collection against an isolated Forge server state and fail on +any drift between the OpenAPI document, the routes actually exposed by that +server, and the Postman collection. A route may not be added, removed or +semantically changed without updating all three artifacts in the same change. +This requirement applies only to a future Forge-owned HTTP API; it neither +redefines the EP API nor makes Forge an EP proxy. + +The autonomy canary requires only installed Forge/EP storage, stable identities, the existing versioned authenticated Forge→EP HTTP seam, a configured/pinned EP binding and restart recovery. Workspace UI, LAN discovery and universal-installer completion are post-canary productization and must not block that proof. From d289b581cda8ca234636b45e568faafa45e34cf0 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Tue, 8 Sep 2026 08:48:51 +0200 Subject: [PATCH 12/14] fix: fail closed on EP terminal evidence --- forge/models/execution_host.py | 6 + forge/operator_identity.py | 20 ++- forge/runtime/runner.py | 40 ++++-- forge/runtime/service.py | 24 +++- forge/scheduler/ep_v11.py | 182 ++++++++++++++++++++----- tests/test_bootstrap_mission_runner.py | 35 ++++- tests/test_ep_http_adapter.py | 33 +++++ tests/test_operator_identity.py | 15 ++ tests/test_runtime_service.py | 36 +++++ 9 files changed, 339 insertions(+), 52 deletions(-) diff --git a/forge/models/execution_host.py b/forge/models/execution_host.py index 1cc2b8d..6adfaac 100644 --- a/forge/models/execution_host.py +++ b/forge/models/execution_host.py @@ -164,6 +164,12 @@ def _default_producer_contract(self) -> ProducerContract: ("repository_id", self.repository_id), ("workspace_id", self.workspace_id), ) + # A rendered runtime prompt is the materialized Mission provenance. + # Carry its actual immutable revision through the Producer Contract; + # adapters must never supply a transport default such as "1". + mission_revision = getattr(self.runtime_prompt, "mission_revision", None) + if isinstance(mission_revision, str) and mission_revision: + metadata += (("mission_revision", mission_revision),) return ProducerContract( producer=Producer(prompt_producer), correlation_id=self.correlation_id, diff --git a/forge/operator_identity.py b/forge/operator_identity.py index 742a5f3..5b31851 100644 --- a/forge/operator_identity.py +++ b/forge/operator_identity.py @@ -84,17 +84,31 @@ def upgrade_legacy_governance_capabilities(self, context, *, decision_source): rows=self.db._connection.execute('SELECT * FROM governance_capability_grants WHERE installation_id=? AND operator_id=? ORDER BY capability',(context.installation_id,operator)).fetchall() authorities=tuple(r['capability'] for r in self.db._connection.execute('SELECT capability FROM governance_authority WHERE installation_id=? AND operator_id=? ORDER BY capability',(context.installation_id,operator))) expected=tuple(sorted((*legacy,'OWNER_PROGRAMME_AUTHORIZATION'))) + legacy_rows=tuple(r for r in rows if r['capability'] in legacy) + def valid_legacy_grants(): + if len(legacy_rows)!=3 or tuple(r['capability'] for r in legacy_rows)!=legacy:return False + for row in legacy_rows: + try: document=json.loads(row['bootstrap_provenance']) + except (TypeError,ValueError): return False + digest='sha256:'+hashlib.sha256(json.dumps(document,sort_keys=True,separators=(',',':')).encode()).hexdigest() + if digest!=row['digest'] or document.get('capability')!=row['capability'] or document.get('installation_id')!=context.installation_id or document.get('operator_id')!=operator:return False + kind=document.get('kind') + if kind=='LOCAL_INSTALLATION_BOOTSTRAP_V1': + if document.get('binding_version')!=binding['version'] or not isinstance(document.get('first_bind_at'),str):return False + elif kind=='EXISTING_G001_GOVERNANCE_ADOPTION_V1': + if document.get('prior_binding_version')!=binding['version'] or not isinstance(document.get('prior_binding_created_at'),str):return False + else:return False + return True if authorities==expected: if len(rows)!=4 or tuple(r['capability'] for r in rows)!=expected: raise PermissionError('upgraded capability grants are incomplete') + if not valid_legacy_grants(): raise PermissionError('legacy capability provenance is invalid') upgraded=next(r for r in rows if r['capability']=='OWNER_PROGRAMME_AUTHORIZATION') document=json.loads(upgraded['bootstrap_provenance']); digest='sha256:'+hashlib.sha256(json.dumps(document,sort_keys=True,separators=(',',':')).encode()).hexdigest() old={r['grant_id']:r['digest'] for r in rows if r['capability']!='OWNER_PROGRAMME_AUTHORIZATION'} if digest!=upgraded['digest'] or document.get('kind')!='G001_OWNER_PROGRAMME_CAPABILITY_UPGRADE_V1' or dict(document.get('legacy_grants',()))!=old or document.get('binding_version')!=binding['version']: raise PermissionError('upgraded capability provenance is invalid') return if authorities!=legacy or tuple(r['capability'] for r in rows)!=legacy: raise PermissionError('state is not the recognized legacy capability set') - for row in rows: - document=json.loads(row['bootstrap_provenance']); digest='sha256:'+hashlib.sha256(json.dumps(document,sort_keys=True,separators=(',',':')).encode()).hexdigest() - if digest!=row['digest'] or document.get('capability')!=row['capability'] or document.get('installation_id')!=context.installation_id or document.get('operator_id')!=operator: raise PermissionError('legacy capability provenance is invalid') + if not valid_legacy_grants(): raise PermissionError('legacy capability provenance is invalid') now=_timestamp(); provenance={'kind':'G001_OWNER_PROGRAMME_CAPABILITY_UPGRADE_V1','upgrade_version':'1','decision_source':decision_source,'legacy_grants':[(r['grant_id'],r['digest']) for r in rows],'binding_version':binding['version'],'added_capability':'OWNER_PROGRAMME_AUTHORIZATION','occurred_at':now,'installation_id':context.installation_id,'operator_id':operator,'capability':'OWNER_PROGRAMME_AUTHORIZATION'}; encoded=json.dumps(provenance,sort_keys=True,separators=(',',':')); digest='sha256:'+hashlib.sha256(encoded.encode()).hexdigest() self.db._insert_governance_grant(digest,context.installation_id,operator,'OWNER_PROGRAMME_AUTHORIZATION',encoded,digest,now); self.db._insert_governance_authority(context.installation_id,operator,'OWNER_PROGRAMME_AUTHORIZATION',now); self._audit(context.installation_id,NamedOperatorIdentity(context.generated_uid,0),'LEGACY_CAPABILITY_UPGRADE',now,'ALLOW') def revoke(self, context): diff --git a/forge/runtime/runner.py b/forge/runtime/runner.py index 17ab213..faf5256 100644 --- a/forge/runtime/runner.py +++ b/forge/runtime/runner.py @@ -144,7 +144,11 @@ def _prompt(document: Mapping[str, Any]) -> RuntimePrompt | CodexCliRuntimePromp def _request(document: Mapping[str, Any]) -> ExecutionRequest: contract_document = document.get("producer_contract") contract = None - if contract_document is not None: + # Absence is an explicit historical compatibility route. Presence is a + # claim that a canonical Producer Contract was persisted, so null or any + # malformed representation must fail closed rather than silently becoming + # a freshly synthesized default contract. + if "producer_contract" in document: if not isinstance(contract_document, Mapping): raise MissionRunnerError("persisted Producer Contract is malformed") try: @@ -154,21 +158,35 @@ def _request(document: Mapping[str, Any]) -> ExecutionRequest: constraints = contract_document["execution_constraints"] receipts = contract_document.get("receipt_references", ()) evidence_references = contract_document.get("execution_evidence_references", ()) - if not isinstance(producer, Mapping) or not isinstance(prompt, Mapping) or not isinstance(metadata, Mapping): + if not isinstance(producer, Mapping) or not isinstance(prompt, Mapping) or not isinstance(metadata, Mapping) or not isinstance(constraints, list): raise TypeError identity = producer["identity"] if not isinstance(identity, Mapping) or not isinstance(receipts, list) or not isinstance(evidence_references, list): raise TypeError + def required(value: Any) -> str: + if not isinstance(value, str) or not value: + raise TypeError + return value + if any(not isinstance(key, str) or not key or not isinstance(value, str) or not value for key, value in metadata.items()): + raise TypeError + if any(not isinstance(item, str) or not item for item in constraints): + raise TypeError + if any(not isinstance(item, str) or not item for item in evidence_references): + raise TypeError + if any(not isinstance(item, Mapping) or set(item) != {"host_id", "receipt_id"} + or not isinstance(item["host_id"], str) or not item["host_id"] + or not isinstance(item["receipt_id"], str) or not item["receipt_id"] for item in receipts): + raise TypeError + mission_id = required(contract_document["mission_id"]) contract = ProducerContract( - Producer(ProducerIdentity(str(identity["id"]), str(identity["type"]), str(identity["version"])), - str(producer["contract_version"])), - str(contract_document["correlation_id"]), str(contract_document["engineering_action_id"]), - RuntimePromptEnvelope(str(prompt["id"]), str(prompt["version"]), str(prompt["format"]), str(prompt["content"]), str(prompt["content_digest"])), - tuple(str(item) for item in constraints), tuple((str(k), str(v)) for k, v in metadata.items()), - mission_id=contract_document.get("mission_id"), - receipt_references=tuple(ExecutionReceiptReference(str(item["host_id"]), str(item["receipt_id"])) for item in receipts), - execution_evidence_references=tuple(str(item) for item in evidence_references), - contract_version=str(contract_document["contract_version"]), + Producer(ProducerIdentity(required(identity["id"]), required(identity["type"]), required(identity["version"])), + required(producer["contract_version"])), + required(contract_document["correlation_id"]), required(contract_document["engineering_action_id"]), + RuntimePromptEnvelope(required(prompt["id"]), required(prompt["version"]), required(prompt["format"]), + required(prompt["content"]), required(prompt["content_digest"])), + tuple(constraints), tuple(metadata.items()), mission_id=mission_id, + receipt_references=tuple(ExecutionReceiptReference(item["host_id"], item["receipt_id"]) for item in receipts), + execution_evidence_references=tuple(evidence_references), contract_version=required(contract_document["contract_version"]), ) except (KeyError, TypeError, ValueError) as error: raise MissionRunnerError("persisted Producer Contract is malformed") from error diff --git a/forge/runtime/service.py b/forge/runtime/service.py index 3e41826..13ce622 100644 --- a/forge/runtime/service.py +++ b/forge/runtime/service.py @@ -4,7 +4,7 @@ from contextlib import contextmanager from dataclasses import dataclass from pathlib import Path -from time import sleep +from threading import Event from typing import Callable, Iterator, Protocol try: # macOS/Linux runtime product path @@ -67,18 +67,28 @@ class ForgeRuntimeService: """ def __init__(self, loop: RuntimeLoop, states: MissionStateStore, *, runtime_database_path: Path | str, - wait: Callable[[float], None] = sleep, minimum_backoff: float = 0.25, + wait: Callable[[float], None] | None = None, minimum_backoff: float = 0.25, maximum_backoff: float = 5.0) -> None: if minimum_backoff <= 0 or maximum_backoff < minimum_backoff: raise ValueError("runtime service backoff bounds are invalid") self._loop, self._states = loop, states self._lock = RuntimeServiceLock(runtime_database_path) self._wait, self._minimum_backoff, self._maximum_backoff = wait, minimum_backoff, maximum_backoff + self._wake, self._stopped = Event(), Event() @property def mutation_lock(self) -> RuntimeServiceLock: return self._lock + def wake(self) -> None: + """Wake a default wait early when new work, pause, or stop arrives.""" + self._wake.set() + + def stop(self) -> None: + """Request a controlled stop without waiting for the current backoff.""" + self._stopped.set() + self.wake() + def tick(self) -> RuntimeServiceTick: with self._lock.acquire(): resumable = {MissionExecutionStatus.READY, MissionExecutionStatus.ACTIVE, @@ -94,7 +104,7 @@ def tick(self) -> RuntimeServiceTick: def serve(self, *, keep_running: Callable[[], bool]) -> None: """Run with bounded interruptible backoff; caller controls pause/stop.""" delay = self._minimum_backoff - while keep_running(): + while keep_running() and not self._stopped.is_set(): try: tick = self.tick() except RuntimeServiceBusy: @@ -104,7 +114,11 @@ def serve(self, *, keep_running: Callable[[], bool]) -> None: continue # Do not sleep through a requested pause/shutdown. The bounded # delay also prevents a WAITING_FOR_EVIDENCE poll from busy-looping. - if not keep_running(): + if not keep_running() or self._stopped.is_set(): break - self._wait(delay) + if self._wait is None: + self._wake.wait(delay) + self._wake.clear() + else: + self._wait(delay) delay = min(self._maximum_backoff, delay * 2) diff --git a/forge/scheduler/ep_v11.py b/forge/scheduler/ep_v11.py index 74e4fa6..329242e 100644 --- a/forge/scheduler/ep_v11.py +++ b/forge/scheduler/ep_v11.py @@ -1,33 +1,151 @@ -"""Strict Forge consumer mapping for the EP producer readback v1.1 fixture.""" +"""Fail-closed Forge consumer mapping for EP producer readback v1.1.""" from __future__ import annotations -import hashlib,json -from typing import Any,Mapping -from forge.models.execution_host import ExecutionEvidenceOutcome,ExecutionHostEvidence,ExecutionRepositoryEvidence - -def terminal_evidence(readback: Mapping[str,Any], artifact: bytes, *, host_id: str) -> ExecutionHostEvidence: - if readback.get('contract_version')!='1.1': raise ValueError('unsupported EP readback contract') - evidence=readback.get('evidence',{}); terminal=evidence.get('terminal_artifact'); run=readback.get('run'); result=readback.get('result',{}); provenance=readback.get('provenance',{}).get('forge_execution') - if not isinstance(terminal,dict) or not isinstance(run,dict) or not isinstance(provenance,dict) or result.get('terminal') is not True: raise ValueError('EP terminal evidence is incomplete') - digest='sha256:'+hashlib.sha256(artifact).hexdigest() - if terminal.get('digest')!=digest: raise ValueError('EP terminal artifact digest mismatch') - try: document=json.loads(artifact) - except (UnicodeDecodeError,json.JSONDecodeError) as e: raise ValueError('EP terminal artifact is invalid JSON') from e - repo=document.get('repository',{}); outcome=ExecutionEvidenceOutcome(str(result.get('outcome')).lower()) - if not isinstance(repo,dict) or not isinstance(document.get('report',{}).get('id'),str): raise ValueError('EP terminal artifact lacks repository evidence') - correlation=readback['correlation']; prompt=provenance['runtime_prompt']; submission=readback['submission'] - # A hash only proves the returned bytes came from *some* registered object. - # Bind their internal subject to the readback before mapping it to Forge. - artifact_correlation=document.get('correlation',{}); artifact_provenance=document.get('provenance',{}); artifact_run=document.get('run',{}); artifact_submission=document.get('submission',{}); artifact_producer=document.get('producer',{}) - if not all(isinstance(value,dict) for value in (artifact_correlation,artifact_provenance,artifact_run,artifact_submission,artifact_producer)): - raise ValueError('EP terminal artifact lacks binding fields') - if artifact_correlation != correlation or artifact_producer != readback.get('producer'): - raise ValueError('EP terminal artifact identity differs from readback') - expected_provenance={key: provenance.get(key) for key in ('action_id','contract_version','correlation_id','host_id','intent_id','intent_revision','mission_id','mission_revision','repository_id','retry_of_correlation_id','runtime_prompt')} - if artifact_provenance != expected_provenance or artifact_run.get('id') != run.get('id') or artifact_submission.get('id') != submission.get('id') or artifact_submission.get('project_id') != submission.get('project_id') or artifact_submission.get('repository_id') != submission.get('repository_id') or artifact_submission.get('accepted_request_digest') != submission.get('accepted_request_digest'): - raise ValueError('EP terminal artifact does not bind readback run and submission') - revision=repo.get('revision') - if revision is not None and not isinstance(revision,str): raise ValueError('EP terminal artifact revision is invalid') - if outcome is ExecutionEvidenceOutcome.COMPLETE and (not revision or document.get('run',{}).get('delivery_qualified') is not True or repo.get('revision_required') is not True): raise ValueError('EP complete terminal artifact lacks qualified delivery revision') - if repo.get('id') != readback.get('evidence',{}).get('repository',{}).get('id') or revision != readback.get('evidence',{}).get('repository',{}).get('revision'): raise ValueError('EP terminal artifact repository differs from readback') - repository=ExecutionRepositoryEvidence(str(correlation['mission_id']),str(provenance['intent_id']),str(provenance['intent_revision']),str(correlation['engineering_action_id']),str(prompt['id']),str(correlation['correlation_id']),str(run['id']),str(repo['id']),revision,str(document['report']['id']),digest) - return ExecutionHostEvidence(host_id,str(correlation['correlation_id']),str(run['id']),str(document['report']['id']),outcome,repository,validation_references=tuple(str(x.get('command')) for x in document.get('references',{}).get('validation',()) if x.get('command'))) + +import hashlib +import json +from typing import Any, Mapping + +from forge.models.execution_host import ( + ExecutionEvidenceOutcome, + ExecutionHostEvidence, + ExecutionRepositoryEvidence, +) + + +_OUTCOMES = frozenset(item.value.upper() for item in ExecutionEvidenceOutcome) + + +def _object(value: Any, name: str) -> Mapping[str, Any]: + if not isinstance(value, Mapping): + raise ValueError(f"EP terminal evidence {name} must be an object") + return value + + +def _string(value: Any, name: str) -> str: + if not isinstance(value, str) or not value: + raise ValueError(f"EP terminal evidence {name} must be a non-empty string") + return value + + +def _sha256(value: Any, name: str) -> str: + value = _string(value, name) + if not value.startswith("sha256:") or len(value.removeprefix("sha256:")) != 64: + raise ValueError(f"EP terminal evidence {name} must be a SHA-256 identity") + try: + int(value.removeprefix("sha256:"), 16) + except ValueError as error: + raise ValueError(f"EP terminal evidence {name} must be a SHA-256 identity") from error + return value + + +def _terminal_outcome(value: Any, name: str) -> str: + value = _string(value, name).upper() + if value not in _OUTCOMES: + raise ValueError(f"EP terminal evidence {name} is not a terminal outcome") + return value + + +def terminal_evidence(readback: Mapping[str, Any], artifact: bytes, *, host_id: str) -> ExecutionHostEvidence: + """Map one immutable EP terminal artifact only when every identity agrees. + + Artifact digest integrity establishes byte origin, not semantic identity; + all readback/artifact fields used below must therefore be explicitly + present, correctly typed, and mutually consistent. + """ + if readback.get("contract_version") != "1.1": + raise ValueError("unsupported EP readback contract") + evidence = _object(readback.get("evidence"), "readback evidence") + terminal = _object(evidence.get("terminal_artifact"), "terminal artifact reference") + run = _object(readback.get("run"), "readback run") + result = _object(readback.get("result"), "readback result") + provenance_root = _object(readback.get("provenance"), "readback provenance") + provenance = _object(provenance_root.get("forge_execution"), "Forge provenance") + correlation = _object(readback.get("correlation"), "readback correlation") + submission = _object(readback.get("submission"), "readback submission") + producer = _object(readback.get("producer"), "readback producer") + repository_readback = _object(evidence.get("repository"), "readback repository") + if result.get("terminal") is not True or run.get("terminal") is not True: + raise ValueError("EP terminal evidence terminal flags are incomplete or contradictory") + + artifact_digest = "sha256:" + hashlib.sha256(artifact).hexdigest() + if _sha256(terminal.get("digest"), "terminal artifact digest") != artifact_digest: + raise ValueError("EP terminal artifact digest mismatch") + try: + document = json.loads(artifact) + except (UnicodeDecodeError, json.JSONDecodeError) as error: + raise ValueError("EP terminal artifact is invalid JSON") from error + document = _object(document, "artifact document") + if document.get("artifact_type") != "EP_TERMINAL_EVIDENCE" or document.get("contract_version") != "1.1": + raise ValueError("unsupported EP terminal artifact contract") + + artifact_correlation = _object(document.get("correlation"), "artifact correlation") + artifact_provenance = _object(document.get("provenance"), "artifact provenance") + artifact_run = _object(document.get("run"), "artifact run") + artifact_report = _object(document.get("report"), "artifact report") + artifact_submission = _object(document.get("submission"), "artifact submission") + artifact_producer = _object(document.get("producer"), "artifact producer") + repository = _object(document.get("repository"), "artifact repository") + + # Required, immutable accepted-request identities cannot be satisfied by + # the accidental equality of two missing values. + accepted_readback = _sha256(submission.get("accepted_request_digest"), "readback accepted request digest") + accepted_artifact = _sha256(artifact_submission.get("accepted_request_digest"), "artifact accepted request digest") + if accepted_readback != accepted_artifact: + raise ValueError("EP terminal artifact accepted request digest differs from readback") + if artifact_correlation != correlation or artifact_producer != producer: + raise ValueError("EP terminal artifact identity differs from readback") + + expected_provenance = {key: provenance.get(key) for key in ( + "action_id", "contract_version", "correlation_id", "host_id", "intent_id", "intent_revision", + "mission_id", "mission_revision", "repository_id", "retry_of_correlation_id", "runtime_prompt", + )} + if artifact_provenance != expected_provenance: + raise ValueError("EP terminal artifact provenance differs from readback") + if artifact_run.get("id") != run.get("id"): + raise ValueError("EP terminal artifact run differs from readback") + if (artifact_submission.get("id"), artifact_submission.get("project_id"), artifact_submission.get("repository_id")) != ( + submission.get("id"), submission.get("project_id"), submission.get("repository_id"), + ): + raise ValueError("EP terminal artifact submission differs from readback") + + outcome = _terminal_outcome(result.get("outcome"), "readback result outcome") + if (_terminal_outcome(artifact_run.get("outcome"), "artifact run outcome"), + _terminal_outcome(artifact_report.get("terminal_state"), "artifact report terminal state"), + _terminal_outcome(run.get("state"), "readback run state")) != (outcome, outcome, outcome): + raise ValueError("EP terminal outcome fields contradict one another") + qualified_readback, qualified_artifact = result.get("delivery_qualified"), artifact_run.get("delivery_qualified") + if not isinstance(qualified_readback, bool) or not isinstance(qualified_artifact, bool) or qualified_readback != qualified_artifact: + raise ValueError("EP terminal delivery qualification fields contradict one another") + + revision = repository.get("revision") + if revision is not None and not isinstance(revision, str): + raise ValueError("EP terminal artifact revision is invalid") + if (repository.get("id"), revision) != (repository_readback.get("id"), repository_readback.get("revision")): + raise ValueError("EP terminal artifact repository differs from readback") + if outcome == "COMPLETE": + if not qualified_readback or not revision or repository.get("revision_required") is not True: + raise ValueError("EP complete terminal artifact lacks qualified delivery revision") + elif qualified_readback: + # A terminal failed/blocked run is never a Forge successful delivery. + raise ValueError("EP non-complete terminal evidence cannot be delivery-qualified") + + prompt = _object(provenance.get("runtime_prompt"), "runtime prompt") + report_id = _string(artifact_report.get("id"), "artifact report id") + repository_evidence = ExecutionRepositoryEvidence( + _string(correlation.get("mission_id"), "mission id"), _string(provenance.get("intent_id"), "intent id"), + _string(provenance.get("intent_revision"), "intent revision"), _string(correlation.get("engineering_action_id"), "action id"), + _string(prompt.get("id"), "runtime prompt id"), _string(correlation.get("correlation_id"), "correlation id"), + _string(run.get("id"), "run id"), _string(repository.get("id"), "repository id"), revision, + report_id, artifact_digest, + ) + references = _object(document.get("references"), "artifact references") + validation = references.get("validation", ()) + if not isinstance(validation, list): + raise ValueError("EP terminal validation references are invalid") + validation_references = tuple( + item["command"] for item in validation + if isinstance(item, Mapping) and isinstance(item.get("command"), str) and item["command"] + ) + return ExecutionHostEvidence(host_id, repository_evidence.correlation_id, repository_evidence.host_run_id, + report_id, ExecutionEvidenceOutcome(outcome.lower()), repository_evidence, + validation_references=validation_references) diff --git a/tests/test_bootstrap_mission_runner.py b/tests/test_bootstrap_mission_runner.py index 063fec2..52bec3f 100644 --- a/tests/test_bootstrap_mission_runner.py +++ b/tests/test_bootstrap_mission_runner.py @@ -12,6 +12,7 @@ ExecutionHostTemporaryUnavailable, IntentApproval, IntentCategory, IntentReference, IntentStatus, IntentTraceability, ProviderPromptDefinition, RuntimePrompt, RuntimePromptSection, RuntimePromptSectionKind, + Producer, ProducerContract, ProducerIdentity, RuntimePromptEnvelope, ExecutionReceiptReference, ) from forge.models.mission import EngineeringMission, MissionIntentMembership, MissionScope from forge.models.codex_runtime_prompt import CodexCliRuntimePromptRequest, ExecutionHostCompatibility, RepositoryState @@ -195,7 +196,7 @@ def test_persisted_codex_prompt_restores_its_type_and_retry_lineage(self) -> Non ) active = EngineeringAction(1, "one", "one", "1", "Run one.", ("repository evidence",), status=EngineeringActionStatus.ACTIVE) rendered = CodexCliRuntimePromptRenderer().render(CodexCliRuntimePromptRequest( - EngineeringMission("mission-1", "1", "Mission", "Complete actions.", MissionScope(("runner",), ("planner",)), (MissionIntentMembership(1, "one", "1"),)), + EngineeringMission("mission-1", "7", "Mission", "Complete actions.", MissionScope(("runner",), ("planner",)), (MissionIntentMembership(1, "one", "1"),)), approved, active, RepositoryState("forge", "abc", "sha256:" + "a" * 64, "now"), ("bounded",), ("test",), ExecutionHostCompatibility("2.4", "GENESIS", ("codex_cli",), "platform>=1.5"), )) @@ -205,6 +206,38 @@ def test_persisted_codex_prompt_restores_its_type_and_retry_lineage(self) -> Non restored = _request(_request_document(request)) self.assertIsInstance(restored.runtime_prompt, type(rendered)) self.assertEqual(restored.original_correlation_id, "retry-1") + self.assertEqual(dict(restored.producer_contract.execution_metadata)["mission_revision"], "7") + + def test_present_null_contract_fails_closed_but_absent_legacy_contract_is_supported(self) -> None: + from forge.runtime.runner import _request, _request_document + request = ExecutionRequest("host", "mission-1", "one", "1", "one", prompt_factory({}, action(1, "one")), + "workspace", "forge", "correlation", "now") + document = _request_document(request) + legacy = dict(document); legacy.pop("producer_contract") + self.assertEqual(_request(legacy).producer_contract.correlation_id, "correlation") + document["producer_contract"] = None + with self.assertRaisesRegex(MissionRunnerError, "Producer Contract"): + _request(document) + + def test_nondefault_producer_contract_roundtrips_losslessly(self) -> None: + from forge.runtime.runner import _request, _request_document + prompt = prompt_factory({}, action(1, "one")) + contract = ProducerContract( + Producer(ProducerIdentity("producer", "FORGE", "1.0"), "1.0"), "correlation", "one", + RuntimePromptEnvelope(prompt.id, "1.0", "text/markdown", "non-default prompt", "sha256:" + "c" * 64), + ("constraint-a", "constraint-b"), + (("custom", "metadata"), ("intent_id", "one"), ("intent_revision", "1"), ("mission_revision", "7")), + mission_id="mission-1", receipt_references=(ExecutionReceiptReference("host", "receipt-a"),), + execution_evidence_references=("evidence-a",), contract_version="1.0", + ) + original = ExecutionRequest("host", "mission-1", "one", "1", "one", prompt, "workspace", "forge", + "correlation", "now", producer_contract=contract) + restored = _request(_request_document(original)) + self.assertEqual(restored.producer_contract.to_dict(), original.producer_contract.to_dict()) + self.assertEqual(restored.producer_contract.digest(), original.producer_contract.digest()) + malformed = _request_document(original); malformed["producer_contract"]["producer"]["identity"]["id"] = None + with self.assertRaisesRegex(MissionRunnerError, "Producer Contract"): + _request(malformed) if __name__ == "__main__": diff --git a/tests/test_ep_http_adapter.py b/tests/test_ep_http_adapter.py index b7aff49..60c51d7 100644 --- a/tests/test_ep_http_adapter.py +++ b/tests/test_ep_http_adapter.py @@ -117,6 +117,7 @@ def test_failed_terminal_evidence_without_delivery_revision_is_preserved(self) - raw = json.dumps(artifact, sort_keys=True, separators=(",", ":")).encode() + b"\n" readback = json.loads(json.dumps(self.readback)) readback["result"].update({"outcome": "FAILED", "delivery_qualified": False}) + readback["run"].update({"state": "FAILED"}) readback["evidence"]["repository"]["revision"] = None readback["evidence"]["terminal_artifact"]["digest"] = "sha256:" + hashlib.sha256(raw).hexdigest() with patch("forge.scheduler.ep_http_adapter.urlopen", self._urlopen([json.dumps(readback).encode(), raw], [])): @@ -124,6 +125,38 @@ def test_failed_terminal_evidence_without_delivery_revision_is_preserved(self) - self.assertEqual(evidence.outcome, ExecutionEvidenceOutcome.FAILED) self.assertIsNone(evidence.repository_evidence.repository_revision) + def _terminal_retrieval(self, readback: dict, artifact: dict): + self.database.save_execution_host_binding(self.request.correlation_id, + {"correlation_id": self.request.correlation_id, "submission_id": "submission-fixture", "host_run_id": "run-fixture"}) + raw = json.dumps(artifact, sort_keys=True, separators=(",", ":")).encode() + b"\n" + readback["evidence"]["terminal_artifact"]["digest"] = "sha256:" + hashlib.sha256(raw).hexdigest() + with patch("forge.scheduler.ep_http_adapter.urlopen", self._urlopen([json.dumps(readback).encode(), raw], [])): + return EngineeringPlatformHttpExecutionHost(self.config, self.database).retrieve_evidence(ExecutionDispatch(self.request, "run-fixture")) + + def test_terminal_outcome_qualification_digest_and_flags_must_have_parity(self) -> None: + cases = ( + ("artifact outcome", lambda r, a: a["run"].update({"outcome": "FAILED"})), + ("report state", lambda r, a: a["report"].update({"terminal_state": "FAILED"})), + ("qualification", lambda r, a: r["result"].update({"delivery_qualified": False})), + ("missing accepted digest", lambda r, a: (r["submission"].pop("accepted_request_digest"), a["submission"].pop("accepted_request_digest"))), + ("invalid accepted digest type", lambda r, a: (r["submission"].update({"accepted_request_digest": []}), a["submission"].update({"accepted_request_digest": []}))), + ("terminal flag", lambda r, a: r["run"].update({"terminal": False})), + ) + for label, mutate in cases: + with self.subTest(label=label): + readback, artifact = json.loads(json.dumps(self.readback)), json.loads(self.artifact) + mutate(readback, artifact) + with self.assertRaises(ValueError): + self._terminal_retrieval(readback, artifact) + self.database._connection.execute("DELETE FROM execution_host_bindings") + + def test_complete_requires_qualified_delivery_revision(self) -> None: + readback, artifact = json.loads(json.dumps(self.readback)), json.loads(self.artifact) + readback["evidence"]["repository"]["revision"] = None + artifact["repository"].update({"revision": None, "revision_required": False}) + with self.assertRaisesRegex(ValueError, "qualified delivery revision"): + self._terminal_retrieval(readback, artifact) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_operator_identity.py b/tests/test_operator_identity.py index 8288dcc..902c953 100644 --- a/tests/test_operator_identity.py +++ b/tests/test_operator_identity.py @@ -74,3 +74,18 @@ class Cursor: def fetchall(self): return [{'capability':capability,'bootstrap_provenance':'not-json','digest':'sha256:'+'0'*64} for capability in ('ARCHITECTURE_APPROVAL','BUSINESS_APPROVAL','OWNER_PROGRAMME_AUTHORIZATION','SECURITY_APPROVAL')] service=InstallationOperatorService(type('Db',(),{'_connection':type('Connection',(),{'execute':lambda *args:Cursor()})()})(),lambda:NamedOperatorIdentity('generated-a',501)) self.assertFalse(service._valid_adoption_provenance(OperatorContext('installation','generated-a',1),{'created_at':'now','version':1})) + def test_legacy_upgrade_reopens_idempotently_and_rejects_tampered_legacy_provenance(self): + with tempfile.TemporaryDirectory() as d: + root=Path(d); path=root/'runtime.db'; identity=NamedOperatorIdentity('generated-a',501); db=RuntimeDatabase(root,path=path); service=InstallationOperatorService(db,lambda:identity); context=service.first_bind() + # Isolated historic-state fixture: remove only the programme capability to + # represent the recognized three-capability predecessor product. + db._connection.executescript('DROP TRIGGER governance_authority_immutable_delete; DROP TRIGGER governance_capability_grants_immutable_delete;') + db._connection.execute("DELETE FROM governance_authority WHERE capability='OWNER_PROGRAMME_AUTHORIZATION'") + db._connection.execute("DELETE FROM governance_capability_grants WHERE capability='OWNER_PROGRAMME_AUTHORIZATION'"); db._connection.commit() + service.upgrade_legacy_governance_capabilities(context,decision_source='owner-decision') + db.close(); db=RuntimeDatabase(root,path=path); service=InstallationOperatorService(db,lambda:identity); service.upgrade_legacy_governance_capabilities(context,decision_source='owner-decision') + db._connection.execute('DROP TRIGGER governance_capability_grants_immutable_update') + db._connection.execute("UPDATE governance_capability_grants SET bootstrap_provenance='{}' WHERE capability='ARCHITECTURE_APPROVAL'"); db._connection.commit() + with self.assertRaisesRegex(PermissionError,'legacy capability provenance'): + service.upgrade_legacy_governance_capabilities(context,decision_source='owner-decision') + db.close() diff --git a/tests/test_runtime_service.py b/tests/test_runtime_service.py index 20e2b23..f3a21c7 100644 --- a/tests/test_runtime_service.py +++ b/tests/test_runtime_service.py @@ -3,6 +3,10 @@ from dataclasses import dataclass from pathlib import Path from tempfile import TemporaryDirectory +from threading import Thread +from time import monotonic, sleep +import subprocess +import sys import unittest from forge.runtime.service import ForgeRuntimeService, RuntimeServiceBusy, RuntimeServiceLock @@ -49,6 +53,38 @@ def test_service_and_mutating_cli_share_one_runtime_lease(self) -> None: with RuntimeServiceLock(path).acquire(): pass + def test_stop_wakes_default_backoff_without_waiting_for_its_cap(self) -> None: + with TemporaryDirectory() as root: + state = _State("mission", MissionExecutionStatus.WAITING_FOR_EVIDENCE, 4) + service = ForgeRuntimeService(_Loop(state, progresses=False), _States(state), + runtime_database_path=Path(root) / "runtime.db", minimum_backoff=5, maximum_backoff=5) + thread = Thread(target=lambda: service.serve(keep_running=lambda: True)) + started = monotonic(); thread.start(); sleep(0.05); service.stop(); thread.join(timeout=1) + self.assertFalse(thread.is_alive()) + self.assertLess(monotonic() - started, 1) + + def test_separate_process_mutator_cannot_share_runtime_lease_and_exit_releases_it(self) -> None: + with TemporaryDirectory() as root: + path, ready = Path(root) / "runtime.db", Path(root) / "ready" + program = ("from pathlib import Path; from forge.runtime.service import RuntimeServiceLock; " + "import sys,time; p=Path(sys.argv[1]); ready=Path(sys.argv[2]); " + "\nwith RuntimeServiceLock(p).acquire():\n ready.touch(); time.sleep(30)") + child = subprocess.Popen([sys.executable, "-c", program, str(path), str(ready)]) + try: + for _ in range(100): + if ready.exists(): break + sleep(0.01) + self.assertTrue(ready.exists()) + with self.assertRaises(RuntimeServiceBusy): + with RuntimeServiceLock(path).acquire(): + pass + finally: + child.kill(); child.wait(timeout=2) + # flock releases on abrupt process termination; a later mutator + # can claim the same canonical runtime without stale ownership. + with RuntimeServiceLock(path).acquire(): + pass + if __name__ == "__main__": unittest.main() From 6a26c5eabacf61004ece390e0c35f51d2730c598 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Tue, 8 Sep 2026 08:51:16 +0200 Subject: [PATCH 13/14] fix: bind runtime service lock to canonical database --- forge/runtime/service.py | 7 +++++-- tests/test_runtime_service.py | 9 +++++++-- 2 files changed, 12 insertions(+), 4 deletions(-) diff --git a/forge/runtime/service.py b/forge/runtime/service.py index 13ce622..8de2b98 100644 --- a/forge/runtime/service.py +++ b/forge/runtime/service.py @@ -66,13 +66,16 @@ class ForgeRuntimeService: work, preserving single-flight dispatch semantics across restarts. """ - def __init__(self, loop: RuntimeLoop, states: MissionStateStore, *, runtime_database_path: Path | str, + def __init__(self, loop: RuntimeLoop, states: MissionStateStore, *, runtime_database, wait: Callable[[float], None] | None = None, minimum_backoff: float = 0.25, maximum_backoff: float = 5.0) -> None: if minimum_backoff <= 0 or maximum_backoff < minimum_backoff: raise ValueError("runtime service backoff bounds are invalid") + runtime_path = getattr(runtime_database, "path", None) + if not isinstance(runtime_path, Path): + raise ValueError("runtime service requires an opened canonical RuntimeDatabase") self._loop, self._states = loop, states - self._lock = RuntimeServiceLock(runtime_database_path) + self._lock = RuntimeServiceLock(runtime_path) self._wait, self._minimum_backoff, self._maximum_backoff = wait, minimum_backoff, maximum_backoff self._wake, self._stopped = Event(), Event() diff --git a/tests/test_runtime_service.py b/tests/test_runtime_service.py index f3a21c7..e010db1 100644 --- a/tests/test_runtime_service.py +++ b/tests/test_runtime_service.py @@ -10,6 +10,7 @@ import unittest from forge.runtime.service import ForgeRuntimeService, RuntimeServiceBusy, RuntimeServiceLock +from forge.runtime.database import RuntimeDatabase from forge.state import MissionExecutionStatus @@ -38,11 +39,13 @@ def test_waiting_evidence_uses_bounded_interruptible_backoff(self) -> None: with TemporaryDirectory() as root: state = _State("mission", MissionExecutionStatus.WAITING_FOR_EVIDENCE, 4) waits: list[float] = [] + database = RuntimeDatabase(".", path=Path(root) / "runtime.db") service = ForgeRuntimeService(_Loop(state, progresses=False), _States(state), - runtime_database_path=Path(root) / "runtime.db", wait=waits.append, minimum_backoff=0.1, maximum_backoff=0.2) + runtime_database=database, wait=waits.append, minimum_backoff=0.1, maximum_backoff=0.2) calls = iter((True, True, True, False)) service.serve(keep_running=lambda: next(calls)) self.assertEqual(waits, [0.1]) + database.close() def test_service_and_mutating_cli_share_one_runtime_lease(self) -> None: with TemporaryDirectory() as root: @@ -56,12 +59,14 @@ def test_service_and_mutating_cli_share_one_runtime_lease(self) -> None: def test_stop_wakes_default_backoff_without_waiting_for_its_cap(self) -> None: with TemporaryDirectory() as root: state = _State("mission", MissionExecutionStatus.WAITING_FOR_EVIDENCE, 4) + database = RuntimeDatabase(".", path=Path(root) / "runtime.db") service = ForgeRuntimeService(_Loop(state, progresses=False), _States(state), - runtime_database_path=Path(root) / "runtime.db", minimum_backoff=5, maximum_backoff=5) + runtime_database=database, minimum_backoff=5, maximum_backoff=5) thread = Thread(target=lambda: service.serve(keep_running=lambda: True)) started = monotonic(); thread.start(); sleep(0.05); service.stop(); thread.join(timeout=1) self.assertFalse(thread.is_alive()) self.assertLess(monotonic() - started, 1) + database.close() def test_separate_process_mutator_cannot_share_runtime_lease_and_exit_releases_it(self) -> None: with TemporaryDirectory() as root: From bab27eacde90a07c71320d56f0ff4c8e40225d3e Mon Sep 17 00:00:00 2001 From: pcvantol Date: Tue, 8 Sep 2026 08:53:21 +0200 Subject: [PATCH 14/14] test: prove producer contract reopen fidelity --- tests/test_bootstrap_mission_runner.py | 25 +++++++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/tests/test_bootstrap_mission_runner.py b/tests/test_bootstrap_mission_runner.py index 52bec3f..2e93d4b 100644 --- a/tests/test_bootstrap_mission_runner.py +++ b/tests/test_bootstrap_mission_runner.py @@ -239,6 +239,31 @@ def test_nondefault_producer_contract_roundtrips_losslessly(self) -> None: with self.assertRaisesRegex(MissionRunnerError, "Producer Contract"): _request(malformed) + def test_nondefault_contract_survives_persisted_state_reopen_losslessly(self) -> None: + from forge.runtime.runner import _request, _request_document + with TemporaryDirectory() as directory: + root = Path(directory) + runtime = RuntimeDatabase(root); store = MissionStateStore(runtime) + path = runtime.path + prompt = prompt_factory({}, action(1, "one")) + contract = ProducerContract(Producer(ProducerIdentity("producer", "FORGE", "1.0")), "correlation", "one", + RuntimePromptEnvelope(prompt.id, "1.0", "text/markdown", "persisted content", "sha256:" + "d" * 64), + ("constraint-a", "constraint-b"), (("intent_id", "one"), ("intent_revision", "1"), ("mission_revision", "7")), + mission_id="mission-1", receipt_references=(ExecutionReceiptReference("host", "receipt-a"),), + execution_evidence_references=("evidence-a",)) + request = ExecutionRequest("host", "mission-1", "one", "1", "one", prompt, "workspace", "forge", "correlation", "now", producer_contract=contract) + state = store.create(mission("one"), (intent("one"),), (action(1, "one"),), occurred_at="now", resume={}) + state = store.transition(state.mission_id, MissionExecutionStatus.READY, occurred_at="now", reason="ready") + state = store.transition(state.mission_id, MissionExecutionStatus.ACTIVE, occurred_at="now", reason="active") + store.transition(state.mission_id, MissionExecutionStatus.WAITING_FOR_EXECUTION, occurred_at="now", reason="persisted", + execution_correlation={"request": _request_document(request), "host_run_id": None}) + store.close(); runtime.close() + reopened = RuntimeDatabase(root); reopened_store = MissionStateStore(reopened) + restored = _request(reopened_store.get("mission-1").execution_correlation["request"]) + self.assertEqual(restored.producer_contract.to_dict(), contract.to_dict()) + self.assertEqual(restored.producer_contract.digest(), contract.digest()) + reopened_store.close(); reopened.close() + if __name__ == "__main__": unittest.main()