butterrobot/internal/queue/queue.go
Felipe M. a2e1196953
Some checks failed
CI / goreleaser-lint (push) Successful in 6s
CI / format (push) Successful in 57s
CI / test (push) Successful in 2m59s
CI / lint (push) Successful in 4m9s
CI / build (push) Successful in 7m17s
Release / release (push) Failing after 7m6s
feat: add Sentry observability support
Report errors, panics and optional performance traces to Sentry. Sentry
stays disabled unless SENTRY_DSN is set.

- slog ERROR records become Sentry issues, with the error attribute
  promoted to an exception so issues group by root cause
- queue worker, reminder worker, cache cleanup and HTTP handler panics
  are captured with a stack trace
- events carry platform/plugin/component tags for filtering
- HTTP requests are traced when SENTRY_TRACES_SAMPLE_RATE is above 0;
  /healthz is never traced

Configured credentials are redacted from every outgoing payload. This is
required rather than defensive: the Telegram webhook embeds the bot token
in its URL path, and failed Telegram API calls quote that URL in their
error text, so events would otherwise carry the token in the clear. A new
platform must register its credential in Config.Secrets().

Also moves the module to Go 1.27, refreshes every dependency and pins
golangci-lint v2.13.2. sentry-go 0.48 removed issue creation from its slog
integration, so the ERROR-to-issue conversion lives in
internal/observability/handler.go instead of relying on the SDK; leaving
it to the SDK would have silently downgraded issues to log lines.

Claude-Session: https://claude.ai/code/session_01W7tcpMTEyk9RrHvT7Be5zZ
2026-09-21 14:10:15 +02:00

154 lines
3 KiB
Go

package queue
import (
"log/slog"
"sync"
"time"
"git.nakama.town/fmartingr/butterrobot/internal/model"
"git.nakama.town/fmartingr/butterrobot/internal/observability"
)
// Item represents a queue item
type Item struct {
Platform string
Request map[string]interface{}
}
// HandlerFunc defines a function that processes queue items
type HandlerFunc func(item Item)
// ReminderHandlerFunc defines a function that processes reminder items
type ReminderHandlerFunc func(reminder *model.Reminder)
// Queue represents a message queue
type Queue struct {
items chan Item
wg sync.WaitGroup
quit chan struct{}
logger *slog.Logger
running bool
runMutex sync.Mutex
reminderTicker *time.Ticker
reminderHandler ReminderHandlerFunc
}
// New creates a new Queue instance
func New(logger *slog.Logger) *Queue {
return &Queue{
items: make(chan Item, 100),
quit: make(chan struct{}),
logger: logger,
}
}
// Start starts processing queue items
func (q *Queue) Start(handler HandlerFunc) {
q.runMutex.Lock()
defer q.runMutex.Unlock()
if q.running {
return
}
q.running = true
// Start worker
q.wg.Add(1)
go q.worker(handler)
}
// StartReminderScheduler starts the reminder scheduler
func (q *Queue) StartReminderScheduler(handler ReminderHandlerFunc) {
q.runMutex.Lock()
defer q.runMutex.Unlock()
if q.reminderTicker != nil {
return
}
q.reminderHandler = handler
// Check for reminders every minute
q.reminderTicker = time.NewTicker(1 * time.Minute)
q.wg.Add(1)
go q.reminderWorker()
}
// Stop stops processing queue items
func (q *Queue) Stop() {
q.runMutex.Lock()
defer q.runMutex.Unlock()
if !q.running {
return
}
q.running = false
// Stop reminder ticker if it exists
if q.reminderTicker != nil {
q.reminderTicker.Stop()
}
close(q.quit)
q.wg.Wait()
}
// Add adds an item to the queue
func (q *Queue) Add(item Item) {
select {
case q.items <- item:
// Item added successfully
default:
// Queue is full
q.logger.Info("Queue is full, dropping message")
}
}
// worker processes queue items
func (q *Queue) worker(handler HandlerFunc) {
defer q.wg.Done()
for {
select {
case item := <-q.items:
// Process item
func() {
defer observability.RecoverAndCapture(q.logger, "queue worker")
handler(item)
}()
case <-q.quit:
// Quit worker
return
}
}
}
// reminderWorker processes reminder items on a schedule
func (q *Queue) reminderWorker() {
defer q.wg.Done()
for {
select {
case <-q.reminderTicker.C:
// This is triggered every minute to check for pending reminders
q.logger.Debug("Checking for pending reminders")
if q.reminderHandler != nil {
// The handler is responsible for fetching and processing reminders
func() {
defer observability.RecoverAndCapture(q.logger, "reminder worker")
// Call the handler with a nil reminder to indicate it should check the database
q.reminderHandler(nil)
}()
}
case <-q.quit:
// Quit worker
return
}
}
}