Files
MOPAC/internal/events/server.go
T
mrcharles 86c39c8fc5 Rename module path to ukrrs.com/mopac/harness
The old path carried the reachableceo org prefix while the repo lives at
git.knownelement.com/ukrrs/MOPAC; the drift was cited as doc rot. go.mod
and every internal import updated; build/vet/test clean in the builder.

💘 Generated with Crush

Assisted-by: Crush:glm-5.2
2026-08-28 22:01:36 -05:00

206 lines
6.1 KiB
Go

package events
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log"
"net/http"
"time"
"ukrrs.com/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) }