Skip to content

Commit 0793a86

Browse files
committed
feat(sdk): migrate deployment logs to logs_v4 and add pod listing
Signed-off-by: Honglin Cao <hocao@nvidia.com>
1 parent ad8e09b commit 0793a86

4 files changed

Lines changed: 234 additions & 67 deletions

File tree

README.md

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,17 @@ Use `python examples/sdk/get_clusters.py` and
5353
Creating the example reserves GPU capacity and may incur usage charges. It does not
5454
delete the deployment automatically.
5555

56+
### Deployment logs SDK example
57+
58+
Logs are read per pod. Discover pod names with `get_deployment_pods()` (terminated
59+
pods still within log retention are included), then fetch that pod's logs oldest-first
60+
with `get_deployment_logs()` — as a flat list, a lazy generator (`stream=True`), or a
61+
live tail that keeps polling for new lines (`follow=True`):
62+
63+
```bash
64+
python examples/sdk/get_deployment_logs.py
65+
```
66+
5667
### Un-installation
5768

5869
To uninstall `centml`, simply do:

centml/sdk/api.py

Lines changed: 53 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,6 @@
1+
import time
12
from contextlib import contextmanager
3+
from typing import Dict, List, Optional
24

35
import platform_api_python_client
46
from platform_api_python_client import (
@@ -20,6 +22,12 @@
2022

2123
STATUS_V3_DEPLOYMENT_TYPES = {DeploymentType.INFERENCE_V3, DeploymentType.CSERVE_V3}
2224

25+
DEFAULT_LOG_PAGE_LINES = 100 # server-side default for max_lines
26+
DEFAULT_LOG_POLL_INTERVAL_SECONDS = 2.0
27+
# The server re-delivers a ~15s look-behind window on follow polls; retain seen event
28+
# ids well past that so long-running follows dedupe correctly without unbounded memory.
29+
LOG_DEDUP_RETENTION_MS = 300_000
30+
2331

2432
class CentMLClient:
2533
def __init__(self, api):
@@ -187,43 +195,69 @@ def get_deployment_revisions(self, deployment_id: int):
187195
deployment_id=deployment_id
188196
).results
189197

