Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions .github/workflows/tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -94,3 +94,18 @@ jobs:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7

- run: bash scripts/release-tests/litellm-existing-secret.sh

litellm-stream-errors:
name: LiteLLM Responses stream errors
runs-on: ubuntu-latest
timeout-minutes: 5
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7

- uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7
with:
python-version: "3.11"

- run: python3 -m unittest discover -s deploy/litellm/tests -p 'test_*.py'

- run: bash scripts/release-tests/litellm-build-context.sh
6 changes: 5 additions & 1 deletion deploy/litellm/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,10 @@
FROM ghcr.io/berriai/litellm:v1.99.0@sha256:570a872d2fde8f1bc4a147634810941103c14697270f0f2918c5aa7d8201cac5

WORKDIR /app

COPY patches/ /tmp/litellm-patches/
RUN python3 /tmp/litellm-patches/apply_responses_stream_errors.py && \
rm -rf /tmp/litellm-patches

# ═══ Inject custom hook files ═══
# In the base image, litellm is pip-installed, so the import path resolves to
Expand Down Expand Up @@ -36,4 +40,4 @@ RUN chmod +x ./docker/entrypoint.sh

EXPOSE 4000/tcp

CMD ["--port", "4000", "--config", "config.yaml"]
CMD ["--port", "4000", "--config", "config.yaml"]
46 changes: 44 additions & 2 deletions deploy/litellm/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -114,11 +114,53 @@ Secret in that namespace and set `serviceMonitor.bearerTokenSecret.name`.

## Building the image

`./build.sh` (needs `REGISTRY`; set `HARBOR_PASSWORD` to build in-cluster with
kaniko when there is no docker daemon). The tag it writes names the LiteLLM
`./build.sh` needs `REGISTRY`. Set `PUSH_SECRET` to select an in-cluster Kaniko
build using an existing registry credential Secret when there is no Docker
daemon. The tag it writes names the LiteLLM
version taken from the Dockerfile, and `deploy.sh` refuses an image whose
version-named tag disagrees with that pin.

## Responses stream errors

The image includes a checked patch for LiteLLM 1.99.0. When an upstream request
fails after a native Responses stream has started, the proxy emits
`response.failed` with the upstream error message and a recognizable error code.
This lets Responses clients report the cause instead of an unexpected end of
stream. The patch preserves the response ID and event sequence when available,
and does not append another failure after a terminal event. Chat Completions
and the Cursor conversion endpoint keep their existing stream formats.

The build verifies SHA-256 hashes of the upstream files before applying the
patch. A different source version fails the build and requires reviewing the
patch against that version. Both Docker and Kaniko include the patch installer
and helper in their build context. Retry, fallback, and cooldown policies are
unchanged by this patch.

From the repository root, run the local checks:

```bash
python3 -m unittest discover -s deploy/litellm/tests -p 'test_*.py'
bash scripts/release-tests/litellm-build-context.sh
```

Verify a built image against a synthetic upstream without provider credentials
or a database:

```bash
docker run --rm --entrypoint python3 \
-v "$PWD/deploy/litellm/tests:/tests:ro" "$LITELLM_TEST_IMAGE" \
/tests/integration_responses_stream_errors.py
```

The integration check exercises the actual LiteLLM HTTP server. Its `--serve`
mode keeps the synthetic service running for a Codex CLI check through a local
port-forward:

```bash
python3 deploy/litellm/tests/integration_responses_stream_errors.py \
--codex-url http://127.0.0.1:4000
```

## Upgrading LiteLLM

**The database migration is one-way.** LiteLLM runs `prisma migrate deploy` on
Expand Down
23 changes: 13 additions & 10 deletions deploy/litellm/build.sh
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,8 @@
# away from the Dockerfile beside it. `claw/deploy/build.sh` had solved the
# same problem with an in-cluster kaniko job; this borrows that shape.
#
# HARBOR_PASSWORD push password for $REGISTRY (presence selects kaniko)
# PUSH_SECRET Secret holding .dockerconfigjson for $REGISTRY (kaniko backend)
# HARBOR_USERNAME push user (default: admin)
# PUSH_SECRET Secret holding .dockerconfigjson; selects the kaniko backend
# HARBOR_PASSWORD legacy selector for kaniko; credentials still use PUSH_SECRET
# REGISTRY e.g. harbor.example.com/primussafe
# NAMESPACE where the build job runs (default: primus-claw)
# TAG default: v<litellm version from the Dockerfile>-<date>
Expand All @@ -26,7 +25,6 @@ SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"

