diff --git a/README.md b/README.md index 733836d2..ad263d11 100644 --- a/README.md +++ b/README.md @@ -44,6 +44,19 @@ function = client.functions.invoke( ) print(function.status, function.version, function.data) +logs = client.logs.search( + "00000000-0000-4000-8000-000000000001", + {"resource": {"type": "function"}, "limit": 100}, +) +for event in logs.data: + print(event["timestamp"], event["body"]) + +activity = client.logs.activity( + "00000000-0000-4000-8000-000000000001", + {"resource": {"type": "function"}, "bucket_count": 24}, +) +print(activity.total) + rows = client.database("main").from_("items").select("*").eq("slug", "a").execute() bucket = client.storage.from_("assets") @@ -141,6 +154,11 @@ then the anonymous key. The immutable result includes the response body, status, headers, and `X-Volcano-Version`. A function's own non-2xx response is returned when the version header proves it ran; platform failures raise typed SDK errors. +`logs.search()` returns an immutable page of retained runtime or deployment log +events. Pass `next_cursor` back as `cursor` to continue a search. `logs.activity()` +returns immutable time buckets using the same resource selector and query syntax. +Both methods require an active user session. + Database builders are immutable, so you can safely reuse a base query. Chain `neq()`, `gt()`, `gte()`, `lt()`, and `lte()` for comparison filters: diff --git a/pyproject.toml b/pyproject.toml index cf36e158..6704e442 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -94,6 +94,7 @@ ignore = [ "src/volcano_sdk/auth.py" = ["SLF001"] "src/volcano_sdk/database.py" = ["SLF001"] "src/volcano_sdk/functions.py" = ["SLF001"] +"src/volcano_sdk/logs.py" = ["SLF001"] "src/volcano_sdk/locks.py" = ["SLF001"] "src/volcano_sdk/realtime.py" = ["ANN401", "SLF001"] "src/volcano_sdk/storage.py" = ["SLF001"] diff --git a/scripts/generate_openapi.py b/scripts/generate_openapi.py index e391d44e..4feb3bd9 100755 --- a/scripts/generate_openapi.py +++ b/scripts/generate_openapi.py @@ -20,6 +20,7 @@ "delete_storage_object.py", "download_storage_object.py", "force_release_project_lock.py", + "get_project_log_activity.py", "get_project_lock.py", "invoke_function.py", "list_storage_objects.py", @@ -28,6 +29,7 @@ "release_project_lock.py", "renew_project_lock.py", "resolve_function_for_invocation.py", + "search_project_logs.py", "update_storage_object_visibility.py", "upload_storage_object.py", } diff --git a/src/volcano_sdk/__init__.py b/src/volcano_sdk/__init__.py index fce3575c..2530026b 100644 --- a/src/volcano_sdk/__init__.py +++ b/src/volcano_sdk/__init__.py @@ -22,6 +22,8 @@ LinkedOAuthProvider, LockLease, LockState, + LogActivityResponse, + LogSearchResponse, OAuthProviderName, OAuthProviderTokenStatus, Session, @@ -48,6 +50,8 @@ "LinkedOAuthProvider", "LockLease", "LockState", + "LogActivityResponse", + "LogSearchResponse", "NotFoundError", "OAuthProviderName", "OAuthProviderTokenStatus", diff --git a/src/volcano_sdk/_transport.py b/src/volcano_sdk/_transport.py index bc388595..231f77b1 100644 --- a/src/volcano_sdk/_transport.py +++ b/src/volcano_sdk/_transport.py @@ -69,6 +69,18 @@ release_project_lock, renew_project_lock, ) +from ._generated.api.logs.get_project_log_activity import ( + _build_response as build_log_activity_response, +) +from ._generated.api.logs.get_project_log_activity import ( + _get_kwargs as log_activity_kwargs, +) +from ._generated.api.logs.search_project_logs import ( + _build_response as build_log_search_response, +) +from ._generated.api.logs.search_project_logs import ( + _get_kwargs as log_search_kwargs, +) from ._generated.api.o_auth_authentication import auth_o_auth_exchange from ._generated.api.o_auth_authentication.auth_link_o_auth_provider import ( _get_kwargs as link_oauth_provider_kwargs, @@ -153,6 +165,8 @@ from ._generated.models.get_o_auth_provider_token_response_200 import ( GetOAuthProviderTokenResponse200, ) +from ._generated.models.log_activity_request import LogActivityRequest +from ._generated.models.log_search_request import LogSearchRequest from ._generated.models.project_lock_lease_request import ProjectLockLeaseRequest from ._generated.models.refresh_o_auth_provider_token_response_200 import ( RefreshOAuthProviderTokenResponse200, @@ -1531,6 +1545,46 @@ def invoke_function( ) return self._raw_response(response) + def search_project_logs( + self, + *, + authorization: str, + project_id: str, + request: Mapping[str, JSONValue], + ) -> TransportResponse: + plain_request = cast("dict[str, Any]", _plain_json(request)) + with self._client(authorization) as client: + request_kwargs = log_search_kwargs( + UUID(project_id), body=LogSearchRequest.from_dict(plain_request) + ) + request_kwargs["json"] = plain_request + raw_response = client.get_httpx_client().request(**request_kwargs) + response = build_log_search_response( + client=client, + response=raw_response, + ) + return self._response(response) + + def get_project_log_activity( + self, + *, + authorization: str, + project_id: str, + request: Mapping[str, JSONValue], + ) -> TransportResponse: + plain_request = cast("dict[str, Any]", _plain_json(request)) + with self._client(authorization) as client: + request_kwargs = log_activity_kwargs( + UUID(project_id), body=LogActivityRequest.from_dict(plain_request) + ) + request_kwargs["json"] = plain_request + raw_response = client.get_httpx_client().request(**request_kwargs) + response = build_log_activity_response( + client=client, + response=raw_response, + ) + return self._response(response) + def acquire_project_lock( self, *, diff --git a/src/volcano_sdk/client.py b/src/volcano_sdk/client.py index bcb4f64b..572054dc 100644 --- a/src/volcano_sdk/client.py +++ b/src/volcano_sdk/client.py @@ -11,6 +11,7 @@ from .database import Database from .functions import Functions from .locks import Locks +from .logs import Logs from .models import ( AuthChangeEvent, AuthStateCallback, @@ -83,6 +84,7 @@ def __init__( ) self.auth = Auth(self) self.functions = Functions(self) + self.logs = Logs(self) self.storage = Storage(self) self.locks = Locks(self) if _realtime_client_factory is None: diff --git a/src/volcano_sdk/logs.py b/src/volcano_sdk/logs.py new file mode 100644 index 00000000..005c92a7 --- /dev/null +++ b/src/volcano_sdk/logs.py @@ -0,0 +1,135 @@ +"""Project runtime log facade.""" + +from __future__ import annotations + +from collections.abc import Mapping +from typing import TYPE_CHECKING, Protocol, cast + +from ._transport import TransportResponse, invoke, response_payload +from .models import JSONValue, LogActivityResponse, LogSearchResponse + +if TYPE_CHECKING: + from .client import VolcanoClient + +_INVALID_PROJECT_ID = "project_id must be a non-empty string" +_INVALID_LOG_REQUEST = "Log request must be a mapping" +_INVALID_LOG_RESPONSE = "Expected a complete log response" + + +class LogsTransport(Protocol): + """Transport operations required by the logs facade.""" + + def search_project_logs( + self, + *, + authorization: str, + project_id: str, + request: Mapping[str, JSONValue], + ) -> TransportResponse: + """Search a project's retained logs.""" + ... + + def get_project_log_activity( + self, + *, + authorization: str, + project_id: str, + request: Mapping[str, JSONValue], + ) -> TransportResponse: + """Get bucketed log activity for a project.""" + ... + + +class Logs: + """Search retained project logs and activity.""" + + def __init__(self, client: VolcanoClient) -> None: + """Bind log reads to a Volcano client.""" + self._client = client + + def search( + self, + project_id: str, + request: Mapping[str, JSONValue], + ) -> LogSearchResponse: + """Search retained logs for one project resource type.""" + project_id, request = _log_request(project_id, request) + transport = cast("LogsTransport", self._client._transport) + response = invoke( + transport.search_project_logs, + authorization=self._client._session_token(), + project_id=project_id, + request=request, + ) + return _search_response(response_payload(response, 200)) + + def activity( + self, + project_id: str, + request: Mapping[str, JSONValue], + ) -> LogActivityResponse: + """Get bucketed activity for one project resource type.""" + project_id, request = _log_request(project_id, request) + transport = cast("LogsTransport", self._client._transport) + response = invoke( + transport.get_project_log_activity, + authorization=self._client._session_token(), + project_id=project_id, + request=request, + ) + return _activity_response(response_payload(response, 200)) + + +def _log_request( + project_id: object, + request: object, +) -> tuple[str, Mapping[str, JSONValue]]: + if not isinstance(project_id, str) or not project_id.strip(): + raise ValueError(_INVALID_PROJECT_ID) + if not isinstance(request, Mapping): + raise TypeError(_INVALID_LOG_REQUEST) + return project_id, cast("Mapping[str, JSONValue]", request) + + +def _response_values(payload: object) -> Mapping[str, object]: + if not isinstance(payload, Mapping): + raise TypeError(_INVALID_LOG_RESPONSE) + return cast("Mapping[str, object]", payload) + + +def _response_data(values: Mapping[str, object]) -> tuple[Mapping[str, JSONValue], ...]: + raw_data = values.get("data") + if not isinstance(raw_data, list): + raise TypeError(_INVALID_LOG_RESPONSE) + data = cast("list[object]", raw_data) + if any(not isinstance(item, Mapping) for item in data): + raise TypeError(_INVALID_LOG_RESPONSE) + return tuple(cast("Mapping[str, JSONValue]", item) for item in data) + + +def _search_response(payload: object) -> LogSearchResponse: + values = _response_values(payload) + limit = values.get("limit") + has_more = values.get("has_more") + next_cursor = values.get("next_cursor") + if ( + not isinstance(limit, int) + or isinstance(limit, bool) + or not isinstance(has_more, bool) + or (next_cursor is not None and not isinstance(next_cursor, str)) + ): + raise TypeError(_INVALID_LOG_RESPONSE) + return LogSearchResponse( + data=_response_data(values), + limit=limit, + has_more=has_more, + next_cursor=next_cursor, + ) + + +def _activity_response(payload: object) -> LogActivityResponse: + values = _response_values(payload) + total = values.get("total") + if not isinstance(total, int) or isinstance(total, bool): + raise TypeError(_INVALID_LOG_RESPONSE) + return LogActivityResponse(data=_response_data(values), total=total) diff --git a/src/volcano_sdk/models.py b/src/volcano_sdk/models.py index f9cf5df8..e1a6cad0 100644 --- a/src/volcano_sdk/models.py +++ b/src/volcano_sdk/models.py @@ -211,6 +211,40 @@ def __post_init__(self) -> None: ) +@dataclass(frozen=True, slots=True) +class LogSearchResponse: + """Immutable page returned by a project log search.""" + + data: tuple[Mapping[str, JSONValue], ...] = field(hash=False) + limit: int + has_more: bool + next_cursor: str | None = None + + def __post_init__(self) -> None: + """Defensively freeze log events owned by this value.""" + object.__setattr__( + self, + "data", + tuple(_freeze_json(event) for event in self.data), + ) + + +@dataclass(frozen=True, slots=True) +class LogActivityResponse: + """Immutable bucketed project log activity.""" + + data: tuple[Mapping[str, JSONValue], ...] = field(hash=False) + total: int + + def __post_init__(self) -> None: + """Defensively freeze activity buckets owned by this value.""" + object.__setattr__( + self, + "data", + tuple(_freeze_json(bucket) for bucket in self.data), + ) + + @dataclass(frozen=True, slots=True) class UploadSession: """Server-created state for a resumable storage upload.""" diff --git a/tests/unit/test_generated_transport.py b/tests/unit/test_generated_transport.py index 2f27be92..8bc4aed6 100644 --- a/tests/unit/test_generated_transport.py +++ b/tests/unit/test_generated_transport.py @@ -103,6 +103,89 @@ def handle(request: httpx.Request) -> httpx.Response: } +def test_generated_transport_reads_project_logs() -> None: + requests: list[httpx.Request] = [] + + def handle(request: httpx.Request) -> httpx.Response: + requests.append(request) + if request.url.path.endswith("/search"): + return httpx.Response( + 200, + json={ + "data": [ + { + "id": "event-1", + "timestamp": "2026-09-02T12:00:00Z", + "body": "ready", + "resource": { + "type": "function", + "id": "00000000-0000-4000-8000-000000000040", + }, + } + ], + "limit": 25, + "has_more": False, + }, + ) + return httpx.Response( + 200, + json={ + "data": [ + { + "start_time": "2026-09-02T12:00:00Z", + "end_time": "2026-09-02T12:05:00Z", + "counts": { + "levels": {"info": 2}, + "regions": {"us-east-1": 2}, + "resource_ids": {"00000000-0000-4000-8000-000000000040": 2}, + }, + "total": 2, + } + ], + "total": 2, + }, + ) + + transport = GeneratedTransport( + api_url="https://api.test.volcano.dev", + httpx_transport=httpx.MockTransport(handle), + ) + project_id = "00000000-0000-4000-8000-000000000001" + resource = {"resource": {"type": "function"}} + + search = transport.search_project_logs( + authorization="access-token", + project_id=project_id, + request={**resource, "limit": 25, "query": "misspelled-filter"}, + ) + activity = transport.get_project_log_activity( + authorization="access-token", + project_id=project_id, + request={ + "resource": {"type": "function", "unknown_selector": True}, + "bucket_count": 12, + }, + ) + + assert search.payload["data"][0]["id"] == "event-1" + assert activity.payload["total"] == 2 + assert [request.url.path for request in requests] == [ + f"/projects/{project_id}/logs/search", + f"/projects/{project_id}/logs/activity", + ] + assert [json.loads(request.content) for request in requests] == [ + {**resource, "limit": 25, "query": "misspelled-filter"}, + { + "resource": {"type": "function", "unknown_selector": True}, + "bucket_count": 12, + }, + ] + assert all( + request.headers["authorization"] == "Bearer access-token" + for request in requests + ) + + def test_generated_transport_exchanges_an_oauth_code() -> None: requests: list[httpx.Request] = [] diff --git a/tests/unit/test_generation.py b/tests/unit/test_generation.py index 48e4bdb9..970b4002 100644 --- a/tests/unit/test_generation.py +++ b/tests/unit/test_generation.py @@ -28,10 +28,12 @@ def test_generate_emits_required_contract_operations(tmp_path: Path) -> None: "update_storage_object_visibility.py", "acquire_project_lock.py", "force_release_project_lock.py", + "get_project_log_activity.py", "get_project_lock.py", "invoke_function.py", "release_project_lock.py", "resolve_function_for_invocation.py", + "search_project_logs.py", "renew_project_lock.py", } @@ -41,6 +43,7 @@ def test_generate_emits_required_contract_operations(tmp_path: Path) -> None: [ "force_release_project_lock.py", "invoke_function.py", + "search_project_logs.py", "update_storage_object_visibility.py", ], ) diff --git a/tests/unit/test_logs.py b/tests/unit/test_logs.py new file mode 100644 index 00000000..4d619267 --- /dev/null +++ b/tests/unit/test_logs.py @@ -0,0 +1,177 @@ +from __future__ import annotations + +from dataclasses import dataclass +from types import MappingProxyType +from typing import TYPE_CHECKING, Any, cast + +import pytest + +from volcano_sdk import ServerError, Session, VolcanoClient + +if TYPE_CHECKING: + from collections.abc import Mapping + + from volcano_sdk._transport import Transport + from volcano_sdk.models import JSONValue + + +@dataclass(frozen=True) +class FakeResponse: + status_code: int + payload: Any + headers: dict[str, str] + content: bytes = b"" + + +class FakeLogsTransport: + def __init__(self) -> None: + self.calls: list[tuple[str, dict[str, Any]]] = [] + self.search_response = FakeResponse( + 200, + { + "data": [ + { + "id": "event-1", + "timestamp": "2026-09-02T12:00:00Z", + "body": {"message": "ready"}, + "resource": {"type": "function", "id": "function-1"}, + } + ], + "limit": 25, + "has_more": True, + "next_cursor": "cursor-2", + }, + {}, + ) + self.activity_response = FakeResponse( + 200, + { + "data": [ + { + "start_time": "2026-09-02T12:00:00Z", + "end_time": "2026-09-02T12:05:00Z", + "counts": { + "levels": {"info": 2}, + "regions": {"us-east-1": 2}, + "resource_ids": {"function-1": 2}, + }, + "total": 2, + } + ], + "total": 2, + }, + {}, + ) + + def search_project_logs(self, **kwargs: Any) -> FakeResponse: + self.calls.append(("searchProjectLogs", kwargs)) + return self.search_response + + def get_project_log_activity(self, **kwargs: Any) -> FakeResponse: + self.calls.append(("getProjectLogActivity", kwargs)) + return self.activity_response + + +def logs_client(transport: FakeLogsTransport) -> VolcanoClient: + client = VolcanoClient( + anon_key="anon-key", + _transport=cast("Transport", transport), + ) + client.auth.set_session( + Session( + access_token="access-token", refresh_token="refresh-token", user_id="user-1" + ) + ) + return client + + +def test_logs_search_returns_an_immutable_page() -> None: + transport = FakeLogsTransport() + request: Mapping[str, JSONValue] = { + "resource": {"type": "function"}, + "limit": 25, + } + + result = logs_client(transport).logs.search("project-1", request) + + assert result.limit == 25 + assert result.has_more is True + assert result.next_cursor == "cursor-2" + assert result.data[0]["body"] == {"message": "ready"} + assert isinstance(result.data[0], MappingProxyType) + assert isinstance(result.data[0]["body"], MappingProxyType) + assert transport.calls == [ + ( + "searchProjectLogs", + { + "authorization": "access-token", + "project_id": "project-1", + "request": request, + }, + ) + ] + + +def test_logs_activity_returns_immutable_buckets() -> None: + transport = FakeLogsTransport() + request: Mapping[str, JSONValue] = { + "resource": {"type": "function"}, + "bucket_count": 12, + } + + result = logs_client(transport).logs.activity("project-1", request) + + assert result.total == 2 + assert result.data[0]["counts"] == { + "levels": {"info": 2}, + "regions": {"us-east-1": 2}, + "resource_ids": {"function-1": 2}, + } + assert isinstance(result.data[0]["counts"], MappingProxyType) + assert transport.calls[0] == ( + "getProjectLogActivity", + { + "authorization": "access-token", + "project_id": "project-1", + "request": request, + }, + ) + + +@pytest.mark.parametrize("project_id", ["", " "]) +def test_logs_rejects_an_empty_project_id(project_id: str) -> None: + transport = FakeLogsTransport() + + with pytest.raises(ValueError, match="project_id"): + logs_client(transport).logs.search( + project_id, + {"resource": {"type": "function"}}, + ) + + assert transport.calls == [] + + +def test_logs_rejects_a_non_mapping_request() -> None: + transport = FakeLogsTransport() + + with pytest.raises(TypeError, match="mapping"): + logs_client(transport).logs.activity("project-1", cast("Any", [])) + + assert transport.calls == [] + + +def test_logs_maps_platform_errors() -> None: + transport = FakeLogsTransport() + transport.search_response = FakeResponse( + 503, + {"error": "logs unavailable", "code": "logs_unavailable"}, + {}, + ) + + with pytest.raises(ServerError, match="logs unavailable") as raised: + logs_client(transport).logs.search( + "project-1", + {"resource": {"type": "function"}}, + ) + + assert raised.value.code == "logs_unavailable"