events: HTTP receiver, CLI wiring, conductor dispatch stub

`harness events` serves POST /hooks/{redmine,discourse,gitea} and
GET /healthz on stdlib net/http until SIGINT/SIGTERM (graceful
shutdown). Flow per delivery: verify (401, generic body) -> normalize
(400) -> append-only store with dedup (200 stored/duplicate + id +
action) -> hand stored actionable events to the conductor via the
Dispatcher interface. Conductor.DispatchEvent is the wiring point and
prints what it will do once phase 3 lands the real event-to-turn
dispatch. Body cap 1 MiB (413); audit log carries normalized fields +
digest only, never headers, secrets or payload. Server refuses to start
without at least one resolvable webhook secret.
This commit is contained in:
2026-08-28 21:38:24 -05:00
parent 05ec1a4142
commit 043e03b830
4 changed files with 622 additions and 1 deletions
+205
View File
@@ -0,0 +1,205 @@
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) }
+328
View File
@@ -0,0 +1,328 @@
package events
import (
"bytes"
"context"
"encoding/json"
"io"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"sync"
"testing"
"git.knownelement.com/reachableceo/MOPAC/harness/internal/config"
)
const (
testRedmineSecret = "redmine-hook-secret"
testDiscourseSecret = "discourse-hook-secret"
testGiteaSecret = "gitea-hook-secret"
)
// recordingDispatcher captures dispatched events for assertions.
type recordingDispatcher struct {
mu sync.Mutex
events []Event
}
func (d *recordingDispatcher) DispatchEvent(ctx context.Context, ev Event) {
d.mu.Lock()
defer d.mu.Unlock()
d.events = append(d.events, ev)
}
func (d *recordingDispatcher) dispatched() []Event {
d.mu.Lock()
defer d.mu.Unlock()
return append([]Event(nil), d.events...)
}
func newTestServer(t *testing.T) (*Server, *recordingDispatcher, string) {
t.Helper()
stateDir := t.TempDir()
cfg := config.EventsConfig{
Listen: ":0",
StateDir: stateDir,
Redmine: config.EventSourceConfig{
SecretRef: "literal:" + testRedmineSecret,
SecretHeader: DefaultRedmineSecretHeader,
},
Discourse: config.EventSourceConfig{
SecretRef: "literal:" + testDiscourseSecret,
SecretHeader: DefaultDiscourseSecretHeader,
},
Gitea: config.EventSourceConfig{
SecretRef: "literal:" + testGiteaSecret,
},
}
disp := &recordingDispatcher{}
srv, err := NewServer(cfg, disp, io.Discard)
if err != nil {
t.Fatalf("NewServer: %v", err)
}
t.Cleanup(func() { srv.Close() })
return srv, disp, stateDir
}
func post(t *testing.T, url string, headers map[string]string, body string) (int, string) {
t.Helper()
req, err := http.NewRequest(http.MethodPost, url, strings.NewReader(body))
if err != nil {
t.Fatal(err)
}
for k, v := range headers {
req.Header.Set(k, v)
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
b, _ := io.ReadAll(resp.Body)
return resp.StatusCode, string(b)
}
func readLog(t *testing.T, stateDir string) []string {
t.Helper()
data, err := os.ReadFile(filepath.Join(stateDir, "events.jsonl"))
if os.IsNotExist(err) {
return nil
}
if err != nil {
t.Fatal(err)
}
var lines []string
for _, ln := range strings.Split(strings.TrimSpace(string(data)), "\n") {
if ln != "" {
lines = append(lines, ln)
}
}
return lines
}
func TestServerVerificationEndToEnd(t *testing.T) {
srv, _, stateDir := newTestServer(t)
ts := httptest.NewServer(srv.Handler())
defer ts.Close()
giteaBody := `{"action":"approved","number":5,"pull_request":{"number":5,"title":"PR"},"repository":{"full_name":"ukrrs/MOPAC"},"sender":{"login":"charles"}}`
giteaHeaders := func(sig string) map[string]string {
return map[string]string{
giteaDeliveryHeader: "delivery-1",
giteaSignatureHeader: sig,
"X-Gitea-Event": "pull_request",
"X-Gitea-Event-Type": "pull_request_approved",
}
}
cases := []struct {
name string
path string
headers map[string]string
body string
wantCode int
wantBody string
}{
{"gitea unsigned rejected", "/hooks/gitea", nil, giteaBody, 401, "unverified webhook"},
{"gitea bad signature rejected", "/hooks/gitea", giteaHeaders(strings.Repeat("ff", 32)), giteaBody, 401, "unverified webhook"},
{"gitea valid accepted", "/hooks/gitea", giteaHeaders(SignGitea([]byte(giteaBody), testGiteaSecret)), giteaBody, 200, `"status":"stored"`},
{"redmine unsigned rejected", "/hooks/redmine", nil, `{}`, 401, "unverified webhook"},
{
"redmine valid accepted",
"/hooks/redmine",
map[string]string{DefaultRedmineSecretHeader: testRedmineSecret},
`{"event_name":"issue_updated","payload":{"issue":{"id":42,"subject":"Ship"},"user":{"login":"charles"}}}`,
200, `"status":"stored"`,
},
{"discourse wrong secret rejected", "/hooks/discourse", map[string]string{DefaultDiscourseSecretHeader: "nope"}, `{}`, 401, "unverified webhook"},
{
"discourse valid accepted",
"/hooks/discourse",
map[string]string{DefaultDiscourseSecretHeader: testDiscourseSecret, "X-Discourse-Event": "post_created", "X-Discourse-Event-Id": "42"},
`{"post":{"id":9,"topic_id":7,"username":"charles","topic_title":"Plan"}}`,
200, `"status":"stored"`,
},
{"gitea malformed body rejected", "/hooks/gitea", giteaHeaders(SignGitea([]byte("not json"), testGiteaSecret)), `not json`, 400, "malformed webhook payload"},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
code, body := post(t, ts.URL+tc.path, tc.headers, tc.body)
if code != tc.wantCode {
t.Errorf("code = %d, want %d (body %s)", code, tc.wantCode, body)
}
if !strings.Contains(body, tc.wantBody) {
t.Errorf("body %q does not contain %q", body, tc.wantBody)
}
for _, secret := range []string{testRedmineSecret, testDiscourseSecret, testGiteaSecret} {
if strings.Contains(body, secret) {
t.Errorf("response body leaks secret: %s", body)
}
}
})
}
if lines := readLog(t, stateDir); len(lines) != 3 {
t.Errorf("event log has %d lines, want 3 (only verified events stored):\n%s",
len(lines), strings.Join(lines, "\n"))
}
}
func TestServerDedupAndDispatch(t *testing.T) {
srv, disp, stateDir := newTestServer(t)
ts := httptest.NewServer(srv.Handler())
defer ts.Close()
body := `{"action":"approved","number":5,"pull_request":{"number":5},"sender":{"login":"charles"}}`
headers := map[string]string{
giteaDeliveryHeader: "delivery-dedup",
giteaSignatureHeader: SignGitea([]byte(body), testGiteaSecret),
}
code, respBody := post(t, ts.URL+"/hooks/gitea", headers, body)
if code != 200 || !strings.Contains(respBody, `"status":"stored"`) {
t.Fatalf("first delivery: code=%d body=%s", code, respBody)
}
code, respBody = post(t, ts.URL+"/hooks/gitea", headers, body)
if code != 200 || !strings.Contains(respBody, `"status":"duplicate"`) {
t.Fatalf("replay: code=%d body=%s", code, respBody)
}
if lines := readLog(t, stateDir); len(lines) != 1 {
t.Errorf("event log has %d lines after replay, want 1", len(lines))
}
dispatched := disp.dispatched()
if len(dispatched) != 1 {
t.Fatalf("dispatcher ran %d times, want 1 (dupes do not dispatch)", len(dispatched))
}
if dispatched[0].Action != ActionPipelineStep {
t.Errorf("dispatched action = %s, want %s", dispatched[0].Action, ActionPipelineStep)
}
}
func TestServerIgnoreActionDoesNotDispatch(t *testing.T) {
srv, disp, stateDir := newTestServer(t)
ts := httptest.NewServer(srv.Handler())
defer ts.Close()
body := `{"action":"opened","number":3,"pull_request":{"number":3},"sender":{"login":"bob"}}`
headers := map[string]string{
giteaDeliveryHeader: "delivery-open",
giteaSignatureHeader: SignGitea([]byte(body), testGiteaSecret),
}
code, respBody := post(t, ts.URL+"/hooks/gitea", headers, body)
if code != 200 || !strings.Contains(respBody, `"action":"ignore"`) {
t.Fatalf("opened PR: code=%d body=%s", code, respBody)
}
if got := disp.dispatched(); len(got) != 0 {
t.Errorf("ignore action dispatched %d events, want 0", len(got))
}
if lines := readLog(t, stateDir); len(lines) != 1 {
t.Errorf("ignored events are still stored: %d lines, want 1", len(lines))
}
}
func TestServerAuxRoutes(t *testing.T) {
srv, _, _ := newTestServer(t)
ts := httptest.NewServer(srv.Handler())
defer ts.Close()
resp, err := http.Get(ts.URL + "/healthz")
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != 200 {
t.Errorf("healthz = %d, want 200", resp.StatusCode)
}
resp, err = http.Get(ts.URL + "/hooks/gitea")
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if resp.StatusCode != http.StatusMethodNotAllowed {
t.Errorf("GET /hooks/gitea = %d, want 405", resp.StatusCode)
}
code, body := post(t, ts.URL+"/hooks/unknown", nil, `{}`)
if code != 404 {
t.Errorf("unknown hook = %d, want 404 (body %s)", code, body)
}
}
func TestServerRejectsOversizedBody(t *testing.T) {
srv, _, _ := newTestServer(t)
ts := httptest.NewServer(srv.Handler())
defer ts.Close()
big := bytes.Repeat([]byte("a"), maxBodyBytes+1)
sig := SignGitea(big, testGiteaSecret)
req, _ := http.NewRequest(http.MethodPost, ts.URL+"/hooks/gitea", bytes.NewReader(big))
req.Header.Set(giteaSignatureHeader, sig)
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusRequestEntityTooLarge {
t.Errorf("oversized body = %d, want 413", resp.StatusCode)
}
}
func TestNewServerRequiresASecret(t *testing.T) {
cfg := config.EventsConfig{Listen: ":0", StateDir: t.TempDir()}
if _, err := NewServer(cfg, nil, io.Discard); err == nil ||
!strings.Contains(err.Error(), "no webhook secrets configured") {
t.Errorf("NewServer without secrets err = %v, want no-webhook-secrets error", err)
}
}
func TestServerLogLineLeaksNothing(t *testing.T) {
var buf bytes.Buffer
cfg := config.EventsConfig{
Listen: ":0",
StateDir: t.TempDir(),
Gitea: config.EventSourceConfig{SecretRef: "literal:" + testGiteaSecret},
}
srv, err := NewServer(cfg, nil, &buf)
if err != nil {
t.Fatal(err)
}
defer srv.Close()
ts := httptest.NewServer(srv.Handler())
defer ts.Close()
// A signed delivery and an unsigned probe both log; neither line may
// contain the secret or any header VALUE.
body := `{"action":"approved","number":1,"pull_request":{"number":1}}`
post(t, ts.URL+"/hooks/gitea", map[string]string{
giteaSignatureHeader: SignGitea([]byte(body), testGiteaSecret),
giteaDeliveryHeader: "leak-check",
}, body)
post(t, ts.URL+"/hooks/gitea", nil, body)
out := buf.String()
if strings.Contains(out, testGiteaSecret) {
t.Errorf("log leaks secret:\n%s", out)
}
if strings.Contains(out, "payload") {
t.Errorf("log echoes payload:\n%s", out)
}
for _, want := range []string{"kind=pr_approved", "provider_id=leak-check", "rejected source=gitea"} {
if !strings.Contains(out, want) {
t.Errorf("log line missing %q:\n%s", want, out)
}
}
// The stored JSONL record must not carry the secret either.
data, _ := os.ReadFile(srv.store.Path())
if strings.Contains(string(data), testGiteaSecret) {
t.Errorf("event log leaks secret: %s", data)
}
var ev Event
if err := json.Unmarshal(data, &ev); err != nil {
t.Errorf("log line is not valid Event JSON: %v (%s)", err, data)
}
}