feat(webhook): trigger brain-sync on Gitea push instead of 15-min poll
Adds POST /webhooks/brain-sync: verifies Gitea's HMAC-SHA256 signature, checks the push is to mathias/brain main, then creates a one-off Job from the existing brain-sync CronJob's template (same script the 15-min poll already runs, just triggered on-demand). Off by default -- opt in via GITEA_WEBHOOK_SECRET, since it needs Job-create RBAC in the "brain" namespace a fresh deploy won't have. 10 new tests (internal/webhook), including a fake-clientset reactor to simulate server-side GenerateName expansion, which the plain fake tracker doesn't do on its own. Needs (follow-up, infra repo): RBAC granting ingestion's ServiceAccount get on cronjobs/brain-sync + create on jobs in the brain namespace, the GITEA_WEBHOOK_SECRET env, and the actual Gitea webhook registration.
This commit is contained in:
@@ -0,0 +1,109 @@
|
||||
// Package webhook triggers an on-demand brain-sync Job when Gitea pushes to
|
||||
// mathias/brain, instead of waiting for the next 15-minute CronJob poll.
|
||||
package webhook
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/hmac"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
|
||||
batchv1 "k8s.io/api/batch/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
)
|
||||
|
||||
// VerifySignature checks a Gitea webhook's X-Gitea-Signature header: a
|
||||
// hex-encoded HMAC-SHA256 of the raw request body, keyed by the shared
|
||||
// webhook secret. Constant-time compare — timing must not leak how much of
|
||||
// the signature matched.
|
||||
func VerifySignature(payload []byte, signatureHeader, secret string) bool {
|
||||
if signatureHeader == "" {
|
||||
return false
|
||||
}
|
||||
mac := hmac.New(sha256.New, []byte(secret))
|
||||
mac.Write(payload)
|
||||
expected := hex.EncodeToString(mac.Sum(nil))
|
||||
return hmac.Equal([]byte(expected), []byte(signatureHeader))
|
||||
}
|
||||
|
||||
// pushEvent is the subset of Gitea's push webhook payload this handler needs.
|
||||
type pushEvent struct {
|
||||
Ref string `json:"ref"`
|
||||
Repo struct {
|
||||
FullName string `json:"full_name"`
|
||||
} `json:"repository"`
|
||||
}
|
||||
|
||||
// TriggerJobFromCronJob reads the named CronJob's job template and creates a
|
||||
// new, uniquely-named Job from it — the same thing `kubectl create job
|
||||
// --from=cronjob/<name>` does. Reuses the CronJob's already-tested script
|
||||
// rather than re-implementing sync logic here.
|
||||
func TriggerJobFromCronJob(ctx context.Context, cs kubernetes.Interface, namespace, cronJobName string) (string, error) {
|
||||
cj, err := cs.BatchV1().CronJobs(namespace).Get(ctx, cronJobName, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("get cronjob %s/%s: %w", namespace, cronJobName, err)
|
||||
}
|
||||
job := &batchv1.Job{
|
||||
ObjectMeta: metav1.ObjectMeta{
|
||||
GenerateName: cronJobName + "-webhook-",
|
||||
Namespace: namespace,
|
||||
Annotations: map[string]string{
|
||||
"triggered-by": "brain-webhook",
|
||||
},
|
||||
},
|
||||
Spec: cj.Spec.JobTemplate.Spec,
|
||||
}
|
||||
created, err := cs.BatchV1().Jobs(namespace).Create(ctx, job, metav1.CreateOptions{})
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("create job from cronjob %s/%s: %w", namespace, cronJobName, err)
|
||||
}
|
||||
return created.Name, nil
|
||||
}
|
||||
|
||||
// Handler is the HTTP handler for Gitea's push webhook on mathias/brain.
|
||||
type Handler struct {
|
||||
Secret string
|
||||
Clientset kubernetes.Interface
|
||||
Namespace string // e.g. "brain"
|
||||
CronJobName string // e.g. "brain-sync"
|
||||
WatchRepo string // e.g. "mathias/brain"
|
||||
Logger *slog.Logger
|
||||
}
|
||||
|
||||
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
body, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, "bad body", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
sig := r.Header.Get("X-Gitea-Signature")
|
||||
if !VerifySignature(body, sig, h.Secret) {
|
||||
http.Error(w, "bad signature", http.StatusUnauthorized)
|
||||
return
|
||||
}
|
||||
var ev pushEvent
|
||||
if err := json.Unmarshal(body, &ev); err != nil {
|
||||
http.Error(w, "bad payload", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if ev.Repo.FullName != h.WatchRepo || ev.Ref != "refs/heads/main" {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
fmt.Fprintf(w, "ignored: repo=%s ref=%s", ev.Repo.FullName, ev.Ref)
|
||||
return
|
||||
}
|
||||
jobName, err := TriggerJobFromCronJob(r.Context(), h.Clientset, h.Namespace, h.CronJobName)
|
||||
if err != nil {
|
||||
h.Logger.Error("webhook: trigger job failed", "err", err)
|
||||
http.Error(w, "trigger failed", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
h.Logger.Info("webhook: triggered brain-sync job", "job", jobName)
|
||||
w.WriteHeader(http.StatusOK)
|
||||
fmt.Fprintf(w, "triggered %s", jobName)
|
||||
}
|
||||
@@ -0,0 +1,222 @@
|
||||
package webhook
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/hmac"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
batchv1 "k8s.io/api/batch/v1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/client-go/kubernetes/fake"
|
||||
k8stesting "k8s.io/client-go/testing"
|
||||
)
|
||||
|
||||
// newFakeClientset simulates the real API server's GenerateName expansion,
|
||||
// which the plain fake clientset tracker does not do on its own -- without
|
||||
// this, every created Job keeps an empty Name (a fake-clientset limitation,
|
||||
// not real API server behavior).
|
||||
func newFakeClientset(objects ...runtime.Object) *fake.Clientset {
|
||||
cs := fake.NewSimpleClientset(objects...)
|
||||
cs.PrependReactor("create", "jobs", func(action k8stesting.Action) (bool, runtime.Object, error) {
|
||||
createAction := action.(k8stesting.CreateAction)
|
||||
job, ok := createAction.GetObject().(*batchv1.Job)
|
||||
if ok && job.Name == "" && job.GenerateName != "" {
|
||||
job.Name = job.GenerateName + "test0001"
|
||||
}
|
||||
return false, nil, nil // not "handled" -- let the default reactor store it
|
||||
})
|
||||
return cs
|
||||
}
|
||||
|
||||
func sign(t *testing.T, payload []byte, secret string) string {
|
||||
t.Helper()
|
||||
mac := hmac.New(sha256.New, []byte(secret))
|
||||
mac.Write(payload)
|
||||
return hex.EncodeToString(mac.Sum(nil))
|
||||
}
|
||||
|
||||
func testLogger() *slog.Logger {
|
||||
return slog.New(slog.NewTextHandler(os.Stderr, nil))
|
||||
}
|
||||
|
||||
func fakeCronJob(namespace, name string) *batchv1.CronJob {
|
||||
return &batchv1.CronJob{
|
||||
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: namespace},
|
||||
Spec: batchv1.CronJobSpec{
|
||||
JobTemplate: batchv1.JobTemplateSpec{
|
||||
Spec: batchv1.JobSpec{
|
||||
Template: corev1.PodTemplateSpec{
|
||||
Spec: corev1.PodSpec{
|
||||
Containers: []corev1.Container{
|
||||
{Name: "sync", Image: "alpine/git:v2.47.2"},
|
||||
},
|
||||
RestartPolicy: corev1.RestartPolicyNever,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// --------------------------------------------------------------------------- #
|
||||
// VerifySignature
|
||||
// --------------------------------------------------------------------------- #
|
||||
func TestVerifySignature_AcceptsCorrectHMAC(t *testing.T) {
|
||||
payload := []byte(`{"ref":"refs/heads/main"}`)
|
||||
secret := "s3cret"
|
||||
sig := sign(t, payload, secret)
|
||||
assert.True(t, VerifySignature(payload, sig, secret))
|
||||
}
|
||||
|
||||
func TestVerifySignature_RejectsWrongSecret(t *testing.T) {
|
||||
payload := []byte(`{"ref":"refs/heads/main"}`)
|
||||
sig := sign(t, payload, "right-secret")
|
||||
assert.False(t, VerifySignature(payload, sig, "wrong-secret"))
|
||||
}
|
||||
|
||||
func TestVerifySignature_RejectsTamperedPayload(t *testing.T) {
|
||||
secret := "s3cret"
|
||||
sig := sign(t, []byte(`{"ref":"refs/heads/main"}`), secret)
|
||||
assert.False(t, VerifySignature([]byte(`{"ref":"refs/heads/evil"}`), sig, secret))
|
||||
}
|
||||
|
||||
func TestVerifySignature_RejectsEmptySignature(t *testing.T) {
|
||||
assert.False(t, VerifySignature([]byte("payload"), "", "secret"))
|
||||
}
|
||||
|
||||
// --------------------------------------------------------------------------- #
|
||||
// TriggerJobFromCronJob
|
||||
// --------------------------------------------------------------------------- #
|
||||
func TestTriggerJobFromCronJob_CreatesJobMatchingTemplate(t *testing.T) {
|
||||
cs := newFakeClientset(fakeCronJob("brain", "brain-sync"))
|
||||
|
||||
jobName, err := TriggerJobFromCronJob(context.Background(), cs, "brain", "brain-sync")
|
||||
require.NoError(t, err)
|
||||
assert.NotEmpty(t, jobName)
|
||||
|
||||
jobs, err := cs.BatchV1().Jobs("brain").List(context.Background(), metav1.ListOptions{})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, jobs.Items, 1)
|
||||
assert.Equal(t, "alpine/git:v2.47.2", jobs.Items[0].Spec.Template.Spec.Containers[0].Image)
|
||||
}
|
||||
|
||||
func TestTriggerJobFromCronJob_ErrorsWhenCronJobMissing(t *testing.T) {
|
||||
cs := fake.NewSimpleClientset()
|
||||
_, err := TriggerJobFromCronJob(context.Background(), cs, "brain", "brain-sync")
|
||||
assert.Error(t, err)
|
||||
}
|
||||
|
||||
// --------------------------------------------------------------------------- #
|
||||
// Handler.ServeHTTP
|
||||
// --------------------------------------------------------------------------- #
|
||||
func newTestHandler(cs *fake.Clientset) *Handler {
|
||||
return &Handler{
|
||||
Secret: "s3cret",
|
||||
Clientset: cs,
|
||||
Namespace: "brain",
|
||||
CronJobName: "brain-sync",
|
||||
WatchRepo: "mathias/brain",
|
||||
Logger: testLogger(),
|
||||
}
|
||||
}
|
||||
|
||||
func pushPayload(t *testing.T, repo, ref string) []byte {
|
||||
t.Helper()
|
||||
body := map[string]any{
|
||||
"ref": ref,
|
||||
"repository": map[string]any{
|
||||
"full_name": repo,
|
||||
},
|
||||
}
|
||||
b, err := json.Marshal(body)
|
||||
require.NoError(t, err)
|
||||
return b
|
||||
}
|
||||
|
||||
func TestHandler_RejectsMissingSignature(t *testing.T) {
|
||||
cs := fake.NewSimpleClientset(fakeCronJob("brain", "brain-sync"))
|
||||
h := newTestHandler(cs)
|
||||
payload := pushPayload(t, "mathias/brain", "refs/heads/main")
|
||||
req := httptest.NewRequest(http.MethodPost, "/webhooks/brain-sync", bytes.NewReader(payload))
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
h.ServeHTTP(rec, req)
|
||||
|
||||
assert.Equal(t, http.StatusUnauthorized, rec.Code)
|
||||
jobs, _ := cs.BatchV1().Jobs("brain").List(context.Background(), metav1.ListOptions{})
|
||||
assert.Empty(t, jobs.Items, "must not trigger a job on an unsigned request")
|
||||
}
|
||||
|
||||
func TestHandler_RejectsWrongSignature(t *testing.T) {
|
||||
cs := fake.NewSimpleClientset(fakeCronJob("brain", "brain-sync"))
|
||||
h := newTestHandler(cs)
|
||||
payload := pushPayload(t, "mathias/brain", "refs/heads/main")
|
||||
req := httptest.NewRequest(http.MethodPost, "/webhooks/brain-sync", bytes.NewReader(payload))
|
||||
req.Header.Set("X-Gitea-Signature", sign(t, payload, "not-the-real-secret"))
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
h.ServeHTTP(rec, req)
|
||||
|
||||
assert.Equal(t, http.StatusUnauthorized, rec.Code)
|
||||
jobs, _ := cs.BatchV1().Jobs("brain").List(context.Background(), metav1.ListOptions{})
|
||||
assert.Empty(t, jobs.Items)
|
||||
}
|
||||
|
||||
func TestHandler_IgnoresOtherRepos(t *testing.T) {
|
||||
cs := fake.NewSimpleClientset(fakeCronJob("brain", "brain-sync"))
|
||||
h := newTestHandler(cs)
|
||||
payload := pushPayload(t, "mathias/some-other-repo", "refs/heads/main")
|
||||
req := httptest.NewRequest(http.MethodPost, "/webhooks/brain-sync", bytes.NewReader(payload))
|
||||
req.Header.Set("X-Gitea-Signature", sign(t, payload, "s3cret"))
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
h.ServeHTTP(rec, req)
|
||||
|
||||
assert.Equal(t, http.StatusOK, rec.Code)
|
||||
jobs, _ := cs.BatchV1().Jobs("brain").List(context.Background(), metav1.ListOptions{})
|
||||
assert.Empty(t, jobs.Items, "must not trigger for a push to an unrelated repo")
|
||||
}
|
||||
|
||||
func TestHandler_IgnoresNonMainBranch(t *testing.T) {
|
||||
cs := fake.NewSimpleClientset(fakeCronJob("brain", "brain-sync"))
|
||||
h := newTestHandler(cs)
|
||||
payload := pushPayload(t, "mathias/brain", "refs/heads/some-feature-branch")
|
||||
req := httptest.NewRequest(http.MethodPost, "/webhooks/brain-sync", bytes.NewReader(payload))
|
||||
req.Header.Set("X-Gitea-Signature", sign(t, payload, "s3cret"))
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
h.ServeHTTP(rec, req)
|
||||
|
||||
assert.Equal(t, http.StatusOK, rec.Code)
|
||||
jobs, _ := cs.BatchV1().Jobs("brain").List(context.Background(), metav1.ListOptions{})
|
||||
assert.Empty(t, jobs.Items, "must not trigger for a push to a non-main branch")
|
||||
}
|
||||
|
||||
func TestHandler_TriggersJobOnValidMainPush(t *testing.T) {
|
||||
cs := fake.NewSimpleClientset(fakeCronJob("brain", "brain-sync"))
|
||||
h := newTestHandler(cs)
|
||||
payload := pushPayload(t, "mathias/brain", "refs/heads/main")
|
||||
req := httptest.NewRequest(http.MethodPost, "/webhooks/brain-sync", bytes.NewReader(payload))
|
||||
req.Header.Set("X-Gitea-Signature", sign(t, payload, "s3cret"))
|
||||
rec := httptest.NewRecorder()
|
||||
|
||||
h.ServeHTTP(rec, req)
|
||||
|
||||
assert.Equal(t, http.StatusOK, rec.Code)
|
||||
jobs, err := cs.BatchV1().Jobs("brain").List(context.Background(), metav1.ListOptions{})
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, jobs.Items, 1, "a valid push to main on the watched repo must trigger exactly one job")
|
||||
}
|
||||
Reference in New Issue
Block a user