mirror of
https://github.com/abhinavxd/libredesk.git
synced 2026-10-04 12:31:32 +00:00
6b8a0f9521
Final review pass before taking the AI agent branch live. Knowledge base: - Text not wrapped in a block tag was never collected, so prose around a table or list never reached the index. The assistant answered "no relevant information" for questions the snippet covered. - Blocks over the token limit were truncated and the remainder dropped. They are split into several chunks now. - Trimming an oversized block ran one rune at a time and re-tokenized the whole string each step. A large table took minutes. It uses a binary search now. - Overlap text was not escaped, so a sentence containing markup swallowed the rest of the chunk. - SVG and template text no longer reaches the index. AI agent: - Verification codes are capped per address and per conversation. The cap was per conversation only, so a customer correcting a mistyped email was told to check an inbox that never got a code. - Livechat verification sends synchronously. A queued send returned nil even when SMTP failed, so a failure counted as a sent code. - Queued jobs drain on shutdown and hand off to a human instead of being dropped with no reply. - Deleting an assistant no longer moves resolved and closed conversations into the fallback team. - Image decode is capped at 25 MP. The old bound allowed a 400 MB decode per attachment. Auth and admin: - A blank OIDC client secret no longer overwrites the stored one. Blank id or secret is rejected instead. - OIDC token exchange uses the SSRF guarded client with a timeout. - Renaming a tool auth header no longer attaches the secret of whichever row now sits at that position. - Clearing embedding dimensions no longer refills 1536 on the next load, which pushed a wrong value to the provider on the next save. - Copilot conversation lookups filter by access before capping at 10.
131 lines
3.1 KiB
Go
131 lines
3.1 KiB
Go
package notifier
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"sync"
|
|
|
|
"github.com/abhinavxd/libredesk/internal/attachment"
|
|
"github.com/zerodha/logf"
|
|
)
|
|
|
|
const (
|
|
ProviderEmail = "email"
|
|
)
|
|
|
|
// Message represents a message to be sent as a notification.
|
|
type Message struct {
|
|
// Email addresses of the recipients
|
|
RecipientEmails []string
|
|
// Subject of the message
|
|
Subject string
|
|
// Body of the message
|
|
Content string
|
|
// Provider to send the message through
|
|
Provider string
|
|
// Attachments to be sent with the message
|
|
Attachments []attachment.Attachment
|
|
// Type of content ("plain" or "html")
|
|
ContentType string
|
|
// Alternative plain text version of the HTML content
|
|
AltContent string
|
|
// Additional email headers
|
|
Headers map[string][]string
|
|
}
|
|
|
|
// Notifier defines the interface for sending notifications through various providers.
|
|
type Notifier interface {
|
|
// Sends the notification message using the specified provider
|
|
Send(message Message) error
|
|
// Returns the name of the provider
|
|
Name() string
|
|
}
|
|
|
|
// Service manages message providers and a worker pool.
|
|
type Service struct {
|
|
providers map[string]Notifier
|
|
messageChannel chan Message
|
|
concurrency int
|
|
lo *logf.Logger
|
|
closed bool
|
|
mu sync.RWMutex
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
// NewService initializes the Service with given concurrency, channel capacity, and logger.
|
|
func NewService(providers map[string]Notifier, concurrency, capacity int, logger *logf.Logger) *Service {
|
|
return &Service{
|
|
providers: providers,
|
|
messageChannel: make(chan Message, capacity),
|
|
concurrency: concurrency,
|
|
lo: logger,
|
|
}
|
|
}
|
|
|
|
// Send sends a message to the message channel.
|
|
func (s *Service) Send(message Message) error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if s.closed {
|
|
return fmt.Errorf("channel closed cannot send message")
|
|
}
|
|
|
|
select {
|
|
case s.messageChannel <- message:
|
|
return nil
|
|
default:
|
|
s.lo.Error("message channel is full")
|
|
return fmt.Errorf("message channel is full")
|
|
}
|
|
}
|
|
|
|
// SendSync sends on the caller's goroutine so delivery failures reach the caller; it holds no lock, so a slow relay cannot stall Send or Close.
|
|
func (s *Service) SendSync(message Message) error {
|
|
provider, exists := s.providers[message.Provider]
|
|
if !exists {
|
|
return fmt.Errorf("unsupported provider: %s", message.Provider)
|
|
}
|
|
return provider.Send(message)
|
|
}
|
|
|
|
// Run starts the worker pool to process messages.
|
|
func (s *Service) Run(ctx context.Context) {
|
|
for range s.concurrency {
|
|
s.wg.Go(func() {
|
|
s.worker(ctx)
|
|
})
|
|
}
|
|
<-ctx.Done()
|
|
s.Close()
|
|
}
|
|
|
|
// worker processes messages from the message channel and sends them using the set provider.
|
|
func (s *Service) worker(ctx context.Context) {
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case message, ok := <-s.messageChannel:
|
|
if !ok {
|
|
return
|
|
}
|
|
if err := s.SendSync(message); err != nil {
|
|
s.lo.Error("error sending message", "error", err)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Close signals service to stop, closes the message channel and
|
|
// waits for all goroutine workers to finish.
|
|
func (s *Service) Close() {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
if s.closed {
|
|
return
|
|
}
|
|
s.closed = true
|
|
close(s.messageChannel)
|
|
s.wg.Wait()
|
|
}
|