Coverage for src/dataknobs_data/pooling/s3.py: 53%
34 statements
« prev ^ index » next coverage.py v7.10.3, created at 2025-08-17 19:59 -0500
« prev ^ index » next coverage.py v7.10.3, created at 2025-08-17 19:59 -0500
1"""S3-specific connection pooling implementation."""
3from dataclasses import dataclass
4from typing import Optional
6from .base import BasePoolConfig
9@dataclass
10class S3PoolConfig(BasePoolConfig):
11 """Configuration for S3 connection pools."""
12 bucket: str
13 prefix: str = ""
14 region_name: Optional[str] = None
15 aws_access_key_id: Optional[str] = None
16 aws_secret_access_key: Optional[str] = None
17 aws_session_token: Optional[str] = None
18 endpoint_url: Optional[str] = None
20 def to_connection_string(self) -> str:
21 """Convert to connection string (not used for S3, but required by base)."""
22 return f"s3://{self.bucket}/{self.prefix}"
24 def to_hash_key(self) -> tuple:
25 """Create a hashable key for this configuration."""
26 return (self.bucket, self.prefix, self.region_name, self.endpoint_url)
28 @classmethod
29 def from_dict(cls, config: dict) -> "S3PoolConfig":
30 """Create from configuration dictionary."""
31 return cls(
32 bucket=config.get("bucket"),
33 prefix=config.get("prefix", ""),
34 region_name=config.get("region_name"),
35 aws_access_key_id=config.get("aws_access_key_id"),
36 aws_secret_access_key=config.get("aws_secret_access_key"),
37 aws_session_token=config.get("aws_session_token"),
38 endpoint_url=config.get("endpoint_url")
39 )
42async def create_aioboto3_session(config: S3PoolConfig):
43 """Create an aioboto3 session for S3 operations."""
44 import aioboto3
46 # Create session with credentials if provided
47 session_config = {}
49 if config.aws_access_key_id:
50 session_config["aws_access_key_id"] = config.aws_access_key_id
51 if config.aws_secret_access_key:
52 session_config["aws_secret_access_key"] = config.aws_secret_access_key
53 if config.aws_session_token:
54 session_config["aws_session_token"] = config.aws_session_token
55 if config.region_name:
56 session_config["region_name"] = config.region_name
58 # Create and return the session
59 return aioboto3.Session(**session_config)
62async def validate_s3_session(session, config: S3PoolConfig) -> None:
63 """Validate an S3 session by checking bucket access."""
64 async with session.client("s3", endpoint_url=config.endpoint_url) as s3:
65 # Try to head the bucket to verify access
66 await s3.head_bucket(Bucket=config.bucket)