diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index ae5583aa66..27a9af3b6c 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -9a4b392262c760fc52256ca00e4d751d9f42e79e +d347295041426d03039b747a148a71e3583c3a6b diff --git a/backend/windmill-api/src/triggers/handler.rs b/backend/windmill-api/src/triggers/handler.rs index 68012a5bcc..6e970fe347 100644 --- a/backend/windmill-api/src/triggers/handler.rs +++ b/backend/windmill-api/src/triggers/handler.rs @@ -359,7 +359,6 @@ pub trait TriggerCrud: Send + Sync + 'static { .sql() .map_err(|e| Error::InternalErr(format!("SQL error: {}", e)))?; - tracing::info!("SQL: {}", sql); let triggers = sqlx::query_as(&sql).fetch_all(&mut *tx).await?; Ok(triggers) diff --git a/backend/windmill-api/src/triggers/listener.rs b/backend/windmill-api/src/triggers/listener.rs index 613774d919..ce4f651ba9 100644 --- a/backend/windmill-api/src/triggers/listener.rs +++ b/backend/windmill-api/src/triggers/listener.rs @@ -587,6 +587,7 @@ async fn listening( let loop_ping_status = Arc::new(RwLock::new(None)); let extra_state = listener.get_extra_state().await; + let path = listening_trigger.path.clone(); tokio::select! { biased; _ = killpill_rx.recv() => { @@ -596,15 +597,18 @@ async fn listening( let _ = listener.cleanup(&db, &listening_trigger, extra_state.as_ref()).await; } consumer = { + tracing::info!("[{}] Getting consumer for trigger {}", T::TRIGGER_KIND, path); listener.get_consumer(&db, &listening_trigger, loop_ping_status.clone(), killpill_rx_get_consumer) } => { tokio::select! { biased; _ = killpill_rx.recv() => { + tracing::info!("[{}] Killing pill received, stopping consumer for trigger {}", T::TRIGGER_KIND, path); let _ = listener.cleanup(&db, &listening_trigger, extra_state.as_ref()).await; return; } _ = listener.loop_ping(&db, &listening_trigger, loop_ping_status.clone(), None) => { + tracing::info!("[{}] Loop ping exited, stopping consumer for trigger {}", T::TRIGGER_KIND, path); let _ = listener.cleanup(&db, &listening_trigger, extra_state.as_ref()).await; return; } @@ -612,14 +616,17 @@ async fn listening( match consumer { Ok(Some(consumer)) => { listener.update_ping_and_loop_ping_status(&db, &listening_trigger, loop_ping_status.clone(), None).await; - let _ = listener.consume(&db, consumer, &listening_trigger, loop_ping_status.clone(), killpill_rx_consumer, extra_state.as_ref()).await; - tracing::debug!("Stopping consumer for trigger"); + tracing::info!("[{}] Starting consumer for trigger {}", T::TRIGGER_KIND, path); + listener.consume(&db, consumer, &listening_trigger, loop_ping_status.clone(), killpill_rx_consumer, extra_state.as_ref()).await; + tracing::info!("[{}] Consumer stopped for trigger {}", T::TRIGGER_KIND, path); } Err(error) => { - tracing::warn!("Disabling trigger due to consumer error: {}", error); + tracing::error!("[{}] Disabling trigger {} due to consumer error: {}", T::TRIGGER_KIND, path, error); listener.disable_with_error(&db, &listening_trigger, error.to_string()).await; } - _ => {} + Ok(None) => { + tracing::error!("[{}] Consumer is None for trigger {}", T::TRIGGER_KIND, path); + } } } => { let _ = listener.cleanup(&db, &listening_trigger, extra_state.as_ref()).await;