mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-02 02:11:58 +00:00
* feat(s3_v2): add s3_partition_granularity option for hourly S3 folders Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(integration): cover s3 v2 partition granularity across surfaces, settings and chaos Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(integration): cover previous_response_id history rebuilt from an hourly cold storage object Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * fix(s3_v2): reuse the cold storage key only when s3_v2 owns cold storage Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): cover hour rollover, postgres outage, in-flight switches, key/team vars and real S3 layout Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * chore(liccheck): authorize libfaketime, the GPLv2 dev-only clock the s3 rollover integration test preloads Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): wait for the rejected-request cell's payloads by id, not by line count Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): declare the postgres outage cell's models in config and trip the relay on burst ids Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): drop the libfaketime hour rollover cell and its dev dependency Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(e2e): deselect the s3_v2 live e2e on the stage-mirror stack The stage-mirror config enables no s3_v2 callback, so every test in test_s3_log_e2e.py fails its readiness check there. The file keeps running in the Buildkite e2e lane, which configures s3_v2 Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(s3_v2): declare the sink outage burst models in config Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * refactor(s3_v2): read cold storage metadata without an empty dict default Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> * test(integration): wait for the proxy to reconnect before the postgres outage recovery request Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: mrinal <mrinal@berri.ai> Co-authored-by: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> Co-authored-by: yucheng <yucheng@berri.ai>
102 lines
3.7 KiB
Python
102 lines
3.7 KiB
Python
import asyncio
|
|
import socket
|
|
import threading
|
|
import time
|
|
from collections.abc import Generator
|
|
from contextlib import contextmanager
|
|
from typing import Final
|
|
from urllib.parse import urlsplit, urlunsplit
|
|
|
|
from pydantic import TypeAdapter
|
|
|
|
PORT: Final = TypeAdapter(int)
|
|
OUTAGE_SECONDS: Final = 10.0
|
|
|
|
|
|
def _free_port() -> int:
|
|
with socket.socket() as reserve:
|
|
reserve.bind(("127.0.0.1", 0))
|
|
return PORT.validate_python(reserve.getsockname()[1])
|
|
|
|
|
|
class DatabaseRelay:
|
|
def __init__(self, upstream_host: str, upstream_port: int, trigger: bytes) -> None:
|
|
self.port: Final = _free_port()
|
|
self._upstream_host: Final = upstream_host
|
|
self._upstream_port: Final = upstream_port
|
|
self._trigger: Final = trigger
|
|
self._loop: Final = asyncio.new_event_loop()
|
|
self._armed: Final = threading.Event()
|
|
self.tripped: Final = threading.Event()
|
|
self.refused = 0
|
|
self.reconnected: Final = threading.Event()
|
|
self._tripped_at = 0.0
|
|
self._writers: tuple[asyncio.StreamWriter, ...] = ()
|
|
self._ready: Final = threading.Event()
|
|
self._thread: Final = threading.Thread(target=self._run, daemon=True)
|
|
|
|
def arm(self) -> None:
|
|
self._armed.set()
|
|
|
|
def start(self) -> None:
|
|
self._thread.start()
|
|
assert self._ready.wait(10), "Database relay did not start"
|
|
|
|
def stop(self) -> None:
|
|
self._loop.call_soon_threadsafe(self._loop.stop)
|
|
self._thread.join(10)
|
|
|
|
def _run(self) -> None:
|
|
asyncio.set_event_loop(self._loop)
|
|
self._loop.run_until_complete(asyncio.start_server(self._serve, "127.0.0.1", self.port))
|
|
self._ready.set()
|
|
self._loop.run_forever()
|
|
|
|
def _drop_all(self) -> None:
|
|
for writer in self._writers:
|
|
writer.close()
|
|
self._writers = ()
|
|
|
|
async def _serve(self, client_reader: asyncio.StreamReader, client_writer: asyncio.StreamWriter) -> None:
|
|
if self.tripped.is_set() and time.monotonic() - self._tripped_at < OUTAGE_SECONDS:
|
|
self.refused += 1
|
|
client_writer.close()
|
|
return
|
|
if self.tripped.is_set():
|
|
self.reconnected.set()
|
|
server_reader, server_writer = await asyncio.open_connection(self._upstream_host, self._upstream_port)
|
|
self._writers = (*self._writers, client_writer, server_writer)
|
|
|
|
async def forward(reader: asyncio.StreamReader, writer: asyncio.StreamWriter, inspect: bool) -> None:
|
|
try:
|
|
while chunk := await reader.read(65536):
|
|
if inspect and self._armed.is_set() and not self.tripped.is_set() and self._trigger in chunk:
|
|
self._tripped_at = time.monotonic()
|
|
self.tripped.set()
|
|
self._drop_all()
|
|
return
|
|
writer.write(chunk)
|
|
await writer.drain()
|
|
except (ConnectionError, asyncio.IncompleteReadError):
|
|
return
|
|
finally:
|
|
writer.close()
|
|
|
|
await asyncio.gather(
|
|
forward(client_reader, server_writer, True),
|
|
forward(server_reader, client_writer, False),
|
|
)
|
|
|
|
|
|
@contextmanager
|
|
def database_relay(database_url: str, trigger: bytes) -> Generator[tuple[DatabaseRelay, str]]:
|
|
parts: Final = urlsplit(database_url)
|
|
assert parts.hostname is not None and parts.port is not None, database_url
|
|
relay: Final = DatabaseRelay(parts.hostname, parts.port, trigger)
|
|
relay.start()
|
|
credentials: Final = f"{parts.username}:{parts.password}@" if parts.username else ""
|
|
relayed: Final = urlunsplit(parts._replace(netloc=f"{credentials}127.0.0.1:{relay.port}"))
|
|
try:
|
|
yield relay, relayed
|
|
finally:
|
|
relay.stop()
|