package httpd import ( "context" "fmt" "os" "slices" "strings" "sync" "time" "github.com/barerepo/server/internal/repo" "github.com/barerepo/server/internal/repocfg" "github.com/barerepo/server/internal/store" "github.com/barerepo/server/internal/webhook" ) // hookBatch bounds one pass, so a quiet server catches up and a busy one does not stall on it. const hookBatch = 100 // HookTick is how often the sender looks, which is the delay a receiver sees on a quiet server. const HookTick = 15 * time.Second // hookPayload is the body a receiver gets, with the same names the event log uses. type hookPayload struct { Event string `json:"event"` Repo string `json:"repo"` URL string `json:"url"` Actor string `json:"actor"` Ref string `json:"ref,omitempty"` Number int `json:"number,omitempty"` Title string `json:"title,omitempty"` At string `json:"at"` } // DeliverHooks sends every event since the cursor to the hooks that named it. Chapter 23. func (s *Server) DeliverHooks(ctx context.Context) { from, err := s.DB.WebhookCursor(ctx) if err != nil { s.warn("webhooks: cursor", "err", err) return } events, err := s.DB.EventsAfter(ctx, from, hookBatch) if err != nil { s.warn("webhooks: events", "err", err) return } for _, e := range events { s.deliverEvent(ctx, e) if err := s.DB.SetWebhookCursor(ctx, e.ID); err != nil { s.warn("webhooks: cursor", "err", err) return } } } // warn logs when there is a logger, because the sender is a goroutine that must not panic. func (s *Server) warn(msg string, args ...any) { if s.Log == nil { return } s.Log.Error(msg, args...) } // splitRepo takes owner/name apart, which is how the event log writes a repository. func splitRepo(full string) (owner, name string, ok bool) { owner, name, ok = strings.Cut(full, "/") return owner, name, ok && owner != "" && name != "" } // deliverEvent posts one event to every hook of its repository that asked for the kind. func (s *Server) deliverEvent(ctx context.Context, e store.Event) { owner, name, ok := splitRepo(e.Repo) if !ok { return } dir, err := repo.Dir(s.Cfg.Paths.Repos, owner, name) if err != nil { return } // The last version that parsed still names the hooks, so a typo does not silently stop them. cfg, _ := repocfg.Load(ctx, dir) hooks := hooksFor(cfg, e.Kind) if len(hooks) == 0 { return } state, err := s.DB.HooksOf(ctx, e.Repo) if err != nil { s.warn("webhooks: state", "repo", e.Repo, "err", err) return } payload := hookPayload{ Event: e.Kind, Repo: e.Repo, URL: s.Cfg.Server.ExternalURL + "/" + e.Repo, Actor: e.Actor, Ref: e.Ref, Number: e.Number, Title: e.Title, At: e.Created.UTC().Format(time.RFC3339), } // One event's hooks are independent, and a slow receiver must not hold up the others. var wg sync.WaitGroup for _, h := range hooks { if state[h.URL].Disabled { continue } wg.Add(1) go func() { defer wg.Done() s.deliverOne(ctx, e.Repo, h, payload, state[h.URL].Failures) }() } wg.Wait() } // deliverOne posts to one hook and records what happened, since the config page is the only report. func (s *Server) deliverOne(ctx context.Context, name string, h repocfg.Webhook, payload hookPayload, failures int) { // A hook that is already failing gets one try, or twenty events take twenty backoffs each. attempts := webhook.Attempts if failures > 0 { attempts = 1 } secret, err := hookSecret(h) if err == nil { err = webhook.Deliver(ctx, h.URL, secret, payload, attempts) } if err != nil { s.warn("webhook failed", "repo", name, "url", h.URL, "err", err) if err := s.DB.HookFailed(ctx, name, h.URL, err.Error()); err != nil { s.warn("webhooks: state", "err", err) } return } if err := s.DB.HookDelivered(ctx, name, h.URL); err != nil { s.warn("webhooks: state", "err", err) } } // hookSecret reads the named value from the server's environment, because the file holds a name. 23.2. func hookSecret(h repocfg.Webhook) (string, error) { if h.SecretEnv == "" { return "", nil } v := os.Getenv(h.SecretEnv) if v == "" { return "", fmt.Errorf("%s is not set on this server, so nothing would sign the body", h.SecretEnv) } return v, nil } // hooksFor is the events filter, and a hook naming no event asked for nothing. func hooksFor(cfg repocfg.Config, kind string) []repocfg.Webhook { var out []repocfg.Webhook for _, h := range cfg.Webhook { if h.URL == "" || !slices.Contains(h.Events, kind) { continue } out = append(out, h) } return out } // DeliverEvery runs the sender until the context ends. func (s *Server) DeliverEvery(ctx context.Context, every time.Duration) { if err := s.DB.StartWebhooksHere(ctx); err != nil { s.warn("webhooks: start", "err", err) return } t := time.NewTicker(every) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: s.DeliverHooks(ctx) } } }