package deliberi import ( "context" "log/slog" "maps" "sync" "time" "github.com/samber/lo" deldb "tangled.org/core/deliberi/db" "tangled.org/core/deliberi/models" "tangled.org/core/orm" ) const ( inviteReconcileInterval = 10 * time.Second inviteFetchTimeout = 2 * time.Second ) var inviteNotificationTypes = []models.NotificationType{ models.NotificationTypeCollaboratorInvited, } type InviteSync struct { db *deldb.DB invites inviteLister logger *slog.Logger mu sync.Mutex lastRun map[string]time.Time swept time.Time } func NewInviteSync(database *deldb.DB, invites inviteLister, logger *slog.Logger) *InviteSync { return &InviteSync{ db: database, invites: invites, logger: logger, lastRun: map[string]time.Time{}, } } func (s *InviteSync) Reconcile(ctx context.Context, recipientDid string) { if !s.claim(recipientDid) { return } answer, err := s.fetch(ctx, recipientDid) if err != nil { s.logger.Warn("listing invites, keeping stored rows", "err", err, "recipient", recipientDid) return } stored, err := deldb.GetNotifications(s.db, recipientDid, 0, orm.FilterIn("type", inviteNotificationTypes)) if err != nil { s.logger.Warn("reading stored invite rows", "err", err, "recipient", recipientDid) return } offered := keySet(answer.Offers, func(o InviteOffer) string { return o.Uri }) unread := keySet(answer.Pending, func(knot string) string { return knot }) held := keySet(stored, func(n *models.Notification) string { return n.AtUri }) withdrawn := lo.Filter(stored, func(n *models.Notification, _ int) bool { _, standing := offered[n.AtUri] _, cold := unread[n.KnotDid] return !standing && !cold && !answer.Truncated }) fresh := lo.Map( lo.Filter(answer.Offers, func(o InviteOffer, _ int) bool { _, cached := held[o.Uri] return !cached && o.AddedBy != recipientDid }), func(o InviteOffer, _ int) *models.Notification { return inviteNotification(recipientDid, o) }, ) if err := s.settle(recipientDid, withdrawn, fresh); err != nil { s.logger.Warn("writing reconciled invite rows", "err", err, "recipient", recipientDid) } } func (s *InviteSync) settle(recipientDid string, withdrawn, fresh []*models.Notification) error { if len(withdrawn) == 0 && len(fresh) == 0 { return nil } tx, err := s.db.Begin() if err != nil { return err } defer func() { _ = tx.Rollback() }() if err := firstErr(withdrawn, func(n *models.Notification) error { return deldb.DeleteNotification(tx, recipientDid, n.AtUri) }); err != nil { return err } if err := firstErr(fresh, func(n *models.Notification) error { return deldb.CreateNotification(tx, n) }); err != nil { return err } return tx.Commit() } func (s *InviteSync) claim(recipientDid string) bool { now := time.Now() s.mu.Lock() defer s.mu.Unlock() if now.Sub(s.swept) >= inviteReconcileInterval { maps.DeleteFunc(s.lastRun, func(_ string, at time.Time) bool { return now.Sub(at) >= inviteReconcileInterval }) s.swept = now } if at, seen := s.lastRun[recipientDid]; seen && now.Sub(at) < inviteReconcileInterval { return false } s.lastRun[recipientDid] = now return true } func (s *InviteSync) fetch(ctx context.Context, recipientDid string) (InviteAnswer, error) { ctx, cancel := context.WithTimeout(ctx, inviteFetchTimeout) defer cancel() return s.invites.ListCollaboratorInvitesBy(ctx, recipientDid) } func keySet[T any, K comparable](items []T, key func(T) K) map[K]struct{} { return lo.SliceToMap(items, func(item T) (K, struct{}) { return key(item), struct{}{} }) } func firstErr[T any](items []T, act func(T) error) error { return lo.Reduce(items, func(err error, item T, _ int) error { if err != nil { return err } return act(item) }, nil) } func inviteNotification(recipientDid string, o InviteOffer) *models.Notification { return &models.Notification{ RecipientDid: recipientDid, AtUri: o.Uri, Type: models.NotificationTypeCollaboratorInvited, ActorDid: o.AddedBy, KnotDid: o.KnotDid, RepoDid: o.RepoDid, Created: o.CreatedAt, } }