Files
libredesk/internal/notification/notification.go
T
Abhinav Raut 6b8a0f9521 fix content loss in knowledge base chunking and harden AI agent limits
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.
2026-07-25 03:52:32 +05:30

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()
}