2000908bac
- Added functionality to list and start Safew tokens in the main server function, ensuring proper initialization of the Safew watcher. - Introduced `ensureSafewWatcher` method in the ChannelHandler to manage Safew tokens during channel creation and updates. - Implemented `StartTokens` method in the Watcher to handle multiple tokens efficiently. - Enhanced error handling and logging for Safew watcher operations to improve observability.
269 lines
5.3 KiB
Go
269 lines
5.3 KiB
Go
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
|
|
}
|
|
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
|
|
}
|