This repository has been archived on 2026-04-03. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
core/server/compaction/scanner.go
2026-03-19 18:50:27 +00:00

362 lines
9.1 KiB
Go

package compaction
import (
"context"
"log"
"sync"
"time"
"chat-switchboard/models"
"chat-switchboard/roles"
"chat-switchboard/store"
"chat-switchboard/treepath"
)
// ── Defaults ────────────────────────────────
const (
DefaultInterval = 5 * time.Minute
DefaultConcurrency = 2
DefaultThreshold = 0.70
DefaultCooldown = 30 * time.Minute
DefaultBudget = 128000 // 128K tokens — conservative fallback
// Candidate query filters
candidateMinMessages = 10
candidateMinChars = 20000 // ~5K tokens
candidateActivityGap = 2 * time.Minute
candidateMaxAge = 7 * 24 * time.Hour
candidateBatchSize = 50
// shouldCompact requires this many post-summary messages
minPostSummaryMessages = 8
)
// ── Scanner Config ──────────────────────────
// ScannerConfig holds startup configuration for the scanner.
// Settings are re-read from global_settings each tick so changes
// take effect without restart.
type ScannerConfig struct {
Interval time.Duration
Concurrency int
}
// ── Scanner ─────────────────────────────────
// Scanner runs a periodic background loop that finds channels exceeding
// their context threshold and auto-compacts them via the utility model role.
//
// Follows the Ingester pattern: goroutine pool with semaphore, WaitGroup
// for graceful shutdown.
type Scanner struct {
service *Service
stores store.Stores
sem chan struct{}
wg sync.WaitGroup
stopCh chan struct{}
// Cooldown: channelID → last compaction time
mu sync.Mutex
lastCompacted map[string]time.Time
// In-flight dedup
inFlight sync.Map // channelID → struct{}
interval time.Duration
}
// NewScanner creates a compaction scanner. Call Start() to begin scanning.
func NewScanner(svc *Service, stores store.Stores, cfg ScannerConfig) *Scanner {
interval := cfg.Interval
if interval <= 0 {
interval = DefaultInterval
}
concurrency := cfg.Concurrency
if concurrency <= 0 {
concurrency = DefaultConcurrency
}
return &Scanner{
service: svc,
stores: stores,
sem: make(chan struct{}, concurrency),
stopCh: make(chan struct{}),
lastCompacted: make(map[string]time.Time),
interval: interval,
}
}
// Start begins the scan loop in a background goroutine.
func (sc *Scanner) Start() {
sc.wg.Add(1)
go func() {
defer sc.wg.Done()
sc.loop()
}()
log.Printf("🔍 compaction scanner started (interval=%s)", sc.interval)
}
// Stop signals the scanner to stop and waits for in-flight compactions
// to drain. Safe to call from main's defer chain.
func (sc *Scanner) Stop() {
close(sc.stopCh)
sc.wg.Wait()
log.Printf("🔍 compaction scanner stopped")
}
// ── Main Loop ───────────────────────────────
func (sc *Scanner) loop() {
ticker := time.NewTicker(sc.interval)
defer ticker.Stop()
for {
select {
case <-sc.stopCh:
return
case <-ticker.C:
sc.tick()
}
}
}
func (sc *Scanner) tick() {
ctx := context.Background()
// Re-read kill switch each tick (no restart required)
if !sc.isEnabled(ctx) {
return
}
candidates := sc.findCandidates(ctx)
if len(candidates) == 0 {
return
}
log.Printf("🔍 compaction: scanning %d candidates", len(candidates))
cooldown := sc.getCooldownDuration(ctx)
for i := range candidates {
ch := candidates[i]
// Skip if in-flight
if _, loaded := sc.inFlight.LoadOrStore(ch.ID, struct{}{}); loaded {
continue
}
// Skip if in cooldown
sc.mu.Lock()
if last, ok := sc.lastCompacted[ch.ID]; ok && time.Since(last) < cooldown {
sc.mu.Unlock()
sc.inFlight.Delete(ch.ID)
remaining := cooldown - time.Since(last)
log.Printf("⏭ compaction: channel %s skipped (cooldown, %s remaining)", ch.ID, remaining.Round(time.Second))
continue
}
sc.mu.Unlock()
// Precise check (loads full path, estimates tokens)
if !sc.shouldCompact(ctx, &ch) {
sc.inFlight.Delete(ch.ID)
continue
}
// Rate limit check (auto-compaction is always org-funded)
if err := sc.service.CheckRateLimit(ctx, ch.UserID); err != nil {
sc.inFlight.Delete(ch.ID)
log.Printf("⏭ compaction: channel %s skipped (rate limit)", ch.ID)
continue
}
// Dispatch compaction
sc.wg.Add(1)
go func(ch models.Channel) {
defer sc.wg.Done()
defer sc.inFlight.Delete(ch.ID)
// Acquire semaphore
sc.sem <- struct{}{}
defer func() { <-sc.sem }()
sc.doCompact(ch)
}(ch)
}
}
// ── Compact One Channel ─────────────────────
func (sc *Scanner) doCompact(ch models.Channel) {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
teamID := sc.service.GetUserTeamID(ctx, ch.UserID)
result, err := sc.service.Compact(ctx, CompactRequest{
ChannelID: ch.ID,
UserID: ch.UserID,
TeamID: teamID,
Trigger: "auto",
})
if err != nil {
log.Printf("⚠ compaction: channel %s failed: %v", ch.ID, err)
return
}
// Record cooldown
sc.mu.Lock()
sc.lastCompacted[ch.ID] = time.Now()
sc.mu.Unlock()
// Audit log
sc.stores.Audit.Log(ctx, &models.AuditEntry{
ActorID: nil, // system action
Action: "compaction.auto",
ResourceType: "channel",
ResourceID: ch.ID,
Metadata: models.JSONMap{
"summarized_count": result.SummarizedCount,
"model": result.Model,
"trigger": "auto",
"user_id": ch.UserID,
},
})
log.Printf("✅ compaction: channel %s done (%d messages → %d chars, model=%s)",
ch.ID, result.SummarizedCount, len(result.Content), result.Model)
}
// ── Candidate Query ─────────────────────────
func (sc *Scanner) findCandidates(ctx context.Context) []models.Channel {
activityGap := time.Now().Add(-candidateActivityGap)
maxAge := time.Now().Add(-candidateMaxAge)
candidates, err := sc.stores.Channels.FindCompactionCandidates(ctx,
activityGap, maxAge, candidateMinMessages, candidateMinChars, candidateBatchSize)
if err != nil {
log.Printf("⚠ compaction: candidate query failed: %v", err)
return nil
}
return candidates
}
// ── Precise Check ───────────────────────────
func (sc *Scanner) shouldCompact(ctx context.Context, ch *models.Channel) bool {
// 1. Channel-level opt-out
if ch.Settings != nil {
if v, ok := ch.Settings["auto_compaction"]; ok {
if b, ok := v.(bool); ok && !b {
return false
}
}
}
// 2. Check utility role is configured (global level — personal overrides
// aren't checked because auto-compaction is org-funded)
if !sc.service.Resolver().IsConfigured(ctx, roles.RoleUtility) {
return false
}
// 3. Load active path
path, err := treepath.GetActivePath(ch.ID, ch.UserID)
if err != nil || len(path) < candidateMinMessages {
return false
}
// 4. Estimate tokens from post-summary messages
tokens, msgCount, ok := EstimatePathFromSummary(path, minPostSummaryMessages)
if !ok {
return false
}
// 5. Compare against threshold
budget := sc.getContextBudget(ctx, ch)
threshold := sc.getThreshold(ch)
ratio := float64(tokens) / float64(budget)
if ratio >= threshold {
log.Printf("🔍 compaction: channel %s qualifies (%.0f%% of %dK context, %d messages post-summary)",
ch.ID, ratio*100, budget/1000, msgCount)
return true
}
return false
}
// ── Settings Helpers ────────────────────────
func (sc *Scanner) isEnabled(ctx context.Context) bool {
val, err := sc.stores.GlobalConfig.Get(ctx, "auto_compaction_enabled")
if err != nil || val == nil {
return false // default: off
}
if v, ok := val["value"]; ok {
switch b := v.(type) {
case bool:
return b
case string:
return b == "true"
}
}
return false
}
func (sc *Scanner) getContextBudget(ctx context.Context, ch *models.Channel) int {
// Try channel's model from catalog
if ch.Model != "" {
entry, err := sc.stores.Catalog.GetByModelIDAny(ctx, ch.Model)
if err == nil && entry.Capabilities.MaxContext > 0 {
return entry.Capabilities.MaxContext
}
}
return DefaultBudget
}
func (sc *Scanner) getThreshold(ch *models.Channel) float64 {
// Channel override
if ch.Settings != nil {
if v, ok := ch.Settings["compaction_threshold"]; ok {
switch f := v.(type) {
case float64:
if f > 0 && f < 1 {
return f
}
}
}
}
// Global override
ctx := context.Background()
val, err := sc.stores.GlobalConfig.Get(ctx, "auto_compaction_threshold")
if err == nil && val != nil {
if v, ok := val["value"]; ok {
if f, ok := v.(float64); ok && f > 0 && f < 1 {
return f
}
}
}
return DefaultThreshold
}
func (sc *Scanner) getCooldownDuration(ctx context.Context) time.Duration {
val, err := sc.stores.GlobalConfig.Get(ctx, "auto_compaction_cooldown_minutes")
if err == nil && val != nil {
if v, ok := val["value"]; ok {
switch n := v.(type) {
case float64:
if n > 0 {
return time.Duration(n) * time.Minute
}
case int:
if n > 0 {
return time.Duration(n) * time.Minute
}
}
}
}
return DefaultCooldown
}