fix(backend): better trigger listening logs (#7392)

* fix(backend): better trigger listening logs

* chore: update ee-repo-ref to d347295041426d03039b747a148a71e3583c3a6b

This commit updates the EE repository reference after PR #362 was merged in windmill-ee-private.

Previous ee-repo-ref: 37b533704e1b40e616ac144bebeff574a5d048e1

New ee-repo-ref: d347295041426d03039b747a148a71e3583c3a6b

Automated by sync-ee-ref workflow.

---------

Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
This commit is contained in:
hugocasa
2025-12-16 16:23:21 +00:00
committed by GitHub
co-authored by windmill-internal-app[bot]
parent 8921b05252
commit 7a05601b11
3 changed files with 12 additions and 6 deletions
+1 -1
View File
@@ -1 +1 @@
9a4b392262c760fc52256ca00e4d751d9f42e79e
d347295041426d03039b747a148a71e3583c3a6b
@@ -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)
+11 -4
View File
@@ -587,6 +587,7 @@ async fn listening<T: Listener>(
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<T: Listener>(
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<T: Listener>(
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;