diff --git a/pyproject.toml b/pyproject.toml index 3b6ad503..ec086624 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "openf1" -version = "1.9.13" +version = "1.9.14" authors = [ { name="Bruno Godefroy" } ] diff --git a/requirements.txt b/requirements.txt index e4155f9b..81b56160 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,10 +1,10 @@ aiohttp==3.12.13 aiomqtt==2.4.0 beautifulsoup4==4.13.4 +boto3==1.43.50 cachetools==5.5.2 fastapi==0.115.13 fastf1_livetiming @ git+https://github.com/br-g/fastf1-livetiming.git#egg=fastf1_livetiming -google-cloud-storage==3.1.1 loguru==0.7.3 lxml==6.0.0 motor==3.7.1 diff --git a/src/openf1/services/ingestor_livetiming/real_time/app.py b/src/openf1/services/ingestor_livetiming/real_time/app.py index 428f46c3..1fe9f4e9 100644 --- a/src/openf1/services/ingestor_livetiming/real_time/app.py +++ b/src/openf1/services/ingestor_livetiming/real_time/app.py @@ -8,10 +8,10 @@ from openf1.services.ingestor_livetiming.core.objects import get_topics from openf1.services.ingestor_livetiming.real_time.processing import ingest_file from openf1.services.ingestor_livetiming.real_time.recording import record_to_file -from openf1.util.gcs import upload_to_gcs_periodically +from openf1.util.storage import upload_to_object_storage_periodically TIMEOUT = 5400 # Terminate job if no data received for 90 minutes (in seconds) -GCS_BUCKET = os.getenv("OPENF1_INGESTOR_LIVETIMING_GCS_BUCKET_RAW") +STORAGE_BUCKET_NAME = os.getenv("OPENF1_INGESTOR_LIVETIMING_STORAGE_BUCKET_RAW") async def main(): @@ -31,15 +31,15 @@ async def main(): ) tasks.append(task_recording) - if GCS_BUCKET: - # Save received raw data to GCS, for debugging - logger.info("Starting periodic GCS upload of raw data") - gcs_filekey = datetime.now(timezone.utc).strftime("%Y/%m/%d/%H:%M:%S.txt") + if STORAGE_BUCKET_NAME: + # Save received raw data to cloud storage, for debugging + logger.info("Starting periodic storage upload of raw data") + filekey = datetime.now(timezone.utc).strftime("%Y/%m/%d/%H:%M:%S.txt") task_upload_raw = asyncio.create_task( - upload_to_gcs_periodically( + upload_to_object_storage_periodically( filepath=temp.name, - bucket=GCS_BUCKET, - destination_key=gcs_filekey, + bucket=STORAGE_BUCKET_NAME, + destination_key=filekey, interval=timedelta(seconds=60), ) ) diff --git a/src/openf1/util/gcs.py b/src/openf1/util/gcs.py deleted file mode 100644 index 6f21f4d2..00000000 --- a/src/openf1/util/gcs.py +++ /dev/null @@ -1,47 +0,0 @@ -import asyncio -from datetime import timedelta -from functools import lru_cache -from pathlib import Path - -from google.auth import default -from google.cloud import storage -from loguru import logger - - -@lru_cache() -def _storage_client() -> storage.Client: - credentials, project_id = default() - client = storage.Client(credentials=credentials, project=project_id) - return client - - -def upload_to_gcs(filepath: Path, bucket: str, destination_key: str): - bucket = _storage_client().bucket(bucket) - blob = bucket.blob(destination_key) - blob.upload_from_filename(filepath) - - -async def upload_to_gcs_periodically( - filepath: Path, - bucket: str, - destination_key: Path, - interval: timedelta, -): - """Periodically uploads a file to Google Cloud Storage (GCS) at specified intervals""" - loop = asyncio.get_running_loop() - - while True: - try: - # Wait for the specified interval - await asyncio.sleep(interval.total_seconds()) - - # Upload the file to GCS using a separate thread to avoid blocking the event loop - await loop.run_in_executor( - None, - upload_to_gcs, - filepath, - bucket, - str(destination_key), - ) - except Exception: - logger.exception("An unexpected error occurred while uploading to GCS") diff --git a/src/openf1/util/storage.py b/src/openf1/util/storage.py new file mode 100644 index 00000000..6f4fc852 --- /dev/null +++ b/src/openf1/util/storage.py @@ -0,0 +1,63 @@ +import asyncio +import os +from datetime import timedelta +from functools import lru_cache +from pathlib import Path + +import boto3 +from loguru import logger + + +@lru_cache() +def _s3_client(): + """ + Client for any S3-compatible object storage. + + Configured entirely through environment variables so the same code runs + against AWS S3, Scaleway, Backblaze B2, MinIO, Cloudflare R2, etc. + + S3_ENDPOINT_URL Endpoint of the provider. Leave unset for AWS S3. + S3_REGION Region name (e.g. fr-par, us-east-1). Optional for some + providers; safe to set to match your bucket. + S3_ACCESS_KEY Access key id. + S3_SECRET_KEY Secret access key. + """ + return boto3.client( + "s3", + endpoint_url=os.environ.get("S3_ENDPOINT_URL") or None, + region_name=os.environ.get("S3_REGION") or None, + aws_access_key_id=os.environ["S3_ACCESS_KEY"], + aws_secret_access_key=os.environ["S3_SECRET_KEY"], + ) + + +def upload_to_object_storage(filepath: Path, bucket: str, destination_key: str): + _s3_client().upload_file(str(filepath), bucket, destination_key) + + +async def upload_to_object_storage_periodically( + filepath: Path, + bucket: str, + destination_key: Path, + interval: timedelta, +): + """Periodically uploads a file to S3-compatible object storage at specified intervals""" + loop = asyncio.get_running_loop() + + while True: + try: + # Wait for the specified interval + await asyncio.sleep(interval.total_seconds()) + + # Upload in a separate thread to avoid blocking the event loop + await loop.run_in_executor( + None, + upload_to_object_storage, + filepath, + bucket, + str(destination_key), + ) + except Exception: + logger.exception( + "An unexpected error occurred while uploading to object storage" + )