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
14 changes: 12 additions & 2 deletions src/pyfwapi/apiconnection.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,13 @@ class APIConnection:
# entity types, like Asset or Rendition.

def __init__(
self, endpoint_url: str, *, client_id: str, client_secret: str
self,
endpoint_url: str,
*,
client_id: str,
client_secret: str,
max_rate: float = 1,
time_period: float = 0.8,
) -> None:
"""
Connect to an instance of the FotoWare API.
Expand All @@ -60,11 +66,15 @@ def __init__(
endpoint_url: URL of the endpoint, e.g. `https://myorg.example.org`
client_id: the registered non-interactive application's `client_id`
client_secret: the application's secret
max_rate: number of requests allowed per `time_period` (client-side
rate limit; defaults to the conservative 1 request / 0.8 s).
time_period: period in seconds over which `max_rate` requests are
allowed.
"""

self.HOST = endpoint_url.removesuffix("/")
self.TOKEN_ENDPOINT = f"{self.HOST}/fotoweb/oauth2/token"
self.rate_limit = aiolimiter.AsyncLimiter(1, 0.8)
self.rate_limit = aiolimiter.AsyncLimiter(max_rate, time_period)

self.client = AsyncOAuth2Client(
client_id=client_id,
Expand Down
75 changes: 61 additions & 14 deletions src/pyfwapi/change/stateful.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import asyncio
import dataclasses
import typing as t
from dataclasses import dataclass, field
Expand Down Expand Up @@ -69,21 +70,45 @@ class BaseChangeManager:
requests, and can keep track of background tasks.

Consider using ChangeManager for a higher-level API.

Args:
max_concurrent_tasks: how many independent changes (uploads, moves,
metadata patches) may be committed in parallel.
max_concurrent_chunks: how many chunks of a single upload may be in
flight at the same time. The FotoWare Upload API explicitly allows
chunks to be uploaded in any order and in parallel.
"""

def __init__(self) -> None:
def __init__(
self, *, max_concurrent_tasks: int = 4, max_concurrent_chunks: int = 4
) -> None:
self.tasks: dict[bytes, ChangeTask] = dict()
self.task_statuslocation: dict[bytes, str] = dict()
self._task_semaphore = asyncio.Semaphore(max_concurrent_tasks)
self._chunk_semaphore = asyncio.Semaphore(max_concurrent_chunks)

def add_task(self, task: ChangeTask):
self.tasks[task.id] = task

async def commit(self, *, conn: APIConnection, await_done: bool = True):
"""Commit changes ready to commit."""
for task in self.tasks.values():
if task.status != "uncommitted":
continue
await self.commit_uncommitted(task, conn=conn)
"""Commit changes ready to commit.

Independent tasks are committed in parallel, bounded by
``max_concurrent_tasks``. All requests still pass through the
connection-level rate limiter.
"""

async def _commit(task: ChangeTask):
async with self._task_semaphore:
await self.commit_uncommitted(task, conn=conn)

await asyncio.gather(
*(
_commit(task)
for task in self.tasks.values()
if task.status == "uncommitted"
)
)

async def commit_uncommitted(self, ch: ChangeTask, *, conn: APIConnection):
"""Commit a single uncommitted ChangeTask."""
Expand All @@ -105,13 +130,16 @@ async def commit_uncommitted(self, ch: ChangeTask, *, conn: APIConnection):
self.task_statuslocation[ch.id] = f"/fotoweb/api/uploads/{task.id}/status"

async def check_submitted(self, *, conn: APIConnection):
"""Check the processing status of backgrounded tasks, like moves and uploads."""
for task in self.tasks.values():
if task.status != "submitted":
continue
"""Check the processing status of backgrounded tasks, like moves and uploads.

All submitted tasks are polled in parallel, bounded by
``max_concurrent_tasks``.
"""

async def _check(task: ChangeTask):
location = self.task_statuslocation.get(task.id)
if location is None:
continue
return

if isinstance(task.change, MoveRequest):
r = await conn.GET(location)
Expand All @@ -133,6 +161,18 @@ async def check_submitted(self, *, conn: APIConnection):
task.status = "failed"
pyfwapiLog.warn(f"Upload failed (fn:{task.change.filename})")

async def _check_bounded(task: ChangeTask):
async with self._task_semaphore:
await _check(task)

await asyncio.gather(
*(
_check_bounded(task)
for task in self.tasks.values()
if task.status == "submitted"
)
)

async def patch_metadata(
self, item: MetadataRequest, *, conn: APIConnection
) -> bool:
Expand Down Expand Up @@ -190,8 +230,15 @@ async def upload_asset(self, item: UploadRequest, *, conn: APIConnection):

upload_info = BatchUploadInfo.model_validate_json(r.content)

for i in range(upload_info.numChunks):
await self._upload_asset_chunk(i, upload_info, item, conn=conn)
# Chunks may be uploaded in any order and in parallel (per the FotoWare
# Upload API docs). Concurrency is bounded so we stay well within
# server-side rate limits; every request also passes through the
# connection-level rate limiter.
async def _upload_chunk(i: int):
async with self._chunk_semaphore:
await self._upload_asset_chunk(i, upload_info, item, conn=conn)

await asyncio.gather(*(_upload_chunk(i) for i in range(upload_info.numChunks)))

return upload_info

Expand All @@ -205,7 +252,7 @@ async def _upload_asset_chunk(
):
"""Upload a chunk of a new asset."""
chunk_offset = i * upload_info.chunkSize
chunk_size = min([upload_info.chunkSize, item.filesize])
chunk_size = min(upload_info.chunkSize, item.filesize - chunk_offset)
chunk_end = chunk_offset + chunk_size

bytes_part = item.contents[chunk_offset:chunk_end]
Expand Down
Loading