feat: persist interrupt contract in run state
This commit is contained in:
@@ -137,6 +137,11 @@ class InterruptRoute:
|
||||
workflow_ref: WorkflowRef
|
||||
|
||||
|
||||
def _object_schema() -> dict[str, object]:
|
||||
"""Default legacy interrupt contract used when older checkpoints are loaded."""
|
||||
return {"type": "object", "additionalProperties": True}
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class InterruptRequest:
|
||||
id: str
|
||||
@@ -146,6 +151,10 @@ class InterruptRequest:
|
||||
payload: dict[str, Any] = field(default_factory=dict)
|
||||
resumable: bool = True
|
||||
route: InterruptRoute | None = None
|
||||
outcomes: list[str] = field(default_factory=lambda: ["submitted"])
|
||||
request_schema: dict[str, object] = field(default_factory=_object_schema)
|
||||
resume_schema: dict[str, object] = field(default_factory=_object_schema)
|
||||
typed: bool = False
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
|
||||
@@ -60,6 +60,10 @@ def build_interrupt_request(
|
||||
kind=node.kind,
|
||||
payload=payload,
|
||||
route=route,
|
||||
outcomes=list(node.outcomes),
|
||||
request_schema=dict(node.request_schema),
|
||||
resume_schema=dict(node.resume_schema),
|
||||
typed=node.has_explicit_contract,
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -81,6 +81,42 @@ def test_run_state_codec_round_trips_child_interrupt_lineage_types() -> None:
|
||||
assert isinstance(restored.interrupt.route.workflow_ref, WorkflowRef)
|
||||
|
||||
|
||||
def test_run_state_codec_round_trips_interrupt_contract_fields() -> None:
|
||||
run = RunState(
|
||||
workflow_name="approval",
|
||||
status=RunStatus.INTERRUPTED,
|
||||
workflow_input={},
|
||||
state={},
|
||||
interrupt=InterruptRequest(
|
||||
id="interrupt:approval",
|
||||
frame_id="frame_1",
|
||||
node_id="approval",
|
||||
kind="approval",
|
||||
payload={"message": "approve?"},
|
||||
outcomes=["submitted", "cancelled"],
|
||||
request_schema={
|
||||
"type": "object",
|
||||
"properties": {"message": {"type": "string"}},
|
||||
"required": ["message"],
|
||||
},
|
||||
resume_schema={
|
||||
"type": "object",
|
||||
"properties": {"approved": {"type": "boolean"}},
|
||||
"required": ["approved"],
|
||||
},
|
||||
typed=True,
|
||||
),
|
||||
)
|
||||
|
||||
restored = load_run_state(dump_run_state(run))
|
||||
|
||||
assert restored.interrupt is not None
|
||||
assert restored.interrupt.outcomes == ["submitted", "cancelled"]
|
||||
assert restored.interrupt.request_schema["required"] == ["message"]
|
||||
assert restored.interrupt.resume_schema["required"] == ["approved"]
|
||||
assert restored.interrupt.typed is True
|
||||
|
||||
|
||||
def test_run_state_codec_restores_root_state_alias_for_resume_writes() -> None:
|
||||
"""Root-scope commits after restore must be visible to final output."""
|
||||
state = {"text": "before"}
|
||||
|
||||
Reference in New Issue
Block a user