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) 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 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.Debug("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 { 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 err := w.store.MergeSafewChats(ctx, token, chats); err != nil { return err } if next != offset { return w.store.SetSafewOffset(ctx, token, next) } return nil } 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 }