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

1"""S3-specific connection pooling implementation.""" 

2 

3from dataclasses import dataclass 

4from typing import Optional 

5 

6from .base import BasePoolConfig 

7 

8 

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 

19 

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}" 

23 

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) 

27 

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 ) 

40 

41 

42async def create_aioboto3_session(config: S3PoolConfig): 

43 """Create an aioboto3 session for S3 operations.""" 

44 import aioboto3 

45 

46 # Create session with credentials if provided 

47 session_config = {} 

48 

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 

57 

58 # Create and return the session 

59 return aioboto3.Session(**session_config) 

60 

61 

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)