package centrifugo import ( "bytes" "context" "encoding/json" "errors" "fmt" "io" "log/slog" "net/http" "strconv" "time" "git.golem15.com/golem15/summercms/modules/wire" ) // requestTimeout bounds every Centrifugo API call, as the WinterCMS client // does. const requestTimeout = 5 * time.Second // maxResponseBytes caps how much of a response body is read. const maxResponseBytes = 1 << 20 // ErrNotConfigured is returned when the secret or key an operation needs // is empty: publishing without an API key, or signing without a token // secret. var ErrNotConfigured = errors.New("centrifugo: not configured") // Client calls the Centrifugo HTTP API. It is safe for concurrent use. type Client struct { apiURL string apiKey string hc *http.Client log *slog.Logger now func() time.Time } // DebugInfo is the connection summary Client.DebugInfo returns. It never // carries the API key. type DebugInfo struct { APIURL string `json:"api_url"` Enabled bool `json:"enabled"` APIKeySet bool `json:"api_key_set"` } // NewClient returns a client for cfg.APIURL and cfg.APIKey. hc may be nil // for a client with a 5s timeout; every request is also bounded by 5s. func NewClient(cfg Config, hc *http.Client) *Client { if hc == nil { hc = &http.Client{Timeout: requestTimeout} } return &Client{ apiURL: cfg.APIURL, apiKey: cfg.APIKey, hc: hc, log: slog.Default(), now: time.Now, } } // Enabled reports whether an API key is configured. func (c *Client) Enabled() bool { return c != nil && c.apiKey != "" } // DebugInfo returns the API URL and whether the client is enabled. func (c *Client) DebugInfo() DebugInfo { if c == nil { return DebugInfo{} } return DebugInfo{APIURL: c.apiURL, Enabled: c.Enabled(), APIKeySet: c.apiKey != ""} } type eventData struct { Event string `json:"event"` Payload json.RawMessage `json:"payload"` Timestamp wire.Time `json:"timestamp"` } type publishRequest struct { Channel string `json:"channel"` Data eventData `json:"data"` } type broadcastRequest struct { Channels []string `json:"channels"` Data eventData `json:"data"` } type presenceRequest struct { Channel string `json:"channel"` } type unsubscribeRequest struct { User string `json:"user"` Channel string `json:"channel"` } // Publish POSTs {"channel":…,"data":{"event":…,"payload":…,"timestamp":…}} // to {api_url}/publish. An empty payload is sent as []. Any 2xx status is // success, including Centrifugo's 200 responses that carry an error body. func (c *Client) Publish(ctx context.Context, channel, event string, payload json.RawMessage) error { if !c.Enabled() { return ErrNotConfigured } body := publishRequest{Channel: channel, Data: c.data(event, payload)} _, err := c.post(ctx, "/publish", body) return err } // Broadcast POSTs the same data with "channels" to {api_url}/broadcast. No // request is sent for an empty channel list. func (c *Client) Broadcast(ctx context.Context, channels []string, event string, payload json.RawMessage) error { if !c.Enabled() { return ErrNotConfigured } if len(channels) == 0 { return nil } body := broadcastRequest{Channels: channels, Data: c.data(event, payload)} _, err := c.post(ctx, "/broadcast", body) return err } // Presence returns result.presence of {api_url}/presence for channel. It // returns an empty map, never nil, when the client is disabled or the call // fails; the error says why. func (c *Client) Presence(ctx context.Context, channel string) (map[string]any, error) { out := map[string]any{} if !c.Enabled() { return out, ErrNotConfigured } raw, err := c.post(ctx, "/presence", presenceRequest{Channel: channel}) if err != nil { return out, err } var resp struct { Result struct { Presence map[string]any `json:"presence"` } `json:"result"` } if err := json.Unmarshal(raw, &resp); err != nil { c.log.Warn("centrifugo: presence response is not JSON", slog.String("channel", channel)) return out, fmt.Errorf("centrifugo: presence: %w", err) } if resp.Result.Presence != nil { out = resp.Result.Presence } return out, nil } // Unsubscribe POSTs {"user":"","channel":…} to {api_url}/unsubscribe. func (c *Client) Unsubscribe(ctx context.Context, userID uint, channel string) error { if !c.Enabled() { return ErrNotConfigured } body := unsubscribeRequest{User: strconv.FormatUint(uint64(userID), 10), Channel: channel} _, err := c.post(ctx, "/unsubscribe", body) return err } // Info POSTs {} to {api_url}/info and returns its result: Centrifugo's node // list and statistics. It is the connectivity probe of websockets:health. // Unlike the publishing calls, an answer with an error body is an error. func (c *Client) Info(ctx context.Context) (map[string]any, error) { out := map[string]any{} if !c.Enabled() { return out, ErrNotConfigured } raw, err := c.post(ctx, "/info", struct{}{}) if err != nil { return out, err } var resp struct { Error *struct { Code int `json:"code"` Message string `json:"message"` } `json:"error"` Result map[string]any `json:"result"` } if err := json.Unmarshal(raw, &resp); err != nil { return out, fmt.Errorf("centrifugo: /info: response is not JSON") } if resp.Error != nil { return out, fmt.Errorf("centrifugo: /info: error %d: %s", resp.Error.Code, resp.Error.Message) } if resp.Result != nil { out = resp.Result } return out, nil } func (c *Client) data(event string, payload json.RawMessage) eventData { if len(bytes.TrimSpace(payload)) == 0 { payload = json.RawMessage("[]") } return eventData{Event: event, Payload: payload, Timestamp: wire.Time{Time: c.now()}} } // post sends body and returns the response body of a 2xx answer. The API // key appears only in the Authorization header, never in logs or errors. func (c *Client) post(ctx context.Context, method string, body any) ([]byte, error) { var buf bytes.Buffer enc := json.NewEncoder(&buf) enc.SetEscapeHTML(false) if err := enc.Encode(body); err != nil { return nil, fmt.Errorf("centrifugo: %s: encode: %w", method, err) } if ctx == nil { ctx = context.Background() } ctx, cancel := context.WithTimeout(ctx, requestTimeout) defer cancel() req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.apiURL+method, bytes.NewReader(bytes.TrimSuffix(buf.Bytes(), []byte("\n")))) if err != nil { return nil, fmt.Errorf("centrifugo: %s: %w", method, err) } req.Header.Set("Authorization", "apikey "+c.apiKey) req.Header.Set("Content-Type", "application/json") resp, err := c.hc.Do(req) if err != nil { return nil, fmt.Errorf("centrifugo: %s: %w", method, err) } defer resp.Body.Close() raw, _ := io.ReadAll(io.LimitReader(resp.Body, maxResponseBytes)) if resp.StatusCode < 200 || resp.StatusCode > 299 { return nil, fmt.Errorf("centrifugo: %s: HTTP %d", method, resp.StatusCode) } if bytes.Contains(raw, []byte(`"error"`)) { // Centrifugo reports API errors as 200 with an error body. The // WinterCMS client counts those as success; so does this one. c.log.Debug("centrifugo: API answered with an error body", slog.String("method", method), slog.Int("status", resp.StatusCode)) } return raw, nil }