File size: 6,339 Bytes
a9df8b1
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
from __future__ import annotations

import pytest

from skillos.layers.skill_runtime.executor import SkillExecutor
from skillos.layers.skill_runtime.planner import ExecutionPlan, PlanStep, StepStatus
from skillos.models.skill_model import Skill, SkillImplementation, SkillInterface, SkillState


def make_skill(name: str, code: str) -> Skill:
    return Skill(
        name=name,
        description=f"{name} test skill",
        state=SkillState.RELEASED,
        interface=SkillInterface(
            input_schema={"type": "object", "properties": {}},
            output_schema={"type": "object", "properties": {}},
        ),
        implementation=SkillImplementation(code=code),
    )


def make_plan(steps: list[PlanStep]) -> ExecutionPlan:
    return ExecutionPlan(
        plan_id="plan-1",
        task_id="plan-1",
        task_description="run stable plan",
        steps=steps,
    )


@pytest.mark.asyncio
async def test_independent_step_can_succeed_after_parallel_step_fails():
    fail_skill = make_skill("fail_skill", "raise RuntimeError('boom')")
    ok_skill = make_skill("ok_skill", "output['ok'] = True")
    failing_step = PlanStep(
        step_index=0,
        skill_id=fail_skill.skill_id,
        skill_name=fail_skill.name,
    )
    success_step = PlanStep(
        step_index=1,
        skill_id=ok_skill.skill_id,
        skill_name=ok_skill.name,
    )
    plan = make_plan([failing_step, success_step])
    events: list[tuple[str, dict]] = []
    executor = SkillExecutor(max_retries=0)
    executor.add_event_callback(lambda event_type, data: events.append((event_type, data)))

    final_state = await executor.execute_plan(
        plan,
        {fail_skill.skill_id: fail_skill, ok_skill.skill_id: ok_skill},
        {},
    )

    assert failing_step.status == StepStatus.FAILED
    assert success_step.status == StepStatus.SUCCESS
    assert final_state["ok_skill_executed"] is True
    assert _event(events, "plan_completed")["status"] == "partial"


@pytest.mark.asyncio
async def test_dependent_step_is_skipped_when_dependency_fails():
    fail_skill = make_skill("fail_skill", "raise RuntimeError('boom')")
    dependent_skill = make_skill("dependent_skill", "output['ok'] = True")
    failing_step = PlanStep(
        step_index=0,
        skill_id=fail_skill.skill_id,
        skill_name=fail_skill.name,
    )
    dependent_step = PlanStep(
        step_index=1,
        skill_id=dependent_skill.skill_id,
        skill_name=dependent_skill.name,
        depends_on=[failing_step.step_id],
    )
    plan = make_plan([failing_step, dependent_step])
    events: list[tuple[str, dict]] = []
    executor = SkillExecutor(max_retries=0)
    executor.add_event_callback(lambda event_type, data: events.append((event_type, data)))

    await executor.execute_plan(
        plan,
        {fail_skill.skill_id: fail_skill, dependent_skill.skill_id: dependent_skill},
        {},
    )

    assert failing_step.status == StepStatus.FAILED
    assert dependent_step.status == StepStatus.SKIPPED
    assert failing_step.step_id in dependent_step.error
    skipped_event = _event(events, "step_skipped")
    assert skipped_event["step_id"] == dependent_step.step_id
    assert skipped_event["failed_dependency"] == failing_step.step_id
    assert _event(events, "plan_completed")["status"] == "failed"


@pytest.mark.asyncio
async def test_skip_cascades_through_dependency_chain():
    fail_skill = make_skill("fail_skill", "raise RuntimeError('boom')")
    middle_skill = make_skill("middle_skill", "output['ok'] = True")
    leaf_skill = make_skill("leaf_skill", "output['ok'] = True")
    first = PlanStep(step_index=0, skill_id=fail_skill.skill_id, skill_name=fail_skill.name)
    middle = PlanStep(
        step_index=1,
        skill_id=middle_skill.skill_id,
        skill_name=middle_skill.name,
        depends_on=[first.step_id],
    )
    leaf = PlanStep(
        step_index=2,
        skill_id=leaf_skill.skill_id,
        skill_name=leaf_skill.name,
        depends_on=[middle.step_id],
    )
    plan = make_plan([first, middle, leaf])
    executor = SkillExecutor(max_retries=0)

    await executor.execute_plan(
        plan,
        {
            fail_skill.skill_id: fail_skill,
            middle_skill.skill_id: middle_skill,
            leaf_skill.skill_id: leaf_skill,
        },
        {},
    )

    assert first.status == StepStatus.FAILED
    assert middle.status == StepStatus.SKIPPED
    assert leaf.status == StepStatus.SKIPPED
    assert middle.step_id in leaf.error


@pytest.mark.asyncio
async def test_missing_skill_failure_has_timestamps_and_event_payload():
    step = PlanStep(step_index=0, skill_id="missing", skill_name="missing_skill")
    plan = make_plan([step])
    events: list[tuple[str, dict]] = []
    executor = SkillExecutor(max_retries=0)
    executor.add_event_callback(lambda event_type, data: events.append((event_type, data)))

    await executor.execute_plan(plan, {}, {})

    assert step.status == StepStatus.FAILED
    assert step.started_at is not None
    assert step.completed_at is not None
    assert step.latency_ms is not None
    failed_event = _event(events, "step_failed")
    assert failed_event["step_index"] == 0
    assert failed_event["skill_id"] == "missing"
    assert failed_event["latency_ms"] is not None


@pytest.mark.asyncio
async def test_step_timeout_fails_and_rolls_back_state():
    slow_skill = make_skill("slow_skill", "output['ok'] = True")
    step = PlanStep(step_index=0, skill_id=slow_skill.skill_id, skill_name=slow_skill.name)
    plan = make_plan([step])
    executor = SlowExecutor(max_retries=0, step_timeout_s=0.01)

    final_state = await executor.execute_plan(plan, {slow_skill.skill_id: slow_skill}, {"before": True})

    assert step.status == StepStatus.FAILED
    assert "timed out" in step.error
    assert final_state == {"before": True}


def _event(events: list[tuple[str, dict]], event_type: str) -> dict:
    matches = [data for event, data in events if event == event_type]
    assert matches, f"missing event {event_type}"
    return matches[-1]


class SlowExecutor(SkillExecutor):
    async def _run_skill(self, skill: Skill, input_data: dict, current_state: dict) -> dict:
        import asyncio

        await asyncio.sleep(0.05)
        return {"ok": True, "_state_changes": {"ok": True}}