Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"

[project]
name = "openf1"
version = "1.9.13"
version = "1.9.14"
authors = [
{ name="Bruno Godefroy" }
]
Expand Down
2 changes: 1 addition & 1 deletion requirements.txt
Original file line number Diff line number Diff line change
@@ -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
Expand Down
18 changes: 9 additions & 9 deletions src/openf1/services/ingestor_livetiming/real_time/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand All @@ -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),
)
)
Expand Down
47 changes: 0 additions & 47 deletions src/openf1/util/gcs.py

This file was deleted.

63 changes: 63 additions & 0 deletions src/openf1/util/storage.py
Original file line number Diff line number Diff line change
@@ -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"
)
Loading