Coverage for packages/dqm-ml-job/src/dqm_ml_job/utils/s3.py: 14%

47 statements  

« prev     ^ index     » next       coverage.py v7.14.1, created at 2026-07-21 08:27 +0000

1"""S3 utilities for DQM-ML job.""" 

2 

3import os 

4from typing import Any 

5 

6from dqm_ml_core.models.global_ import StorageConfig 

7import pyarrow as pa 

8 

9_SIMPLE_KWARGS: list[tuple[str, str]] = [ 

10 ("region", "region"), 

11 ("session_token", "session_token"), 

12 ("anonymous", "anonymous"), 

13 ("request_timeout", "request_timeout"), 

14 ("connect_timeout", "connect_timeout"), 

15 ("scheme", "scheme"), 

16 ("proxy_options", "proxy_options"), 

17 ("tls_ca_file_path", "tls_ca_file_path"), 

18] 

19 

20_ROLE_KWARGS: list[tuple[str, str]] = [ 

21 ("session_name", "session_name"), 

22 ("external_id", "external_id"), 

23] 

24 

25 

26def _build_retry_strategy( 

27 retry_cfg: Any, 

28) -> pa.fs.S3RetryStrategy | None: 

29 """Build an S3 retry strategy from configuration. 

30 

31 Args: 

32 retry_cfg: RetryConfig instance or dict with mode and max_attempts. 

33 

34 Returns: 

35 Configured PyArrow S3 retry strategy, or None if no config. 

36 """ 

37 from dqm_ml_core.models.global_ import RetryConfig 

38 

39 if retry_cfg is None: 

40 return None 

41 if isinstance(retry_cfg, dict): 

42 retry_cfg = RetryConfig.model_validate(retry_cfg) 

43 max_attempts = retry_cfg.max_attempts 

44 if retry_cfg.mode == "default": 

45 return pa.fs.AwsDefaultS3RetryStrategy(max_attempts=max_attempts) 

46 return pa.fs.AwsStandardS3RetryStrategy(max_attempts=max_attempts) 

47 

48 

49def _apply_simple_kwargs(storage_config: StorageConfig, kwargs: dict[str, Any]) -> None: 

50 """Apply simple 1-to-1 config-attribute-to-kwargs mappings.""" 

51 for attr, kwarg in _SIMPLE_KWARGS: 

52 val = getattr(storage_config, attr, None) 

53 if kwarg == "anonymous": 

54 if val: 

55 kwargs[kwarg] = True 

56 elif val is not None: 

57 kwargs[kwarg] = val 

58 

59 

60def _apply_role_kwargs(storage_config: StorageConfig, kwargs: dict[str, Any]) -> None: 

61 """Apply role-based access kwargs.""" 

62 kwargs["role_arn"] = storage_config.role_arn 

63 for attr, kwarg in _ROLE_KWARGS: 

64 val = getattr(storage_config, attr, None) 

65 if val is not None: 

66 kwargs[kwarg] = val 

67 kwargs["load_frequency"] = storage_config.load_frequency 

68 

69 

70def get_s3_filesystem( 

71 storage_config: StorageConfig, 

72) -> pa.fs.S3FileSystem | None: 

73 """Create and return an S3 filesystem instance from StorageConfig. 

74 

75 Credentials fall back to environment variables: 

76 S3_ACCESS_KEY, S3_SECRET_KEY, S3_ENDPOINT, S3_REGION 

77 

78 Args: 

79 storage_config: Validated StorageConfig instance. 

80 

81 Returns: 

82 Configured S3 filesystem or None if credentials not available. 

83 """ 

84 os.environ["AWS_RESPONSE_CHECKSUM_VALIDATION"] = storage_config.checksum_validation 

85 

86 access_key = storage_config.access_key or os.getenv("S3_ACCESS_KEY") 

87 secret_key = storage_config.secret_key or os.getenv("S3_SECRET_KEY") 

88 

89 if not access_key or not secret_key: 

90 return None 

91 

92 kwargs: dict[str, Any] = { 

93 "access_key": access_key, 

94 "secret_key": secret_key, 

95 "endpoint_override": storage_config.endpoint or os.getenv("S3_ENDPOINT", ""), 

96 } 

97 

98 _apply_simple_kwargs(storage_config, kwargs) 

99 

100 if storage_config.role_arn: 

101 _apply_role_kwargs(storage_config, kwargs) 

102 

103 if storage_config.retry: 

104 kwargs["retry_strategy"] = _build_retry_strategy(storage_config.retry) 

105 

106 if os.getenv("ENVIRONMENT") == "mock": 

107 kwargs.setdefault("background_writes", False) 

108 kwargs.setdefault("retry_strategy", pa.fs.AwsStandardS3RetryStrategy(max_attempts=10)) 

109 

110 return pa.fs.S3FileSystem(**kwargs)