Combine message dispatching into a dispatch function

This commit is contained in:
binwiederhier
2026-08-03 08:11:51 +02:00
parent dc11655153
commit 0ecba37334
3 changed files with 68 additions and 44 deletions
+56 -43
View File
@@ -833,6 +833,41 @@ func (s *Server) handleMatrixDiscovery(w http.ResponseWriter) error {
return writeMatrixDiscoveryResponse(w) return writeMatrixDiscoveryResponse(w)
} }
// dispatch delivers m to local subscribers and fires the requested side-effect targets. It is
// the single choke point through which every published message must pass; t may be nil when
// the topic has no local subscribers (delayed sender).
func (s *Server) dispatch(v *visitor, t *topic, m *model.Message, opts dispatchOpts) error {
// Deliver to local subscribers
if t != nil {
if opts.async {
go func() {
if err := t.Publish(v, m); err != nil {
logvm(v, m).Err(err).Warn("Unable to publish message")
}
}()
} else if err := t.Publish(v, m); err != nil {
return err
}
}
// Fire the requested side-effect targets
if s.firebaseClient != nil && opts.firebase {
go s.sendToFirebase(v, m)
}
if s.mailer != nil && opts.email != "" {
go s.sendEmail(v, m, opts.email)
}
if s.config.TwilioAccount != "" && opts.call != "" {
go s.callPhone(v, m, opts.call)
}
if s.config.UpstreamBaseURL != "" && opts.upstream {
go s.forwardPollRequest(v, m)
}
if s.config.WebPushPublicKey != "" && opts.webPush {
go s.publishToWebPushEndpoints(v, m)
}
return nil
}
func (s *Server) handlePublishInternal(r *http.Request, v *visitor) (*model.Message, error) { func (s *Server) handlePublishInternal(r *http.Request, v *visitor) (*model.Message, error) {
start := time.Now() start := time.Now()
t, err := fromContext[*topic](r, contextTopic) t, err := fromContext[*topic](r, contextTopic)
@@ -911,24 +946,16 @@ func (s *Server) handlePublishInternal(r *http.Request, v *visitor) (*model.Mess
ev.Debug("Received message") ev.Debug("Received message")
} }
if !delayed { if !delayed {
if err := t.Publish(v, m); err != nil { err := s.dispatch(v, t, m, dispatchOpts{
firebase: firebase,
email: email,
call: call,
upstream: !unifiedpush, // UP messages are not sent to upstream
webPush: true,
})
if err != nil {
return nil, err return nil, err
} }
if s.firebaseClient != nil && firebase {
go s.sendToFirebase(v, m)
}
if s.mailer != nil && email != "" {
go s.sendEmail(v, m, email)
}
if s.config.TwilioAccount != "" && call != "" {
go s.callPhone(v, m, call)
}
if s.config.UpstreamBaseURL != "" && !unifiedpush { // UP messages are not sent to upstream
go s.forwardPollRequest(v, m)
}
if s.config.WebPushPublicKey != "" {
go s.publishToWebPushEndpoints(v, m)
}
} else { } else {
logvrm(v, r, m).Tag(tagPublish).Debug("Message delayed, will process later") logvrm(v, r, m).Tag(tagPublish).Debug("Message delayed, will process later")
} }
@@ -1027,18 +1054,10 @@ func (s *Server) handleActionMessage(w http.ResponseWriter, r *http.Request, v *
m.Sender = v.IP() m.Sender = v.IP()
m.User = v.MaybeUserID() m.User = v.MaybeUserID()
m.Expires = time.Unix(m.Time, 0).Add(v.Limits().MessageExpiryDuration).Unix() m.Expires = time.Unix(m.Time, 0).Add(v.Limits().MessageExpiryDuration).Unix()
// Publish to subscribers // Publish to subscribers, Firebase (for Android clients), and web push endpoints
if err := t.Publish(v, m); err != nil { if err := s.dispatch(v, t, m, dispatchOpts{firebase: true, webPush: true}); err != nil {
return err return err
} }
// Send to Firebase for Android clients
if s.firebaseClient != nil {
go s.sendToFirebase(v, m)
}
// Send to web push endpoints
if s.config.WebPushPublicKey != "" {
go s.publishToWebPushEndpoints(v, m)
}
if event == model.MessageDeleteEvent { if event == model.MessageDeleteEvent {
// Delete any existing scheduled message with the same sequence ID // Delete any existing scheduled message with the same sequence ID
deletedIDs, err := s.messageCache.DeleteScheduledBySequenceID(t.ID, sequenceID) deletedIDs, err := s.messageCache.DeleteScheduledBySequenceID(t.ID, sequenceID)
@@ -1987,24 +2006,18 @@ func (s *Server) sendDelayedMessages() error {
func (s *Server) sendDelayedMessage(v *visitor, m *model.Message) error { func (s *Server) sendDelayedMessage(v *visitor, m *model.Message) error {
logvm(v, m).Debug("Sending delayed message") logvm(v, m).Debug("Sending delayed message")
s.mu.RLock() s.mu.RLock()
t, ok := s.topics[m.Topic] // If no subscribers, just mark message as published t := s.topics[m.Topic] // May be nil if there are no local subscribers; dispatch handles that
s.mu.RUnlock() s.mu.RUnlock()
if ok { // We do not rate-limit messages here, since we've rate limited them in the PUT/POST handler.
go func() { // Firebase subscribers may not show up in the topics map, so side effects fire regardless.
// We do not rate-limit messages here, since we've rate limited them in the PUT/POST handler err := s.dispatch(v, t, m, dispatchOpts{
if err := t.Publish(v, m); err != nil { firebase: true,
logvm(v, m).Err(err).Warn("Unable to publish message") upstream: true,
} webPush: true,
}() async: true,
} })
if s.firebaseClient != nil { // Firebase subscribers may not show up in topics map if err != nil {
go s.sendToFirebase(v, m) return err
}
if s.config.UpstreamBaseURL != "" {
go s.forwardPollRequest(v, m)
}
if s.config.WebPushPublicKey != "" {
go s.publishToWebPushEndpoints(v, m)
} }
if err := s.messageCache.MarkPublished(m); err != nil { if err := s.messageCache.MarkPublished(m); err != nil {
return err return err
+1 -1
View File
@@ -985,7 +985,7 @@ func (s *Server) publishSyncEventForUser(v *visitor, u *user.User) error {
return err return err
} }
m := model.NewDefaultMessage(syncTopic.ID, string(messageBytes)) m := model.NewDefaultMessage(syncTopic.ID, string(messageBytes))
if err := syncTopic.Publish(v, m); err != nil { if err := s.dispatch(v, syncTopic, m, dispatchOpts{}); err != nil {
return err return err
} }
return nil return nil
+11
View File
@@ -29,6 +29,17 @@ type publishMessage struct {
Delay string `json:"delay"` Delay string `json:"delay"`
} }
// dispatchOpts selects which delivery targets fire for a published message, beyond delivery
// to local subscribers (see Server.dispatch)
type dispatchOpts struct {
firebase bool // Send to Firebase (if configured)
email string // Send an email to this address (if a mailer is configured)
call string // Call this phone number (if Twilio is configured)
upstream bool // Forward a poll request to the upstream server (if configured)
webPush bool // Publish to web push endpoints (if configured)
async bool // Deliver to local subscribers in a goroutine, logging errors instead of returning them
}
// messageEncoder is a function that knows how to encode a message // messageEncoder is a function that knows how to encode a message
type messageEncoder func(msg *model.Message) (string, error) type messageEncoder func(msg *model.Message) (string, error)