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
206 lines
6.1 KiB
Go
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) }
|