NAMESPACE="${NAMESPACE:-primus-claw}"
REGISTRY="${REGISTRY:?REGISTRY is required, e.g. harbor.example.com/primussafe}"
HARBOR_USERNAME="${HARBOR_USERNAME:-admin}"

# Read the pinned version out of the Dockerfile so the tag cannot disagree.
BASE_VERSION="$(grep -oE '^FROM .*litellm:v[0-9.]+' "$SCRIPT_DIR/Dockerfile" | grep -oE 'v[0-9.]+' | head -1)"
Expand All @@ -36,8 +34,8 @@ IMG="$REGISTRY/litellm:$TAG"

echo "[litellm-build] building $IMG (base $BASE_VERSION)"

if [ -z "${HARBOR_PASSWORD:-}" ]; then
command -v docker >/dev/null || { echo "ERROR: no HARBOR_PASSWORD for the kaniko backend and no docker daemon" >&2; exit 1; }
if [ -z "${PUSH_SECRET:-}" ] && [ -z "${HARBOR_PASSWORD:-}" ]; then
command -v docker >/dev/null || { echo "ERROR: set PUSH_SECRET for kaniko, or install docker" >&2; exit 1; }
docker build -t "$IMG" "$SCRIPT_DIR"
docker push "$IMG"
echo "[litellm-build] pushed $IMG"
Expand All @@ -46,12 +44,17 @@ fi

JOB="litellm-build-$(date +%s)"
CTX="litellm-build-ctx-$(date +%s)"
cleanup() { kubectl -n "$NAMESPACE" delete cm "$CTX" --ignore-not-found >/dev/null 2>&1 || true; }
CONTEXT_ARCHIVE="$(mktemp)"
cleanup() {
rm -f "$CONTEXT_ARCHIVE"
kubectl -n "$NAMESPACE" delete cm "$CTX" --ignore-not-found >/dev/null 2>&1 || true
}
trap cleanup EXIT

tar -C "$SCRIPT_DIR" --exclude='__pycache__' --exclude='*.pyc' \
-czf "$CONTEXT_ARCHIVE" Dockerfile apim_key_hook.py patches
kubectl -n "$NAMESPACE" create cm "$CTX" \
--from-file=Dockerfile="$SCRIPT_DIR/Dockerfile" \
--from-file=apim_key_hook.py="$SCRIPT_DIR/apim_key_hook.py" \
--from-file=build-context.tar.gz="$CONTEXT_ARCHIVE" \
--dry-run=client -o yaml | kubectl apply -f - >/dev/null

# Registry auth is read from an existing pull/push secret rather than inlined
Expand All @@ -75,7 +78,7 @@ spec:
initContainers:
- name: ctx
image: busybox:1.36
command: ["sh","-c","cp /cm/Dockerfile /cm/apim_key_hook.py /workspace/"]
command: ["sh","-c","tar -xzf /cm/build-context.tar.gz -C /workspace"]
volumeMounts:
- {name: ws, mountPath: /workspace}
- {name: cm, mountPath: /cm}
Expand Down
173 changes: 173 additions & 0 deletions deploy/litellm/patches/apply_responses_stream_errors.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
# Copyright Advanced Micro Devices, Inc.
# SPDX-License-Identifier: MIT

from __future__ import annotations

import argparse
import ast
import hashlib
import importlib.util
from pathlib import Path


BASE_SHA256 = {
"proxy/proxy_server.py": "f63cd83c5c4459d84dbfbd350a460caf9e2907389f19d2c3090fe24ad937f47b",
"proxy/response_api_endpoints/endpoints.py": "563b462e7c36e729869d0acdc016a54c2e99ec07f70a3f6ea17482b1dbf11e4f",
"responses/streaming_iterator.py": "e25f32c3bf3e1815f08a5f2e328d4b32359d189d386780a7b9b4640bdbe4e56a",
}
HELPER_NAME = "responses_stream_errors.py"


def _replace_once(source: str, before: str, after: str) -> str:
count = source.count(before)
if count != 1:
raise ValueError(f"Expected exactly one patch context, found {count}: {before[:80]!r}")
return source.replace(before, after, 1)


def _native_responses_endpoint(source: str) -> str:
function = next(node for node in ast.parse(source).body if isinstance(node, ast.AsyncFunctionDef) and node.name == "responses_api")
lines = source.splitlines(keepends=True)
start, end = function.lineno - 1, function.end_lineno
block = "".join(lines[start:end])
block = _replace_once(block, " select_data_generator,\n", " select_responses_data_generator as select_data_generator,\n")
return "".join(lines[:start]) + block + "".join(lines[end:])


