Community-driven AWS SDK for Python. Manage AWS services like a capo
See the codeCommunity-driven AWS SDK for Python. Manage AWS services like a capo.
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.
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.
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.
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)
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.
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.
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.
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_TOKENrole_arn (STS AssumeRole through source_profile or credential_source)AWS_WEB_IDENTITY_TOKEN_FILE + AWS_ROLE_ARN, or the profile's web_identity_token_file (EKS IRSA)sso_* (IAM Identity Center, after aws sso login)~/.aws/credentials and ~/.aws/config static keysAWS_CONTAINER_CREDENTIALS_*)AWS_EC2_METADATA_DISABLED=true)Steps 2, 3 and 4 need the sso extra:
uv add "capo-s3[sso]"
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
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 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)
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)
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 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.
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())
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})
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:
GET and HEAD responses are cached. Every other operation goes straight to S3.range=) bypass the cache.304 is still a GET request, billed per thousand requests. Only the transfer is free.AsyncSqliteStorage(default_ttl=86400) drops them after that many seconds.if_none_match= yourself is not a substitute: a 304 has no S3 error body, so the operation raises UnknownServiceError.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:
Async* clients work; S3Client() raises RuntimeError. asyncio only, no trio.AsyncPyodideHandler, picked automatically.Body.async_from_path is unavailable. Pass bytes, an async iterator, or a Body with your own opener.157 commits
18 commits
Python
100.0%
Community-driven AWS SDK for Python. Manage AWS services like a capo
See the codeCommunity-driven AWS SDK for Python. Manage AWS services like a capo.
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.
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.
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.
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)
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.
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.
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.
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_TOKENrole_arn (STS AssumeRole through source_profile or credential_source)AWS_WEB_IDENTITY_TOKEN_FILE + AWS_ROLE_ARN, or the profile's web_identity_token_file (EKS IRSA)sso_* (IAM Identity Center, after aws sso login)~/.aws/credentials and ~/.aws/config static keysAWS_CONTAINER_CREDENTIALS_*)AWS_EC2_METADATA_DISABLED=true)Steps 2, 3 and 4 need the sso extra:
uv add "capo-s3[sso]"
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
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 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)
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)
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 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.
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())
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})
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:
GET and HEAD responses are cached. Every other operation goes straight to S3.range=) bypass the cache.304 is still a GET request, billed per thousand requests. Only the transfer is free.AsyncSqliteStorage(default_ttl=86400) drops them after that many seconds.if_none_match= yourself is not a substitute: a 304 has no S3 error body, so the operation raises UnknownServiceError.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:
Async* clients work; S3Client() raises RuntimeError. asyncio only, no trio.AsyncPyodideHandler, picked automatically.Body.async_from_path is unavailable. Pass bytes, an async iterator, or a Body with your own opener.157 commits
18 commits
Python
100.0%