diff --git a/pkg/channels/dispatcher.go b/pkg/channels/dispatcher.go new file mode 100644 index 000000000..320ba54c7 --- /dev/null +++ b/pkg/channels/dispatcher.go @@ -0,0 +1,97 @@ +// PicoClaw - Ultra-lightweight personal AI agent +// License: MIT +// Copyright (c) 2026 PicoClaw contributors + +package channels + +import ( + "context" + + "jane/pkg/bus" + "jane/pkg/constants" + "jane/pkg/logger" +) + +func dispatchLoop[M any]( + ctx context.Context, + m *Manager, + subscribe func(context.Context) (M, bool), + getChannel func(M) string, + enqueue func(context.Context, *channelWorker, M) bool, + startMsg, stopMsg, unknownMsg, noWorkerMsg string, +) { + logger.InfoC("channels", startMsg) + + for { + msg, ok := subscribe(ctx) + if !ok { + logger.InfoC("channels", stopMsg) + return + } + + channel := getChannel(msg) + + // Silently skip internal channels + if constants.IsInternalChannel(channel) { + continue + } + + m.mu.RLock() + _, exists := m.channels[channel] + w, wExists := m.workers[channel] + m.mu.RUnlock() + + if !exists { + logger.WarnCF("channels", unknownMsg, map[string]any{"channel": channel}) + continue + } + + if wExists && w != nil { + if !enqueue(ctx, w, msg) { + return + } + } else if exists { + logger.WarnCF("channels", noWorkerMsg, map[string]any{"channel": channel}) + } + } +} + +func (m *Manager) dispatchOutbound(ctx context.Context) { + dispatchLoop( + ctx, m, + m.bus.SubscribeOutbound, + func(msg bus.OutboundMessage) string { return msg.Channel }, + func(ctx context.Context, w *channelWorker, msg bus.OutboundMessage) bool { + select { + case w.queue <- msg: + return true + case <-ctx.Done(): + return false + } + }, + "Outbound dispatcher started", + "Outbound dispatcher stopped", + "Unknown channel for outbound message", + "Channel has no active worker, skipping message", + ) +} + +func (m *Manager) dispatchOutboundMedia(ctx context.Context) { + dispatchLoop( + ctx, m, + m.bus.SubscribeOutboundMedia, + func(msg bus.OutboundMediaMessage) string { return msg.Channel }, + func(ctx context.Context, w *channelWorker, msg bus.OutboundMediaMessage) bool { + select { + case w.mediaQueue <- msg: + return true + case <-ctx.Done(): + return false + } + }, + "Outbound media dispatcher started", + "Outbound media dispatcher stopped", + "Unknown channel for outbound media message", + "Channel has no active worker, skipping media message", + ) +} diff --git a/pkg/channels/http.go b/pkg/channels/http.go new file mode 100644 index 000000000..b4d5ef6b5 --- /dev/null +++ b/pkg/channels/http.go @@ -0,0 +1,52 @@ +// PicoClaw - Ultra-lightweight personal AI agent +// License: MIT +// Copyright (c) 2026 PicoClaw contributors + +package channels + +import ( + "net/http" + "time" + + "jane/pkg/health" + "jane/pkg/logger" +) + +// SetupHTTPServer creates a shared HTTP server with the given listen address. +// It registers health endpoints from the health server and discovers channels +// that implement WebhookHandler and/or HealthChecker to register their handlers. +func (m *Manager) SetupHTTPServer(addr string, healthServer *health.Server) { + m.mux = http.NewServeMux() + + // Register health endpoints + if healthServer != nil { + healthServer.RegisterOnMux(m.mux) + } + + // Discover and register webhook handlers and health checkers + for name, ch := range m.channels { + if wh, ok := ch.(WebhookHandler); ok { + m.mux.Handle(wh.WebhookPath(), wh) + logger.InfoCF("channels", "Webhook handler registered", map[string]any{ + "channel": name, + "path": wh.WebhookPath(), + }) + } + if hc, ok := ch.(HealthChecker); ok { + m.mux.HandleFunc(hc.HealthPath(), hc.HealthHandler) + logger.InfoCF("channels", "Health endpoint registered", map[string]any{ + "channel": name, + "path": hc.HealthPath(), + }) + } + } + + m.httpServer = &http.Server{ + Addr: addr, + Handler: m.mux, + ReadTimeout: 30 * time.Second, + ReadHeaderTimeout: 10 * time.Second, + WriteTimeout: 30 * time.Second, + IdleTimeout: 120 * time.Second, + } +} diff --git a/pkg/channels/janitor.go b/pkg/channels/janitor.go new file mode 100644 index 000000000..ba12a33e2 --- /dev/null +++ b/pkg/channels/janitor.go @@ -0,0 +1,54 @@ +// PicoClaw - Ultra-lightweight personal AI agent +// License: MIT +// Copyright (c) 2026 PicoClaw contributors + +package channels + +import ( + "context" + "time" +) + +// runTTLJanitor periodically scans the typingStops and placeholders maps +// and evicts entries that have exceeded their TTL. This prevents memory +// accumulation when outbound paths fail to trigger preSend (e.g. LLM errors). +func (m *Manager) runTTLJanitor(ctx context.Context) { + ticker := time.NewTicker(janitorInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case now := <-ticker.C: + m.typingStops.Range(func(key, value any) bool { + if entry, ok := value.(typingEntry); ok { + if now.Sub(entry.createdAt) > typingStopTTL { + if _, loaded := m.typingStops.LoadAndDelete(key); loaded { + entry.stop() // idempotent, safe + } + } + } + return true + }) + m.reactionUndos.Range(func(key, value any) bool { + if entry, ok := value.(reactionEntry); ok { + if now.Sub(entry.createdAt) > typingStopTTL { + if _, loaded := m.reactionUndos.LoadAndDelete(key); loaded { + entry.undo() // idempotent, safe + } + } + } + return true + }) + m.placeholders.Range(func(key, value any) bool { + if entry, ok := value.(placeholderEntry); ok { + if now.Sub(entry.createdAt) > placeholderTTL { + m.placeholders.Delete(key) + } + } + return true + }) + } + } +} diff --git a/pkg/channels/manager.go b/pkg/channels/manager.go index bf08399f5..72164e32e 100644 --- a/pkg/channels/manager.go +++ b/pkg/channels/manager.go @@ -9,73 +9,16 @@ package channels import ( "context" "errors" - "fmt" - "math" "net/http" "sync" "time" - "golang.org/x/time/rate" - "jane/pkg/bus" "jane/pkg/config" - "jane/pkg/constants" - "jane/pkg/health" "jane/pkg/logger" "jane/pkg/media" ) -const ( - defaultChannelQueueSize = 16 - defaultRateLimit = 10 // default 10 msg/s - maxRetries = 3 - rateLimitDelay = 1 * time.Second - baseBackoff = 500 * time.Millisecond - maxBackoff = 8 * time.Second - - janitorInterval = 10 * time.Second - typingStopTTL = 5 * time.Minute - placeholderTTL = 10 * time.Minute -) - -// typingEntry wraps a typing stop function with a creation timestamp for TTL eviction. -type typingEntry struct { - stop func() - createdAt time.Time -} - -// reactionEntry wraps a reaction undo function with a creation timestamp for TTL eviction. -type reactionEntry struct { - undo func() - createdAt time.Time -} - -// placeholderEntry wraps a placeholder ID with a creation timestamp for TTL eviction. -type placeholderEntry struct { - id string - createdAt time.Time -} - -// channelRateConfig maps channel name to per-second rate limit. -var channelRateConfig = map[string]float64{ - "telegram": 20, - "discord": 1, - "slack": 1, - "matrix": 2, - "line": 10, - "qq": 5, - "irc": 2, -} - -type channelWorker struct { - ch Channel - queue chan bus.OutboundMessage - mediaQueue chan bus.OutboundMediaMessage - done chan struct{} - mediaDone chan struct{} - limiter *rate.Limiter -} - type Manager struct { channels map[string]Channel workers map[string]*channelWorker @@ -91,91 +34,6 @@ type Manager struct { reactionUndos sync.Map // "channel:chatID" → reactionEntry } -type asyncTask struct { - cancel context.CancelFunc -} - -// RecordPlaceholder registers a placeholder message for later editing. -// Implements PlaceholderRecorder. -func (m *Manager) RecordPlaceholder(channel, chatID, placeholderID string) { - key := channel + ":" + chatID - m.placeholders.Store(key, placeholderEntry{id: placeholderID, createdAt: time.Now()}) -} - -// SendPlaceholder sends a "Thinking…" placeholder for the given channel/chatID -// and records it for later editing. Returns true if a placeholder was sent. -func (m *Manager) SendPlaceholder(ctx context.Context, channel, chatID string) bool { - m.mu.RLock() - ch, ok := m.channels[channel] - m.mu.RUnlock() - if !ok { - return false - } - pc, ok := ch.(PlaceholderCapable) - if !ok { - return false - } - phID, err := pc.SendPlaceholder(ctx, chatID) - if err != nil || phID == "" { - return false - } - m.RecordPlaceholder(channel, chatID, phID) - return true -} - -// RecordTypingStop registers a typing stop function for later invocation. -// Implements PlaceholderRecorder. -func (m *Manager) RecordTypingStop(channel, chatID string, stop func()) { - key := channel + ":" + chatID - entry := typingEntry{stop: stop, createdAt: time.Now()} - if previous, loaded := m.typingStops.Swap(key, entry); loaded { - if oldEntry, ok := previous.(typingEntry); ok && oldEntry.stop != nil { - oldEntry.stop() - } - } -} - -// RecordReactionUndo registers a reaction undo function for later invocation. -// Implements PlaceholderRecorder. -func (m *Manager) RecordReactionUndo(channel, chatID string, undo func()) { - key := channel + ":" + chatID - m.reactionUndos.Store(key, reactionEntry{undo: undo, createdAt: time.Now()}) -} - -// preSend handles typing stop, reaction undo, and placeholder editing before sending a message. -// Returns true if the message was edited into a placeholder (skip Send). -func (m *Manager) preSend(ctx context.Context, name string, msg bus.OutboundMessage, ch Channel) bool { - key := name + ":" + msg.ChatID - - // 1. Stop typing - if v, loaded := m.typingStops.LoadAndDelete(key); loaded { - if entry, ok := v.(typingEntry); ok { - entry.stop() // idempotent, safe - } - } - - // 2. Undo reaction - if v, loaded := m.reactionUndos.LoadAndDelete(key); loaded { - if entry, ok := v.(reactionEntry); ok { - entry.undo() // idempotent, safe - } - } - - // 3. Try editing placeholder - if v, loaded := m.placeholders.LoadAndDelete(key); loaded { - if entry, ok := v.(placeholderEntry); ok && entry.id != "" { - if editor, ok := ch.(MessageEditor); ok { - if err := editor.EditMessage(ctx, msg.ChatID, entry.id, msg.Content); err == nil { - return true // edited successfully, skip Send - } - // edit failed → fall through to normal Send - } - } - } - - return false -} - func NewManager(cfg *config.Config, messageBus *bus.MessageBus, store media.MediaStore) (*Manager, error) { m := &Manager{ channels: make(map[string]Channel), @@ -298,45 +156,6 @@ func (m *Manager) initChannels() error { return nil } -// SetupHTTPServer creates a shared HTTP server with the given listen address. -// It registers health endpoints from the health server and discovers channels -// that implement WebhookHandler and/or HealthChecker to register their handlers. -func (m *Manager) SetupHTTPServer(addr string, healthServer *health.Server) { - m.mux = http.NewServeMux() - - // Register health endpoints - if healthServer != nil { - healthServer.RegisterOnMux(m.mux) - } - - // Discover and register webhook handlers and health checkers - for name, ch := range m.channels { - if wh, ok := ch.(WebhookHandler); ok { - m.mux.Handle(wh.WebhookPath(), wh) - logger.InfoCF("channels", "Webhook handler registered", map[string]any{ - "channel": name, - "path": wh.WebhookPath(), - }) - } - if hc, ok := ch.(HealthChecker); ok { - m.mux.HandleFunc(hc.HealthPath(), hc.HealthHandler) - logger.InfoCF("channels", "Health endpoint registered", map[string]any{ - "channel": name, - "path": hc.HealthPath(), - }) - } - } - - m.httpServer = &http.Server{ - Addr: addr, - Handler: m.mux, - ReadTimeout: 30 * time.Second, - ReadHeaderTimeout: 10 * time.Second, - WriteTimeout: 30 * time.Second, - IdleTimeout: 120 * time.Second, - } -} - func (m *Manager) StartAll(ctx context.Context) error { m.mu.Lock() defer m.mu.Unlock() @@ -458,322 +277,6 @@ func (m *Manager) StopAll(ctx context.Context) error { return nil } -// newChannelWorker creates a channelWorker with a rate limiter configured -// for the given channel name. -func newChannelWorker(name string, ch Channel) *channelWorker { - rateVal := float64(defaultRateLimit) - if r, ok := channelRateConfig[name]; ok { - rateVal = r - } - burst := int(math.Max(1, math.Ceil(rateVal/2))) - - return &channelWorker{ - ch: ch, - queue: make(chan bus.OutboundMessage, defaultChannelQueueSize), - mediaQueue: make(chan bus.OutboundMediaMessage, defaultChannelQueueSize), - done: make(chan struct{}), - mediaDone: make(chan struct{}), - limiter: rate.NewLimiter(rate.Limit(rateVal), burst), - } -} - -// runWorker processes outbound messages for a single channel, splitting -// messages that exceed the channel's maximum message length. -func (m *Manager) runWorker(ctx context.Context, name string, w *channelWorker) { - defer close(w.done) - for { - select { - case msg, ok := <-w.queue: - if !ok { - return - } - maxLen := 0 - if mlp, ok := w.ch.(MessageLengthProvider); ok { - maxLen = mlp.MaxMessageLength() - } - if maxLen > 0 && len([]rune(msg.Content)) > maxLen { - chunks := SplitMessage(msg.Content, maxLen) - for _, chunk := range chunks { - chunkMsg := msg - chunkMsg.Content = chunk - m.sendWithRetry(ctx, name, w, chunkMsg) - } - } else { - m.sendWithRetry(ctx, name, w, msg) - } - case <-ctx.Done(): - return - } - } -} - -// sendWithRetry sends a message through the channel with rate limiting and -// retry logic. It classifies errors to determine the retry strategy: -// - ErrNotRunning / ErrSendFailed: permanent, no retry -// - ErrRateLimit: fixed delay retry -// - ErrTemporary / unknown: exponential backoff retry -func (m *Manager) sendWithRetry(ctx context.Context, name string, w *channelWorker, msg bus.OutboundMessage) { - // Rate limit: wait for token - if err := w.limiter.Wait(ctx); err != nil { - // ctx canceled, shutting down - return - } - - // Pre-send: stop typing and try to edit placeholder - if m.preSend(ctx, name, msg, w.ch) { - return // placeholder was edited successfully, skip Send - } - - var lastErr error - for attempt := 0; attempt <= maxRetries; attempt++ { - lastErr = w.ch.Send(ctx, msg) - if lastErr == nil { - return - } - - // Permanent failures — don't retry - if errors.Is(lastErr, ErrNotRunning) || errors.Is(lastErr, ErrSendFailed) { - break - } - - // Last attempt exhausted — don't sleep - if attempt == maxRetries { - break - } - - // Rate limit error — fixed delay - if errors.Is(lastErr, ErrRateLimit) { - select { - case <-time.After(rateLimitDelay): - continue - case <-ctx.Done(): - return - } - } - - // ErrTemporary or unknown error — exponential backoff - backoff := min(time.Duration(float64(baseBackoff)*math.Pow(2, float64(attempt))), maxBackoff) - select { - case <-time.After(backoff): - case <-ctx.Done(): - return - } - } - - // All retries exhausted or permanent failure - logger.ErrorCF("channels", "Send failed", map[string]any{ - "channel": name, - "chat_id": msg.ChatID, - "error": lastErr.Error(), - "retries": maxRetries, - }) -} - -func dispatchLoop[M any]( - ctx context.Context, - m *Manager, - subscribe func(context.Context) (M, bool), - getChannel func(M) string, - enqueue func(context.Context, *channelWorker, M) bool, - startMsg, stopMsg, unknownMsg, noWorkerMsg string, -) { - logger.InfoC("channels", startMsg) - - for { - msg, ok := subscribe(ctx) - if !ok { - logger.InfoC("channels", stopMsg) - return - } - - channel := getChannel(msg) - - // Silently skip internal channels - if constants.IsInternalChannel(channel) { - continue - } - - m.mu.RLock() - _, exists := m.channels[channel] - w, wExists := m.workers[channel] - m.mu.RUnlock() - - if !exists { - logger.WarnCF("channels", unknownMsg, map[string]any{"channel": channel}) - continue - } - - if wExists && w != nil { - if !enqueue(ctx, w, msg) { - return - } - } else if exists { - logger.WarnCF("channels", noWorkerMsg, map[string]any{"channel": channel}) - } - } -} - -func (m *Manager) dispatchOutbound(ctx context.Context) { - dispatchLoop( - ctx, m, - m.bus.SubscribeOutbound, - func(msg bus.OutboundMessage) string { return msg.Channel }, - func(ctx context.Context, w *channelWorker, msg bus.OutboundMessage) bool { - select { - case w.queue <- msg: - return true - case <-ctx.Done(): - return false - } - }, - "Outbound dispatcher started", - "Outbound dispatcher stopped", - "Unknown channel for outbound message", - "Channel has no active worker, skipping message", - ) -} - -func (m *Manager) dispatchOutboundMedia(ctx context.Context) { - dispatchLoop( - ctx, m, - m.bus.SubscribeOutboundMedia, - func(msg bus.OutboundMediaMessage) string { return msg.Channel }, - func(ctx context.Context, w *channelWorker, msg bus.OutboundMediaMessage) bool { - select { - case w.mediaQueue <- msg: - return true - case <-ctx.Done(): - return false - } - }, - "Outbound media dispatcher started", - "Outbound media dispatcher stopped", - "Unknown channel for outbound media message", - "Channel has no active worker, skipping media message", - ) -} - -// runMediaWorker processes outbound media messages for a single channel. -func (m *Manager) runMediaWorker(ctx context.Context, name string, w *channelWorker) { - defer close(w.mediaDone) - for { - select { - case msg, ok := <-w.mediaQueue: - if !ok { - return - } - m.sendMediaWithRetry(ctx, name, w, msg) - case <-ctx.Done(): - return - } - } -} - -// sendMediaWithRetry sends a media message through the channel with rate limiting and -// retry logic. If the channel does not implement MediaSender, it silently skips. -func (m *Manager) sendMediaWithRetry(ctx context.Context, name string, w *channelWorker, msg bus.OutboundMediaMessage) { - ms, ok := w.ch.(MediaSender) - if !ok { - logger.DebugCF("channels", "Channel does not support MediaSender, skipping media", map[string]any{ - "channel": name, - }) - return - } - - // Rate limit: wait for token - if err := w.limiter.Wait(ctx); err != nil { - return - } - - var lastErr error - for attempt := 0; attempt <= maxRetries; attempt++ { - lastErr = ms.SendMedia(ctx, msg) - if lastErr == nil { - return - } - - // Permanent failures — don't retry - if errors.Is(lastErr, ErrNotRunning) || errors.Is(lastErr, ErrSendFailed) { - break - } - - // Last attempt exhausted — don't sleep - if attempt == maxRetries { - break - } - - // Rate limit error — fixed delay - if errors.Is(lastErr, ErrRateLimit) { - select { - case <-time.After(rateLimitDelay): - continue - case <-ctx.Done(): - return - } - } - - // ErrTemporary or unknown error — exponential backoff - backoff := min(time.Duration(float64(baseBackoff)*math.Pow(2, float64(attempt))), maxBackoff) - select { - case <-time.After(backoff): - case <-ctx.Done(): - return - } - } - - // All retries exhausted or permanent failure - logger.ErrorCF("channels", "SendMedia failed", map[string]any{ - "channel": name, - "chat_id": msg.ChatID, - "error": lastErr.Error(), - "retries": maxRetries, - }) -} - -// runTTLJanitor periodically scans the typingStops and placeholders maps -// and evicts entries that have exceeded their TTL. This prevents memory -// accumulation when outbound paths fail to trigger preSend (e.g. LLM errors). -func (m *Manager) runTTLJanitor(ctx context.Context) { - ticker := time.NewTicker(janitorInterval) - defer ticker.Stop() - - for { - select { - case <-ctx.Done(): - return - case now := <-ticker.C: - m.typingStops.Range(func(key, value any) bool { - if entry, ok := value.(typingEntry); ok { - if now.Sub(entry.createdAt) > typingStopTTL { - if _, loaded := m.typingStops.LoadAndDelete(key); loaded { - entry.stop() // idempotent, safe - } - } - } - return true - }) - m.reactionUndos.Range(func(key, value any) bool { - if entry, ok := value.(reactionEntry); ok { - if now.Sub(entry.createdAt) > typingStopTTL { - if _, loaded := m.reactionUndos.LoadAndDelete(key); loaded { - entry.undo() // idempotent, safe - } - } - } - return true - }) - m.placeholders.Range(func(key, value any) bool { - if entry, ok := value.(placeholderEntry); ok { - if now.Sub(entry.createdAt) > placeholderTTL { - m.placeholders.Delete(key) - } - } - return true - }) - } - } -} - func (m *Manager) GetChannel(name string) (Channel, bool) { m.mu.RLock() defer m.mu.RUnlock() @@ -824,66 +327,3 @@ func (m *Manager) UnregisterChannel(name string) { delete(m.workers, name) delete(m.channels, name) } - -// SendMessage sends an outbound message synchronously through the channel -// worker's rate limiter and retry logic. It blocks until the message is -// delivered (or all retries are exhausted), which preserves ordering when -// a subsequent operation depends on the message having been sent. -func (m *Manager) SendMessage(ctx context.Context, msg bus.OutboundMessage) error { - m.mu.RLock() - _, exists := m.channels[msg.Channel] - w, wExists := m.workers[msg.Channel] - m.mu.RUnlock() - - if !exists { - return fmt.Errorf("channel %s not found", msg.Channel) - } - if !wExists || w == nil { - return fmt.Errorf("channel %s has no active worker", msg.Channel) - } - - maxLen := 0 - if mlp, ok := w.ch.(MessageLengthProvider); ok { - maxLen = mlp.MaxMessageLength() - } - if maxLen > 0 && len([]rune(msg.Content)) > maxLen { - for _, chunk := range SplitMessage(msg.Content, maxLen) { - chunkMsg := msg - chunkMsg.Content = chunk - m.sendWithRetry(ctx, msg.Channel, w, chunkMsg) - } - } else { - m.sendWithRetry(ctx, msg.Channel, w, msg) - } - return nil -} - -func (m *Manager) SendToChannel(ctx context.Context, channelName, chatID, content string) error { - m.mu.RLock() - _, exists := m.channels[channelName] - w, wExists := m.workers[channelName] - m.mu.RUnlock() - - if !exists { - return fmt.Errorf("channel %s not found", channelName) - } - - msg := bus.OutboundMessage{ - Channel: channelName, - ChatID: chatID, - Content: content, - } - - if wExists && w != nil { - select { - case w.queue <- msg: - return nil - case <-ctx.Done(): - return ctx.Err() - } - } - - // Fallback: direct send (should not happen) - channel, _ := m.channels[channelName] - return channel.Send(ctx, msg) -} diff --git a/pkg/channels/sender.go b/pkg/channels/sender.go new file mode 100644 index 000000000..d3663f052 --- /dev/null +++ b/pkg/channels/sender.go @@ -0,0 +1,284 @@ +// PicoClaw - Ultra-lightweight personal AI agent +// License: MIT +// Copyright (c) 2026 PicoClaw contributors + +package channels + +import ( + "context" + "errors" + "fmt" + "math" + "time" + + "jane/pkg/bus" + "jane/pkg/logger" +) + +// RecordPlaceholder registers a placeholder message for later editing. +// Implements PlaceholderRecorder. +func (m *Manager) RecordPlaceholder(channel, chatID, placeholderID string) { + key := channel + ":" + chatID + m.placeholders.Store(key, placeholderEntry{id: placeholderID, createdAt: time.Now()}) +} + +// SendPlaceholder sends a "Thinking…" placeholder for the given channel/chatID +// and records it for later editing. Returns true if a placeholder was sent. +func (m *Manager) SendPlaceholder(ctx context.Context, channel, chatID string) bool { + m.mu.RLock() + ch, ok := m.channels[channel] + m.mu.RUnlock() + if !ok { + return false + } + pc, ok := ch.(PlaceholderCapable) + if !ok { + return false + } + phID, err := pc.SendPlaceholder(ctx, chatID) + if err != nil || phID == "" { + return false + } + m.RecordPlaceholder(channel, chatID, phID) + return true +} + +// RecordTypingStop registers a typing stop function for later invocation. +// Implements PlaceholderRecorder. +func (m *Manager) RecordTypingStop(channel, chatID string, stop func()) { + key := channel + ":" + chatID + entry := typingEntry{stop: stop, createdAt: time.Now()} + if previous, loaded := m.typingStops.Swap(key, entry); loaded { + if oldEntry, ok := previous.(typingEntry); ok && oldEntry.stop != nil { + oldEntry.stop() + } + } +} + +// RecordReactionUndo registers a reaction undo function for later invocation. +// Implements PlaceholderRecorder. +func (m *Manager) RecordReactionUndo(channel, chatID string, undo func()) { + key := channel + ":" + chatID + m.reactionUndos.Store(key, reactionEntry{undo: undo, createdAt: time.Now()}) +} + +// preSend handles typing stop, reaction undo, and placeholder editing before sending a message. +// Returns true if the message was edited into a placeholder (skip Send). +func (m *Manager) preSend(ctx context.Context, name string, msg bus.OutboundMessage, ch Channel) bool { + key := name + ":" + msg.ChatID + + // 1. Stop typing + if v, loaded := m.typingStops.LoadAndDelete(key); loaded { + if entry, ok := v.(typingEntry); ok { + entry.stop() // idempotent, safe + } + } + + // 2. Undo reaction + if v, loaded := m.reactionUndos.LoadAndDelete(key); loaded { + if entry, ok := v.(reactionEntry); ok { + entry.undo() // idempotent, safe + } + } + + // 3. Try editing placeholder + if v, loaded := m.placeholders.LoadAndDelete(key); loaded { + if entry, ok := v.(placeholderEntry); ok && entry.id != "" { + if editor, ok := ch.(MessageEditor); ok { + if err := editor.EditMessage(ctx, msg.ChatID, entry.id, msg.Content); err == nil { + return true // edited successfully, skip Send + } + // edit failed → fall through to normal Send + } + } + } + + return false +} + +// sendWithRetry sends a message through the channel with rate limiting and +// retry logic. It classifies errors to determine the retry strategy: +// - ErrNotRunning / ErrSendFailed: permanent, no retry +// - ErrRateLimit: fixed delay retry +// - ErrTemporary / unknown: exponential backoff retry +func (m *Manager) sendWithRetry(ctx context.Context, name string, w *channelWorker, msg bus.OutboundMessage) { + // Rate limit: wait for token + if err := w.limiter.Wait(ctx); err != nil { + // ctx canceled, shutting down + return + } + + // Pre-send: stop typing and try to edit placeholder + if m.preSend(ctx, name, msg, w.ch) { + return // placeholder was edited successfully, skip Send + } + + var lastErr error + for attempt := 0; attempt <= maxRetries; attempt++ { + lastErr = w.ch.Send(ctx, msg) + if lastErr == nil { + return + } + + // Permanent failures — don't retry + if errors.Is(lastErr, ErrNotRunning) || errors.Is(lastErr, ErrSendFailed) { + break + } + + // Last attempt exhausted — don't sleep + if attempt == maxRetries { + break + } + + // Rate limit error — fixed delay + if errors.Is(lastErr, ErrRateLimit) { + select { + case <-time.After(rateLimitDelay): + continue + case <-ctx.Done(): + return + } + } + + // ErrTemporary or unknown error — exponential backoff + backoff := min(time.Duration(float64(baseBackoff)*math.Pow(2, float64(attempt))), maxBackoff) + select { + case <-time.After(backoff): + case <-ctx.Done(): + return + } + } + + // All retries exhausted or permanent failure + logger.ErrorCF("channels", "Send failed", map[string]any{ + "channel": name, + "chat_id": msg.ChatID, + "error": lastErr.Error(), + "retries": maxRetries, + }) +} + +// sendMediaWithRetry sends a media message through the channel with rate limiting and +// retry logic. If the channel does not implement MediaSender, it silently skips. +func (m *Manager) sendMediaWithRetry(ctx context.Context, name string, w *channelWorker, msg bus.OutboundMediaMessage) { + ms, ok := w.ch.(MediaSender) + if !ok { + logger.DebugCF("channels", "Channel does not support MediaSender, skipping media", map[string]any{ + "channel": name, + }) + return + } + + // Rate limit: wait for token + if err := w.limiter.Wait(ctx); err != nil { + return + } + + var lastErr error + for attempt := 0; attempt <= maxRetries; attempt++ { + lastErr = ms.SendMedia(ctx, msg) + if lastErr == nil { + return + } + + // Permanent failures — don't retry + if errors.Is(lastErr, ErrNotRunning) || errors.Is(lastErr, ErrSendFailed) { + break + } + + // Last attempt exhausted — don't sleep + if attempt == maxRetries { + break + } + + // Rate limit error — fixed delay + if errors.Is(lastErr, ErrRateLimit) { + select { + case <-time.After(rateLimitDelay): + continue + case <-ctx.Done(): + return + } + } + + // ErrTemporary or unknown error — exponential backoff + backoff := min(time.Duration(float64(baseBackoff)*math.Pow(2, float64(attempt))), maxBackoff) + select { + case <-time.After(backoff): + case <-ctx.Done(): + return + } + } + + // All retries exhausted or permanent failure + logger.ErrorCF("channels", "SendMedia failed", map[string]any{ + "channel": name, + "chat_id": msg.ChatID, + "error": lastErr.Error(), + "retries": maxRetries, + }) +} + +// SendMessage sends an outbound message synchronously through the channel +// worker's rate limiter and retry logic. It blocks until the message is +// delivered (or all retries are exhausted), which preserves ordering when +// a subsequent operation depends on the message having been sent. +func (m *Manager) SendMessage(ctx context.Context, msg bus.OutboundMessage) error { + m.mu.RLock() + _, exists := m.channels[msg.Channel] + w, wExists := m.workers[msg.Channel] + m.mu.RUnlock() + + if !exists { + return fmt.Errorf("channel %s not found", msg.Channel) + } + if !wExists || w == nil { + return fmt.Errorf("channel %s has no active worker", msg.Channel) + } + + maxLen := 0 + if mlp, ok := w.ch.(MessageLengthProvider); ok { + maxLen = mlp.MaxMessageLength() + } + if maxLen > 0 && len([]rune(msg.Content)) > maxLen { + for _, chunk := range SplitMessage(msg.Content, maxLen) { + chunkMsg := msg + chunkMsg.Content = chunk + m.sendWithRetry(ctx, msg.Channel, w, chunkMsg) + } + } else { + m.sendWithRetry(ctx, msg.Channel, w, msg) + } + return nil +} + +// SendToChannel sends a message to the specified channel. +func (m *Manager) SendToChannel(ctx context.Context, channelName, chatID, content string) error { + m.mu.RLock() + _, exists := m.channels[channelName] + w, wExists := m.workers[channelName] + m.mu.RUnlock() + + if !exists { + return fmt.Errorf("channel %s not found", channelName) + } + + msg := bus.OutboundMessage{ + Channel: channelName, + ChatID: chatID, + Content: content, + } + + if wExists && w != nil { + select { + case w.queue <- msg: + return nil + case <-ctx.Done(): + return ctx.Err() + } + } + + // Fallback: direct send (should not happen) + channel, _ := m.channels[channelName] + return channel.Send(ctx, msg) +} diff --git a/pkg/channels/types.go b/pkg/channels/types.go new file mode 100644 index 000000000..791820173 --- /dev/null +++ b/pkg/channels/types.go @@ -0,0 +1,56 @@ +// PicoClaw - Ultra-lightweight personal AI agent +// License: MIT +// Copyright (c) 2026 PicoClaw contributors + +package channels + +import ( + "context" + "time" +) + +const ( + defaultChannelQueueSize = 16 + defaultRateLimit = 10 // default 10 msg/s + maxRetries = 3 + rateLimitDelay = 1 * time.Second + baseBackoff = 500 * time.Millisecond + maxBackoff = 8 * time.Second + + janitorInterval = 10 * time.Second + typingStopTTL = 5 * time.Minute + placeholderTTL = 10 * time.Minute +) + +// typingEntry wraps a typing stop function with a creation timestamp for TTL eviction. +type typingEntry struct { + stop func() + createdAt time.Time +} + +// reactionEntry wraps a reaction undo function with a creation timestamp for TTL eviction. +type reactionEntry struct { + undo func() + createdAt time.Time +} + +// placeholderEntry wraps a placeholder ID with a creation timestamp for TTL eviction. +type placeholderEntry struct { + id string + createdAt time.Time +} + +// channelRateConfig maps channel name to per-second rate limit. +var channelRateConfig = map[string]float64{ + "telegram": 20, + "discord": 1, + "slack": 1, + "matrix": 2, + "line": 10, + "qq": 5, + "irc": 2, +} + +type asyncTask struct { + cancel context.CancelFunc +} diff --git a/pkg/channels/worker.go b/pkg/channels/worker.go new file mode 100644 index 000000000..6f08a57f2 --- /dev/null +++ b/pkg/channels/worker.go @@ -0,0 +1,88 @@ +// PicoClaw - Ultra-lightweight personal AI agent +// License: MIT +// Copyright (c) 2026 PicoClaw contributors + +package channels + +import ( + "context" + "math" + + "golang.org/x/time/rate" + + "jane/pkg/bus" +) + +type channelWorker struct { + ch Channel + queue chan bus.OutboundMessage + mediaQueue chan bus.OutboundMediaMessage + done chan struct{} + mediaDone chan struct{} + limiter *rate.Limiter +} + +// newChannelWorker creates a channelWorker with a rate limiter configured +// for the given channel name. +func newChannelWorker(name string, ch Channel) *channelWorker { + rateVal := float64(defaultRateLimit) + if r, ok := channelRateConfig[name]; ok { + rateVal = r + } + burst := int(math.Max(1, math.Ceil(rateVal/2))) + + return &channelWorker{ + ch: ch, + queue: make(chan bus.OutboundMessage, defaultChannelQueueSize), + mediaQueue: make(chan bus.OutboundMediaMessage, defaultChannelQueueSize), + done: make(chan struct{}), + mediaDone: make(chan struct{}), + limiter: rate.NewLimiter(rate.Limit(rateVal), burst), + } +} + +// runWorker processes outbound messages for a single channel, splitting +// messages that exceed the channel's maximum message length. +func (m *Manager) runWorker(ctx context.Context, name string, w *channelWorker) { + defer close(w.done) + for { + select { + case msg, ok := <-w.queue: + if !ok { + return + } + maxLen := 0 + if mlp, ok := w.ch.(MessageLengthProvider); ok { + maxLen = mlp.MaxMessageLength() + } + if maxLen > 0 && len([]rune(msg.Content)) > maxLen { + chunks := SplitMessage(msg.Content, maxLen) + for _, chunk := range chunks { + chunkMsg := msg + chunkMsg.Content = chunk + m.sendWithRetry(ctx, name, w, chunkMsg) + } + } else { + m.sendWithRetry(ctx, name, w, msg) + } + case <-ctx.Done(): + return + } + } +} + +// runMediaWorker processes outbound media messages for a single channel. +func (m *Manager) runMediaWorker(ctx context.Context, name string, w *channelWorker) { + defer close(w.mediaDone) + for { + select { + case msg, ok := <-w.mediaQueue: + if !ok { + return + } + m.sendMediaWithRetry(ctx, name, w, msg) + case <-ctx.Done(): + return + } + } +}