def _responses_iterator(source: str) -> str:
before = """ if 400 <= status_code < 500 and status_code != 429:
raise mapped_exception
"""
after = """ from litellm.proxy.responses_stream_errors import preserve_upstream_error

preserve_upstream_error(mapped_exception, error_obj, result)
if 400 <= status_code < 500 and status_code != 429:
raise mapped_exception
"""
return _replace_once(source, before, after)


def _proxy_generator(source: str) -> str:
before = """async def async_data_generator(
response,
user_api_key_dict: UserAPIKeyAuth,
request_data: dict,
request: Request | None = None,
):
verbose_proxy_logger.debug("inside generator")
stream_completed = False
client_disconnected = False
"""
after = """async def async_data_generator(
response,
user_api_key_dict: UserAPIKeyAuth,
request_data: dict,
request: Request | None = None,
*,
responses_stream_errors: bool = False,
):
from litellm.proxy.responses_stream_errors import ResponsesStreamErrorState

verbose_proxy_logger.debug("inside generator")
stream_completed = False
client_disconnected = False
error_state = ResponsesStreamErrorState() if responses_stream_errors else None
"""
source = _replace_once(source, before, after)
source = _replace_once(source, " raw_passthrough = False\n", " if error_state is not None:\n error_state.observe_chunk(chunk, emitted=False)\n raw_passthrough = False\n")
before = """ if isinstance(e, HTTPException):
raise e
elif isinstance(e, StreamingCallbackError):
"""
after = """ if error_state is not None:
stream_completed = True
error_frame = error_state.format_failure(e)
if error_frame is not None:
yield error_frame
return
if isinstance(e, HTTPException):
raise e
elif isinstance(e, StreamingCallbackError):
"""
source = _replace_once(source, before, after)
before = ''' yield _format_streaming_sse_chunk(chunk=chunk)
except Exception as e:
yield f"data: {e}\\n\\n"
'''
after = ''' formatted_chunk = _format_streaming_sse_chunk(chunk=chunk)
if error_state is not None:
error_state.mark_emitted()
yield formatted_chunk
except Exception as e:
if error_state is not None:
raise
yield f"data: {e}\\n\\n"
'''
return _replace_once(source, before, after)


def _responses_selector(source: str) -> str:
before = "\ndef select_data_generator(\n"
after = """
def select_responses_data_generator(
response,
user_api_key_dict: UserAPIKeyAuth,
request_data: dict,
request: Request | None = None,
):
return async_data_generator(
response=response,
user_api_key_dict=user_api_key_dict,
request_data=request_data,
request=request,
responses_stream_errors=True,
)


def select_data_generator(
"""
return _replace_once(source, before, after)


def apply_patch(package_root: Path) -> None:
sources: dict[str, str] = {}
for relative, expected in BASE_SHA256.items():
content = (package_root / relative).read_bytes()
actual = hashlib.sha256(content).hexdigest()
if actual != expected:
raise ValueError(f"Unsupported LiteLLM source {relative}: expected SHA-256 {expected}, found {actual}")
sources[relative] = content.decode("utf-8")
helper = Path(__file__).with_name(HELPER_NAME).read_text()
target_helper = package_root / "proxy" / HELPER_NAME
if target_helper.exists():
raise ValueError(f"Refusing to replace an existing {target_helper}")
sources["proxy/proxy_server.py"] = _responses_selector(_proxy_generator(sources["proxy/proxy_server.py"]))
sources["proxy/response_api_endpoints/endpoints.py"] = _native_responses_endpoint(sources["proxy/response_api_endpoints/endpoints.py"])
sources["responses/streaming_iterator.py"] = _responses_iterator(sources["responses/streaming_iterator.py"])
for relative, source in sources.items():
compile(source, relative, "exec")
compile(helper, HELPER_NAME, "exec")
for relative, source in sources.items():
(package_root / relative).write_text(source)
target_helper.write_text(helper)


def main() -> None:
parser = argparse.ArgumentParser(description="Apply the checked Responses stream error patch to LiteLLM 1.99.0.")
parser.add_argument("package_root", nargs="?", type=Path)
parser.add_argument("--package-dir", type=Path)
args = parser.parse_args()
if args.package_root is not None and args.package_dir is not None:
parser.error("Specify either package_root or --package-dir, not both")
package_root = args.package_dir or args.package_root
if package_root is None:
spec = importlib.util.find_spec("litellm")
if spec is None or spec.origin is None:
raise RuntimeError("Cannot locate the installed litellm package")
package_root = Path(spec.origin).parent
apply_patch(package_root)
print("Applied the Responses stream error patch to LiteLLM 1.99.0.")


if __name__ == "__main__":
main()
Loading
Loading