capo (caporegime)
Community-driven AWS SDK for Python. Manage AWS services like a capo.
Features
- async and sync — Use the same API for both async and sync code, with support for asyncio and trio.
- WASM support — Runs in WASM environments via Pyodide.
- typed and documented — Fully typed and documented for a better developer experience.
- simple inputs — No need to import separate parameter classes; nested inputs are fully typed via TypedDicts.
- generated from Smithy models — The SDK is generated from the official Smithy models describing AWS APIs, ensuring accuracy and consistency.
- zero runtime overhead — Codegen produces dedicated serialization and deserialization code for each operation, avoiding reflection.
- interchangeable input/output — Input and output types use the same TypedDicts, so you can pass a response directly as input when appropriate.
- built on zapros — A modern HTTP client for Python that abstracts HTTP semantics from the transport implementation.
- fast to import — import time stays flat even for huge services like EC2.
- interceptors — Hook into the operation request/response lifecycle to inspect, log, or modify calls.
Warning The API should be mostly stable, but some breaking changes may occur as the SDK is still in early development. We strongly recommend pinning the version before the first major release.
Installation
Every service is a standalone package, so you only install the ones you actually use. Packages are published under the capo- prefix:
uv add capo-ec2 # Amazon EC2 (also covers Amazon VPC) uv add capo-s3 # Amazon S3 uv add capo-iam # AWS IAM uv add capo-rds # Amazon RDS uv add capo-lambda # AWS Lambda uv add capo-cloudwatch # Amazon CloudWatch uv add capo-elastic-load-balancing # Elastic Load Balancing (ELB) uv add capo-route-53 # Amazon Route 53 uv add capo-cloudfront # Amazon CloudFront
Note that Amazon VPC has no separate package — its API is part of EC2, so capo-ec2 covers it.
Services not yet on PyPI
Not every service is on PyPI yet. We publish incrementally because of PyPI's limits on new projects and total upload size. If the service you need isn't published, please open an issue so we can prioritize publishing it.
Async/Sync Usage
All the services have both async and sync client. The async client is simply prefixed with Async (e.g. AsyncS3Client).
from capo_s3 import AsyncS3Client async def main(): async with AsyncS3Client() as s3: response = await s3.create_bucket("capo") print(response)
Sync usage is the same, just without the Async prefix and without await:
from capo_s3 import S3Client with S3Client() as s3: response = s3.create_bucket("capo") print(response)
Input/Output Types
No matter how nested the input types are, you don't need to import any additional classes. All input and output types are fully typed via TypedDicts.
from capo_s3 import AsyncS3Client async with AsyncS3Client() as s3_client: response = await s3_client.create_bucket( bucket="some_bucket", create_bucket_configuration={ "location": { "name": "location_name" } } ) print(response["location"])
The output types are also TypeDicts, we do this to make the input/output types interchangeable. You can pass the output member to another operation that accepts the same type as input.
Error Handling
The SDK raises exceptions for errors returned by the API. Catch them to handle failures gracefully.
from capo_s3 import AsyncS3Client from capo_s3.errors import NoSuchUpload async with AsyncS3Client() as s3: try: await s3.abort_multipart_upload() except NoSuchUpload as e: print(f"Error: {e}") print(e.data) # additional error data
Note that the service errors (errors returned by AWS) might have additional data of any shape stored in the data attribute, which is a TypedDict. You can access it to get more information about the error.
All the service errors that an operation can raise are documented in that operation method's docstring.
Configuration
Everything is a client argument:
from capo_s3 import AsyncS3Client AsyncS3Client(region="eu-central-1") # S3-compatible server (RustFS, MinIO, LocalStack, ...) AsyncS3Client( endpoint="http://localhost:9000", force_path_style=True, region="us-east-1", credentials={"access_key": "rustfsadmin", "secret_key": "rustfsadmin"}, ) AsyncS3Client(use_fips=True) AsyncS3Client(use_dual_stack=True) AsyncS3Client(accelerate=True) # S3 Transfer Acceleration AsyncS3Client(retry_max_attempts=5)
Or a per-call override:
async with AsyncS3Client(region="eu-central-1") as s3: await s3.head_bucket("us-bucket", config_overrides={"region": "us-west-2"}) await s3.head_bucket( "local", config_overrides={"endpoint": "http://localhost:9000", "force_path_style": True}, )
Or environment variables:
export AWS_REGION=eu-central-1 export AWS_ENDPOINT_URL_S3=http://localhost:9000 # this service only export AWS_ENDPOINT_URL=http://localhost:4566 # every service export AWS_USE_FIPS_ENDPOINT=true export AWS_USE_DUALSTACK_ENDPOINT=true export AWS_MAX_ATTEMPTS=5 export AWS_PROFILE=dev # or AWS_DEFAULT_PROFILE export AWS_CONFIG_FILE=~/.aws/config
Or the active profile in ~/.aws/config:
[profile dev] region = eu-central-1 endpoint_url = http://localhost:4566 use_fips_endpoint = false use_dualstack_endpoint = false max_attempts = 5 services = local [services local] s3 = endpoint_url = http://localhost:9000
First match wins: config_overrides, then the client argument, then the environment variable, then the profile. For the endpoint the order within the environment and the profile is AWS_ENDPOINT_URL_S3, AWS_ENDPOINT_URL, the [services] entry, then the profile's endpoint_url.
Authentication
With no arguments, credentials are resolved by the default chain when the first operation runs:
from capo_s3 import AsyncS3Client async with AsyncS3Client() as s3: await s3.list_buckets()
The chain, in order:
AWS_ACCESS_KEY_ID,AWS_SECRET_ACCESS_KEY,AWS_SESSION_TOKEN- Profile with
role_arn(STS AssumeRole throughsource_profileorcredential_source) AWS_WEB_IDENTITY_TOKEN_FILE+AWS_ROLE_ARN, or the profile'sweb_identity_token_file(EKS IRSA)- Profile with
sso_*(IAM Identity Center, afteraws sso login) ~/.aws/credentialsand~/.aws/configstatic keys- ECS / EKS Pod Identity container endpoint (
AWS_CONTAINER_CREDENTIALS_*) - EC2 instance metadata (IMDSv2, off with
AWS_EC2_METADATA_DISABLED=true)
Steps 2, 3 and 4 need the sso extra:
Static credentials:
AsyncS3Client(credentials={"access_key": "AKIA...", "secret_key": "..."}) AsyncS3Client(credentials={"access_key": "ASIA...", "secret_key": "...", "session_token": "..."})
A specific profile, either via the environment or a provider:
AWS_PROFILE=dev python app.py
from zapros import AsyncClient from capo_s3 import ( AssumeRoleCredentialsProvider, AsyncS3Client, ProfileCredentialsProvider, SsoCredentialsProvider, ) # static keys in the profile AsyncS3Client(credentials_provider=ProfileCredentialsProvider(profile="dev")) # sso_* profile AsyncS3Client(credentials_provider=SsoCredentialsProvider(AsyncClient(), profile="dev")) # role_arn profile AsyncS3Client(credentials_provider=AssumeRoleCredentialsProvider(AsyncClient(), profile="dev"))
Providers that call AWS take a zapros client: AsyncClient() for Async* clients, Client() for sync ones.
Your own chain:
from zapros import AsyncClient from capo_s3 import ( AsyncS3Client, CachedProvider, ChainedProvider, EnvCredentialsProvider, WebIdentityCredentialsProvider, ) chain = ChainedProvider(EnvCredentialsProvider(), WebIdentityCredentialsProvider(AsyncClient())) AsyncS3Client(credentials_provider=CachedProvider(chain))
Credentials are resolved on every request, so wrap providers that do I/O in CachedProvider. It re-resolves 60 seconds before expiration.
Your own provider:
from datetime import datetime, timedelta, timezone from capo_s3 import AsyncS3Client, CachedProvider, Credentials, CredentialsProvider class VaultCredentials(CredentialsProvider): def resolve_identity(self) -> Credentials: lease = vault.read("aws/creds/app") return { "access_key": lease["access_key"], "secret_key": lease["secret_key"], "session_token": lease["security_token"], "expiration": datetime.now(timezone.utc) + timedelta(seconds=lease["lease_duration"]), } # optional: override `aresolve_identity` to do async I/O; the default calls `resolve_identity` AsyncS3Client(credentials_provider=CachedProvider(VaultCredentials()))
Resolve once, share across services. Every package exports the same providers and the same Credentials shape:
from zapros import AsyncClient from capo_s3 import AsyncS3Client, SsoCredentialsProvider from capo_sts import AsyncSTSClient credentials = await SsoCredentialsProvider(AsyncClient()).aresolve_identity() async with AsyncS3Client(credentials=credentials) as s3, AsyncSTSClient(credentials=credentials) as sts: print(await sts.get_caller_identity()) print(await s3.list_buckets())
Per call:
await s3.get_object( "bucket", "key", config_overrides={"credentials_provider": ProfileCredentialsProvider(profile="other")} )
Errors:
from capo_s3 import AssumeRoleError, IdentityNotFound, MissingDependencyError, SSOError try: await s3.list_buckets() except IdentityNotFound: ... # no provider in the chain found credentials except SSOError: ... # SSO profile found but the token is unusable: run `aws sso login` except AssumeRoleError: ... # role_arn profile found but AssumeRole failed except MissingDependencyError: ... # the profile needs the `sso` extra
Streaming
If the operation's input is a streaming blob, you can pass any AsyncIterator[bytes] or just a bytes object.
from capo_s3 import AsyncS3Client s3_client = AsyncS3Client() response = await s3_client.put_object("bucket_name", "key", body=b"some binary data")
Or, if you don't want to load the entire blob into memory, you can pass an AsyncIterator[bytes]:
from capo_s3 import AsyncS3Client async def async_iterator(): yield b"capo" response = await s3_client.put_object("bucket_name", "key", body=async_iterator(), content_length=4)
As you might have noticed, we also passed the content_length. That's an AWS requirement when using streaming inputs; it must always know the length of the blob before sending it to AWS.
Note that the stream can be any iterator of bytes; it need not be the file's content. You can stream any data you want, for example, directly from the HTTP response of another service, or from a database, etc.
The catch with a plain iterator is that it can be sent only once. If the request fails after the body was transmitted (a throttling error, a dropped connection), there is nothing left to resend, so the operation is not retried. To get retries for streamed uploads, pass a Body instead: it wraps a source that can be reopened, and every attempt streams a fresh copy. Body.from_path (sync client) and Body.async_from_path (async client) stream a file from disk and take the content_length from the file size, so you don't need to pass it:
from capo_s3 import AsyncS3Client, Body s3_client = AsyncS3Client() response = await s3_client.put_object("bucket_name", "key", body=Body.async_from_path("data.bin"))
For sources other than files, build a Body from an opener — a context manager that yields a (stream, length) pair each time it is entered. The SDK enters it before every attempt and exits it when the operation finishes:
from contextlib import asynccontextmanager from capo_s3 import AsyncS3Client, Body @asynccontextmanager async def open_rows(): rows = await db.fetch_all() # re-read from your data source on every attempt async def chunks(): for row in rows: yield row yield chunks(), sum(len(row) for row in rows) response = await s3_client.put_object("bucket_name", "key", body=Body(open_rows))
The output as mentioned before also can be a stream, in such case, the operation will return a context manager that yield the response, ensuring that the resource is properly closed after the response is consumed.
from capo_s3 import AsyncS3Client s3_client = AsyncS3Client() async with s3_client.get_object("bucket_name", "key") as response: async for chunk in response["body"]: print(chunk)
The event streaming operations are similar, but instead of using AsyncIterator[bytes], they use AsyncIterator[Event], where Event is a TypedDict that represents the event type.
from capo_s3 import AsyncS3Client s3_client = AsyncS3Client() async def main(): async with s3_client.select_object_content( "bucket_name", "key", expression="SELECT * FROM S3Object s WHERE s._1 > 100", expression_type="SQL", input_serialization={ "csv": {"file_header_info": "NONE"}, "compression_type": "NONE", }, output_serialization={"csv": {}}, ) as response: async for event in response["payload"]: if "End" in event: print(event["End"])
Waiters
Waiters poll an operation until a resource reaches a desired state. If the operation supports waiters it will have a wait_until_ prefixed method.
from capo_s3 import AsyncS3Client async with AsyncS3Client() as s3: # Example: wait for bucket_exists await s3.wait_until_bucket_exists(max_wait_time=300)
Pagination
Some operations in this SDK support pagination. If the operation supports pagination it will have an iter_ prefixed method that returns an async iterator.
from capo_s3 import AsyncS3Client async with AsyncS3Client() as s3: # Example: paginate over list_buckets async for item in s3.iter_list_buckets(): print(item)
Presigning
Some operations support presigning, which generates a URL that can be used without credentials. Use the presigned_ prefixed method on the client to get a presigned URL.
from capo_s3 import AsyncS3Client async def main(): async with AsyncS3Client() as s3: # Example: get a presigned URL for delete_object url = s3.presigned_delete_object() print(url)
Note Smithy models don't indicate which operations support presigning, so presigned methods are added by maintainers rather than the code generator. If you notice an operation that should support presigning but has no
presigned_method, please open an issue.
Interceptors
Interceptors let you hook into the operation lifecycle. An operation interceptor receives the operation request and a next callable that invokes the rest of the chain, returning the operation response. You can inspect or modify the request before calling next, and inspect or modify the response after.
import asyncio from typing import Any, Awaitable, Callable from capo_s3 import AsyncOperationRequest, AsyncOperationResponse, AsyncS3Client async def debug_interceptor( request: AsyncOperationRequest[Any], next: Callable[[AsyncOperationRequest[Any]], Awaitable[AsyncOperationResponse]], ): print(request) response = await next(request) print(response) return response async def main(): async with AsyncS3Client(operation_interceptors=[debug_interceptor]) as client: async for item in client.iter_list_buckets(): print(item) asyncio.run(main())
Note An
operation_interceptorsits between the operation request and the operation response — not at the HTTP layer. For HTTP-level interceptors/middlewares, see the Configuring the HTTP client section below.
Configuring the HTTP client
The SDK is built on zapros. You can pass your own HTTP handler via the http_handler argument, including any zapros middleware wrapping a network handler. This is the right place for HTTP-level concerns such as caching, mocks, or custom transports. See the zapros handlers documentation for the available handlers and middlewares.
import asyncio from zapros import AsyncStdNetworkHandler, CacheMiddleware from capo_s3 import AsyncS3Client async def main(): async with AsyncS3Client( http_handler=CacheMiddleware(AsyncStdNetworkHandler()) ) as client: async for item in client.iter_list_buckets(): print(item) asyncio.run(main())
Retrying
The SDK retries failed operations automatically. Retry behaviour follows the Smithy specification: errors are retried based on their is_retryable and is_throttling_error attributes. Throttling errors use a longer base delay. Network-level failures (connection errors and timeouts) are also retried. Non-retryable errors, such as client errors without the @retryable trait, are raised immediately without further attempts.
The number of attempts defaults to 3 and can be changed at the client level via retry_max_attempts, or per call via config_overrides.
from capo_s3 import AsyncS3Client async def main(): async with AsyncS3Client() as s3: # Default: 3 attempts for every operation response = await s3.abort_multipart_upload() # Override per operation response = await s3.abort_multipart_upload(config_overrides={"retry_max_attempts": 5}) # Disable retries for this call response = await s3.abort_multipart_upload(config_overrides={"retry_max_attempts": 1})
Caching
Every download from S3 is billed as data transfer out, about 0.09 USD per GB, so a 5 GB object costs about 0.45 USD each time it is fetched. S3 supports HTTP conditional requests: send the ETag you already have as If-None-Match, and if the object has not changed S3 answers 304 Not Modified with no body and no transfer charge. With the zapros cache middleware the SDK does this for you: the first download is stored locally, every later download is a cheap revalidation, and the body is transferred again only when the object changed.
uv add "zapros[caching]" # sync clients uv add "zapros[caching]" "hishel[async]" # Async* clients
from hishel import AsyncSqliteStorage, CacheOptions, SpecificationPolicy from zapros import AsyncStdNetworkHandler, CacheMiddleware from capo_s3 import AsyncS3Client cache = CacheMiddleware( AsyncStdNetworkHandler(), # S3 requests carry an Authorization header, which a shared cache must not store: use a private one policy=SpecificationPolicy(CacheOptions(shared=False)), # defaults to hishel_cache.db in the working directory storage=AsyncSqliteStorage(database_path="s3-cache.db"), ) async with AsyncS3Client(http_handler=cache) as s3: async with s3.get_object("bucket", "big.bin") as response: # 200: downloaded and stored async for chunk in response["body"]: ... async with s3.get_object("bucket", "big.bin") as response: # 304: body served from the cache async for chunk in response["body"]: ...
On the wire:
| Download | Request | Response | Transfer billed |
|---|---|---|---|
| first | GET /big.bin |
200, 5 GB body |
5 GB |
| later, object unchanged | GET /big.bin + If-None-Match: "<etag>" |
304, no body |
none |
| later, object changed | GET /big.bin + If-None-Match: "<etag>" |
200, new body, stored |
new size |
Sync:
from hishel import CacheOptions, SpecificationPolicy, SyncSqliteStorage from zapros import CacheMiddleware, StdNetworkHandler from capo_s3 import S3Client cache = CacheMiddleware( StdNetworkHandler(), policy=SpecificationPolicy(CacheOptions(shared=False)), storage=SyncSqliteStorage(database_path="s3-cache.db"), ) with S3Client(http_handler=cache) as s3: with s3.get_object("bucket", "big.bin") as response: for chunk in response["body"]: ...
Make sure the cache asks S3 on every download. Without a Cache-Control header on the object, the cache may serve a stored copy without checking. Set the header on the object:
await s3.put_object("bucket", "big.bin", body=data, cache_control="no-cache")
See what the cache did for each operation:
async def cache_log(request, next): response = await next(request) print(response.response.context["caching"]) # {'from_cache': True, 'revalidated': True, 'stored': False, ...} return response async with AsyncS3Client(http_handler=cache, operation_interceptors=[cache_log]) as s3: ...
Notes:
- Only
GETandHEADresponses are cached. Every other operation goes straight to S3. - The cache key is the URL. Range requests (
range=) bypass the cache. - A
304is still a GET request, billed per thousand requests. Only the transfer is free. - Entries stay until evicted.
AsyncSqliteStorage(default_ttl=86400)drops them after that many seconds. - Passing
if_none_match=yourself is not a substitute: a304has no S3 error body, so the operation raisesUnknownServiceError.
WASM (Pyodide)
The packages are pure Python wheels, so micropip installs them:
import micropip await micropip.install("capo-s3")
from capo_s3 import AsyncS3Client async with AsyncS3Client( region="eu-central-1", credentials={"access_key": "ASIA...", "secret_key": "...", "session_token": "..."}, ) as s3: await s3.put_object("bucket", "hello.txt", body=b"hi from the browser") async with s3.get_object("bucket", "hello.txt") as response: async for chunk in response["body"]: print(chunk)
In a page:
<script src="https://cdn.jsdelivr.net/pyodide/v314.0.7/full/pyodide.js"></script> <script type="module"> const pyodide = await loadPyodide(); await pyodide.loadPackage("micropip"); await pyodide.runPythonAsync(` import micropip await micropip.install("capo-s3") from capo_s3 import AsyncS3Client async with AsyncS3Client(region="eu-central-1", credentials={...}) as s3: async for bucket in s3.iter_list_buckets(): print(bucket) `); </script>
There is no ~/.aws and no metadata service in a browser, so pass credentials= or fetch them from your backend:
from datetime import datetime from zapros import AsyncClient from capo_s3 import AsyncS3Client, CachedProvider, Credentials, CredentialsProvider class BackendCredentials(CredentialsProvider): def resolve_identity(self) -> Credentials: raise NotImplementedError("async only") async def aresolve_identity(self) -> Credentials: response = await AsyncClient().get("https://api.example.com/aws-credentials") data = response.json return { "access_key": data["access_key"], "secret_key": data["secret_key"], "session_token": data["session_token"], "expiration": datetime.fromisoformat(data["expiration"]), } AsyncS3Client(region="eu-central-1", credentials_provider=CachedProvider(BackendCredentials()))
Requests go through the browser's fetch, so the bucket needs a CORS rule for your origin:
await s3.put_bucket_cors( "bucket", cors_configuration={ "cors_rules": [ { "allowed_origins": ["https://app.example.com"], "allowed_methods": ["GET", "PUT"], "allowed_headers": ["*"], "expose_headers": ["ETag"], } ] }, )
What is different under Pyodide:
- Only the
Async*clients work;S3Client()raisesRuntimeError. asyncio only, no trio. - The HTTP handler is zapros'
AsyncPyodideHandler, picked automatically. - No threads, so
Body.async_from_pathis unavailable. Passbytes, an async iterator, or aBodywith your own opener. - Node works too; that is how the integration tests run in CI.