Skip to content
Merged
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
18 changes: 18 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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:

Expand Down
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
Expand Down
2 changes: 2 additions & 0 deletions scripts/generate_openapi.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
}
Expand Down
4 changes: 4 additions & 0 deletions src/volcano_sdk/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
LinkedOAuthProvider,
LockLease,
LockState,
LogActivityResponse,
LogSearchResponse,
OAuthProviderName,
OAuthProviderTokenStatus,
Session,
Expand All @@ -48,6 +50,8 @@
"LinkedOAuthProvider",
"LockLease",
"LockState",
"LogActivityResponse",
"LogSearchResponse",
"NotFoundError",
"OAuthProviderName",
"OAuthProviderTokenStatus",
Expand Down
54 changes: 54 additions & 0 deletions src/volcano_sdk/_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Avoid parsing malformed requests before sending them

When a request omits the required resource field, LogSearchRequest.from_dict() raises a bare KeyError before the raw-body override or HTTP request runs, so Hosting cannot return its 400 response and the facade cannot map it to ValidationError; malformed known timestamps or resource IDs similarly leak ValueError. The post-fix override preserves unknown fields only for bodies that the generated parser already accepts, and the activity path has the same ordering, so construct the generated URL without parsing the caller body or translate these local failures consistently.

Useful? React with 👍 / 👎.

)
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,
*,
Expand Down
2 changes: 2 additions & 0 deletions src/volcano_sdk/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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:
Expand Down
135 changes: 135 additions & 0 deletions src/volcano_sdk/logs.py
Original file line number Diff line number Diff line change
@@ -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)
34 changes: 34 additions & 0 deletions src/volcano_sdk/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down
Loading
Loading