package events import ( "context" "encoding/json" "errors" "fmt" "io" "log" "net/http" "time" "git.knownelement.com/reachableceo/MOPAC/harness/internal/config" ) // maxBodyBytes bounds webhook payloads; bigger deliveries are rejected // before parsing (413). const maxBodyBytes = 1 << 20 // Dispatcher is the conductor-side event sink. Turn dispatch is a stub // call for now: the conductor implements it, real event->turn wiring lands // after the skeleton (phase 3). type Dispatcher interface { DispatchEvent(ctx context.Context, ev Event) } // Server is the `harness events` webhook receiver. Routes: // // POST /hooks/redmine shared-secret header // POST /hooks/discourse shared-secret header // POST /hooks/gitea HMAC-SHA256 X-Gitea-Signature // GET /healthz // // Unsigned or unverified requests get 401 without any detail beyond // "unverified"; secrets and header values are never logged. type Server struct { cfg config.EventsConfig store *Store disp Dispatcher logger *log.Logger secrets map[string]string // resolved at startup; never logged, never persisted } // NewServer resolves the webhook secrets (fail-fast: at least one source // must be configured) and opens the event store. func NewServer(cfg config.EventsConfig, disp Dispatcher, out io.Writer) (*Server, error) { secrets := map[string]string{ SourceRedmine: "", SourceDiscourse: "", SourceGitea: "", } for source, sc := range map[string]config.EventSourceConfig{ SourceRedmine: cfg.Redmine, SourceDiscourse: cfg.Discourse, SourceGitea: cfg.Gitea, } { if sc.SecretRef == "" { continue } v, err := config.ResolveKeyRef(sc.SecretRef) if err != nil { return nil, fmt.Errorf("[events.%s]: resolve secret: %w", source, err) } secrets[source] = v } if secrets[SourceRedmine] == "" && secrets[SourceDiscourse] == "" && secrets[SourceGitea] == "" { return nil, fmt.Errorf("[events]: no webhook secrets configured; set secret_ref for at least one of redmine/discourse/gitea") } store, err := OpenStore(cfg.StateDir) if err != nil { return nil, err } return &Server{ cfg: cfg, store: store, disp: disp, logger: log.New(out, "events: ", log.LstdFlags|log.Lmsgprefix), secrets: secrets, }, nil } // Handler builds the receiver's HTTP routes. func (s *Server) Handler() http.Handler { mux := http.NewServeMux() mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") fmt.Fprint(w, `{"status":"ok"}`) }) mux.HandleFunc("/hooks/redmine", s.hook(SourceRedmine)) mux.HandleFunc("/hooks/discourse", s.hook(SourceDiscourse)) mux.HandleFunc("/hooks/gitea", s.hook(SourceGitea)) return mux } // Close releases the event store. func (s *Server) Close() error { return s.store.Close() } func (s *Server) hook(source string) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { w.Header().Set("Allow", http.MethodPost) s.writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "POST only"}) return } body, ok := s.readBody(w, r) if !ok { return } if err := s.verify(source, r.Header, body); err != nil { // Log the class of failure only; never header values. s.logger.Printf("rejected source=%s remote=%s reason=%v", source, r.RemoteAddr, err) s.writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "unverified webhook"}) return } ev, err := Normalize(source, r.Header, body, time.Now()) if err != nil { s.logger.Printf("malformed source=%s remote=%s", source, r.RemoteAddr) s.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "malformed webhook payload"}) return } stored, err := s.store.Append(ev) if err != nil { s.logger.Printf("store error source=%s: %v", source, err) s.writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "event log write failed"}) return } status := "stored" if !stored { status = "duplicate" } // One audit line per delivery: normalized fields + digest, no // headers, no secrets, no payload contents. s.logger.Printf( "source=%s kind=%s actor=%s subject_id=%s repo=%s action=%s provider_id=%s digest=%s %s", ev.Source, ev.Kind, ev.Actor, ev.SubjectID, ev.Repo, ev.Action, ev.ProviderID, shortDigest(ev.PayloadDigest), status, ) if stored && ev.Action != ActionIgnore && s.disp != nil { s.disp.DispatchEvent(r.Context(), ev) } s.writeJSON(w, http.StatusOK, map[string]string{ "status": status, "id": ev.DedupKey(), "action": ev.Action, }) } } func (s *Server) verify(source string, header http.Header, body []byte) error { secret := s.secrets[source] if secret == "" { return fmt.Errorf("%w: no secret configured for source %s", ErrUnverified, source) } switch source { case SourceGitea: return VerifyGiteaHMAC(body, header.Get(giteaSignatureHeader), secret) case SourceRedmine: name := s.cfg.Redmine.SecretHeader if name == "" { name = DefaultRedmineSecretHeader } return VerifySharedSecret(name, header.Get(name), secret) case SourceDiscourse: name := s.cfg.Discourse.SecretHeader if name == "" { name = DefaultDiscourseSecretHeader } return VerifySharedSecret(name, header.Get(name), secret) } return fmt.Errorf("%w: unknown source %s", ErrUnverified, source) } func (s *Server) readBody(w http.ResponseWriter, r *http.Request) ([]byte, bool) { body, err := io.ReadAll(io.LimitReader(r.Body, maxBodyBytes+1)) if err != nil { s.writeJSON(w, http.StatusBadRequest, map[string]string{"error": "unreadable body"}) return nil, false } if len(body) > maxBodyBytes { s.writeJSON(w, http.StatusRequestEntityTooLarge, map[string]string{"error": "payload too large"}) return nil, false } return body, true } func (s *Server) writeJSON(w http.ResponseWriter, code int, v any) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(code) _ = json.NewEncoder(w).Encode(v) } func shortDigest(d string) string { if len(d) > 16 { return d[:16] } return d } // IsUnverified reports whether err is a verification failure (CLI/tests). func IsUnverified(err error) bool { return errors.Is(err, ErrUnverified) }