From 45212d1415097fc5bbbc1fe5764cbb65890f2052 Mon Sep 17 00:00:00 2001 From: Honglin Cao Date: Fri, 14 Aug 2026 13:28:28 -0400 Subject: [PATCH 1/3] feat(sdk): migrate deployment logs to logs_v4 and add pod listing Adds get_deployment_pods, a page-based get_deployment_logs anchored on previously fetched events or an epoch-ms boundary (before/after), a get_deployment_logs_range time-window read that can merge all pods of a revision, and a stateful DeploymentLogSession for backfill, tailing, and cross-process resume. Signed-off-by: Honglin Cao --- README.md | 16 ++ centml/sdk/api.py | 209 +++++++++++++--- examples/sdk/get_deployment_logs.py | 100 ++++---- tests/test_sdk_api.py | 357 ++++++++++++++++++++++++++++ 4 files changed, 594 insertions(+), 88 deletions(-) diff --git a/README.md b/README.md index c99ede76..44405297 100644 --- a/README.md +++ b/README.md @@ -53,6 +53,22 @@ Use `python examples/sdk/get_clusters.py` and Creating the example reserves GPU capacity and may incur usage charges. It does not delete the deployment automatically. +### Deployment logs SDK example + +Logs are read per pod. Discover pod names with `get_deployment_pods()` (terminated +pods still within log retention are included), then read with a +`deployment_log_session()`: `fetch_older()` pages toward the beginning of history and +`fetch_newer()` returns only new lines, while the session keeps the merged, ordered +log in `.events`. `get_deployment_logs_range()` fetches a specific time window +(epoch-millisecond bounds, both optional) and, with `pod=None`, merges every pod's +stream chronologically. The same paging is available statelessly through +`get_deployment_logs(before=..., after=...)`, anchored on events you already hold +or on a bare epoch-millisecond boundary: + +```bash +python examples/sdk/get_deployment_logs.py +``` + ### Un-installation To uninstall `centml`, simply do: diff --git a/centml/sdk/api.py b/centml/sdk/api.py index 4a4bea49..9f9b554b 100644 --- a/centml/sdk/api.py +++ b/centml/sdk/api.py @@ -1,4 +1,7 @@ +from bisect import insort from contextlib import contextmanager +from dataclasses import dataclass +from typing import List, Optional, Union import platform_api_python_client from platform_api_python_client import ( @@ -21,6 +24,23 @@ STATUS_V3_DEPLOYMENT_TYPES = {DeploymentType.INFERENCE_V3, DeploymentType.CSERVE_V3} +DEFAULT_LOG_PAGE_LINES = 100 # server-side default for max_lines +MAX_LOG_PAGE_LINES = 5000 # server-side ceiling for max_lines +# The server re-delivers a ~15s look-behind window on fetch-newer requests; only the +# caller's events within this generous margin of the boundary can be re-delivered. +LOG_DEDUP_RETENTION_MS = 300_000 + + +@dataclass(frozen=True) +class DeploymentLogEvent: + """One log line with its pod attached — logs_v4 events carry no pod name, so + merged multi-pod views need the SDK to attribute each line itself.""" + + id: str + timestamp: int + message: str + pod: str + class CentMLClient: def __init__(self, api): @@ -191,46 +211,175 @@ def get_deployment_revisions(self, deployment_id: int): deployment_id=deployment_id ).results + def get_deployment_pods(self, deployment_id: int, revision_number: int) -> List[str]: + """List pods that have logged for a deployment revision, including terminated + pods still within log retention. A fresh deployment may return an empty list.""" + return self._api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get( + deployment_id=deployment_id, revision_number=revision_number + ).pods + + # pylint: disable=R0917 def get_deployment_logs( self, deployment_id: int, revision_number: int, - start_time: int, - end_time: int, - line_count: int = 100, - start_from_head: bool = True, - stream: bool = False, - ): - """Fetch logs for a deployment within a time window, handling pagination automatically. - - start_time and end_time are Unix timestamps in milliseconds. - Use get_deployment_revisions() to find the current revision number. - - If stream=True, returns a generator that yields events as each page is fetched. - If stream=False (default), returns a flat list of all events. + pod: str, + before: Optional[Union[list, int]] = None, + after: Optional[Union[list, int]] = None, + max_lines: int = DEFAULT_LOG_PAGE_LINES, + ) -> list: + """Fetch one page of a pod's logs, oldest-first. Use get_deployment_pods() to + discover pod names and get_deployment_revisions() for the revision number. + + before and after anchor the page to events a previous call returned for the + same pod (pass your accumulated list; only the relevant boundary is used): + - neither: the newest page (tail). + - before=: the page strictly older than the oldest of them; + an empty result means the beginning of history is reached. + - after=: lines strictly newer than the newest of them; an empty + result means nothing new yet — call again later to keep tailing. Late + lines still landing near that boundary are included on top of max_lines + and may sort below events you already hold (order by id if that matters). + Either anchor also accepts a bare epoch-millisecond int as the (exclusive) + boundary itself; an int after anchor holds no event ids, so the re-delivered + span at the boundary comes through undeduplicated. """ + if before is not None and after is not None: + raise ValueError("before and after are mutually exclusive") + + fetch_newer = after is not None + anchor = after if fetch_newer else before + anchor_events = None + boundary_timestamp = None + if isinstance(anchor, int): + boundary_timestamp = anchor + elif anchor: + anchor_events = anchor + timestamps = [event.timestamp for event in anchor_events] + boundary_timestamp = max(timestamps) if fetch_newer else min(timestamps) + + response = self._api.get_deployment_logs_v4_logs_deployment_id_revision_number_get( + deployment_id=deployment_id, + revision_number=revision_number, + pod=pod, + fetch_newer=fetch_newer, + timestamp=boundary_timestamp, + max_lines=max_lines, + ) + if not fetch_newer or not anchor_events: + return response.events + + # fetch_newer re-delivers a look-behind window at and before the boundary + # (late-arrival protection); drop the lines the caller already holds by id. + cutoff = max(event.timestamp for event in anchor_events) - LOG_DEDUP_RETENTION_MS + held_event_ids = {event.id for event in anchor_events if event.timestamp >= cutoff} + return [event for event in response.events if event.id not in held_event_ids] - def _iter_events(): - next_page_token = None + # pylint: disable=R0917 + def get_deployment_logs_range( + self, + deployment_id: int, + revision_number: int, + pod: Optional[str] = None, + start_time: Optional[int] = None, + end_time: Optional[int] = None, + ) -> List[DeploymentLogEvent]: + """Fetch every log line in [start_time, end_time] (epoch ms, inclusive; both + optional — an open end reads to the beginning or the present), oldest first. + pod=None reads all pods of the revision and merges the streams + chronologically; each returned event carries its pod name.""" + if start_time is not None and end_time is not None and start_time > end_time: + raise ValueError("start_time must not exceed end_time") + + pods = [pod] if pod is not None else self.get_deployment_pods(deployment_id, revision_number) + merged: List[DeploymentLogEvent] = [] + for pod_name in pods: + events: list = [] while True: - response = self._api.get_deployment_logs_v3_deployments_logs_v3_deployment_id_revision_number_get( - deployment_id=deployment_id, - revision_number=revision_number, - start_time=start_time, - end_time=end_time, - next_page_token=next_page_token, - start_from_head=start_from_head, - line_count=line_count, + # after is exclusive, so start_time - 1 admits lines at start_time itself; + # start_time 0 (or None) means the whole window — scan from the head. + anchor: Union[list, int] = events if events else (start_time - 1 if start_time else []) + page = self.get_deployment_logs( + deployment_id, revision_number, pod_name, after=anchor, max_lines=MAX_LOG_PAGE_LINES ) - yield from response.events - next_page_token = response.next_page_token - if not next_page_token: + if not page: break + events += page + if end_time is not None and page[-1].timestamp > end_time: + break + merged += [ + DeploymentLogEvent(id=event.id, timestamp=event.timestamp, message=event.message, pod=pod_name) + for event in events + if (start_time is None or event.timestamp >= start_time) + and (end_time is None or event.timestamp <= end_time) + ] + merged.sort(key=lambda event: event.id) + return merged + + def deployment_log_session( + self, deployment_id: int, revision_number: int, pod: str, events: Optional[list] = None + ) -> "DeploymentLogSession": + """Stateful reader for one pod's logs that tracks fetched pages and anchors + every request itself — see DeploymentLogSession. Seed events with logs a + previous session (or get_deployment_logs) returned for the same pod.""" + return DeploymentLogSession(self, deployment_id, revision_number, pod, events) + + +class DeploymentLogSession: + """Maintains a contiguous, ordered window of one pod's logs across fetches. + + Every fetch is anchored on the window itself, so pages can never overlap or + leave gaps inside it (within log retention; an undetectable gap forms if the + session idles past retention before fetching newer lines). + """ - if stream: - return _iter_events() - - return list(_iter_events()) + # pylint: disable=R0917 + def __init__(self, client: CentMLClient, deployment_id: int, revision_number: int, pod: str, events=None): + self._client = client + self._deployment_id = deployment_id + self._revision_number = revision_number + self._pod = pod + # Seeded events come from outside the session: canonicalize to unique ids in + # chronological order (id order == time order at nanosecond precision). + unique_events = {event.id: event for event in events or []} + self._events = [unique_events[event_id] for event_id in sorted(unique_events)] + + @property + def events(self) -> list: + """Copy of the window fetched so far, oldest first. Complete from the beginning + of history only once fetch_older() has returned an empty list.""" + return list(self._events) + + def fetch_older(self, max_lines: int = DEFAULT_LOG_PAGE_LINES) -> list: + """Fetch the page older than the window and prepend it; on an empty session + fetches the newest page (tail). Returns the page; empty list = no older + lines exist (yet).""" + page = self._client.get_deployment_logs( + self._deployment_id, self._revision_number, self._pod, before=self._events, max_lines=max_lines + ) + self._events[:0] = page + return page + + def fetch_newer(self, max_lines: int = DEFAULT_LOG_PAGE_LINES) -> list: + """Fetch lines newer than the window and merge them in; on an empty session + fetches the newest page (tail) — to read from the beginning of history + instead, loop fetch_older() until it returns an empty list. Returns only + the new lines; empty list = nothing new yet, call again later to keep + tailing. Rare late arrivals sort into the window below its newest lines.""" + if not self._events: + return self.fetch_older(max_lines=max_lines) + delta = self._client.get_deployment_logs( + self._deployment_id, self._revision_number, self._pod, after=self._events, max_lines=max_lines + ) + for event in delta: + if event.id > self._events[-1].id: + self._events.append(event) + else: + # A late arrival may even precede the window's oldest line (tail page + # cut inside the look-behind span); the server delivers that span + # completely on top of max_lines, so the window stays contiguous. + insort(self._events, event, key=lambda held: held.id) + return delta @contextmanager diff --git a/examples/sdk/get_deployment_logs.py b/examples/sdk/get_deployment_logs.py index 6c7d7ea8..f1e1bdbb 100644 --- a/examples/sdk/get_deployment_logs.py +++ b/examples/sdk/get_deployment_logs.py @@ -1,74 +1,58 @@ -from datetime import datetime, timezone, timedelta +import time +from datetime import datetime, timezone from centml.sdk.api import get_centml_client # --- Configuration --- DEPLOYMENT_ID = 1234 # Replace with your deployment ID REVISION_NUMBER = 10 -HOURS_BACK = 1 # Fetch logs from the last N hours +TAIL_SECONDS = 30 # How long to keep polling for new lines after reading history -def format_event(event: dict) -> str: - timestamp_ms = ( - event.get("timestamp") - or event.get("time") - or event.get("ts") - or "" - ) - message = ( - event.get("message") - or event.get("msg") - or event.get("log") - or str(event) - ) - if timestamp_ms: - ts = datetime.fromtimestamp(int(timestamp_ms) / 1000, tz=timezone.utc).isoformat() - return f"[{ts}] {message}" - return message +def format_event(event) -> str: + ts = datetime.fromtimestamp(event.timestamp / 1000, tz=timezone.utc).isoformat() + return f"[{ts}] {event.message}" def main(): - stream = True - end_time = int(datetime.now(timezone.utc).timestamp() * 1000) - start_time = end_time - int(timedelta(hours=HOURS_BACK).total_seconds() * 1000) - - print(f"Fetching logs for deployment {DEPLOYMENT_ID}") - print( - f"Time window: " - f"{datetime.fromtimestamp(start_time / 1000, tz=timezone.utc).isoformat()} → " - f"{datetime.fromtimestamp(end_time / 1000, tz=timezone.utc).isoformat()}" - ) - print() - with get_centml_client() as cclient: - if stream: - # Streaming: print events as each page arrives - for event in cclient.get_deployment_logs( - deployment_id=DEPLOYMENT_ID, - revision_number=REVISION_NUMBER, - start_time=start_time, - end_time=end_time, - start_from_head=False, - stream=stream, - ): - print(format_event(event)) - else: - # Batch: collect all events then process - events = cclient.get_deployment_logs( - deployment_id=DEPLOYMENT_ID, - revision_number=REVISION_NUMBER, - start_time=start_time, - end_time=end_time, - start_from_head=False, - ) - - if not events: - print("No logs found in the given time window.") - return - - print(f"Found {len(events)} log entries:\n") - for event in events: + # Logs are read per pod: discover the pods that have logged for this revision + # (terminated pods within log retention are included). + pods = cclient.get_deployment_pods(DEPLOYMENT_ID, REVISION_NUMBER) + if not pods: + print("No pods have logged for this revision yet.") + return + + pod = pods[0] + print(f"Reading logs for deployment {DEPLOYMENT_ID} revision {REVISION_NUMBER}, pod {pod}\n") + + # The session tracks what it has fetched and anchors every request itself. + session = cclient.deployment_log_session(DEPLOYMENT_ID, REVISION_NUMBER, pod) + + # Read the full history: newest page first, then page back to the beginning. + while session.fetch_older(): + pass + print(f"Found {len(session.events)} log entries:\n") + for event in session.events: + print(format_event(event)) + + # Keep tailing: each call returns only the lines the session does not hold yet. + print(f"\nPolling for new lines for {TAIL_SECONDS}s...") + deadline = time.monotonic() + TAIL_SECONDS + while time.monotonic() < deadline: + for event in session.fetch_newer(): print(format_event(event)) + time.sleep(2) + + # The same paging is available statelessly via get_deployment_logs, anchored + # on events you already hold — useful when you manage storage yourself: + # page = cclient.get_deployment_logs(DEPLOYMENT_ID, REVISION_NUMBER, pod=pod) # tail + # older = cclient.get_deployment_logs(..., pod=pod, before=page) # [] = beginning + # newer = cclient.get_deployment_logs(..., pod=pod, after=page) # [] = nothing new + # A specific time window (all pods merged, oldest first, pod on each event): + # window = cclient.get_deployment_logs_range( + # DEPLOYMENT_ID, REVISION_NUMBER, start_time=t1_ms, end_time=t2_ms + # ) if __name__ == "__main__": diff --git a/tests/test_sdk_api.py b/tests/test_sdk_api.py index 03ae917c..598fabcf 100644 --- a/tests/test_sdk_api.py +++ b/tests/test_sdk_api.py @@ -1,6 +1,8 @@ from types import SimpleNamespace from unittest.mock import MagicMock, patch +import pytest + import platform_api_python_client from platform_api_python_client import ( CreateDynamoDeploymentRequest, @@ -228,3 +230,358 @@ def test_delete_hardware_instance_delegates_to_platform_client(): assert response is expected_response api.delete_hardware_instance_hardware_instances_hardware_instance_id_delete.assert_called_once_with(123) + + +def _log_event(event_id, timestamp, message="line"): + return SimpleNamespace(id=event_id, timestamp=timestamp, message=message) + + +def _log_page(*events): + return SimpleNamespace(events=list(events)) + + +def test_generated_client_exposes_logs_v4_contract(): + assert hasattr( + platform_api_python_client.EXTERNALApi, "get_deployment_logs_v4_logs_deployment_id_revision_number_get" + ) + assert hasattr( + platform_api_python_client.EXTERNALApi, "get_deployment_pods_deployments_pods_deployment_id_revision_number_get" + ) + + +def test_get_deployment_pods_returns_pod_names(): + api = MagicMock() + api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get.return_value = SimpleNamespace( + pods=["pod-a", "pod-b"] + ) + client = CentMLClient(api) + + assert client.get_deployment_pods(123, 2) == ["pod-a", "pod-b"] + + api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get.assert_called_once_with( + deployment_id=123, revision_number=2 + ) + + +def test_get_deployment_logs_returns_tail_page_when_unanchored(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page( + _log_event("1-a", 1000), _log_event("2-b", 2000) + ) + client = CentMLClient(api) + + events = client.get_deployment_logs(123, 2, pod="pod-a") + + assert [e.id for e in events] == ["1-a", "2-b"] + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_called_once_with( + deployment_id=123, revision_number=2, pod="pod-a", fetch_newer=False, timestamp=None, max_lines=100 + ) + + +def test_get_deployment_logs_before_pages_older_from_oldest_anchor(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page(_log_event("1-a", 1000)) + client = CentMLClient(api) + held = [_log_event("2-b", 2000), _log_event("3-c", 3000)] + + events = client.get_deployment_logs(123, 2, pod="pod-a", before=held) + + assert [e.id for e in events] == ["1-a"] + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + # The boundary is the oldest held timestamp, passed verbatim (exclusive server-side). + assert call.kwargs["fetch_newer"] is False + assert call.kwargs["timestamp"] == 2000 + + +def test_get_deployment_logs_before_empty_page_signals_beginning_of_history(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page() + client = CentMLClient(api) + + assert client.get_deployment_logs(123, 2, pod="pod-a", before=[_log_event("1-a", 1000)]) == [] + + +def test_get_deployment_logs_after_fetches_newer_from_newest_anchor(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page(_log_event("3-c", 3000)) + client = CentMLClient(api) + held = [_log_event("1-a", 1000), _log_event("2-b", 2000)] + + events = client.get_deployment_logs(123, 2, pod="pod-a", after=held) + + assert [e.id for e in events] == ["3-c"] + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is True + assert call.kwargs["timestamp"] == 2000 + + +def test_get_deployment_logs_after_drops_redelivered_lines_but_keeps_late_arrivals(): + api = MagicMock() + # The look-behind window re-delivers 2-b (already held) plus a late arrival at 1500. + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page( + _log_event("15-l", 1500), _log_event("2-b", 2000), _log_event("3-c", 3000) + ) + client = CentMLClient(api) + held = [_log_event("1-a", 1000), _log_event("2-b", 2000)] + + events = client.get_deployment_logs(123, 2, pod="pod-a", after=held) + + assert [e.id for e in events] == ["15-l", "3-c"] + + +def test_get_deployment_logs_after_empty_anchor_reads_from_head(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page(_log_event("1-a", 1000)) + client = CentMLClient(api) + + events = client.get_deployment_logs(123, 2, pod="pod-a", after=[]) + + assert [e.id for e in events] == ["1-a"] + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is True + assert call.kwargs["timestamp"] is None + + +def test_get_deployment_logs_after_empty_page_signals_nothing_new(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page() + client = CentMLClient(api) + + assert client.get_deployment_logs(123, 2, pod="pod-a", after=[_log_event("1-a", 1000)]) == [] + + +def test_get_deployment_logs_passes_max_lines_through(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page() + client = CentMLClient(api) + + client.get_deployment_logs(123, 2, pod="pod-a", max_lines=7) + + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["max_lines"] == 7 + + +def test_get_deployment_logs_rejects_before_and_after_together(): + api = MagicMock() + client = CentMLClient(api) + + with pytest.raises(ValueError): + client.get_deployment_logs(123, 2, pod="pod-a", before=[_log_event("1-a", 1000)], after=[]) + + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_not_called() + + +def _session(api, events=None): + return CentMLClient(api).deployment_log_session(123, 2, "pod-a", events=events) + + +def test_log_session_first_fetch_is_tail_for_both_directions(): + for method in ("fetch_older", "fetch_newer"): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page( + _log_event("1-a", 1000), _log_event("2-b", 2000) + ) + session = _session(api) + + page = getattr(session, method)() + + assert [e.id for e in page] == ["1-a", "2-b"] + assert [e.id for e in session.events] == ["1-a", "2-b"] + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is False and call.kwargs["timestamp"] is None + + +def test_log_session_fetch_older_prepends_until_beginning(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + _log_page(_log_event("3-c", 3000), _log_event("4-d", 4000)), # tail + _log_page(_log_event("1-a", 1000), _log_event("2-b", 2000)), # older page + _log_page(), # beginning reached + ] + session = _session(api) + session.fetch_older() + + older = session.fetch_older() + assert [e.id for e in older] == ["1-a", "2-b"] + assert [e.id for e in session.events] == ["1-a", "2-b", "3-c", "4-d"] + # The older fetch anchors on the window's oldest timestamp. + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is False and call.kwargs["timestamp"] == 3000 + + assert session.fetch_older() == [] + assert [e.id for e in session.events] == ["1-a", "2-b", "3-c", "4-d"] + + +def test_log_session_fetch_newer_merges_delta_and_sorts_late_arrivals(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + _log_page(_log_event("1-a", 1000), _log_event("2-b", 2000)), # tail + # Look-behind re-delivers 2-b (held) plus a late arrival at 1500 and a fresh line. + _log_page(_log_event("15-l", 1500), _log_event("2-b", 2000), _log_event("3-c", 3000)), + ] + session = _session(api) + session.fetch_newer() + + delta = session.fetch_newer() + + assert [e.id for e in delta] == ["15-l", "3-c"] + assert [e.id for e in session.events] == ["1-a", "15-l", "2-b", "3-c"] + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is True and call.kwargs["timestamp"] == 2000 + + +def test_log_session_fetch_newer_empty_delta_leaves_window_unchanged(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + _log_page(_log_event("1-a", 1000)), + _log_page(), + ] + session = _session(api) + session.fetch_newer() + + assert session.fetch_newer() == [] + assert [e.id for e in session.events] == ["1-a"] + + +def test_log_session_seed_canonicalizes_and_anchors_follow_up_fetches(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page() + seed = [_log_event("2-b", 2000), _log_event("1-a", 1000), _log_event("2-b", 2000)] + session = _session(api, events=seed) + + assert [e.id for e in session.events] == ["1-a", "2-b"] + + session.fetch_newer() + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is True and call.kwargs["timestamp"] == 2000 + + +def test_log_session_events_returns_a_copy(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page(_log_event("1-a", 1000)) + session = _session(api) + session.fetch_older() + + view = session.events + view.append(_log_event("9-z", 9000)) + + assert [e.id for e in session.events] == ["1-a"] + + +def test_log_session_passes_max_lines_through(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page() + session = _session(api) + + session.fetch_older(max_lines=7) + assert api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args.kwargs["max_lines"] == 7 + + +def test_get_deployment_logs_accepts_timestamp_anchors(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page(_log_event("2-b", 2000)) + client = CentMLClient(api) + + events = client.get_deployment_logs(123, 2, pod="pod-a", after=1999) + assert [e.id for e in events] == ["2-b"] + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is True and call.kwargs["timestamp"] == 1999 + + client.get_deployment_logs(123, 2, pod="pod-a", before=5000) + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is False and call.kwargs["timestamp"] == 5000 + + +def test_get_deployment_logs_range_trims_to_window(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + # First page anchors at start_time-1; the look-behind may re-deliver older lines. + _log_page(_log_event("05-x", 500), _log_event("1-a", 1000), _log_event("2-b", 2000)), + _log_page(_log_event("3-c", 3000), _log_event("4-d", 4000)), # newest passes end_time + ] + client = CentMLClient(api) + + events = client.get_deployment_logs_range(123, 2, pod="pod-a", start_time=1000, end_time=3000) + + assert [(e.id, e.pod) for e in events] == [("1-a", "pod-a"), ("2-b", "pod-a"), ("3-c", "pod-a")] + calls = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args_list + assert calls[0].kwargs["timestamp"] == 999 + assert calls[1].kwargs["timestamp"] == 2000 + assert len(calls) == 2 # stops once a page reaches past end_time + + +def test_get_deployment_logs_range_open_ended_reads_full_history(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + _log_page(_log_event("1-a", 1000)), + _log_page(_log_event("2-b", 2000)), + _log_page(), + ] + client = CentMLClient(api) + + events = client.get_deployment_logs_range(123, 2, pod="pod-a") + + assert [e.id for e in events] == ["1-a", "2-b"] + calls = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args_list + assert calls[0].kwargs["timestamp"] is None and calls[0].kwargs["fetch_newer"] is True + + +def test_get_deployment_logs_range_merges_all_pods_by_id(): + api = MagicMock() + api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get.return_value = SimpleNamespace( + pods=["pod-a", "pod-b"] + ) + + def pages(**kwargs): + if kwargs["timestamp"] is not None: + return _log_page() + if kwargs["pod"] == "pod-a": + return _log_page(_log_event("1-a", 1000), _log_event("3-a", 3000)) + return _log_page(_log_event("2-b", 2000), _log_event("4-b", 4000)) + + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = pages + client = CentMLClient(api) + + events = client.get_deployment_logs_range(123, 2) + + assert [(e.id, e.pod, e.message) for e in events] == [ + ("1-a", "pod-a", "line"), + ("2-b", "pod-b", "line"), + ("3-a", "pod-a", "line"), + ("4-b", "pod-b", "line"), + ] + + +def test_get_deployment_logs_range_returns_empty_when_no_pod_has_logged(): + api = MagicMock() + api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get.return_value = SimpleNamespace(pods=[]) + client = CentMLClient(api) + + assert client.get_deployment_logs_range(123, 2) == [] + + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_not_called() + + +def test_get_deployment_logs_range_single_millisecond_window(): + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + _log_page(_log_event("1-a", 1000), _log_event("2-b", 2000), _log_event("3-c", 3000)), + _log_page(), + ] + client = CentMLClient(api) + + events = client.get_deployment_logs_range(123, 2, pod="pod-a", start_time=2000, end_time=2000) + + assert [e.id for e in events] == ["2-b"] + calls = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args_list + assert calls[0].kwargs["timestamp"] == 1999 + + +def test_get_deployment_logs_range_rejects_inverted_window(): + api = MagicMock() + client = CentMLClient(api) + + with pytest.raises(ValueError): + client.get_deployment_logs_range(123, 2, pod="pod-a", start_time=2000, end_time=1000) + + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_not_called() From ee074f52eeae8f1f8ea295e57b976b6f1752ae3b Mon Sep 17 00:00:00 2001 From: Honglin Cao Date: Tue, 1 Sep 2026 14:40:08 -0400 Subject: [PATCH 2/3] fix(sdk): reject empty log anchors and slim session poll anchors Empty before/after lists now raise ValueError (after=[] silently meant a head scan while before=[] meant the tail); head reads use after=0. The session passes only the boundary event / trailing retention window to the primitive instead of its whole window, and the example prints a count plus the last lines instead of the full history. Signed-off-by: Honglin Cao --- centml/sdk/api.py | 32 ++++++++++++--- examples/sdk/get_deployment_logs.py | 6 ++- tests/test_sdk_api.py | 60 ++++++++++++++++++++++++++--- 3 files changed, 85 insertions(+), 13 deletions(-) diff --git a/centml/sdk/api.py b/centml/sdk/api.py index 9f9b554b..fc388da2 100644 --- a/centml/sdk/api.py +++ b/centml/sdk/api.py @@ -241,8 +241,10 @@ def get_deployment_logs( lines still landing near that boundary are included on top of max_lines and may sort below events you already hold (order by id if that matters). Either anchor also accepts a bare epoch-millisecond int as the (exclusive) - boundary itself; an int after anchor holds no event ids, so the re-delivered - span at the boundary comes through undeduplicated. + boundary itself — after=0 scans from the head of the log window; an int after + anchor holds no event ids, so the re-delivered span at the boundary comes + through undeduplicated. An empty anchor list raises ValueError. Pages never + split a millisecond, so a delivered boundary millisecond is always complete. """ if before is not None and after is not None: raise ValueError("before and after are mutually exclusive") @@ -253,6 +255,11 @@ def get_deployment_logs( boundary_timestamp = None if isinstance(anchor, int): boundary_timestamp = anchor + elif anchor is not None and len(anchor) == 0: + raise ValueError( + "anchor events must be non-empty; omit the anchor for the tail page, " + "or pass an epoch-ms boundary (after=0 reads from the head)" + ) elif anchor: anchor_events = anchor timestamps = [event.timestamp for event in anchor_events] @@ -298,7 +305,7 @@ def get_deployment_logs_range( while True: # after is exclusive, so start_time - 1 admits lines at start_time itself; # start_time 0 (or None) means the whole window — scan from the head. - anchor: Union[list, int] = events if events else (start_time - 1 if start_time else []) + anchor: Union[list, int] = events if events else (start_time - 1 if start_time else 0) page = self.get_deployment_logs( deployment_id, revision_number, pod_name, after=anchor, max_lines=MAX_LOG_PAGE_LINES ) @@ -355,7 +362,11 @@ def fetch_older(self, max_lines: int = DEFAULT_LOG_PAGE_LINES) -> list: fetches the newest page (tail). Returns the page; empty list = no older lines exist (yet).""" page = self._client.get_deployment_logs( - self._deployment_id, self._revision_number, self._pod, before=self._events, max_lines=max_lines + self._deployment_id, + self._revision_number, + self._pod, + before=[self._events[0]] if self._events else None, + max_lines=max_lines, ) self._events[:0] = page return page @@ -368,8 +379,19 @@ def fetch_newer(self, max_lines: int = DEFAULT_LOG_PAGE_LINES) -> list: tailing. Rare late arrivals sort into the window below its newest lines.""" if not self._events: return self.fetch_older(max_lines=max_lines) + # Only the trailing retention window matters to the primitive (boundary + + # look-behind dedup ids); passing it keeps long tails from scanning the + # whole window on every poll. + cutoff = self._events[-1].timestamp - LOG_DEDUP_RETENTION_MS + first_recent = len(self._events) + while first_recent > 0 and self._events[first_recent - 1].timestamp >= cutoff: + first_recent -= 1 delta = self._client.get_deployment_logs( - self._deployment_id, self._revision_number, self._pod, after=self._events, max_lines=max_lines + self._deployment_id, + self._revision_number, + self._pod, + after=self._events[first_recent:], + max_lines=max_lines, ) for event in delta: if event.id > self._events[-1].id: diff --git a/examples/sdk/get_deployment_logs.py b/examples/sdk/get_deployment_logs.py index f1e1bdbb..7339a6c7 100644 --- a/examples/sdk/get_deployment_logs.py +++ b/examples/sdk/get_deployment_logs.py @@ -7,6 +7,7 @@ DEPLOYMENT_ID = 1234 # Replace with your deployment ID REVISION_NUMBER = 10 TAIL_SECONDS = 30 # How long to keep polling for new lines after reading history +TAIL_LINES = 20 # How much history to print before tailing def format_event(event) -> str: @@ -32,8 +33,9 @@ def main(): # Read the full history: newest page first, then page back to the beginning. while session.fetch_older(): pass - print(f"Found {len(session.events)} log entries:\n") - for event in session.events: + events = session.events + print(f"Found {len(events)} log entries; showing the last {TAIL_LINES}:\n") + for event in events[-TAIL_LINES:]: print(format_event(event)) # Keep tailing: each call returns only the lines the session does not hold yet. diff --git a/tests/test_sdk_api.py b/tests/test_sdk_api.py index 598fabcf..ab1c55dc 100644 --- a/tests/test_sdk_api.py +++ b/tests/test_sdk_api.py @@ -329,17 +329,28 @@ def test_get_deployment_logs_after_drops_redelivered_lines_but_keeps_late_arriva assert [e.id for e in events] == ["15-l", "3-c"] -def test_get_deployment_logs_after_empty_anchor_reads_from_head(): +def test_get_deployment_logs_zero_after_anchor_reads_from_head(): api = MagicMock() api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page(_log_event("1-a", 1000)) client = CentMLClient(api) - events = client.get_deployment_logs(123, 2, pod="pod-a", after=[]) + events = client.get_deployment_logs(123, 2, pod="pod-a", after=0) assert [e.id for e in events] == ["1-a"] call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args assert call.kwargs["fetch_newer"] is True - assert call.kwargs["timestamp"] is None + assert call.kwargs["timestamp"] == 0 + + +def test_get_deployment_logs_rejects_empty_anchor_lists(): + api = MagicMock() + client = CentMLClient(api) + + for kwargs in ({"before": []}, {"after": []}): + with pytest.raises(ValueError): + client.get_deployment_logs(123, 2, pod="pod-a", **kwargs) + + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_not_called() def test_get_deployment_logs_after_empty_page_signals_nothing_new(): @@ -366,7 +377,7 @@ def test_get_deployment_logs_rejects_before_and_after_together(): client = CentMLClient(api) with pytest.raises(ValueError): - client.get_deployment_logs(123, 2, pod="pod-a", before=[_log_event("1-a", 1000)], after=[]) + client.get_deployment_logs(123, 2, pod="pod-a", before=[_log_event("1-a", 1000)], after=0) api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_not_called() @@ -523,7 +534,7 @@ def test_get_deployment_logs_range_open_ended_reads_full_history(): assert [e.id for e in events] == ["1-a", "2-b"] calls = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args_list - assert calls[0].kwargs["timestamp"] is None and calls[0].kwargs["fetch_newer"] is True + assert calls[0].kwargs["timestamp"] == 0 and calls[0].kwargs["fetch_newer"] is True def test_get_deployment_logs_range_merges_all_pods_by_id(): @@ -533,7 +544,7 @@ def test_get_deployment_logs_range_merges_all_pods_by_id(): ) def pages(**kwargs): - if kwargs["timestamp"] is not None: + if kwargs["timestamp"]: return _log_page() if kwargs["pod"] == "pod-a": return _log_page(_log_event("1-a", 1000), _log_event("3-a", 3000)) @@ -585,3 +596,40 @@ def test_get_deployment_logs_range_rejects_inverted_window(): client.get_deployment_logs_range(123, 2, pod="pod-a", start_time=2000, end_time=1000) api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_not_called() + + +def test_log_session_merges_by_production_shaped_ids(): + def _real_id(ms, suffix): + return f"{ms * 10**6:019d}-{suffix}" + + api = MagicMock() + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + _log_page(_log_event(_real_id(1000, "aa"), 1000), _log_event(_real_id(2000, "bb"), 2000)), + _log_page(_log_event(_real_id(1500, "ll"), 1500), _log_event(_real_id(3000, "cc"), 3000)), + ] + session = _session(api) + session.fetch_newer() + session.fetch_newer() + + assert [e.timestamp for e in session.events] == [1000, 1500, 2000, 3000] + + +def test_log_session_long_window_keeps_boundary_and_dedup_correct(): + api = MagicMock() + old_events = [_log_event(f"{i}-x", i) for i in range(1, 4)] # far older than the retention window + recent = _log_event("900000-y", 900_000) + api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [ + # Look-behind re-delivers the recent held line next to a fresh one. + _log_page(_log_event("900000-y", 900_000), _log_event("901000-z", 901_000)), + _log_page(_log_event("0-w", 500)), + ] + session = _session(api, events=old_events + [recent]) + + delta = session.fetch_newer() + assert [e.id for e in delta] == ["901000-z"] + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is True and call.kwargs["timestamp"] == 900_000 + + session.fetch_older() + call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args + assert call.kwargs["fetch_newer"] is False and call.kwargs["timestamp"] == 1 From f714c3aba66aceab057a1c97f7b50f3c8656854c Mon Sep 17 00:00:00 2001 From: Honglin Cao Date: Tue, 1 Sep 2026 15:11:32 -0400 Subject: [PATCH 3/3] refactor(sdk): share the trailing-anchor slice with the range walk MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit get_deployment_logs_range paged with the full accumulated list, the same O(window) pattern just removed from the session; both now anchor via a shared _recent_anchor helper. Example comment reworded — the [] shorthand read as an input since empty anchor lists became a ValueError. Signed-off-by: Honglin Cao --- centml/sdk/api.py | 22 +++++++++++++--------- examples/sdk/get_deployment_logs.py | 4 ++-- 2 files changed, 15 insertions(+), 11 deletions(-) diff --git a/centml/sdk/api.py b/centml/sdk/api.py index fc388da2..afa8b26d 100644 --- a/centml/sdk/api.py +++ b/centml/sdk/api.py @@ -31,6 +31,17 @@ LOG_DEDUP_RETENTION_MS = 300_000 +def _recent_anchor(events: list) -> list: + """Trailing slice within LOG_DEDUP_RETENTION_MS of the newest event — everything + an after anchor contributes (the exclusive boundary and the look-behind dedup + ids), without rescanning the whole accumulated window on every page.""" + cutoff = events[-1].timestamp - LOG_DEDUP_RETENTION_MS + first_recent = len(events) + while first_recent > 0 and events[first_recent - 1].timestamp >= cutoff: + first_recent -= 1 + return events[first_recent:] + + @dataclass(frozen=True) class DeploymentLogEvent: """One log line with its pod attached — logs_v4 events carry no pod name, so @@ -305,7 +316,7 @@ def get_deployment_logs_range( while True: # after is exclusive, so start_time - 1 admits lines at start_time itself; # start_time 0 (or None) means the whole window — scan from the head. - anchor: Union[list, int] = events if events else (start_time - 1 if start_time else 0) + anchor: Union[list, int] = _recent_anchor(events) if events else (start_time - 1 if start_time else 0) page = self.get_deployment_logs( deployment_id, revision_number, pod_name, after=anchor, max_lines=MAX_LOG_PAGE_LINES ) @@ -379,18 +390,11 @@ def fetch_newer(self, max_lines: int = DEFAULT_LOG_PAGE_LINES) -> list: tailing. Rare late arrivals sort into the window below its newest lines.""" if not self._events: return self.fetch_older(max_lines=max_lines) - # Only the trailing retention window matters to the primitive (boundary + - # look-behind dedup ids); passing it keeps long tails from scanning the - # whole window on every poll. - cutoff = self._events[-1].timestamp - LOG_DEDUP_RETENTION_MS - first_recent = len(self._events) - while first_recent > 0 and self._events[first_recent - 1].timestamp >= cutoff: - first_recent -= 1 delta = self._client.get_deployment_logs( self._deployment_id, self._revision_number, self._pod, - after=self._events[first_recent:], + after=_recent_anchor(self._events), max_lines=max_lines, ) for event in delta: diff --git a/examples/sdk/get_deployment_logs.py b/examples/sdk/get_deployment_logs.py index 7339a6c7..0dd7ec59 100644 --- a/examples/sdk/get_deployment_logs.py +++ b/examples/sdk/get_deployment_logs.py @@ -49,8 +49,8 @@ def main(): # The same paging is available statelessly via get_deployment_logs, anchored # on events you already hold — useful when you manage storage yourself: # page = cclient.get_deployment_logs(DEPLOYMENT_ID, REVISION_NUMBER, pod=pod) # tail - # older = cclient.get_deployment_logs(..., pod=pod, before=page) # [] = beginning - # newer = cclient.get_deployment_logs(..., pod=pod, after=page) # [] = nothing new + # older = cclient.get_deployment_logs(..., pod=pod, before=page) # empty return = beginning + # newer = cclient.get_deployment_logs(..., pod=pod, after=page) # empty return = nothing new # A specific time window (all pods merged, oldest first, pod on each event): # window = cclient.get_deployment_logs_range( # DEPLOYMENT_ID, REVISION_NUMBER, start_time=t1_ms, end_time=t2_ms