mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-30 01:52:18 +00:00
* 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>
74 lines
2.3 KiB
Python
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()
|