Merge pull request #47 from hobbyistlabs-coder/refactor-channel-manager-1242404581752262388
🧹 refactor pkg/channels/manager.go into multiple files
This commit is contained in:
commit
fec84e473b
7 changed files with 631 additions and 560 deletions
97
pkg/channels/dispatcher.go
Normal file
97
pkg/channels/dispatcher.go
Normal file
|
|
@ -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",
|
||||
)
|
||||
}
|
||||
52
pkg/channels/http.go
Normal file
52
pkg/channels/http.go
Normal file
|
|
@ -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,
|
||||
}
|
||||
}
|
||||
54
pkg/channels/janitor.go
Normal file
54
pkg/channels/janitor.go
Normal file
|
|
@ -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
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
284
pkg/channels/sender.go
Normal file
284
pkg/channels/sender.go
Normal file
|
|
@ -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)
|
||||
}
|
||||
56
pkg/channels/types.go
Normal file
56
pkg/channels/types.go
Normal file
|
|
@ -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
|
||||
}
|
||||
88
pkg/channels/worker.go
Normal file
88
pkg/channels/worker.go
Normal file
|
|
@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue