diff --git a/litellm/__init__.py b/litellm/__init__.py index 9a4f4605519..1e9e7037477 100644 --- a/litellm/__init__.py +++ b/litellm/__init__.py @@ -2282,7 +2282,6 @@ if TYPE_CHECKING: # Track if async client cleanup has been registered (for lazy loading) _async_client_cleanup_registered = False -# litellm.agent() entrypoints, resolved lazily from litellm.harness by __getattr__. _AGENT_EXPORTS: Final = frozenset( { "agent", @@ -2333,7 +2332,6 @@ def __getattr__(name: str) -> Any: handler_func: Final = registry[name] return handler_func(name) - # litellm.agent() and friends: imported on first access (not needed for completion calls) if name == "harness" or name in _AGENT_EXPORTS: import importlib diff --git a/litellm/harness/endpoint.py b/litellm/harness/endpoint.py index c492e4794ba..ce3586735c6 100644 --- a/litellm/harness/endpoint.py +++ b/litellm/harness/endpoint.py @@ -131,11 +131,6 @@ class UsageTracker: ) -# --------------------------------------------------------------------------- -# Usage parsing -# --------------------------------------------------------------------------- - - def _as_int(value: object) -> int: if isinstance(value, bool): return 0 @@ -230,11 +225,6 @@ class SSEUsageParser: self.output_tokens = output_tokens -# --------------------------------------------------------------------------- -# Cost + helpers -# --------------------------------------------------------------------------- - - def compute_cost(model: str | None, input_tokens: int, output_tokens: int) -> float: """Cost from LiteLLM's price map. Never raises; unknown models cost 0.0.""" if not model or not (input_tokens or output_tokens): @@ -373,11 +363,6 @@ def _noop() -> None: return None -# --------------------------------------------------------------------------- -# ModelEndpoint -# --------------------------------------------------------------------------- - - class ModelEndpoint: """Local HTTP endpoint for one harness session. Use as an async context manager.""" @@ -401,7 +386,6 @@ class ModelEndpoint: self.token = secrets.token_urlsafe(HARNESS_SESSION_TOKEN_BYTES) self.usage = UsageTracker() self.port = 0 - # Injected client (tests); production uses LiteLLM's shared cached client. self._injected_client = client self._deps: _ServerDeps | None = None self._client: httpx.AsyncClient | None = None @@ -412,8 +396,6 @@ class ModelEndpoint: def url(self) -> str: return f"http://{HARNESS_ENDPOINT_HOST}:{self.port}" - # -- lifecycle ---------------------------------------------------------- - async def __aenter__(self) -> ModelEndpoint: await self.start() return self @@ -500,8 +482,6 @@ class ModelEndpoint: routes=[*post_routes, *get_routes] # mutable-ok: Starlette takes a routes list ) - # -- request handling --------------------------------------------------- - @property def _responses(self) -> ModuleType: if self._deps is None: @@ -565,8 +545,6 @@ class ModelEndpoint: cost = compute_cost(model, input_tokens, output_tokens) self.usage.add(input_tokens, output_tokens, cost) - # -- gateway mode ------------------------------------------------------- - async def _forward(self, request: Request, route: str, body: Mapping[str, Any]) -> Response: if self._client is None or self.gateway is None: raise HarnessError("gateway client is not started") @@ -622,8 +600,6 @@ class ModelEndpoint: tokens = (0, 0) self._record(model, tokens[0], tokens[1], header_cost(upstream.headers)) - # -- SDK mode ----------------------------------------------------------- - def _sdk_kwargs( self, body: Mapping[str, Any] ) -> dict[str, Any]: # mutable-ok: SDK call kwargs, mutated by _invoke_sdk then splatted diff --git a/litellm/harness/runtime.py b/litellm/harness/runtime.py index 53a06e7d50d..4ec167a06b4 100644 --- a/litellm/harness/runtime.py +++ b/litellm/harness/runtime.py @@ -68,11 +68,6 @@ verbose_logger: Final = logging.getLogger("LiteLLM") PROPAGATED_ERRORS: Final = (HarnessInstallFailed, CapabilityUnsupported) -# --------------------------------------------------------------------------- -# Configuration + validation -# --------------------------------------------------------------------------- - - @dataclass(frozen=True) class SessionConfig: """Every per-session parameter a caller can pass, already normalized.""" @@ -249,11 +244,6 @@ def _context_for(config: SessionConfig) -> SessionContext: ) -# --------------------------------------------------------------------------- -# Structured output -# --------------------------------------------------------------------------- - - def parse_output(output: type[BaseModel], output_json: str | None, text: str) -> tuple[BaseModel | None, str | None]: """Return (model, None) on success or (None, error message) on failure.""" raw = output_json or last_json_object(text) @@ -265,11 +255,6 @@ def parse_output(output: type[BaseModel], output_json: str | None, text: str) -> return None, str(e) -# --------------------------------------------------------------------------- -# Turn machinery -# --------------------------------------------------------------------------- - - @dataclass class _End: """Sentinel the producer puts on the queue when the handler turn is over.""" @@ -392,8 +377,6 @@ class _Turn: self.usage_before: tuple[int, int, int, float] = (0, 0, 0, 0.0) self.deadline: float | None = None - # -- setup / teardown --------------------------------------------------- - async def _begin(self) -> None: sandbox = self.ctx.sandbox self.before = await sandbox.snapshot() @@ -419,8 +402,6 @@ class _Turn: if pending: await asyncio.wait(pending) - # -- event loop --------------------------------------------------------- - async def _next_item(self) -> Event | _End: if self.deadline is None: return await self.queue.get() @@ -476,8 +457,6 @@ class _Turn: # The consumer asked for the next event without answering. item.deny("approval not answered") - # -- results ------------------------------------------------------------ - async def _file_changes(self) -> list[FileChange]: # mutable-ok: becomes the public Result.files list sandbox = self.ctx.sandbox after = await sandbox.snapshot() @@ -535,8 +514,6 @@ class _Turn: parsed, error = parse_output(output_type, self.ctx.output_json, text) return parsed, self.ctx.output_json or text, error - # -- entry -------------------------------------------------------------- - async def run(self) -> AsyncIterator[Event]: await self._begin() self._start_producer() @@ -564,11 +541,6 @@ class _Turn: ) -# --------------------------------------------------------------------------- -# Streams -# --------------------------------------------------------------------------- - - class AsyncEventStream: """Async iterator of events for one turn. `.result` is set once Done is seen.""" @@ -609,11 +581,6 @@ async def _one_shot(session: AsyncSession, prompt: str, control: TurnControl) -> await session.aclose() -# --------------------------------------------------------------------------- -# Sessions -# --------------------------------------------------------------------------- - - class AsyncSession: """A multi-turn conversation with one harness. Use `async with` or `await`.""" @@ -641,8 +608,6 @@ class AsyncSession: self._busy = False self._restart_needed = False - # -- lifecycle ---------------------------------------------------------- - def __await__(self) -> Generator[object, None, AsyncSession]: return self.start().__await__() @@ -752,8 +717,6 @@ class AsyncSession: detach = adetach stop = astop - # -- turns -------------------------------------------------------------- - def usage_counters(self) -> tuple[int, int, int, float]: """(input_tokens, output_tokens, calls, cost) so far, from endpoint or handler.""" endpoint = self.ctx.endpoint @@ -823,11 +786,6 @@ async def _collect(events: AsyncIterator[Event]) -> Result: return result -# --------------------------------------------------------------------------- -# Public async API -# --------------------------------------------------------------------------- - - def aagent_session( harness: Harness, *, diff --git a/litellm/harness/sandbox/docker.py b/litellm/harness/sandbox/docker.py index 32dbfc27ddc..ac8f357200b 100644 --- a/litellm/harness/sandbox/docker.py +++ b/litellm/harness/sandbox/docker.py @@ -77,8 +77,6 @@ class DockerSandbox: def __repr__(self) -> str: return f"DockerSandbox({self.image!r}, workdir={self.workdir!r})" - # -- docker CLI plumbing (tests monkeypatch these two) --------------------- - def _docker_binary(self) -> str: binary = shutil.which("docker") if binary is None: @@ -116,8 +114,6 @@ class DockerSandbox: await handle.kill() raise SandboxError(f"docker {args[0]} timed out after {timeout}s") - # -- command construction -------------------------------------------------- - def run_args( self, ) -> list[str]: # mutable-ok: argv is returned as a list, the shape callers and tests compare against @@ -169,8 +165,6 @@ class DockerSandbox: joined = path if posixpath.isabs(path) else posixpath.join(self.workdir, path) return posixpath.normpath(joined) - # -- lifecycle --------------------------------------------------------------- - async def start(self) -> str: """Start the container if needed and return its id.""" if self._closed: @@ -191,8 +185,6 @@ class DockerSandbox: container_id = await self.start() return await self._docker(self.exec_args(container_id, cmd), input=input) - # -- Sandbox protocol -------------------------------------------------------- - async def exec( self, cmd: Sequence[str], diff --git a/litellm/harness/sync.py b/litellm/harness/sync.py index 4583a98bf54..1788543f9b5 100644 --- a/litellm/harness/sync.py +++ b/litellm/harness/sync.py @@ -62,9 +62,6 @@ class _LoopThread: self._thread.start() return self._loop - def in_loop_thread(self) -> bool: - return self._thread is not None and threading.current_thread() is self._thread - def submit(self, coro: Coroutine[Any, Any, T]) -> Future[T]: return asyncio.run_coroutine_threadsafe(coro, self.loop()) diff --git a/litellm/llms/deepagents/harness/sandbox_backend.py b/litellm/llms/deepagents/harness/sandbox_backend.py index 49610014e21..d39690ccdd1 100644 --- a/litellm/llms/deepagents/harness/sandbox_backend.py +++ b/litellm/llms/deepagents/harness/sandbox_backend.py @@ -129,8 +129,6 @@ class SandboxBackend(SandboxBackendProtocol): # pyright: ignore[reportUntypedBa def id(self) -> str: return self._id - # -- paths -------------------------------------------------------------- - def to_real(self, path: str) -> str: """Sandbox path for a virtual path (or an absolute path already under workdir).""" normalized = posixpath.normpath("/" + path.lstrip("/")) @@ -186,8 +184,6 @@ class SandboxBackend(SandboxBackendProtocol): # pyright: ignore[reportUntypedBa async def _run(self, cmd: Sequence[str], timeout: float | None = DEEPAGENTS_FS_TIMEOUT_SECONDS) -> CompletedRun: return await self._sandbox.run(cmd, timeout=timeout) - # -- ls ----------------------------------------------------------------- - async def als(self, path: str) -> LsResult: try: real = await self.to_confined(path) @@ -207,8 +203,6 @@ class SandboxBackend(SandboxBackendProtocol): # pyright: ignore[reportUntypedBa def ls(self, path: str) -> LsResult: return self._sync(self.als(path)) - # -- read / write / edit ------------------------------------------------ - async def _read_bytes(self, path: str) -> bytes: return await self._sandbox.read(await self.to_confined(path)) @@ -296,8 +290,6 @@ class SandboxBackend(SandboxBackendProtocol): # pyright: ignore[reportUntypedBa def delete(self, file_path: str) -> DeleteResult: return self._sync(self.adelete(file_path)) - # -- glob / grep -------------------------------------------------------- - def _find_cmd(self, root: str) -> tuple[str, ...]: prune = tuple( itertools.chain.from_iterable( @@ -378,8 +370,6 @@ class SandboxBackend(SandboxBackendProtocol): # pyright: ignore[reportUntypedBa ) -> GrepResult: return self._sync(self.agrep(pattern, path, glob, max_count=max_count)) - # -- upload / download -------------------------------------------------- - async def _upload_one(self, path: str, data: bytes) -> FileUploadResponse: if not self._writable: return FileUploadResponse(path=path, error="permission_denied") @@ -425,8 +415,6 @@ class SandboxBackend(SandboxBackendProtocol): # pyright: ignore[reportUntypedBa ) -> list[FileDownloadResponse]: # mutable-ok: return type fixed by deepagents BackendProtocol return self._sync(self.adownload_files(paths)) - # -- execute ------------------------------------------------------------ - async def aexecute(self, command: str, *, timeout: int | None = None) -> ExecuteResponse: if not self._allow_execute: return ExecuteResponse(output=_NO_EXECUTE_ERROR, exit_code=1) diff --git a/litellm/llms/deepagents/harness/transformation.py b/litellm/llms/deepagents/harness/transformation.py index 339a6274d33..82b2c044eac 100644 --- a/litellm/llms/deepagents/harness/transformation.py +++ b/litellm/llms/deepagents/harness/transformation.py @@ -70,11 +70,6 @@ APPROVAL_TOOLS: Final = WRITE_TOOLS | EXECUTE_TOOLS _APPROVAL_DECISIONS: Final = ("approve", "reject") -# --------------------------------------------------------------------------- -# Pure helpers (unit tested directly; kept module-level so they port cleanly) -# --------------------------------------------------------------------------- - - def gateway_headers( ctx: SessionContext, ) -> dict[str, str]: # mutable-ok: ChatLiteLLM.extra_headers is a pydantic dict field