From c2c53c5ea6a63d8a3e32bdf8a99ad44709e01bda Mon Sep 17 00:00:00 2001 From: yashnevatia Date: Tue, 18 Aug 2026 11:50:01 +0100 Subject: [PATCH 1/6] Adding metrics --- pkg/capabilities/base_trigger.go | 126 +++++++++++++++++++---- pkg/capabilities/base_trigger_metrics.go | 93 +++++++++++++++-- 2 files changed, 189 insertions(+), 30 deletions(-) diff --git a/pkg/capabilities/base_trigger.go b/pkg/capabilities/base_trigger.go index 20424e523e..335016dcb3 100644 --- a/pkg/capabilities/base_trigger.go +++ b/pkg/capabilities/base_trigger.go @@ -30,6 +30,26 @@ const ( ackMemoryOutcomeMissNoEvent = "miss_no_event" ackMemoryOutcomeMissNilRecord = "miss_nil_record" ackMemoryOutcomePreAckDeliverySkipped = "pre_ack_delivery_skipped" + + // b.mu acquisition sites, used as the "op" label on mu_wait_ms so lock + // contention can be attributed to the path that is waiting. + muOpStart = "start" + muOpRegister = "register" + muOpUnregister = "unregister" + muOpDeliverPre = "deliver_pre_persist" + muOpDeliverPost = "deliver_post_persist" + muOpSendToInbox = "send_to_inbox" + muOpAck = "ack" + muOpScanPending = "scan_pending" + muOpTrySend = "try_send" + muOpPrune = "prune" + + // EventStore operations, used as the "op" label on store_op_duration_ms. + storeOpList = "list" + storeOpInsert = "insert" + storeOpUpdateDelivery = "update_delivery" + storeOpDeleteEvent = "delete_event" + storeOpDeleteEventsForTrigger = "delete_events_for_trigger" ) // triggerReg holds all per-triggerID runtime state for a [BaseTriggerCapability]. @@ -77,6 +97,18 @@ type BaseTriggerMetrics interface { AddPendingEvents(delta int64) // IncStoppedResending records a gauge (Unix seconds) when the node exhausted max retries and stopped resending. IncStoppedResending(triggerID string, attempts int) + // ObserveMuWait records how long a caller blocked acquiring b.mu. op is one of the muOp* constants. + ObserveMuWait(op string, d time.Duration) + // ObserveScanPendingLockHeld records how long scanPending held b.mu. The retransmit loop + // re-arms 100ms after scanPending returns, so held/(held+100ms) is the share of wall-clock + // the scan owns the lock — i.e. how much it starves AckEvent and DeliverEvent. + ObserveScanPendingLockHeld(d time.Duration) + // SetPreAckedEntries reports the live count of pre-ACK tombstones across all triggers. + // This is the collection expirePreAcked walks on every scan, so it drives lock hold time. + SetPreAckedEntries(n int64) + // ObserveStoreOp records EventStore call latency. op is one of the storeOp* constants; + // outcome is "success" or "error". + ObserveStoreOp(op string, d time.Duration, outcome string) } // BaseTriggerCapability keeps track of trigger registrations and handles resending events until @@ -131,6 +163,27 @@ func NewBaseTriggerCapability[T proto.Message]( } } +// lockMu acquires b.mu and records how long the acquisition blocked, attributed to op. +// b.mu is a single global lock shared by every path (ACK, delivery, the 10Hz retransmit +// scan, prune), so this is the primary signal for whether a slow AckEvent is waiting on +// the lock rather than on the store. +func (b *BaseTriggerCapability[T]) lockMu(op string) { + start := time.Now() + b.mu.Lock() + b.metrics.ObserveMuWait(op, time.Since(start)) +} + +// observeStoreOp records the latency and outcome of an EventStore call. The store is +// backed by Postgres over gRPC, so these calls are unbounded in the absence of a +// deadline and can hold an ACK executor slot indefinitely. +func (b *BaseTriggerCapability[T]) observeStoreOp(op string, start time.Time, err error) { + outcome := "success" + if err != nil { + outcome = "error" + } + b.metrics.ObserveStoreOp(op, time.Since(start), outcome) +} + // getOrCreateRegLocked returns the registration record for triggerID, creating it if needed. // Both pending and preAcked are always initialized; inbox remains nil until RegisterTrigger. // Must be called with b.mu held. @@ -229,14 +282,16 @@ func (b *BaseTriggerCapability[T]) pruneAge(ctx context.Context) time.Duration { func (b *BaseTriggerCapability[T]) Start(ctx context.Context) error { b.lggr.Info("starting base trigger") + listStart := time.Now() recs, err := b.store.List(ctx) + b.observeStoreOp(storeOpList, listStart, err) if err != nil { b.lggr.Errorf("failed to load persisted trigger events") return err } // Initialize in-memory persistence - b.mu.Lock() + b.lockMu(muOpStart) for i := range recs { r := &recs[i] reg := b.getOrCreateRegLocked(r.TriggerId) @@ -266,7 +321,7 @@ func (b *BaseTriggerCapability[T]) Stop() { } func (b *BaseTriggerCapability[T]) RegisterTrigger(triggerID string, sendCh chan<- TriggerAndId[T]) { - b.mu.Lock() + b.lockMu(muOpRegister) // getOrCreateRegLocked returns the existing record if AckEvent already created // one (preserving any pre-ACKed entries), or a fresh one. In either case we // only overwrite the inbox field, so pre-ACK state is never erased. @@ -281,7 +336,7 @@ func (b *BaseTriggerCapability[T]) RegisterTrigger(triggerID string, sendCh chan } func (b *BaseTriggerCapability[T]) UnregisterTrigger(triggerID string) { - b.mu.Lock() + b.lockMu(muOpUnregister) reg, ok := b.byTrigger[triggerID] var pendingCount int64 var existed bool @@ -299,7 +354,10 @@ func (b *BaseTriggerCapability[T]) UnregisterTrigger(triggerID string) { b.metrics.AddPendingEvents(-pendingCount) } - if err := b.store.DeleteEventsForTrigger(b.ctx, triggerID); err != nil { + deleteStart := time.Now() + err := b.store.DeleteEventsForTrigger(b.ctx, triggerID) + b.observeStoreOp(storeOpDeleteEventsForTrigger, deleteStart, err) + if err != nil { b.lggr.Errorf("Failed to delete events for trigger (TriggerID=%s): %v", triggerID, err) } } @@ -323,7 +381,7 @@ func (b *BaseTriggerCapability[T]) DeliverEvent( // Already pending: the EVM trigger re-delivers after finalization // while the event is still awaiting ACK. var reg *triggerReg[T] - b.mu.Lock() + b.lockMu(muOpDeliverPre) reg = b.byTrigger[triggerID] if reg != nil { if _, wasAcked := reg.preAcked[te.ID]; wasAcked { @@ -353,7 +411,9 @@ func (b *BaseTriggerCapability[T]) DeliverEvent( OrgID: orgID, } + insertStart := time.Now() if err := b.store.Insert(ctx, rec); err != nil { + b.observeStoreOp(storeOpInsert, insertStart, err) if isDuplicateKeyError(err) { b.lggr.Debugw("base trigger DeliverEvent: event already in store (re-delivery after give-up), skipping", "capabilityID", b.capabilityId, "triggerID", triggerID, "eventID", te.ID) @@ -363,6 +423,7 @@ func (b *BaseTriggerCapability[T]) DeliverEvent( "capabilityID", b.capabilityId, "triggerID", triggerID, "eventID", te.ID, "err", err) return err } + b.observeStoreOp(storeOpInsert, insertStart, nil) b.lggr.Infow("base trigger persisted pending event for ACK tracking", "capabilityID", b.capabilityId, "triggerID", triggerID, "eventID", te.ID) @@ -370,7 +431,7 @@ func (b *BaseTriggerCapability[T]) DeliverEvent( // An ACK may have arrived during the store.Insert call above. Without // this second check, the event would be retransmitted forever because // the first preAcked check (before Insert) narrowly missed the ACK. - b.mu.Lock() + b.lockMu(muOpDeliverPost) reg = b.getOrCreateRegLocked(triggerID) if _, wasAcked := reg.preAcked[te.ID]; wasAcked { delete(reg.preAcked, te.ID) @@ -378,7 +439,10 @@ func (b *BaseTriggerCapability[T]) DeliverEvent( b.lggr.Infow("base trigger DeliverEvent skipped after persist: event was ACKed during store write (pre-ACK double-check)", "capabilityID", b.capabilityId, "triggerID", triggerID, "eventID", te.ID) b.metrics.IncAckMemoryOutcome(ackMemoryOutcomePreAckDeliverySkipped) - if err := b.store.DeleteEvent(ctx, triggerID, te.ID); err != nil { + deleteStart := time.Now() + err := b.store.DeleteEvent(ctx, triggerID, te.ID) + b.observeStoreOp(storeOpDeleteEvent, deleteStart, err) + if err != nil { b.lggr.Errorw("base trigger failed to delete pre-ACKed event from store", "capabilityID", b.capabilityId, "triggerID", triggerID, "eventID", te.ID, "err", err) } @@ -402,7 +466,7 @@ func (b *BaseTriggerCapability[T]) DeliverEvent( // sendToInbox unmarshals the payload and delivers it to the registered inbox channel. func (b *BaseTriggerCapability[T]) sendToInbox(triggerID, eventID string, payload []byte) error { - b.mu.Lock() + b.lockMu(muOpSendToInbox) reg := b.byTrigger[triggerID] var sendCh chan<- TriggerAndId[T] if reg != nil { @@ -443,7 +507,7 @@ func (b *BaseTriggerCapability[T]) AckEvent(ctx context.Context, triggerId strin hadNilPendingRecord bool ) - b.mu.Lock() + b.lockMu(muOpAck) reg := b.byTrigger[triggerId] eventWasInPending := false if reg != nil { @@ -510,7 +574,10 @@ func (b *BaseTriggerCapability[T]) AckEvent(ctx context.Context, triggerId strin "hadNilPendingRecord", hadNilPendingRecord) } - if err := b.store.DeleteEvent(ctx, triggerId, eventId); err != nil { + deleteStart := time.Now() + err := b.store.DeleteEvent(ctx, triggerId, eventId) + b.observeStoreOp(storeOpDeleteEvent, deleteStart, err) + if err != nil { b.lggr.Errorw("base trigger ACK failed to delete event from store", "capabilityID", b.capabilityId, "triggerID", triggerId, "eventID", eventId, "foundInMemory", found, "err", err) @@ -572,9 +639,10 @@ func (b *BaseTriggerCapability[T]) scanPending() { maxRetries := b.maxRetries(ctx) - b.mu.Lock() + b.lockMu(muOpScanPending) + lockedAt := time.Now() - b.expirePreAcked(now) + preAckedRemaining := b.expirePreAcked(now) toResend := make([]PendingEvent, 0, defaultMaxSendsPerTick) var toStop []stoppedResendingEvent @@ -591,6 +659,11 @@ func (b *BaseTriggerCapability[T]) scanPending() { } } b.mu.Unlock() + // Recorded outside the lock. The retransmit loop re-arms a 100ms timer after + // scanPending returns, so held/(held+100ms) is the share of wall-clock this scan + // holds b.mu — i.e. how much it starves AckEvent. + b.metrics.ObserveScanPendingLockHeld(time.Since(lockedAt)) + b.metrics.SetPreAckedEntries(int64(preAckedRemaining)) for _, ev := range toStop { b.emitStoppedResending(ev, maxRetries) @@ -623,20 +696,27 @@ func (b *BaseTriggerCapability[T]) scanPending() { } } -// expirePreAcked removes old preAcked entries so the cache doesn't grow unbounded. +// expirePreAcked removes old preAcked entries so the cache doesn't grow unbounded, +// returning the number of entries that survived. The count is free here because this +// function already walks every entry, and it is the size of the collection driving +// this scan's lock hold time. // The TTL is generous (24h) because entries are tiny (two strings + timestamp) and // the cost of expiring too early is severe: a slow node would persist and retransmit // an already-ACKed event forever. UnregisterTrigger already clears entries per trigger. // Must be called under b.mu. -func (b *BaseTriggerCapability[T]) expirePreAcked(now time.Time) { +func (b *BaseTriggerCapability[T]) expirePreAcked(now time.Time) int { preAckTTL := 24 * time.Hour + remaining := 0 for _, reg := range b.byTrigger { for eventID, ackedAt := range reg.preAcked { if now.Sub(ackedAt) > preAckTTL { delete(reg.preAcked, eventID) + continue } + remaining++ } } + return remaining } // collectStoppedResending removes the event from pending, returning the metadata @@ -720,7 +800,9 @@ func (b *BaseTriggerCapability[T]) pruneStaleEvents() { } cutoff := time.Now().Add(-age) + listStart := time.Now() recs, err := b.store.List(b.ctx) + b.observeStoreOp(storeOpList, listStart, err) if err != nil { b.lggr.Errorw("prune: failed to list events from store", "capabilityID", b.capabilityId, "err", err) return @@ -737,7 +819,7 @@ func (b *BaseTriggerCapability[T]) pruneStaleEvents() { continue } - b.mu.Lock() + b.lockMu(muOpPrune) if reg, ok := b.byTrigger[rec.TriggerId]; ok { delete(reg.pending, rec.EventId) } @@ -746,7 +828,10 @@ func (b *BaseTriggerCapability[T]) pruneStaleEvents() { b.lggr.Infow("prune: removing stale event from store", "capabilityID", b.capabilityId, "triggerID", rec.TriggerId, "eventID", rec.EventId, "firstAt", rec.FirstAt, "lastSentAt", rec.LastSentAt, "attempts", rec.Attempts, "pruneAge", age) - if err := b.store.DeleteEvent(b.ctx, rec.TriggerId, rec.EventId); err != nil { + deleteStart := time.Now() + deleteErr := b.store.DeleteEvent(b.ctx, rec.TriggerId, rec.EventId) + b.observeStoreOp(storeOpDeleteEvent, deleteStart, deleteErr) + if deleteErr != nil { b.lggr.Errorw("prune: failed to delete stale event", "capabilityID", b.capabilityId, "triggerID", rec.TriggerId, "eventID", rec.EventId, "err", err) } @@ -761,7 +846,7 @@ func (b *BaseTriggerCapability[T]) trySend(event PendingEvent) { return } - b.mu.Lock() + b.lockMu(muOpTrySend) reg := b.byTrigger[event.TriggerId] if reg == nil { b.mu.Unlock() @@ -788,8 +873,11 @@ func (b *BaseTriggerCapability[T]) trySend(event PendingEvent) { b.mu.Unlock() b.metrics.IncRetry(event.TriggerId) - if err := b.store.UpdateDelivery(b.ctx, event.TriggerId, event.EventId, lastSent, attempts); err != nil { - b.lggr.Errorf("failed to persist delivery update for trigger=%s event=%s: %v", event.TriggerId, event.EventId, err) + updateStart := time.Now() + updateErr := b.store.UpdateDelivery(b.ctx, event.TriggerId, event.EventId, lastSent, attempts) + b.observeStoreOp(storeOpUpdateDelivery, updateStart, updateErr) + if updateErr != nil { + b.lggr.Errorf("failed to persist delivery update for trigger=%s event=%s: %v", event.TriggerId, event.EventId, updateErr) } if err := b.sendToInbox(event.TriggerId, event.EventId, payloadCopy); err != nil { diff --git a/pkg/capabilities/base_trigger_metrics.go b/pkg/capabilities/base_trigger_metrics.go index aa0f26f75b..09d70cd935 100644 --- a/pkg/capabilities/base_trigger_metrics.go +++ b/pkg/capabilities/base_trigger_metrics.go @@ -23,8 +23,20 @@ type BaseTriggerBeholderMetrics struct { activeRegistrations metric.Int64UpDownCounter pendingEvents metric.Int64UpDownCounter stoppedResendingTimestamp metric.Int64Gauge + muWaitMs metric.Int64Histogram + scanPendingLockHeldMs metric.Int64Histogram + preAckedEntries metric.Int64Gauge + storeOpDurationMs metric.Int64Histogram } +// contentionBuckets span sub-millisecond fast paths through multi-minute stalls. +// The upper bound must stay well above any observed stall: values past the last +// boundary land in +Inf, where histogram_quantile clamps and silently censors the +// tail (a p95 pinned exactly at the top bucket means "at least this", not "this"). +var contentionBuckets = metric.WithExplicitBucketBoundaries( + 1, 5, 10, 25, 50, 100, 250, 500, 1_000, 5_000, 30_000, 120_000, +) + var _ BaseTriggerMetrics = &BaseTriggerBeholderMetrics{} func NewBaseTriggerBeholderMetrics(capabilityID string) (BaseTriggerMetrics, error) { @@ -82,6 +94,26 @@ func NewBaseTriggerBeholderMetrics(capabilityID string) (BaseTriggerMetrics, err return nil, err } + muWaitMs, err := beholder.GetMeter().Int64Histogram("capabilities_base_trigger_mu_wait_ms", contentionBuckets) + if err != nil { + return nil, err + } + + scanPendingLockHeldMs, err := beholder.GetMeter().Int64Histogram("capabilities_base_trigger_scan_pending_lock_held_ms", contentionBuckets) + if err != nil { + return nil, err + } + + preAckedEntries, err := beholder.GetMeter().Int64Gauge("capabilities_base_trigger_preacked_entries") + if err != nil { + return nil, err + } + + storeOpDurationMs, err := beholder.GetMeter().Int64Histogram("capabilities_base_trigger_store_op_duration_ms", contentionBuckets) + if err != nil { + return nil, err + } + return &BaseTriggerBeholderMetrics{ capabilityID: capabilityID, retryCount: retryCount, @@ -95,6 +127,10 @@ func NewBaseTriggerBeholderMetrics(capabilityID string) (BaseTriggerMetrics, err activeRegistrations: activeRegistrations, pendingEvents: pendingEvents, stoppedResendingTimestamp: stoppedResendingTimestamp, + muWaitMs: muWaitMs, + scanPendingLockHeldMs: scanPendingLockHeldMs, + preAckedEntries: preAckedEntries, + storeOpDurationMs: storeOpDurationMs, }, nil } @@ -190,18 +226,53 @@ func (m *BaseTriggerBeholderMetrics) IncStoppedResending(triggerID string, attem ) } +func (m *BaseTriggerBeholderMetrics) ObserveMuWait(op string, d time.Duration) { + m.muWaitMs.Record(context.Background(), d.Milliseconds(), + metric.WithAttributes( + attribute.String("capability_id", m.capabilityID), + attribute.String("op", op), + ), + ) +} + +func (m *BaseTriggerBeholderMetrics) ObserveScanPendingLockHeld(d time.Duration) { + m.scanPendingLockHeldMs.Record(context.Background(), d.Milliseconds(), + metric.WithAttributes(attribute.String("capability_id", m.capabilityID)), + ) +} + +func (m *BaseTriggerBeholderMetrics) SetPreAckedEntries(n int64) { + m.preAckedEntries.Record(context.Background(), n, + metric.WithAttributes(attribute.String("capability_id", m.capabilityID)), + ) +} + +func (m *BaseTriggerBeholderMetrics) ObserveStoreOp(op string, d time.Duration, outcome string) { + m.storeOpDurationMs.Record(context.Background(), d.Milliseconds(), + metric.WithAttributes( + attribute.String("capability_id", m.capabilityID), + attribute.String("op", op), + attribute.String("outcome", outcome), + ), + ) +} + type noopBaseTriggerMetrics struct{} var _ BaseTriggerMetrics = &noopBaseTriggerMetrics{} -func (noopBaseTriggerMetrics) IncActiveTriggers() {} -func (noopBaseTriggerMetrics) DecActiveTriggers() {} -func (noopBaseTriggerMetrics) IncRetry(string) {} -func (noopBaseTriggerMetrics) IncAck(string) {} -func (noopBaseTriggerMetrics) ObserveTimeToAck(string, time.Duration, int) {} -func (noopBaseTriggerMetrics) IncInboxMissing(string) {} -func (noopBaseTriggerMetrics) IncInboxFull(string) {} -func (noopBaseTriggerMetrics) IncAckError(string) {} -func (noopBaseTriggerMetrics) IncAckMemoryOutcome(string) {} -func (noopBaseTriggerMetrics) AddPendingEvents(int64) {} -func (noopBaseTriggerMetrics) IncStoppedResending(string, int) {} +func (noopBaseTriggerMetrics) IncActiveTriggers() {} +func (noopBaseTriggerMetrics) DecActiveTriggers() {} +func (noopBaseTriggerMetrics) IncRetry(string) {} +func (noopBaseTriggerMetrics) IncAck(string) {} +func (noopBaseTriggerMetrics) ObserveTimeToAck(string, time.Duration, int) {} +func (noopBaseTriggerMetrics) IncInboxMissing(string) {} +func (noopBaseTriggerMetrics) IncInboxFull(string) {} +func (noopBaseTriggerMetrics) IncAckError(string) {} +func (noopBaseTriggerMetrics) IncAckMemoryOutcome(string) {} +func (noopBaseTriggerMetrics) AddPendingEvents(int64) {} +func (noopBaseTriggerMetrics) IncStoppedResending(string, int) {} +func (noopBaseTriggerMetrics) ObserveMuWait(string, time.Duration) {} +func (noopBaseTriggerMetrics) ObserveScanPendingLockHeld(time.Duration) {} +func (noopBaseTriggerMetrics) SetPreAckedEntries(int64) {} +func (noopBaseTriggerMetrics) ObserveStoreOp(string, time.Duration, string) {} From 396e0625a089987280b98890bbcdd1d21d440508 Mon Sep 17 00:00:00 2001 From: Yashvardhan Nevatia Date: Tue, 18 Aug 2026 12:02:53 +0100 Subject: [PATCH 2/6] Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- pkg/capabilities/base_trigger.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/capabilities/base_trigger.go b/pkg/capabilities/base_trigger.go index 335016dcb3..1c86aed4ea 100644 --- a/pkg/capabilities/base_trigger.go +++ b/pkg/capabilities/base_trigger.go @@ -833,7 +833,7 @@ func (b *BaseTriggerCapability[T]) pruneStaleEvents() { b.observeStoreOp(storeOpDeleteEvent, deleteStart, deleteErr) if deleteErr != nil { b.lggr.Errorw("prune: failed to delete stale event", - "capabilityID", b.capabilityId, "triggerID", rec.TriggerId, "eventID", rec.EventId, "err", err) + "capabilityID", b.capabilityId, "triggerID", rec.TriggerId, "eventID", rec.EventId, "err", deleteErr) } } } From 5a3fddb49cf044c7c6256def7313cbc121a711f7 Mon Sep 17 00:00:00 2001 From: Yashvardhan Nevatia Date: Tue, 18 Aug 2026 12:03:37 +0100 Subject: [PATCH 3/6] Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- pkg/capabilities/base_trigger_metrics.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/capabilities/base_trigger_metrics.go b/pkg/capabilities/base_trigger_metrics.go index 09d70cd935..a5fc55d087 100644 --- a/pkg/capabilities/base_trigger_metrics.go +++ b/pkg/capabilities/base_trigger_metrics.go @@ -29,7 +29,7 @@ type BaseTriggerBeholderMetrics struct { storeOpDurationMs metric.Int64Histogram } -// contentionBuckets span sub-millisecond fast paths through multi-minute stalls. +// contentionBuckets span millisecond-scale fast paths (>=1ms) through multi-minute stalls. // The upper bound must stay well above any observed stall: values past the last // boundary land in +Inf, where histogram_quantile clamps and silently censors the // tail (a p95 pinned exactly at the top bucket means "at least this", not "this"). From 9697734230b7625f5f67bffbd53ef6324910c784 Mon Sep 17 00:00:00 2001 From: yashnevatia Date: Wed, 19 Aug 2026 10:14:57 +0100 Subject: [PATCH 4/6] adding logs as metrics are fishy --- pkg/capabilities/base_trigger.go | 19 +++++++++++++++++-- 1 file changed, 17 insertions(+), 2 deletions(-) diff --git a/pkg/capabilities/base_trigger.go b/pkg/capabilities/base_trigger.go index 1c86aed4ea..9407ce3632 100644 --- a/pkg/capabilities/base_trigger.go +++ b/pkg/capabilities/base_trigger.go @@ -168,9 +168,12 @@ func NewBaseTriggerCapability[T proto.Message]( // scan, prune), so this is the primary signal for whether a slow AckEvent is waiting on // the lock rather than on the store. func (b *BaseTriggerCapability[T]) lockMu(op string) { + b.lggr.Infow("waiting for lock", "capabilityID", b.capabilityId, "op", op) start := time.Now() b.mu.Lock() - b.metrics.ObserveMuWait(op, time.Since(start)) + waited := time.Since(start) + b.metrics.ObserveMuWait(op, waited) + b.lggr.Infow("got lock", "capabilityID", b.capabilityId, "op", op, "waitedMs", waited.Milliseconds()) } // observeStoreOp records the latency and outcome of an EventStore call. The store is @@ -629,6 +632,12 @@ func isDuplicateKeyError(err error) bool { } func (b *BaseTriggerCapability[T]) scanPending() { + scanStart := time.Now() + defer func() { + b.lggr.Infow("scanPending finished", "capabilityID", b.capabilityId, + "durationMs", time.Since(scanStart).Milliseconds()) + }() + now := time.Now() ctx := b.ctx @@ -662,8 +671,14 @@ func (b *BaseTriggerCapability[T]) scanPending() { // Recorded outside the lock. The retransmit loop re-arms a 100ms timer after // scanPending returns, so held/(held+100ms) is the share of wall-clock this scan // holds b.mu — i.e. how much it starves AckEvent. - b.metrics.ObserveScanPendingLockHeld(time.Since(lockedAt)) + lockHeld := time.Since(lockedAt) + b.metrics.ObserveScanPendingLockHeld(lockHeld) b.metrics.SetPreAckedEntries(int64(preAckedRemaining)) + b.lggr.Infow("scanPending released lock", "capabilityID", b.capabilityId, + "lockHeldMs", lockHeld.Milliseconds(), + "preAckedEntries", preAckedRemaining, + "triggers", len(b.byTrigger), + "toResend", len(toResend)) for _, ev := range toStop { b.emitStoppedResending(ev, maxRetries) From 58b918f6f9856f9f6b2fd00124a9d68e8d0e4099 Mon Sep 17 00:00:00 2001 From: yashnevatia Date: Thu, 3 Sep 2026 12:18:15 +0100 Subject: [PATCH 5/6] adding more logging --- .../core/services/capability/capabilities.go | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/pkg/loop/internal/core/services/capability/capabilities.go b/pkg/loop/internal/core/services/capability/capabilities.go index 20ec9f1485..5ec495d8f6 100644 --- a/pkg/loop/internal/core/services/capability/capabilities.go +++ b/pkg/loop/internal/core/services/capability/capabilities.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "sync" + "time" "google.golang.org/grpc" "google.golang.org/grpc/connectivity" @@ -207,9 +208,18 @@ func newTriggerExecutableServer(brokerExt *net.BrokerExt, impl capabilities.Trig var _ pb.TriggerExecutableServer = (*triggerExecutableServer)(nil) func (t *triggerExecutableServer) AckEvent(ctx context.Context, req *pb.AckEventRequest) (*emptypb.Empty, error) { + // Temporary instrumentation for the AckEvent stall investigation. Pair on eventId + // with the client-side lines to split total AckEvent latency into three segments: + // outbound transport, handler, and return transport. + t.Logger.Infow("grpc AckEvent server: entered", "eventId", req.EventId, "triggerId", req.TriggerId) + handlerStart := time.Now() if err := t.impl.AckEvent(ctx, req.TriggerId, req.EventId, req.Method); err != nil { + t.Logger.Infow("grpc AckEvent server: failed", "eventId", req.EventId, + "handlerMs", time.Since(handlerStart).Milliseconds(), "err", err) return nil, fmt.Errorf("error acking event: %w", err) } + t.Logger.Infow("grpc AckEvent server: done", "eventId", req.EventId, + "handlerMs", time.Since(handlerStart).Milliseconds()) return &emptypb.Empty{}, nil } @@ -306,7 +316,13 @@ func (t *triggerExecutableClient) AckEvent(ctx context.Context, triggerId string EventId: eventId, Method: method, } + // Temporary instrumentation: brackets the node->plugin gRPC hop. The gap between + // this line and "grpc AckEvent server: entered" is pure outbound transport time. + t.Logger.Infow("grpc AckEvent client: sending", "eventId", eventId, "triggerId", triggerId, "method", method) + callStart := time.Now() _, err := t.grpc.AckEvent(ctx, req) + t.Logger.Infow("grpc AckEvent client: returned", "eventId", eventId, + "callMs", time.Since(callStart).Milliseconds(), "err", err) if err != nil { return fmt.Errorf("failed to call AckEvent: %w", err) } From 1c7703887d99e710896d7ef590a5115c3e06cf8b Mon Sep 17 00:00:00 2001 From: yashnevatia Date: Fri, 4 Sep 2026 12:02:41 +0100 Subject: [PATCH 6/6] remove lock hold from AckEvent --- pkg/capabilities/registry/base.go | 27 +++++++++++++++++++-------- 1 file changed, 19 insertions(+), 8 deletions(-) diff --git a/pkg/capabilities/registry/base.go b/pkg/capabilities/registry/base.go index 6c2d7348e0..bf554d7268 100644 --- a/pkg/capabilities/registry/base.go +++ b/pkg/capabilities/registry/base.go @@ -384,13 +384,20 @@ func (a *atomicTriggerCapability) GetState() connectivity.State { return connectivity.State(-1) // unknown } +// AckEvent reads a.cap only; unlike RegisterTrigger/UnregisterTrigger it does not +// touch a.registrations, so it takes a read lock and releases it before the call. +// Holding the write lock across the call serialises every ACK for this capability: +// the downstream call crosses gRPC to the LOOP plugin and hits Postgres, so a single +// in-flight ACK blocks all others regardless of how much parallelism the caller has. +// Mirrors Execute. func (a *atomicTriggerCapability) AckEvent(ctx context.Context, triggerID string, eventID string, method string) error { - a.mu.Lock() - defer a.mu.Unlock() - if a.cap == nil { + a.mu.RLock() + cap := a.cap + a.mu.RUnlock() + if cap == nil { return errors.New("capability unavailable") } - return a.cap.AckEvent(ctx, triggerID, eventID, method) + return cap.AckEvent(ctx, triggerID, eventID, method) } func (a *atomicTriggerCapability) Load() *capabilities.TriggerCapability { @@ -560,13 +567,17 @@ func (a *atomicExecuteAndTriggerCapability) GetState() connectivity.State { return connectivity.State(-1) // unknown } +// AckEvent reads a.cap only; unlike RegisterTrigger/UnregisterTrigger it does not +// touch a.registrations, so it takes a read lock and releases it before the call. +// See atomicTriggerCapability.AckEvent. Mirrors Execute. func (a *atomicExecuteAndTriggerCapability) AckEvent(ctx context.Context, triggerID string, eventID string, method string) error { - a.mu.Lock() - defer a.mu.Unlock() - if a.cap == nil { + a.mu.RLock() + cap := a.cap + a.mu.RUnlock() + if cap == nil { return errors.New("capability unavailable") } - return a.cap.AckEvent(ctx, triggerID, eventID, method) + return cap.AckEvent(ctx, triggerID, eventID, method) } func (a *atomicExecuteAndTriggerCapability) Load() *capabilities.ExecutableAndTriggerCapability {