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
123 changes: 107 additions & 16 deletions app/services/upload_manager.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
"""Upload manager for orchestrating file uploads to S3."""

import logging
import os
import shutil
import threading
import uuid
from collections.abc import Callable
from concurrent.futures import ThreadPoolExecutor, as_completed
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor, as_completed
from dataclasses import dataclass, field
from datetime import UTC, datetime
from enum import Enum
Expand All @@ -23,6 +24,18 @@
EPOCH_CUTOFF = datetime(1980, 1, 1, tzinfo=UTC)


def _extract_start_time_worker(local_path: str) -> datetime | str:
"""Worker function for ProcessPoolExecutor — must be top-level for pickling.

Returns:
datetime on success, or error message string on failure.
"""
try:
return mcap_service.extract_start_time(local_path)
except Exception as e:
return str(e)


class UploadStatus(Enum):
"""Status of a file upload."""

Expand Down Expand Up @@ -348,6 +361,39 @@ def analyze_job(

return job

def _check_duplicate(
self,
file_state: FileUploadState,
s3_client: Any,
s3_bucket: str,
use_cache: bool = True,
) -> None:
"""Check if a file already exists in S3 (I/O-bound, safe for threads)."""
s3_path = file_state.s3_path
if not s3_path:
return

cache_result: bool | None = None
if use_cache:
cache = get_cache_service()
cache_result = cache.check_exists_cached(s3_bucket, s3_path)

if cache_result is not None:
file_state.is_duplicate = cache_result
else:
file_state.is_duplicate = s3_service.check_file_exists(
s3_client, s3_bucket, s3_path
)
if use_cache:
cache = get_cache_service()
cache.update_cache(
s3_bucket,
s3_path,
file_state.is_duplicate,
file_state.filename,
file_state.file_size,
)

def _analyze_single_file(
self,
file_state: FileUploadState,
Expand Down Expand Up @@ -497,34 +543,79 @@ def analyze_job_async(
progress_callback(job, file_state)
return job

# Analyze files in parallel
with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
futures = {}
for file_state in job.files:
future = executor.submit(
self._analyze_single_file,
# Phase 1: MCAP parsing (CPU-bound) — use ProcessPoolExecutor for true
# parallelism across cores, bypassing the GIL.
cpu_workers = max(1, (os.cpu_count() or 4) - 1)
for file_state in job.files:
file_state.status = UploadStatus.ANALYZING
if progress_callback:
progress_callback(job, file_state)

with ProcessPoolExecutor(max_workers=cpu_workers) as proc_executor:
parse_futures = {
proc_executor.submit(
_extract_start_time_worker, file_state.local_path
): file_state
for file_state in job.files
}
for future in as_completed(parse_futures):
file_state = parse_futures[future]
result = future.result()
if isinstance(result, str):
# Error message returned from worker
file_state.status = UploadStatus.FAILED
file_state.error_message = result
log.error(
"analysis",
"file_analysis_failed",
f"Failed to analyze {file_state.filename}: {result}",
{"job_id": job_id, "filename": file_state.filename, "error": result},
)
else:
file_state.start_time = result
naive_start = mcap_service.to_naive_utc(result)
file_state.is_valid = naive_start >= EPOCH_CUTOFF.replace(tzinfo=None)
file_state.s3_path = mcap_service.generate_s3_path(result, file_state.filename)
if progress_callback:
progress_callback(job, file_state)

# Phase 2: S3 duplicate checks (I/O-bound) — threads are fine here.
parsed_files = [f for f in job.files if f.status != UploadStatus.FAILED]
with ThreadPoolExecutor(max_workers=self.max_workers) as io_executor:
dup_futures: dict[Any, FileUploadState] = {}
for file_state in parsed_files:
fut = io_executor.submit(
self._check_duplicate,
file_state,
s3_client,
s3_bucket,
use_cache,
job_id,
progress_callback,
job,
)
futures[future] = file_state
dup_futures[fut] = file_state

# Process results as they complete
for future in as_completed(futures):
file_state = futures[future]
for fut in as_completed(dup_futures):
file_state = dup_futures[fut]
try:
# Result is already updated in-place, but get it to handle exceptions
future.result()
file_state.status = UploadStatus.READY
log.info(
"analysis",
"file_analysis_completed",
f"Analyzed {file_state.filename}",
{
"job_id": job_id,
"filename": file_state.filename,
"file_size": file_state.file_size,
"s3_path": file_state.s3_path,
"is_duplicate": file_state.is_duplicate,
"is_valid": file_state.is_valid,
},
)
except Exception as e:
with job.lock:
file_state.status = UploadStatus.FAILED
file_state.error_message = str(e)

# Call progress callback
if progress_callback:
progress_callback(job, file_state)

Expand Down
2 changes: 1 addition & 1 deletion deploy/modaq-upload.service
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ WorkingDirectory=/opt/modaq-upload
Environment="PATH=/opt/modaq-upload/venv/bin"
ExecStart=/opt/modaq-upload/venv/bin/gunicorn \
--bind 0.0.0.0:8080 \
--workers 2 \
--workers 1 \
--threads 4 \
--timeout 300 \
--access-logfile /var/log/modaq-upload/access.log \
Expand Down
13 changes: 11 additions & 2 deletions launch.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,15 @@ def open_browser(url: str) -> None:
import shutil

if shutil.which("firefox"):
subprocess.Popen(["firefox", url], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
subprocess.Popen(
["firefox", url],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, start_new_session=True,
)
elif shutil.which("xdg-open"):
subprocess.Popen(["xdg-open", url], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
subprocess.Popen(
["xdg-open", url],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, start_new_session=True,
)
else:
log(f"Could not find a browser. Open {url} manually.")

Expand Down Expand Up @@ -98,6 +104,7 @@ def main() -> None:
if port_in_use(PORT):
log(f"Port {PORT} is already in use — opening browser to existing instance.")
open_browser(BROWSER_URL)
time.sleep(2)
sys.exit(0)

# ── Register signal handlers ─────────────────────────────
Expand All @@ -113,6 +120,8 @@ def main() -> None:
"--bind",
f"127.0.0.1:{PORT}",
"--workers",
"1",
"--threads",
"4",
"--timeout",
"300",
Expand Down
4 changes: 2 additions & 2 deletions modaq-upload.desktop.template
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
[Desktop Entry]
Version=0.2.1
Version=0.2.2
Type=Application
Name=MODAQ Upload
Comment=MODAQ File Uploader
Exec=python3 {{PROJECT_DIR}}/launch.py
Exec={{PROJECT_DIR}}/venv/bin/python {{PROJECT_DIR}}/launch.py
Icon={{PROJECT_DIR}}/app/static/images/modaq-logo.png
Terminal=true
Categories=Science;Utility;
Expand Down