diff --git a/src/pyfwapi/apiconnection.py b/src/pyfwapi/apiconnection.py index d6f774a..702b1aa 100644 --- a/src/pyfwapi/apiconnection.py +++ b/src/pyfwapi/apiconnection.py @@ -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. @@ -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, diff --git a/src/pyfwapi/change/stateful.py b/src/pyfwapi/change/stateful.py index eed6996..742b33d 100644 --- a/src/pyfwapi/change/stateful.py +++ b/src/pyfwapi/change/stateful.py @@ -1,3 +1,4 @@ +import asyncio import dataclasses import typing as t from dataclasses import dataclass, field @@ -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.""" @@ -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) @@ -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: @@ -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 @@ -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]