Forgejo Webhook push 이벤트를 branch.updated 후생성자로 연결한다. WebhookAdapter에 VerifyWebhook, NormalizeWebhook 인터페이스 구현. 컨트롤 플레임 라우터에서 webhook 경로를 adapter와 연동.
515 lines
17 KiB
Go
515 lines
17 KiB
Go
package controlplane
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"io"
|
|
"log/slog"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
|
|
"git.toki-labs.com/toki/gito/services/core/internal/config"
|
|
"git.toki-labs.com/toki/gito/services/core/internal/protosocket"
|
|
"git.toki-labs.com/toki/gito/services/core/internal/provider"
|
|
"git.toki-labs.com/toki/gito/services/core/internal/provider/forgejo"
|
|
"git.toki-labs.com/toki/gito/services/core/internal/storage"
|
|
)
|
|
|
|
func NewRouter(cfg config.Config, logger *slog.Logger) http.Handler {
|
|
return newRouterWithStore(cfg, logger, nil)
|
|
}
|
|
|
|
// NewRouterWithStore constructs a router backed by the given durable store.
|
|
// Pass nil to fall back to in-memory runtime (same as NewRouter).
|
|
func NewRouterWithStore(cfg config.Config, logger *slog.Logger, store storage.Store) http.Handler {
|
|
return newRouterWithStore(cfg, logger, store)
|
|
}
|
|
|
|
func newRouterWithStore(cfg config.Config, logger *slog.Logger, store storage.Store) http.Handler {
|
|
mux := http.NewServeMux()
|
|
protoServer := protosocket.NewServer(protosocket.Config{
|
|
HeartbeatIntervalSec: cfg.ProtoSocketHeartbeatSec,
|
|
HeartbeatWaitSec: cfg.ProtoSocketHeartbeatWait,
|
|
}, logger)
|
|
runtime := NewRuntimeWithStore(protoServer, store)
|
|
registerProtoSocketHandlers(protoServer.Dispatcher(), runtime, protoServer)
|
|
|
|
mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) {
|
|
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
|
|
})
|
|
mux.HandleFunc("/readyz", func(w http.ResponseWriter, _ *http.Request) {
|
|
writeJSON(w, http.StatusOK, map[string]string{
|
|
"status": "ready",
|
|
"env": cfg.AppEnv,
|
|
})
|
|
})
|
|
mux.HandleFunc("/api/listeners/branches", handleBranchListeners(runtime))
|
|
mux.HandleFunc("/api/webhook-subscriptions", handleWebhookSubscriptions(runtime))
|
|
mux.HandleFunc("/api/events", handleEvents(runtime))
|
|
mux.HandleFunc("/callbacks/providers/", handleProviderWebhook(runtime))
|
|
mux.HandleFunc("/callbacks/forgejo/push", handleForgejoPush(cfg, runtime))
|
|
mux.HandleFunc(cfg.ProtoSocketPath, func(w http.ResponseWriter, r *http.Request) {
|
|
if isWebSocketUpgrade(r) {
|
|
protoServer.ServeHTTP(w, r)
|
|
return
|
|
}
|
|
if r.Method != http.MethodGet {
|
|
w.Header().Set("Allow", http.MethodGet)
|
|
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{
|
|
"error": "method not allowed",
|
|
})
|
|
return
|
|
}
|
|
logger.Info("proto-socket registry scaffold", "path", r.URL.Path)
|
|
writeJSON(w, http.StatusOK, DefaultProtoSocketRegistry())
|
|
})
|
|
return mux
|
|
}
|
|
|
|
func registerProtoSocketHandlers(dispatcher *protosocket.Dispatcher, runtime *Runtime, subscriber protosocket.EventSubscriber) {
|
|
dispatcher.Register("event.subscribe", func(_ context.Context, env protosocket.Envelope) protosocket.Envelope {
|
|
connectionID := stringValue(env.Meta["connection_id"])
|
|
if connectionID == "" {
|
|
return protosocket.ErrorResponse(env, "event.missing_connection", "event.subscribe requires a proto-socket connection id", false)
|
|
}
|
|
subscription := eventSubscriptionFromPayload(env.Payload)
|
|
registered := subscriber.Subscribe(connectionID, subscription)
|
|
return protosocket.SuccessResponse(env, map[string]any{
|
|
"status": "subscribed",
|
|
"subscription": registered.Payload(),
|
|
})
|
|
})
|
|
dispatcher.Register("event.list", func(_ context.Context, env protosocket.Envelope) protosocket.Envelope {
|
|
return protosocket.SuccessResponse(env, map[string]any{
|
|
"events": eventRecordsPayload(runtime.ListEvents()),
|
|
})
|
|
})
|
|
}
|
|
|
|
func eventSubscriptionFromPayload(payload map[string]any) protosocket.EventSubscription {
|
|
events := stringListValue(payload["events"])
|
|
if len(events) == 0 {
|
|
events = []string{"branch.updated"}
|
|
}
|
|
return protosocket.EventSubscription{
|
|
Events: events,
|
|
RepoID: stringValue(payload["repo_id"]),
|
|
Branch: stringValue(payload["branch"]),
|
|
}
|
|
}
|
|
|
|
func stringListValue(value any) []string {
|
|
items, ok := value.([]any)
|
|
if !ok {
|
|
if values, ok := value.([]string); ok {
|
|
return values
|
|
}
|
|
return nil
|
|
}
|
|
result := make([]string, 0, len(items))
|
|
for _, item := range items {
|
|
if text := stringValue(item); text != "" {
|
|
result = append(result, text)
|
|
}
|
|
}
|
|
return result
|
|
}
|
|
|
|
func stringValue(value any) string {
|
|
text, ok := value.(string)
|
|
if !ok {
|
|
return ""
|
|
}
|
|
return strings.TrimSpace(text)
|
|
}
|
|
|
|
func handleBranchListeners(runtime *Runtime) http.HandlerFunc {
|
|
type request struct {
|
|
RepoID string `json:"repo_id"`
|
|
Branch string `json:"branch"`
|
|
Provider string `json:"provider"`
|
|
}
|
|
return func(w http.ResponseWriter, r *http.Request) {
|
|
switch r.Method {
|
|
case http.MethodGet:
|
|
writeJSON(w, http.StatusOK, map[string]any{"listeners": runtime.ListBranchWatches()})
|
|
case http.MethodPost:
|
|
var input request
|
|
if err := json.NewDecoder(r.Body).Decode(&input); err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid listener payload"})
|
|
return
|
|
}
|
|
watch, err := runtime.RegisterBranchWatch(input.RepoID, input.Branch, input.Provider)
|
|
if err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()})
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusCreated, map[string]any{"listener": watch})
|
|
default:
|
|
w.Header().Set("Allow", "GET, POST")
|
|
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
|
|
}
|
|
}
|
|
}
|
|
|
|
func handleWebhookSubscriptions(runtime *Runtime) http.HandlerFunc {
|
|
type request struct {
|
|
Name string `json:"name"`
|
|
TargetURL string `json:"target_url"`
|
|
Events []string `json:"events"`
|
|
RepoID string `json:"repo_id"`
|
|
Branch string `json:"branch"`
|
|
SecretRef string `json:"secret_ref"`
|
|
}
|
|
return func(w http.ResponseWriter, r *http.Request) {
|
|
switch r.Method {
|
|
case http.MethodGet:
|
|
writeJSON(w, http.StatusOK, map[string]any{"subscriptions": webhookSubscriptionsPayload(runtime.ListWebhookSubscriptions())})
|
|
case http.MethodPost:
|
|
var input request
|
|
if err := json.NewDecoder(r.Body).Decode(&input); err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid webhook subscription payload"})
|
|
return
|
|
}
|
|
subscription, err := runtime.RegisterWebhookSubscription(
|
|
input.Name,
|
|
input.TargetURL,
|
|
input.Events,
|
|
input.RepoID,
|
|
input.Branch,
|
|
input.SecretRef,
|
|
)
|
|
if err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()})
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusCreated, map[string]any{"subscription": webhookSubscriptionPayload(subscription)})
|
|
default:
|
|
w.Header().Set("Allow", "GET, POST")
|
|
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
|
|
}
|
|
}
|
|
}
|
|
|
|
func handleEvents(runtime *Runtime) http.HandlerFunc {
|
|
return func(w http.ResponseWriter, r *http.Request) {
|
|
if r.Method != http.MethodGet {
|
|
w.Header().Set("Allow", http.MethodGet)
|
|
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, map[string]any{"events": eventRecordsPayload(runtime.ListEvents())})
|
|
}
|
|
}
|
|
|
|
func handleProviderWebhook(runtime *Runtime) http.HandlerFunc {
|
|
const prefix = "/callbacks/providers/"
|
|
return func(w http.ResponseWriter, r *http.Request) {
|
|
if r.Method != http.MethodPost {
|
|
w.Header().Set("Allow", http.MethodPost)
|
|
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
|
|
return
|
|
}
|
|
providerID := strings.Trim(strings.TrimPrefix(r.URL.Path, prefix), "/")
|
|
if providerID == "" || strings.Contains(providerID, "/") {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "provider is required"})
|
|
return
|
|
}
|
|
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
|
|
if err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "failed to read payload"})
|
|
return
|
|
}
|
|
record, err := runtime.HandleProviderWebhook(r.Context(), provider.WebhookRequest{
|
|
Provider: provider.ProviderID(providerID),
|
|
EventType: providerWebhookEventName(r),
|
|
ExternalID: providerWebhookDeliveryID(r),
|
|
Payload: body,
|
|
Headers: map[string][]string(r.Header),
|
|
Query: r.URL.Query(),
|
|
ReceivedAt: time.Now().UTC(),
|
|
})
|
|
if err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()})
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusAccepted, map[string]any{
|
|
"accepted": true,
|
|
"event": eventRecordPayload(record),
|
|
})
|
|
}
|
|
}
|
|
|
|
func handleForgejoPush(cfg config.Config, runtime *Runtime) http.HandlerFunc {
|
|
return func(w http.ResponseWriter, r *http.Request) {
|
|
if r.Method != http.MethodPost {
|
|
w.Header().Set("Allow", http.MethodPost)
|
|
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
|
|
return
|
|
}
|
|
body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20))
|
|
if err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "failed to read payload"})
|
|
return
|
|
}
|
|
|
|
adapter := forgejo.WebhookAdapter{Secret: cfg.ForgejoWebhookSecret}
|
|
req := provider.WebhookRequest{
|
|
Provider: provider.ProviderForgejo,
|
|
EventType: webhookEventName(r),
|
|
ExternalID: r.Header.Get("X-Forgejo-Delivery"),
|
|
Payload: body,
|
|
Headers: map[string][]string(r.Header),
|
|
Query: r.URL.Query(),
|
|
ReceivedAt: time.Now().UTC(),
|
|
}
|
|
|
|
// Verify webhook signature.
|
|
if err := adapter.VerifyWebhook(r.Context(), req); err != nil {
|
|
writeJSON(w, http.StatusUnauthorized, map[string]string{"error": "invalid forgejo signature"})
|
|
return
|
|
}
|
|
|
|
// Normalize to provider-neutral candidates.
|
|
candidates, err := adapter.NormalizeWebhook(r.Context(), req)
|
|
if err != nil {
|
|
writeJSON(w, http.StatusBadRequest, map[string]string{"error": "failed to normalize payload"})
|
|
return
|
|
}
|
|
|
|
// No candidates = unsupported event or non-branch ref.
|
|
if len(candidates) == 0 {
|
|
writeJSON(w, http.StatusAccepted, map[string]any{"accepted": true, "matched": false, "reason": "unsupported event"})
|
|
return
|
|
}
|
|
|
|
// Process each candidate (typically one branch.updated).
|
|
for _, candidate := range candidates {
|
|
if candidate.Type == "branch.updated" && candidate.Revision != nil {
|
|
record, matched, err := runtime.HandleRevision(r.Context(), "forgejo", candidate.ExternalID, *candidate.Revision)
|
|
if err != nil {
|
|
writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "failed to publish event"})
|
|
return
|
|
}
|
|
resp := map[string]any{
|
|
"accepted": true,
|
|
"matched": matched,
|
|
}
|
|
if matched {
|
|
resp["event"] = eventRecordPayload(record)
|
|
}
|
|
if record.Duplicate {
|
|
resp["duplicate"] = true
|
|
}
|
|
writeJSON(w, http.StatusAccepted, resp)
|
|
return
|
|
}
|
|
}
|
|
|
|
// No matching candidate type.
|
|
writeJSON(w, http.StatusAccepted, map[string]any{"accepted": true, "matched": false, "reason": "non-branch ref"})
|
|
}
|
|
}
|
|
|
|
func webhookEventName(r *http.Request) string {
|
|
for _, header := range []string{"X-Forgejo-Event", "X-Gitea-Event", "X-Gogs-Event", "X-GitHub-Event"} {
|
|
if value := strings.TrimSpace(r.Header.Get(header)); value != "" {
|
|
return strings.ToLower(value)
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func providerWebhookEventName(r *http.Request) string {
|
|
for _, header := range []string{"X-Gito-Provider-Event", "X-GitHub-Event", "X-Forgejo-Event", "X-Gitea-Event", "X-Gogs-Event", "X-GitLab-Event"} {
|
|
if value := strings.TrimSpace(r.Header.Get(header)); value != "" {
|
|
return strings.ToLower(value)
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func providerWebhookDeliveryID(r *http.Request) string {
|
|
for _, header := range []string{"X-Gito-Delivery", "X-GitHub-Delivery", "X-Forgejo-Delivery", "X-Gitea-Delivery", "X-Gogs-Delivery", "X-GitLab-Event-UUID"} {
|
|
if value := strings.TrimSpace(r.Header.Get(header)); value != "" {
|
|
return value
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func isWebSocketUpgrade(r *http.Request) bool {
|
|
return strings.EqualFold(r.Header.Get("Upgrade"), "websocket") &&
|
|
strings.Contains(strings.ToLower(r.Header.Get("Connection")), "upgrade")
|
|
}
|
|
|
|
type ProtoSocketRegistry struct {
|
|
Transport string `json:"transport"`
|
|
Status string `json:"status"`
|
|
Channels []ProtoSocketChannel `json:"channels"`
|
|
}
|
|
|
|
type ProtoSocketChannel struct {
|
|
Name string `json:"name"`
|
|
Purpose string `json:"purpose"`
|
|
Status string `json:"status"`
|
|
Actions []ProtoSocketAction `json:"actions"`
|
|
}
|
|
|
|
type ProtoSocketAction struct {
|
|
Name string `json:"name"`
|
|
Status string `json:"status"`
|
|
}
|
|
|
|
func DefaultProtoSocketRegistry() ProtoSocketRegistry {
|
|
return ProtoSocketRegistry{
|
|
Transport: "proto-socket",
|
|
Status: "registry-placeholder",
|
|
Channels: []ProtoSocketChannel{
|
|
{
|
|
Name: "repo",
|
|
Purpose: "register, inspect, and list managed repositories",
|
|
Status: "placeholder",
|
|
Actions: []ProtoSocketAction{
|
|
{Name: "register", Status: "placeholder"},
|
|
{Name: "get", Status: "placeholder"},
|
|
{Name: "list", Status: "placeholder"},
|
|
},
|
|
},
|
|
{
|
|
Name: "workspace",
|
|
Purpose: "lease, release, and inspect workspace slots",
|
|
Status: "placeholder",
|
|
Actions: []ProtoSocketAction{
|
|
{Name: "lease", Status: "placeholder"},
|
|
{Name: "release", Status: "placeholder"},
|
|
{Name: "get", Status: "placeholder"},
|
|
},
|
|
},
|
|
{
|
|
Name: "operation",
|
|
Purpose: "create, cancel, inspect, and stream operations",
|
|
Status: "placeholder",
|
|
Actions: []ProtoSocketAction{
|
|
{Name: "create", Status: "placeholder"},
|
|
{Name: "cancel", Status: "placeholder"},
|
|
{Name: "get", Status: "placeholder"},
|
|
{Name: "stream", Status: "placeholder"},
|
|
},
|
|
},
|
|
{
|
|
Name: "git",
|
|
Purpose: "request platformless Git operations",
|
|
Status: "placeholder",
|
|
Actions: []ProtoSocketAction{
|
|
{Name: "clone", Status: "placeholder"},
|
|
{Name: "fetch", Status: "placeholder"},
|
|
{Name: "status", Status: "placeholder"},
|
|
{Name: "diff", Status: "placeholder"},
|
|
{Name: "commit", Status: "placeholder"},
|
|
{Name: "push", Status: "placeholder"},
|
|
},
|
|
},
|
|
{
|
|
Name: "change_request",
|
|
Purpose: "manage provider-neutral change requests",
|
|
Status: "placeholder",
|
|
Actions: []ProtoSocketAction{
|
|
{Name: "open", Status: "placeholder"},
|
|
{Name: "update", Status: "placeholder"},
|
|
{Name: "list", Status: "placeholder"},
|
|
},
|
|
},
|
|
{
|
|
Name: "agent_shell",
|
|
Purpose: "shell heartbeat, command dispatch, and log streaming",
|
|
Status: "placeholder",
|
|
Actions: []ProtoSocketAction{
|
|
{Name: "heartbeat", Status: "placeholder"},
|
|
{Name: "dispatch", Status: "placeholder"},
|
|
{Name: "stream_logs", Status: "placeholder"},
|
|
},
|
|
},
|
|
{
|
|
Name: "event",
|
|
Purpose: "subscribe to normalized events",
|
|
Status: "mvp",
|
|
Actions: []ProtoSocketAction{
|
|
{Name: "subscribe", Status: "mvp"},
|
|
{Name: "list", Status: "mvp"},
|
|
{Name: "ack", Status: "placeholder"},
|
|
},
|
|
},
|
|
},
|
|
}
|
|
}
|
|
|
|
func writeJSON(w http.ResponseWriter, status int, value any) {
|
|
w.Header().Set("Content-Type", "application/json")
|
|
w.WriteHeader(status)
|
|
_ = json.NewEncoder(w).Encode(value)
|
|
}
|
|
|
|
func eventRecordsPayload(records []EventRecord) []any {
|
|
payload := make([]any, 0, len(records))
|
|
for _, record := range records {
|
|
payload = append(payload, eventRecordPayload(record))
|
|
}
|
|
return payload
|
|
}
|
|
|
|
func webhookSubscriptionsPayload(subscriptions []WebhookSubscription) []any {
|
|
payload := make([]any, 0, len(subscriptions))
|
|
for _, subscription := range subscriptions {
|
|
payload = append(payload, webhookSubscriptionPayload(subscription))
|
|
}
|
|
return payload
|
|
}
|
|
|
|
func webhookSubscriptionPayload(subscription WebhookSubscription) map[string]any {
|
|
return map[string]any{
|
|
"id": subscription.ID,
|
|
"name": subscription.Name,
|
|
"events": subscription.Events,
|
|
"repo_id": subscription.RepoID,
|
|
"branch": subscription.Branch,
|
|
"created_at": subscription.CreatedAt.UTC().Format(time.RFC3339Nano),
|
|
}
|
|
}
|
|
|
|
func eventRecordPayload(record EventRecord) map[string]any {
|
|
payload := map[string]any{
|
|
"id": record.ID,
|
|
"type": record.Type,
|
|
"provider": record.Provider,
|
|
"delivery_id": record.DeliveryID,
|
|
"created_at": record.CreatedAt.UTC().Format(time.RFC3339Nano),
|
|
}
|
|
if record.Type == "branch.updated" {
|
|
changedFiles := make([]any, 0, len(record.Revision.ChangedFiles))
|
|
for _, file := range record.Revision.ChangedFiles {
|
|
changedFiles = append(changedFiles, map[string]any{
|
|
"path": file.Path,
|
|
"change_type": file.ChangeType,
|
|
})
|
|
}
|
|
payload["revision"] = map[string]any{
|
|
"repo_id": record.Revision.RepoID,
|
|
"branch": record.Revision.Branch,
|
|
"before": record.Revision.Before,
|
|
"after": record.Revision.After,
|
|
"changed_files": changedFiles,
|
|
"observed_at": record.Revision.ObservedAt.UTC().Format(time.RFC3339Nano),
|
|
}
|
|
}
|
|
if record.Type == "provider.webhook.received" && record.Webhook != nil {
|
|
payload["webhook"] = map[string]any{
|
|
"event_type": record.Webhook.EventType,
|
|
"external_id": record.Webhook.ExternalID,
|
|
"payload_size": record.Webhook.PayloadSize,
|
|
"received_at": record.Webhook.ReceivedAt.UTC().Format(time.RFC3339Nano),
|
|
}
|
|
}
|
|
return payload
|
|
}
|