mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-03 02:22:24 +00:00
chore(harness): remove banner comments, restating comments and dead in_loop_thread (#44161)
Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
7d50a31eb5
commit
62a1b7fdf3
7 changed files with 0 additions and 96 deletions
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
*,
|
||||
|
|
|
|||
|
|
@ -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],
|
||||
|
|
|
|||
|
|
@ -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())
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue