Some checks failed
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
154 lines
3 KiB
Go
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
|
|
}
|
|
}
|
|
}
|