From d380a290ad3df07b8a618f2110986af9faea18b8 Mon Sep 17 00:00:00 2001 From: Redmer Kronemeijer <12477216+redmer@users.noreply.github.com> Date: Thu, 24 Sep 2026 12:50:29 +0200 Subject: [PATCH] feat: parallelize ChangeManager operations with bounded concurrency Commit tasks, poll submitted tasks, and upload chunks concurrently via semaphores, and make the connection rate limit configurable. Also fix the upload chunk size calculation for final chunks. --- src/pyfwapi/apiconnection.py | 14 ++++++- src/pyfwapi/change/stateful.py | 75 +++++++++++++++++++++++++++------- 2 files changed, 73 insertions(+), 16 deletions(-) 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]