Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
126 changes: 107 additions & 19 deletions pkg/capabilities/base_trigger.go
Original file line number Diff line number Diff line change
Expand Up @@ -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].
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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.
Expand All @@ -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
Expand All @@ -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)
}
}
Expand All @@ -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 {
Expand Down Expand Up @@ -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)
Expand All @@ -363,22 +423,26 @@ 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)

// Double-check preAcked under the same lock as adding to pending.
// 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)
b.mu.Unlock()
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)
}
Expand All @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Comment on lines +651 to +654

toResend := make([]PendingEvent, 0, defaultMaxSendsPerTick)
var toStop []stoppedResendingEvent
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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)
}
Expand All @@ -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)
}
Comment thread
Copilot marked this conversation as resolved.
Expand All @@ -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()
Expand All @@ -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 {
Expand Down
Loading
Loading