litellm/tests/unit/rust_bridge/test_lifecycle.py
devin-ai-integration[bot] e4190d86a6
refactor(rust): centralize host execution and compose callbacks (#43515)
* refactor(rust): extract litellm-host-native as the shared Rust host driver

Move service and hook dispatch out of host-http into a Driver that owns the
machine and Rust handlers, returning at completion or a stream boundary and
holding the demand reply until the consumer advances. Move the in-process
runner onto the same driver. host-http now layers encoding, SSE, body polling
and lifecycle observation over it. host-python keeps driving litellm-host
directly

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* fix(rust): interrupt the machine when the in-process stream consumer fails

Restores the pre-refactor interruption path for StreamConsumer errors via
Driver::fail and ports the generic run lifecycle tests into host-native.

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* refactor(rust): separate the machine contract from coroutine execution

* auth update

* refactor(rust): use standard flow control for host requests

* style(rust): keep host driver imports formatted

* chores

* mostly relocation

* refactor(rust): separate interceptors from queued observers

* refactor(rust): centralize legacy callback mappings and lifecycle

* docs: define Python host boundaries and migration plan

* refactor: enforce Python host and bridge boundaries

* refactor(rust): separate operations from callback composition

* refactor(rust): compose SDK policy through call hooks

---------

Co-authored-by: Yujong Lee <yujong@berri.ai>
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-09-28 19:20:27 +00:00

74 lines
2.3 KiB
Python

from __future__ import annotations
import asyncio
from collections.abc import Sequence
from typing import Final
import pytest
from litellm.rust_bridge.lifecycle import Await, Complete, Execution, Open, Step, drive
class ScriptedExecution:
"""Plays scripted steps and records how it was resumed and whether it was closed."""
def __init__(self, steps: Sequence[Step]) -> None:
self._steps: Final = list(steps)
self.resumed: list[tuple[str, object]] = []
self.closed = False
def start(self) -> Step:
return self._steps.pop(0)
def resume_value(self, value: object) -> Step:
self.resumed.append(("value", value))
return self._steps.pop(0)
def resume_error(self, error: BaseException) -> Step:
self.resumed.append(("error", type(error)))
return self._steps.pop(0)
def close(self) -> None:
self.closed = True
async def ready(value: object) -> object:
return value
async def failing() -> object:
raise ValueError("boom")
def test_drive_resumes_each_await_with_its_result_or_error_and_returns_the_completed_value() -> None:
execution: Final = ScriptedExecution([Await(ready(1)), Await(failing()), Complete("done")])
assert asyncio.run(drive(execution)) == "done"
assert execution.resumed == [("value", 1), ("error", ValueError)]
assert execution.closed
@pytest.mark.parametrize("factory_fails", (False, True))
def test_stream_handoff_preserves_head_identity_and_closes_on_construction_failure(factory_fails: bool) -> None:
head: Final = object()
stream: Final = object()
execution: Final = ScriptedExecution([Open(head)])
failure: Final = ValueError("stream construction failed")
def construct(owner: Execution, received: object) -> object:
assert owner is execution
assert received is head
if factory_fails:
raise failure
return stream
if factory_fails:
with pytest.raises(ValueError, match="stream construction failed") as caught:
asyncio.run(drive(execution, construct))
assert caught.value is failure
assert execution.closed
else:
assert asyncio.run(drive(execution, construct)) is stream
assert not execution.closed
execution.close()