feat(receiver): ingest server and admin dashboard
POST /webhook (bearer-token auth, strict validation, 1MB cap) and /healthz stay open; everything under /admin requires basic auth (RECEIVER_ADMIN_USER/PASSWORD, refuses to start without a password, constant-time compares). Admin JSON APIs for the traversal log (filters, pagination) and stats, plus an embedded dashboard: stat cards, four Chart.js charts (lazy CDN load with graceful degradation), filterable log with expandable summaries, auto-refresh. Fetches resolve against location.origin so credentialed bookmark URLs work. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5
parent
e8f56a1ef4
commit
beb595442f
@@ -0,0 +1,222 @@
|
||||
// Package server exposes the webhook telemetry receiver over HTTP:
|
||||
// - POST /webhook — ingest start/complete events from the main app
|
||||
// - GET /healthz — health check
|
||||
// - GET /admin — basic-auth-protected admin UI and JSON API
|
||||
package server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net"
|
||||
"net/http"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gitea.hansenits.com.au/hits/ExploreDNS/internal/receiver/store"
|
||||
)
|
||||
|
||||
// maxBodyBytes caps webhook request bodies, matching the sender-side API
|
||||
// request cap.
|
||||
const maxBodyBytes = 1 << 20
|
||||
|
||||
// Server is the receiver HTTP server.
|
||||
type Server struct {
|
||||
addr string
|
||||
version string
|
||||
token string
|
||||
adminUser string
|
||||
adminPass string
|
||||
st *store.Store
|
||||
srv *http.Server
|
||||
}
|
||||
|
||||
// New creates a Server that listens on addr and writes events to st.
|
||||
func New(addr string, st *store.Store) *Server {
|
||||
return &Server{addr: addr, version: "dev", st: st}
|
||||
}
|
||||
|
||||
// SetVersion records the build version reported by GET /healthz.
|
||||
// Call before Start; empty values are ignored.
|
||||
func (s *Server) SetVersion(v string) {
|
||||
if v != "" {
|
||||
s.version = v
|
||||
}
|
||||
}
|
||||
|
||||
// SetIngestToken enables bearer authentication on POST /webhook. Empty
|
||||
// leaves the endpoint open. Call before Start.
|
||||
func (s *Server) SetIngestToken(t string) { s.token = t }
|
||||
|
||||
// SetAdminAuth sets the basic-auth credentials for /admin. Call before
|
||||
// Start; Start refuses to run without a password so the admin interface
|
||||
// can never be exposed unprotected.
|
||||
func (s *Server) SetAdminAuth(user, pass string) {
|
||||
s.adminUser = user
|
||||
s.adminPass = pass
|
||||
}
|
||||
|
||||
// Start begins listening. Call Shutdown to stop gracefully.
|
||||
func (s *Server) Start() error {
|
||||
if s.adminPass == "" {
|
||||
return errors.New("admin password not set (RECEIVER_ADMIN_PASSWORD): refusing to expose /admin unprotected")
|
||||
}
|
||||
s.srv = &http.Server{
|
||||
Addr: s.addr,
|
||||
Handler: newHandler(s.st, s.version, s.token, s.adminUser, s.adminPass),
|
||||
ReadHeaderTimeout: 10 * time.Second,
|
||||
ReadTimeout: 30 * time.Second,
|
||||
WriteTimeout: 30 * time.Second,
|
||||
IdleTimeout: 120 * time.Second,
|
||||
}
|
||||
|
||||
ln, err := net.Listen("tcp", s.addr)
|
||||
if err != nil {
|
||||
return fmt.Errorf("listen %s: %w", s.addr, err)
|
||||
}
|
||||
s.addr = ln.Addr().String()
|
||||
|
||||
go func() {
|
||||
if err := s.srv.Serve(ln); err != nil && err != http.ErrServerClosed {
|
||||
log.Printf("receiver server: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Addr returns the address the server is listening on. Valid after Start.
|
||||
func (s *Server) Addr() string { return s.addr }
|
||||
|
||||
// Shutdown gracefully stops the server, waiting up to timeout for in-flight
|
||||
// requests to complete.
|
||||
func (s *Server) Shutdown(timeout time.Duration) error {
|
||||
if s.srv == nil {
|
||||
return nil
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), timeout)
|
||||
defer cancel()
|
||||
return s.srv.Shutdown(ctx)
|
||||
}
|
||||
|
||||
// handler routes receiver HTTP requests.
|
||||
type handler struct {
|
||||
st *store.Store
|
||||
version string
|
||||
token string
|
||||
adminUser string
|
||||
adminPass string
|
||||
mux *http.ServeMux
|
||||
}
|
||||
|
||||
// newHandler builds the receiver's HTTP handler. token "" leaves /webhook
|
||||
// open; adminPass "" leaves /admin permanently locked (every request 401s).
|
||||
func newHandler(st *store.Store, version, token, adminUser, adminPass string) http.Handler {
|
||||
h := &handler{st: st, version: version, token: token,
|
||||
adminUser: adminUser, adminPass: adminPass, mux: http.NewServeMux()}
|
||||
h.mux.HandleFunc("POST /webhook", h.ingest)
|
||||
h.mux.HandleFunc("GET /healthz", h.healthz)
|
||||
|
||||
// The whole /admin subtree sits behind basic auth, including paths the
|
||||
// inner mux will 404.
|
||||
admin := http.NewServeMux()
|
||||
admin.HandleFunc("GET /admin", h.adminPage)
|
||||
admin.HandleFunc("GET /admin/{$}", h.adminPage)
|
||||
admin.HandleFunc("GET /admin/api/traversals", h.adminTraversals)
|
||||
admin.HandleFunc("GET /admin/api/stats", h.adminStats)
|
||||
protected := h.requireAdmin(admin)
|
||||
h.mux.Handle("/admin", protected)
|
||||
h.mux.Handle("/admin/", protected)
|
||||
return h.mux
|
||||
}
|
||||
|
||||
// healthz handles GET /healthz.
|
||||
func (h *handler) healthz(w http.ResponseWriter, _ *http.Request) {
|
||||
writeJSON(w, http.StatusOK, map[string]string{"status": "ok", "version": h.version})
|
||||
}
|
||||
|
||||
// ingest handles POST /webhook: it routes on the payload's "event" field
|
||||
// and upserts the event into the store. The payload structs live in the
|
||||
// store package and mirror web/api/webhook.go exactly.
|
||||
func (h *handler) ingest(w http.ResponseWriter, r *http.Request) {
|
||||
if !h.authorized(r) {
|
||||
w.Header().Set("WWW-Authenticate", "Bearer")
|
||||
writeError(w, http.StatusUnauthorized, "missing or invalid bearer token")
|
||||
return
|
||||
}
|
||||
|
||||
r.Body = http.MaxBytesReader(w, r.Body, maxBodyBytes)
|
||||
body, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
var tooLarge *http.MaxBytesError
|
||||
if errors.As(err, &tooLarge) {
|
||||
writeError(w, http.StatusRequestEntityTooLarge,
|
||||
fmt.Sprintf("body exceeds %d bytes", tooLarge.Limit))
|
||||
return
|
||||
}
|
||||
writeError(w, http.StatusBadRequest, "read body: "+err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
var probe struct {
|
||||
Event string `json:"event"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &probe); err != nil {
|
||||
writeError(w, http.StatusBadRequest, "invalid JSON: "+err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
switch probe.Event {
|
||||
case "start":
|
||||
var ev store.StartEvent
|
||||
if err := json.Unmarshal(body, &ev); err != nil {
|
||||
writeError(w, http.StatusBadRequest, "invalid start event: "+err.Error())
|
||||
return
|
||||
}
|
||||
err = h.st.RecordStart(r.Context(), ev)
|
||||
case "complete":
|
||||
var ev store.CompleteEvent
|
||||
if err := json.Unmarshal(body, &ev); err != nil {
|
||||
writeError(w, http.StatusBadRequest, "invalid complete event: "+err.Error())
|
||||
return
|
||||
}
|
||||
err = h.st.RecordComplete(r.Context(), ev)
|
||||
default:
|
||||
writeError(w, http.StatusBadRequest, fmt.Sprintf("unknown event %q", probe.Event))
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
log.Printf("receiver: store %s event: %v", probe.Event, err)
|
||||
writeError(w, http.StatusInternalServerError, "store event failed")
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
// authorized checks the bearer token when one is configured.
|
||||
func (h *handler) authorized(r *http.Request) bool {
|
||||
if h.token == "" {
|
||||
return true
|
||||
}
|
||||
const prefix = "Bearer "
|
||||
auth := r.Header.Get("Authorization")
|
||||
if !strings.HasPrefix(auth, prefix) {
|
||||
return false
|
||||
}
|
||||
return secretEqual(strings.TrimPrefix(auth, prefix), h.token)
|
||||
}
|
||||
|
||||
// writeJSON encodes v as JSON and writes it to w with the given status code.
|
||||
func writeJSON(w http.ResponseWriter, status int, v any) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.WriteHeader(status)
|
||||
_ = json.NewEncoder(w).Encode(v)
|
||||
}
|
||||
|
||||
// writeError writes a JSON error response.
|
||||
func writeError(w http.ResponseWriter, status int, msg string) {
|
||||
writeJSON(w, status, map[string]string{"error": msg})
|
||||
}
|
||||
Reference in New Issue
Block a user