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
3 changes: 3 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ print(public_url)
state = client.locks.get("build")
print(state.held)
lease = client.locks.acquire("build", ttl=30)
lease = client.locks.renew("build", lease, ttl=30)
client.locks.release("build", lease)
```

Expand All @@ -112,6 +113,8 @@ a best-effort abort and raises the original error.
part number to replace that part.
`locks.get()` returns immutable lock availability, expiry, and fencing-token
state without acquiring the lock.
`locks.renew()` returns a new immutable lease and leaves the previous value
unchanged.
`get_upload_session()` returns immutable progress and uploaded-part metadata for
resuming an interrupted upload.
`complete_upload_session()` assembles the uploaded parts and returns the stored
Expand Down
19 changes: 19 additions & 0 deletions src/volcano_sdk/_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@
acquire_project_lock,
get_project_lock,
release_project_lock,
renew_project_lock,
)
from ._generated.api.o_auth_authentication import auth_o_auth_exchange
from ._generated.api.o_auth_authentication.auth_link_o_auth_provider import (
Expand Down Expand Up @@ -1520,6 +1521,24 @@ def get_project_lock(
)
return self._response(response)

def renew_project_lock(
self,
*,
authorization: str,
key: str,
ttl: int,
token: str,
) -> TransportResponse:
with self._client(authorization) as client:
response = renew_project_lock.sync_detailed(
key,
client=client,
body=ProjectLockLeaseRequest(ttl_seconds=ttl),
x_volcano_lock_token=cast("UUID", token),
x_volcano_request_id=cast("UUID", str(uuid4())),
)
return self._response(response)

def release_project_lock(
self,
*,
Expand Down
33 changes: 33 additions & 0 deletions src/volcano_sdk/locks.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,21 @@ def get_project_lock(
...


class LockRenewTransport(Protocol):
"""Transport capability required to renew a lock."""

def renew_project_lock(
self,
*,
authorization: str,
key: str,
ttl: int,
token: str,
) -> object:
"""Renew one project-scoped lock."""
...


def _parse_datetime(value: object) -> datetime | None:
if value is None:
return None
Expand Down Expand Up @@ -77,6 +92,24 @@ def acquire(self, key: str, *, ttl: int) -> LockLease:
fencing_token=payload.get("fencing_token"),
)

def renew(self, key: str, lease: LockLease, *, ttl: int) -> LockLease:
"""Renew a lock lease and return its immutable replacement."""
transport = cast("LockRenewTransport", self._client._transport)
response = invoke(
transport.renew_project_lock,
authorization=self._client._service_token(),
key=key,
ttl=ttl,
token=lease.token,
)
payload = response_payload(response, 200)
return LockLease(
key=key,
token=lease.token,
expires_at=_parse_datetime(payload.get("expires_at")),
fencing_token=payload.get("fencing_token"),
)

def release(self, key: str, lease: LockLease) -> None:
"""Release a lock lease."""
response = invoke(
Expand Down
43 changes: 43 additions & 0 deletions tests/unit/test_facade.py
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,13 @@ def get_project_lock(self, **kwargs: Any) -> FakeResponse:
},
)

def renew_project_lock(self, **kwargs: Any) -> FakeResponse:
self.calls.append(("renewProjectLock", kwargs))
return FakeResponse(
200,
{"expires_at": "2026-08-26T12:01:00Z", "fencing_token": 7},
)


def anon_key_with_project_id(project_id: str | None) -> str:
payload = {} if project_id is None else {"project_id": project_id}
Expand Down Expand Up @@ -452,6 +459,42 @@ def test_locks_gets_immutable_current_state() -> None:
]


def test_locks_renews_a_lease_without_mutating_the_original() -> None:
transport = FakeTransport()
client = VolcanoClient(
anon_key="anon-key",
service_key="service-key",
_transport=transport,
)
lease = LockLease(
key="build",
token="00000000-0000-4000-8000-000000000001",
expires_at=datetime(2026, 8, 26, 12, 0, 30, tzinfo=UTC),
fencing_token=7,
)

renewed = client.locks.renew("build", lease, ttl=60)

assert renewed == LockLease(
key="build",
token=lease.token,
expires_at=datetime(2026, 8, 26, 12, 1, tzinfo=UTC),
fencing_token=7,
)
assert lease.expires_at == datetime(2026, 8, 26, 12, 0, 30, tzinfo=UTC)
assert transport.calls == [
(
"renewProjectLock",
{
"authorization": "service-key",
"key": "build",
"ttl": 60,
"token": lease.token,
},
)
]


def test_storage_remove_accepts_one_path() -> None:
transport = FakeTransport()
client = VolcanoClient(anon_key="anon-key", _transport=transport)
Expand Down
34 changes: 34 additions & 0 deletions tests/unit/test_generated_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -1531,6 +1531,40 @@ def handle(request: httpx.Request) -> httpx.Response:
assert requests[0].headers["x-volcano-request-id"]


def test_generated_transport_renews_a_project_lock() -> None:
requests: list[httpx.Request] = []

def handle(request: httpx.Request) -> httpx.Response:
requests.append(request)
return httpx.Response(
200,
json={"expires_at": "2026-08-26T12:01:00Z", "fencing_token": 7},
)

transport = GeneratedTransport(
api_url="https://api.test.volcano.dev",
httpx_transport=httpx.MockTransport(handle),
)

response = transport.renew_project_lock(
authorization="service-key",
key="build:queue",
ttl=60,
token="00000000-0000-4000-8000-000000000001",
)

assert response.payload["fencing_token"] == 7
assert len(requests) == 1
assert requests[0].method == "PATCH"
assert requests[0].url.path == "/locks/build:queue/lease"
assert json.loads(requests[0].content) == {"ttl_seconds": 60}
assert requests[0].headers["authorization"] == "Bearer service-key"
assert requests[0].headers["x-volcano-lock-token"] == (
"00000000-0000-4000-8000-000000000001"
)
assert requests[0].headers["x-volcano-request-id"]


def test_generated_transport_moves_a_storage_object() -> None:
requests: list[httpx.Request] = []

Expand Down
Loading