from datetime import timedelta from typing import Optional, Tuple import boto3 from config import S3_ARTIFACTS_BUCKET from constants import LOCKFILE_NAME, LOCKFILE_TTL from .key import Key __all__ = ["S3Lock", "LockError"] class LockError(Exception): pass class S3Lock: _s3_artifacts_bucket: str = S3_ARTIFACTS_BUCKET _lock_ttl: timedelta = timedelta(minutes=LOCKFILE_TTL) _s3_key: str = LOCKFILE_NAME def __init__(self, key: Key): self._s3_client = boto3.client("s3") self._key = key def __enter__(self): self.get_lock() def __exit__(self, exc_type, exc_val, exc_tb): self.release_lock() def _check(self) -> Tuple[bool, Optional[Key]]: try: obj = self._s3_client.get_object(Bucket=self._s3_artifacts_bucket, Key=self._s3_key) except self._s3_client.exceptions.NoSuchKey: return True, None except Exception as e: raise LockError("Something went wrong while acquiring a lock") from e else: return False, self._decode_key(obj["Body"].read()) @staticmethod def _encode_key(key: Key) -> bytes: return str(key).encode() @staticmethod def _decode_key(key: bytes) -> Key: return Key.from_string(key.decode()) def get_lock(self): is_available, old_key = self._check() if not is_available and self._key.timestamp.to_datetime() - old_key.timestamp.to_datetime() < self._lock_ttl: raise LockError(f"Lock object already acquired by '{old_key}'") try: self._s3_client.put_object( Bucket=self._s3_artifacts_bucket, Key=self._s3_key, Body=self._encode_key(self._key) ) except Exception as e: raise LockError("Something went wrong while acquiring a lock") from e def release_lock(self): try: self._s3_client.delete_object(Bucket=self._s3_artifacts_bucket, Key=self._s3_key) except Exception as e: raise LockError("Something went wrong while releasing the lock") from e