Files
crewAI/lib/crewai/tests/test_execution_uuid.py
Lorenze Jay 1818792fb1 feat(execution): introduce execution context management with UUID sup… (#6988)
* feat(execution): introduce execution context management with UUID support

- Added `execution.py` to manage execution UUIDs for tracking nested execution contexts.
- Implemented `begin_execution` and `end_execution` functions to handle the lifecycle of execution contexts.
- Updated `Crew` and `Flow` classes to utilize the new execution context management, ensuring proper tracking during execution.
- Added tests for execution UUID creation, inheritance, and lifecycle management to ensure functionality and correctness.

* feat(execution): enhance execution context management in Crew class

- Introduced `begin_execution` and `end_execution` calls in the `Crew` class to manage execution tokens effectively.
- Updated the `akickoff` method to ensure proper lifecycle handling of execution contexts.
- Added tests to verify the creation and clearing of execution UUIDs during the `akickoff` process, ensuring correct behavior in various scenarios.

* feat(execution): add execution_uuid to PendingFeedbackContext for flow resumption

- Introduced `execution_uuid` to the `PendingFeedbackContext` class to maintain the UUID across flow pauses and resumes, ensuring traceability of execution contexts.
- Updated the `Flow` class to utilize the new `execution_uuid` during execution management, enhancing the handling of paused flows.
- Added tests to verify that the execution UUID is correctly persisted and restored during flow operations, ensuring consistent behavior across sessions.

* refactor(execution): streamline execution UUID management and update tests

- Removed the `ensure_execution_uuid` function to simplify UUID handling, consolidating logic into `begin_execution` and `end_execution`.
- Updated the `clear_execution_uuid` function to ensure it correctly restores previous UUIDs using context tokens.
- Modified tests to reflect changes in execution UUID management, ensuring proper creation, inheritance, and clearing of UUIDs during execution contexts.
- Enhanced the `PendingFeedbackContext` documentation to clarify the handling of `execution_uuid` for pending rows.
2026-08-13 16:29:46 -07:00

239 lines
7.1 KiB
Python

"""Tests for OSS execution uuid creation and nesting inheritance."""
from __future__ import annotations
import pytest
from crewai.execution import (
_current_execution_uuid,
begin_execution,
clear_execution_uuid,
end_execution,
get_execution_uuid,
set_execution_uuid,
)
@pytest.fixture(autouse=True)
def _isolate_execution_uuid() -> None:
token = _current_execution_uuid.set(None)
yield
_current_execution_uuid.reset(token)
def test_begin_creates_when_empty() -> None:
assert get_execution_uuid() is None
first_token = begin_execution()
first = get_execution_uuid()
second_token = begin_execution()
assert first
assert get_execution_uuid() == first
assert second_token is None
end_execution(second_token)
end_execution(first_token)
def test_begin_does_not_overwrite_existing() -> None:
token = set_execution_uuid("enterprise-kickoff-id")
try:
assert begin_execution("should-not-win") is None
assert get_execution_uuid() == "enterprise-kickoff-id"
finally:
clear_execution_uuid(token)
def test_set_rejects_empty() -> None:
with pytest.raises(ValueError, match="non-empty"):
set_execution_uuid("")
def test_nested_execution_inherits_and_only_owner_clears() -> None:
parent_token = begin_execution()
parent_id = get_execution_uuid()
child_token = begin_execution()
assert parent_id
assert child_token is None
assert get_execution_uuid() == parent_id
end_execution(child_token)
assert get_execution_uuid() == parent_id
end_execution(parent_token)
assert get_execution_uuid() is None
def test_two_sequential_outer_runs_get_distinct_uuids() -> None:
first_token = begin_execution()
first_id = get_execution_uuid()
end_execution(first_token)
second_token = begin_execution()
second_id = get_execution_uuid()
end_execution(second_token)
assert first_id != second_id
def test_flow_kickoff_creates_and_clears_execution_uuid() -> None:
from crewai.flow.flow import Flow, start
seen: dict[str, str | None] = {}
class ProbeFlow(Flow):
@start()
def begin(self) -> str:
seen["during"] = get_execution_uuid()
return "ok"
flow = ProbeFlow()
assert get_execution_uuid() is None
flow.kickoff()
assert seen["during"]
assert get_execution_uuid() is None
def test_flow_kickoff_inherits_enterprise_execution_uuid() -> None:
from crewai.flow.flow import Flow, start
seen: dict[str, str | None] = {}
class ProbeFlow(Flow):
@start()
def begin(self) -> str:
seen["during"] = get_execution_uuid()
return "ok"
token = set_execution_uuid("celery-kickoff-id")
try:
ProbeFlow().kickoff()
assert seen["during"] == "celery-kickoff-id"
# Owner was enterprise set(), not the flow — still set here.
assert get_execution_uuid() == "celery-kickoff-id"
finally:
clear_execution_uuid(token)
assert get_execution_uuid() is None
@pytest.mark.asyncio
async def test_crew_akickoff_creates_and_clears_execution_uuid() -> None:
from unittest.mock import patch
from crewai import Agent, Crew, Task
from crewai.tasks.task_output import TaskOutput
seen: dict[str, str | None] = {}
agent = Agent(role="r", goal="g", backstory="b", llm="gpt-4o-mini")
task = Task(description="d", expected_output="o", agent=agent)
crew = Crew(agents=[agent], tasks=[task])
async def capture(*_args: object, **_kwargs: object) -> TaskOutput:
seen["during"] = get_execution_uuid()
return TaskOutput(description="d", raw="ok", agent="r")
with patch("crewai.task.Task.aexecute_sync", side_effect=capture):
assert get_execution_uuid() is None
await crew.akickoff()
assert seen["during"]
assert get_execution_uuid() is None
@pytest.mark.asyncio
async def test_crew_akickoff_inherits_enterprise_execution_uuid() -> None:
from unittest.mock import patch
from crewai import Agent, Crew, Task
from crewai.tasks.task_output import TaskOutput
seen: dict[str, str | None] = {}
agent = Agent(role="r", goal="g", backstory="b", llm="gpt-4o-mini")
task = Task(description="d", expected_output="o", agent=agent)
crew = Crew(agents=[agent], tasks=[task])
async def capture(*_args: object, **_kwargs: object) -> TaskOutput:
seen["during"] = get_execution_uuid()
return TaskOutput(description="d", raw="ok", agent="r")
token = set_execution_uuid("celery-kickoff-id")
try:
with patch("crewai.task.Task.aexecute_sync", side_effect=capture):
await crew.akickoff()
assert seen["during"] == "celery-kickoff-id"
assert get_execution_uuid() == "celery-kickoff-id"
finally:
clear_execution_uuid(token)
assert get_execution_uuid() is None
def test_flow_pause_persists_execution_uuid_and_resume_restores_it() -> None:
import os
import tempfile
from crewai.flow import Flow, human_feedback, listen, start
from crewai.flow.async_feedback.types import (
HumanFeedbackPending,
PendingFeedbackContext,
)
from crewai.flow.persistence import SQLiteFlowPersistence
seen: dict[str, str | None] = {}
class PausingProvider:
def request_feedback(
self, context: PendingFeedbackContext, flow: Flow
) -> str:
raise HumanFeedbackPending(context=context)
class ReviewFlow(Flow):
@start()
@human_feedback(message="Review:", provider=PausingProvider())
def generate(self) -> str:
seen["during_kickoff"] = get_execution_uuid()
return "draft"
@listen(generate)
def process(self, result: object) -> str:
seen["during_resume"] = get_execution_uuid()
return "done"
with tempfile.TemporaryDirectory() as tmpdir:
persistence = SQLiteFlowPersistence(os.path.join(tmpdir, "test.db"))
pending = ReviewFlow(persistence=persistence).kickoff()
assert isinstance(pending, HumanFeedbackPending)
paused_id = seen["during_kickoff"]
assert paused_id
assert pending.context.execution_uuid == paused_id
assert get_execution_uuid() is None
ReviewFlow.from_pending(pending.context.flow_id, persistence).resume("ok")
assert seen["during_resume"] == paused_id
assert get_execution_uuid() is None
def test_pending_feedback_context_roundtrips_execution_uuid() -> None:
from crewai.flow.async_feedback.types import PendingFeedbackContext
original = PendingFeedbackContext(
flow_id="flow-1",
flow_class="test.Flow",
method_name="review",
method_output="draft",
message="Review:",
execution_uuid="kickoff-uuid",
)
restored = PendingFeedbackContext.from_dict(original.to_dict())
assert restored.execution_uuid == "kickoff-uuid"
legacy = PendingFeedbackContext.from_dict(
{
"flow_id": "flow-1",
"flow_class": "test.Flow",
"method_name": "review",
"method_output": "draft",
"message": "Review:",
}
)
assert legacy.execution_uuid is None