package safew import ( "context" "errors" "log/slog" "strings" "sync" "time" "aiaa-notification-service/internal/adapter" ) type Poller func(token string, offset int64, timeout int) ([]adapter.SafewChat, int64, error) const pollLockTTL = 60 * time.Second type runner struct { cancel context.CancelFunc firstDone chan struct{} mu sync.Mutex firstErr error } func (r *runner) setFirst(err error) { r.mu.Lock() r.firstErr = err r.mu.Unlock() select { case <-r.firstDone: default: close(r.firstDone) } } func (r *runner) waitFirst(ctx context.Context) error { select { case <-ctx.Done(): return ctx.Err() case <-r.firstDone: r.mu.Lock() defer r.mu.Unlock() return r.firstErr } } type Watcher struct { store ChatStore poll Poller bgIdle time.Duration mu sync.Mutex running map[string]*runner stopped bool } func NewWatcher(store ChatStore, poll Poller) *Watcher { return &Watcher{ store: store, poll: poll, bgIdle: time.Second, running: map[string]*runner{}, } } func (w *Watcher) Ensure(token string) { w.ensure(token) } func (w *Watcher) StartTokens(tokens []string) { for _, token := range tokens { w.Ensure(strings.TrimSpace(token)) } } func (w *Watcher) ensure(token string) *runner { if token == "" || w == nil { return nil } w.mu.Lock() defer w.mu.Unlock() if w.stopped { return nil } if r, ok := w.running[token]; ok { return r } ctx, cancel := context.WithCancel(context.Background()) r := &runner{cancel: cancel, firstDone: make(chan struct{})} w.running[token] = r slog.Info("safew watcher started") go w.loop(ctx, token, r) return r } func (w *Watcher) Stop() { w.mu.Lock() w.stopped = true for _, r := range w.running { r.cancel() } w.running = map[string]*runner{} w.mu.Unlock() } func (w *Watcher) Refresh(ctx context.Context, token string) error { r := w.ensure(token) if r == nil { return nil } return r.waitFirst(ctx) } func (w *Watcher) List(ctx context.Context, token, q string) ([]adapter.SafewChat, error) { chats, err := w.store.ListSafewChats(ctx, token) if err != nil { return nil, err } return adapter.FilterSafewChats(chats, q), nil } func (w *Watcher) loop(ctx context.Context, token string, r *runner) { for { err := w.pollOnce(ctx, token, 0) if ctx.Err() != nil { r.setFirst(ctx.Err()) return } if isSafewAuth(err) { r.setFirst(err) return } if isGetUpdatesConflict(err) { if !w.sleep(ctx, conflictBackoff(w.bgIdle)) { r.setFirst(ctx.Err()) return } continue } r.setFirst(err) break } for { if err := w.pollOnce(ctx, token, 30); err != nil { if ctx.Err() != nil { return } if isSafewAuth(err) { slog.Warn("safew watcher poll", "error", err) return } if isGetUpdatesConflict(err) { slog.Info("safew watcher poll conflict", "error", err) if !w.sleep(ctx, conflictBackoff(w.bgIdle)) { return } continue } slog.Warn("safew watcher poll", "error", err) } if !w.sleep(ctx, w.bgIdle) { return } } } func (w *Watcher) sleep(ctx context.Context, d time.Duration) bool { if d <= 0 { d = time.Second } select { case <-ctx.Done(): return false case <-time.After(d): return true } } func (w *Watcher) pollOnce(ctx context.Context, token string, timeout int) error { ok, err := w.store.TrySafewPollLock(ctx, token, pollLockTTL) if err != nil { return err } if !ok { slog.Info("safew poll skipped", "reason", "lock held", "timeout", timeout) return nil } defer func() { _ = w.store.UnlockSafewPoll(ctx, token) }() offset, err := w.store.GetSafewOffset(ctx, token) if err != nil { return err } chats, next, err := w.poll(token, offset, timeout) if err != nil { return err } if len(chats) > 0 { slog.Info("safew poll", "timeout", timeout, "groups", len(chats), "offset", next) } if err := w.mergeAndLog(ctx, token, chats); err != nil { return err } if next != offset { return w.store.SetSafewOffset(ctx, token, next) } return nil } func (w *Watcher) mergeAndLog(ctx context.Context, token string, chats []adapter.SafewChat) error { if len(chats) == 0 { return nil } existing, err := w.store.ListSafewChats(ctx, token) if err != nil { return w.store.MergeSafewChats(ctx, token, chats) } known := map[string]struct{}{} for _, c := range existing { known[c.ID] = struct{}{} } if err := w.store.MergeSafewChats(ctx, token, chats); err != nil { return err } for _, c := range newSafewChats(known, chats) { slog.Info("safew chat discovered", "id", c.ID, "type", c.Type, "title", c.Title) } return nil } func newSafewChats(known map[string]struct{}, chats []adapter.SafewChat) []adapter.SafewChat { var out []adapter.SafewChat seen := map[string]struct{}{} for _, c := range chats { if _, ok := known[c.ID]; ok { continue } if _, ok := seen[c.ID]; ok { continue } seen[c.ID] = struct{}{} out = append(out, c) } return out } func isSafewAuth(err error) bool { var e *adapter.SafewAuthError return errors.As(err, &e) } func isGetUpdatesConflict(err error) bool { if err == nil { return false } s := strings.ToLower(err.Error()) return strings.Contains(s, "conflict") && strings.Contains(s, "getupdates") } func conflictBackoff(idle time.Duration) time.Duration { if idle <= 0 { return 3 * time.Second } d := idle * 3 if d < 3*time.Second { return 3 * time.Second } return d }