modified internal/atproto/sync.go +0 -6
| @@ -36,12 +36,6 @@ func (s *Sync) Run(ctx context.Context, userDID string) error { |
| 36 | 36 | if err := s.syncSubscriptions(ctx, userDID); err != nil { |
| 37 | 37 | s.logger.Error("sync subscriptions failed", "error", err, "did", userDID) |
| 38 | 38 | } |
| 39 | | - // Recompute subscriber_count from subscriptions table rather than |
| 40 | | - // maintaining it incrementally (done here to avoid drift from race conditions |
| 41 | | - // between stream handler events and sync operations). |
| 42 | | - if err := s.articles.RecountSubscriberCounts(ctx); err != nil { |
| 43 | | - s.logger.Error("recount subscriber counts failed", "error", err, "did", userDID) |
| 44 | | - } |
| 45 | 39 | if err := s.syncLikes(ctx, userDID); err != nil { |
| 46 | 40 | s.logger.Error("sync likes failed", "error", err, "did", userDID) |
| 47 | 41 | } |
| @@ -36,12 +36,6 @@ func (s *Sync) Run(ctx context.Context, userDID string) error { |
| 36 | if err := s.syncSubscriptions(ctx, userDID); err != nil { | 36 | if err := s.syncSubscriptions(ctx, userDID); err != nil { |
| 37 | s.logger.Error("sync subscriptions failed", "error", err, "did", userDID) | 37 | s.logger.Error("sync subscriptions failed", "error", err, "did", userDID) |
| 38 | } | 38 | } |
| 39 | - // Recompute subscriber_count from subscriptions table rather than | | |
| 40 | - // maintaining it incrementally (done here to avoid drift from race conditions | | |
| 41 | - // between stream handler events and sync operations). | | |
| 42 | - if err := s.articles.RecountSubscriberCounts(ctx); err != nil { | | |
| 43 | - s.logger.Error("recount subscriber counts failed", "error", err, "did", userDID) | | |
| 44 | - } | | |
| 45 | if err := s.syncLikes(ctx, userDID); err != nil { | 39 | if err := s.syncLikes(ctx, userDID); err != nil { |
| 46 | s.logger.Error("sync likes failed", "error", err, "did", userDID) | 40 | s.logger.Error("sync likes failed", "error", err, "did", userDID) |
| 47 | } | 41 | } |
modified internal/server/server.go +5 -0
| @@ -520,6 +520,11 @@ func (s *Server) runSyncAll(ctx context.Context) { |
| 520 | 520 | |
| 521 | 521 | metrics.SyncRuns.Inc() |
| 522 | 522 | } |
| 523 | + |
| 524 | + // Recompute subscriber_count once after all users are synced. |
| 525 | + if err := s.dbs.Articles.RecountSubscriberCounts(ctx); err != nil { |
| 526 | + s.logger.Error("recount subscriber counts failed", "error", err) |
| 527 | + } |
| 523 | 528 | } |
| 524 | 529 | |
| 525 | 530 | func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { |
| @@ -520,6 +520,11 @@ func (s *Server) runSyncAll(ctx context.Context) { |
| 520 | | 520 | |
| 521 | metrics.SyncRuns.Inc() | 521 | metrics.SyncRuns.Inc() |
| 522 | } | 522 | } |
| | 523 | + |
| | 524 | + // Recompute subscriber_count once after all users are synced. |
| | 525 | + if err := s.dbs.Articles.RecountSubscriberCounts(ctx); err != nil { |
| | 526 | + s.logger.Error("recount subscriber counts failed", "error", err) |
| | 527 | + } |
| 523 | } | 528 | } |
| 524 | | 529 | |
| 525 | func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { | 530 | func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) { |