198+
def get_deployment_pods(self, deployment_id: int, revision_number: int) -> List[str]:
199+
"""List pods that have logged for a deployment revision, including terminated
200+
pods still within log retention. A fresh deployment may return an empty list."""
201+
return self._api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get(
202+
deployment_id=deployment_id, revision_number=revision_number
203+
).pods
204+
205+
# pylint: disable=R0917
190206
def get_deployment_logs(
191207
self,
192208
deployment_id: int,
193209
revision_number: int,
194-
start_time: int,
195-
end_time: int,
196-
line_count: int = 100,
197-
start_from_head: bool = True,
210+
pod: str,
211+
max_lines: int = DEFAULT_LOG_PAGE_LINES,
198212
stream: bool = False,
213+
follow: bool = False,
214+
poll_interval: float = DEFAULT_LOG_POLL_INTERVAL_SECONDS,
199215
):
200-
"""Fetch logs for a deployment within a time window, handling pagination automatically.
216+
"""Fetch one pod's logs oldest-first, handling pagination and deduplication automatically.
201217
202-
start_time and end_time are Unix timestamps in milliseconds.
203-
Use get_deployment_revisions() to find the current revision number.
218+
Use get_deployment_pods() to discover pod names and get_deployment_revisions()
219+
to find the current revision number.
204220
205221
If stream=True, returns a generator that yields events as each page is fetched.
206-
If stream=False (default), returns a flat list of all events.
222+
If follow=True, returns an endless generator that keeps polling for new lines
223+
every poll_interval seconds (implies streaming); rare late-arriving lines may
224+
then be yielded slightly out of timestamp order, as with `kubectl logs -f`.
225+
If neither (default), returns a flat list of all events logged so far.
207226
"""
208227

209228
def _iter_events():
210-
next_page_token = None
229+
seen_event_ids: Dict[str, int] = {}
230+
newest_timestamp: Optional[int] = None
211231
while True:
212-
response = self._api.get_deployment_logs_v3_deployments_logs_v3_deployment_id_revision_number_get(
232+
response = self._api.get_deployment_logs_v4_logs_deployment_id_revision_number_get(
213233
deployment_id=deployment_id,
214234
revision_number=revision_number,
215-
start_time=start_time,
216-
end_time=end_time,
217-
next_page_token=next_page_token,
218-
start_from_head=start_from_head,
219-
line_count=line_count,
235+
pod=pod,
236+
fetch_newer=True,
237+
timestamp=newest_timestamp,
238+
max_lines=max_lines,
220239
)
221-
yield from response.events
222-
next_page_token = response.next_page_token
223-
if not next_page_token:
240+
# Follow polls re-deliver a look-behind window of already-seen lines
241+
# (late-arrival protection); the event id is the deduplication key.
242+
has_newer_events = False
243+
for event in response.events:
244+
if event.id in seen_event_ids:
245+
continue
246+
seen_event_ids[event.id] = event.timestamp
247+
if newest_timestamp is None or event.timestamp > newest_timestamp:
248+
has_newer_events = True
249+
yield event
250+
if has_newer_events:
251+
newest_timestamp = max(seen_event_ids.values())
252+
cutoff = newest_timestamp - LOG_DEDUP_RETENTION_MS
253+
seen_event_ids = {i: t for i, t in seen_event_ids.items() if t >= cutoff}
254+
elif follow:
255+
time.sleep(poll_interval)
256+
else:
257+
# No event newer than the boundary: history is exhausted.
224258
break
225259

226-
if stream:
260+
if stream or follow:
227261
return _iter_events()
228262

229263
return list(_iter_events())

examples/sdk/get_deployment_logs.py

Lines changed: 21 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -1,69 +1,42 @@
1-
from datetime import datetime, timezone, timedelta
1+
from datetime import datetime, timezone
22

33
from centml.sdk.api import get_centml_client
44

55
# --- Configuration ---
66
DEPLOYMENT_ID = 1234 # Replace with your deployment ID
77
REVISION_NUMBER = 10
8-
HOURS_BACK = 1 # Fetch logs from the last N hours
8+
FOLLOW = False # Set True to keep polling for new log lines (like `kubectl logs -f`)
99

1010

11-
def format_event(event: dict) -> str:
12-
timestamp_ms = (
13-
event.get("timestamp")
14-
or event.get("time")
15-
or event.get("ts")
16-
or ""
17-
)
18-
message = (
19-
event.get("message")
20-
or event.get("msg")
21-
or event.get("log")
22-
or str(event)
23-
)
24-
if timestamp_ms:
25-
ts = datetime.fromtimestamp(int(timestamp_ms) / 1000, tz=timezone.utc).isoformat()
26-
return f"[{ts}] {message}"
27-
return message
11+
def format_event(event) -> str:
12+
ts = datetime.fromtimestamp(event.timestamp / 1000, tz=timezone.utc).isoformat()
13+
return f"[{ts}] {event.message}"
2814

2915

3016
def main():
31-
stream = True
32-
end_time = int(datetime.now(timezone.utc).timestamp() * 1000)
33-
start_time = end_time - int(timedelta(hours=HOURS_BACK).total_seconds() * 1000)
34-
35-
print(f"Fetching logs for deployment {DEPLOYMENT_ID}")
36-
print(
37-
f"Time window: "
38-
f"{datetime.fromtimestamp(start_time / 1000, tz=timezone.utc).isoformat()} → "
39-
f"{datetime.fromtimestamp(end_time / 1000, tz=timezone.utc).isoformat()}"
40-
)
41-
print()
42-
4317
with get_centml_client() as cclient:
44-
if stream:
45-
# Streaming: print events as each page arrives
18+
# Logs are read per pod: discover the pods that have logged for this revision
19+
# (terminated pods within log retention are included).
20+
pods = cclient.get_deployment_pods(DEPLOYMENT_ID, REVISION_NUMBER)
21+
if not pods:
22+
print("No pods have logged for this revision yet.")
23+
return
24+
25+
pod = pods[0]
26+
print(f"Fetching logs for deployment {DEPLOYMENT_ID} revision {REVISION_NUMBER}, pod {pod}\n")
27+
28+
if FOLLOW:
29+
# Endless generator: yields history oldest-first, then keeps polling for new lines
4630
for event in cclient.get_deployment_logs(
47-
deployment_id=DEPLOYMENT_ID,
48-
revision_number=REVISION_NUMBER,
49-
start_time=start_time,
50-
end_time=end_time,
51-
start_from_head=False,
52-
stream=stream,
31+
deployment_id=DEPLOYMENT_ID, revision_number=REVISION_NUMBER, pod=pod, follow=True
5332
):
5433
print(format_event(event))
5534
else:
56-
# Batch: collect all events then process
57-
events = cclient.get_deployment_logs(
58-
deployment_id=DEPLOYMENT_ID,
59-
revision_number=REVISION_NUMBER,
60-
start_time=start_time,
61-
end_time=end_time,
62-
start_from_head=False,
63-
)
35+
# Batch: collect all events logged so far, oldest first
36+
events = cclient.get_deployment_logs(deployment_id=DEPLOYMENT_ID, revision_number=REVISION_NUMBER, pod=pod)
6437

6538
if not events:
66-
print("No logs found in the given time window.")
39+
print("No logs found for this pod.")
6740
return
6841

6942
print(f"Found {len(events)} log entries:\n")

tests/test_sdk_api.py

Lines changed: 149 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,6 @@
1+
from itertools import islice
12
from types import SimpleNamespace
3+
from typing import Generator
24
from unittest.mock import MagicMock, patch
35

46
import platform_api_python_client
@@ -209,3 +211,150 @@ def test_delete_hardware_instance_delegates_to_platform_client():
209211

210212
assert response is expected_response
211213
api.delete_hardware_instance_hardware_instances_hardware_instance_id_delete.assert_called_once_with(123)
214+
215+
216+
def _log_event(event_id, timestamp, message="line"):
217+
return SimpleNamespace(id=event_id, timestamp=timestamp, message=message)
218+
219+
220+
def _log_page(*events):
221+
return SimpleNamespace(events=list(events))
222+
223+
224+
def test_generated_client_exposes_logs_v4_contract():
225+
assert hasattr(
226+
platform_api_python_client.EXTERNALApi, "get_deployment_logs_v4_logs_deployment_id_revision_number_get"
227+
)
228+
assert hasattr(
229+
platform_api_python_client.EXTERNALApi, "get_deployment_pods_deployments_pods_deployment_id_revision_number_get"
230+
)
231+
232+
233+
def test_get_deployment_pods_returns_pod_names():
234+
api = MagicMock()
235+
api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get.return_value = SimpleNamespace(
236+
pods=["pod-a", "pod-b"]
237+
)
238+
client = CentMLClient(api)
239+
240+
assert client.get_deployment_pods(123, 2) == ["pod-a", "pod-b"]
241+
242+
api.get_deployment_pods_deployments_pods_deployment_id_revision_number_get.assert_called_once_with(
243+
deployment_id=123, revision_number=2
244+
)
245+
246+
247+
def test_get_deployment_logs_paginates_oldest_first_until_no_newer_events():
248+
api = MagicMock()
249+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [
250+
_log_page(_log_event("1-a", 1000), _log_event("2-b", 2000)),
251+
_log_page(_log_event("3-c", 3000)),
252+
_log_page(),
253+
]
254+
client = CentMLClient(api)
255+
256+
events = client.get_deployment_logs(123, 2, pod="pod-a")
257+
258+
assert [e.id for e in events] == ["1-a", "2-b", "3-c"]
259+
calls = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args_list
260+
assert len(calls) == 3
261+
# First page scans from the head of the log window; later pages pass the newest
262+
# seen timestamp verbatim (the boundary is exclusive server-side).
263+
assert [c.kwargs["timestamp"] for c in calls] == [None, 2000, 3000]
264+
assert all(c.kwargs["fetch_newer"] is True for c in calls)
265+
assert all(c.kwargs["deployment_id"] == 123 and c.kwargs["revision_number"] == 2 for c in calls)
266+
assert all(c.kwargs["pod"] == "pod-a" for c in calls)
267+
268+
269+
def test_get_deployment_logs_returns_empty_list_when_pod_has_no_logs():
270+
api = MagicMock()
271+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page()
272+
client = CentMLClient(api)
273+
274+
assert client.get_deployment_logs(123, 2, pod="pod-a") == []
275+
276+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.assert_called_once()
277+
278+
279+
def test_get_deployment_logs_deduplicates_look_behind_redelivery():
280+
api = MagicMock()
281+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [
282+
_log_page(_log_event("1-a", 1000), _log_event("2-b", 2000)),
283+
# Look-behind re-delivers 2000 alongside fresh lines; a late arrival shows up at 1500.
284+
_log_page(_log_event("15-l", 1500), _log_event("2-b", 2000), _log_event("3-c", 3000)),
285+
# Re-delivery-only page: no event newer than the boundary means history is done.
286+
_log_page(_log_event("3-c", 3000)),
287+
]
288+
client = CentMLClient(api)
289+
290+
events = client.get_deployment_logs(123, 2, pod="pod-a")
291+
292+
assert [e.id for e in events] == ["1-a", "2-b", "15-l", "3-c"]
293+
assert api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_count == 3
294+
295+
296+
def test_get_deployment_logs_passes_max_lines_through():
297+
api = MagicMock()
298+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.return_value = _log_page()
299+
client = CentMLClient(api)
300+
301+
client.get_deployment_logs(123, 2, pod="pod-a", max_lines=7)
302+
303+
call = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args
304+
assert call.kwargs["max_lines"] == 7
305+
306+
307+
def test_get_deployment_logs_stream_fetches_pages_lazily():
308+
api = MagicMock()
309+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [
310+
_log_page(_log_event("1-a", 1000), _log_event("2-b", 2000)),
311+
_log_page(_log_event("3-c", 3000)),
312+
_log_page(),
313+
]
314+
client = CentMLClient(api)
315+
316+
events = client.get_deployment_logs(123, 2, pod="pod-a", stream=True)
317+
318+
assert isinstance(events, Generator)
319+
assert [e.id for e in islice(events, 2)] == ["1-a", "2-b"]
320+
assert api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_count == 1
321+
assert [e.id for e in events] == ["3-c"]
322+
assert api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_count == 3
323+
324+
325+
def test_get_deployment_logs_follow_polls_for_new_events_and_dedupes():
326+
api = MagicMock()
327+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [
328+
_log_page(_log_event("1-a", 1000)),
329+
_log_page(), # nothing new yet -> sleep, poll again with the same boundary
330+
_log_page(_log_event("1-a", 1000), _log_event("2-b", 2000)),
331+
]
332+
client = CentMLClient(api)
333+
334+
with patch("centml.sdk.api.time.sleep") as mock_sleep:
335+
events = client.get_deployment_logs(123, 2, pod="pod-a", follow=True, poll_interval=0.5)
336+
assert [e.id for e in islice(events, 2)] == ["1-a", "2-b"]
337+
338+
calls = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args_list
339+
assert [c.kwargs["timestamp"] for c in calls] == [None, 1000, 1000]
340+
mock_sleep.assert_called_once_with(0.5)
341+
342+
343+
def test_get_deployment_logs_follow_sleeps_on_redelivery_only_page():
344+
api = MagicMock()
345+
api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.side_effect = [
346+
_log_page(_log_event("1-a", 1000)),
347+
# Look-behind re-delivery with nothing new: the boundary must not advance
348+
# and the poll must sleep instead of spinning.
349+
_log_page(_log_event("1-a", 1000)),
350+
_log_page(_log_event("1-a", 1000), _log_event("2-b", 2000)),
351+
]
352+
client = CentMLClient(api)
353+
354+
with patch("centml.sdk.api.time.sleep") as mock_sleep:
355+
events = client.get_deployment_logs(123, 2, pod="pod-a", follow=True, poll_interval=0.5)
356+
assert [e.id for e in islice(events, 2)] == ["1-a", "2-b"]
357+
358+
calls = api.get_deployment_logs_v4_logs_deployment_id_revision_number_get.call_args_list
359+
assert [c.kwargs["timestamp"] for c in calls] == [None, 1000, 1000]
360+
mock_sleep.assert_called_once_with(0.5)

0 commit comments

Comments
 (0)