ifieryarrows commited on
Commit
f6973c1
·
verified ·
1 Parent(s): a3d6c68

Sync from GitHub (tests passed)

Browse files
Files changed (1) hide show
  1. worker/tasks.py +15 -16
worker/tasks.py CHANGED
@@ -40,19 +40,16 @@ logger = logging.getLogger(__name__)
40
  # Helper functions for metrics tracking
41
  # =============================================================================
42
 
43
- def _create_pipeline_session() -> tuple[Session, Optional[Any]]:
44
- """Return a work session pinned to one PostgreSQL backend connection.
45
-
46
- PostgreSQL advisory locks are session-scoped. A normal SQLAlchemy Session
47
- may return its connection to the pool after each commit, leaking the lock
48
- on the original backend and attempting the final unlock on another one.
49
- Binding to an explicitly held Connection keeps every stage commit and the
50
- eventual unlock on the same physical database session.
51
  """
52
  if get_db_type() == "postgresql":
53
- connection = get_engine().connect()
54
- return Session(bind=connection, autoflush=False), connection
55
- return SessionLocal(), None
56
 
57
  def create_run_metrics(
58
  session: Session,
@@ -233,7 +230,9 @@ async def run_pipeline(
233
 
234
  # Get a dedicated session for this pipeline run
235
  # IMPORTANT: This session holds the advisory lock
236
- session, pinned_connection = _create_pipeline_session()
 
 
237
  quality_state = "ok"
238
  result = {}
239
 
@@ -250,7 +249,7 @@ async def run_pipeline(
250
  session.commit()
251
 
252
  # 1. Acquire distributed lock
253
- if not try_acquire_lock(session, PIPELINE_LOCK_KEY):
254
  logger.warning(f"[run_id={run_id}] Pipeline skipped: lock held by another process")
255
  finalize_run_metrics(session, run_id, status="skipped_locked", quality_state="skipped")
256
  session.commit()
@@ -328,15 +327,15 @@ async def run_pipeline(
328
  finally:
329
  # Always release lock and cleanup
330
  try:
331
- release_lock(session, PIPELINE_LOCK_KEY)
332
  clear_lock_visibility(session, PIPELINE_LOCK_KEY)
333
  session.commit()
334
  except Exception:
335
  session.rollback()
336
  finally:
337
  session.close()
338
- if pinned_connection is not None:
339
- pinned_connection.close()
340
 
341
 
342
  async def _execute_pipeline_stages_v2(
 
40
  # Helper functions for metrics tracking
41
  # =============================================================================
42
 
43
+ def _create_pipeline_lock_connection() -> Optional[Any]:
44
+ """Hold the production advisory lock on a dedicated DB connection.
45
+
46
+ Stage code intentionally commits and rolls back many transactions. Keeping
47
+ the session-scoped advisory lock on a separate physical connection makes
48
+ its lifetime independent of ORM transaction/pool behavior.
 
 
49
  """
50
  if get_db_type() == "postgresql":
51
+ return get_engine().connect()
52
+ return None
 
53
 
54
  def create_run_metrics(
55
  session: Session,
 
230
 
231
  # Get a dedicated session for this pipeline run
232
  # IMPORTANT: This session holds the advisory lock
233
+ session: Session = SessionLocal()
234
+ lock_connection = _create_pipeline_lock_connection()
235
+ lock_handle = lock_connection if lock_connection is not None else session
236
  quality_state = "ok"
237
  result = {}
238
 
 
249
  session.commit()
250
 
251
  # 1. Acquire distributed lock
252
+ if not try_acquire_lock(lock_handle, PIPELINE_LOCK_KEY):
253
  logger.warning(f"[run_id={run_id}] Pipeline skipped: lock held by another process")
254
  finalize_run_metrics(session, run_id, status="skipped_locked", quality_state="skipped")
255
  session.commit()
 
327
  finally:
328
  # Always release lock and cleanup
329
  try:
330
+ release_lock(lock_handle, PIPELINE_LOCK_KEY)
331
  clear_lock_visibility(session, PIPELINE_LOCK_KEY)
332
  session.commit()
333
  except Exception:
334
  session.rollback()
335
  finally:
336
  session.close()
337
+ if lock_connection is not None:
338
+ lock_connection.close()
339
 
340
 
341
  async def _execute_pipeline_stages_v2(