Skip to content

Commit e027a86

Browse files
committed
fix: improve setup, logging, and task resiliency [no ci]
1 parent 0ce7ae0 commit e027a86

4 files changed

Lines changed: 19 additions & 9 deletions

File tree

install.sh

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -95,7 +95,11 @@ clone_repository() {
9595
setup_virtualenv() {
9696
log "Setting up virtual environment"
9797
cd "$INSTALL_DIR"
98-
[ -d "venv" ] || python3 -m venv venv || error "Failed to create venv"
98+
if [ -d "venv" ]; then
99+
log "Removing existing venv for a clean rebuild"
100+
rm -rf venv
101+
fi
102+
python3 -m venv venv || error "Failed to create venv"
99103
venv/bin/pip install --quiet --upgrade pip || error "Failed to upgrade pip"
100104
venv/bin/pip install --quiet -r requirements.txt || error "Failed to install Python dependencies"
101105
}

logutils.py

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,10 +13,14 @@
1313
if not isinstance(numeric_level, int):
1414
raise ValueError(f"Invalid log level: {LOG_LEVEL}")
1515

16-
logging.basicConfig(
17-
level=numeric_level, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s"
16+
_LOG_FORMAT = (
17+
"%(name)s - %(levelname)s - %(message)s"
18+
if os.getenv("JOURNAL_STREAM")
19+
else "%(asctime)s - %(name)s - %(levelname)s - %(message)s"
1820
)
1921

22+
logging.basicConfig(level=numeric_level, format=_LOG_FORMAT)
23+
2024

2125
def get_logger(name: str = None) -> logging.Logger:
2226
"""

publications.py

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -308,18 +308,19 @@ def _store_segment_and_try_join(
308308

309309
try:
310310
joined = rrs.V1Payloads.join(segment_data)
311-
except Exception:
312-
logger.exception("Failed to join segments for session %s.", sess_id)
311+
except Exception as exc:
313312
logger.warning(
314-
"Session %s incomplete: %d segments saved.", sess_id, len(segment_data)
313+
"Session %s not ready yet (%d segment(s) stored so far): %s: %s",
314+
sess_id,
315+
len(segment_data),
316+
type(exc).__name__,
317+
exc or "-",
315318
)
316319
return None
317320

318321
delete_session(payload_session, self.session)
319322
logger.info(
320-
"Session %s assembled successfully from %d segments.",
321-
sess_id,
322-
len(segment_data),
323+
"Session %s assembled from %d segments.", sess_id, len(segment_data)
323324
)
324325
return joined
325326

tasks/celery_app.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,7 @@ def make_celery() -> Celery:
8585
worker_enable_remote_control=False,
8686
worker_concurrency=concurrency,
8787
task_acks_late=True,
88+
task_reject_on_worker_lost=True,
8889
worker_prefetch_multiplier=1,
8990
task_ignore_result=True,
9091
beat_schedule_filename=schedule_path,

0 commit comments

Comments
 (0)