nomadcode/services/core/internal/http/handlers.go
toki 210f67bd9a feat: complete G07 trigger dispatch and refactor webhook handling
- Move G07 task artifacts to archive (plan, code review, complete.log)
- Refactor main.go to use new webhook router
- Update plane webhook handlers with improved error handling
- Add config for Plane webhook integration
- Update docker-compose for webhook testing
- Update plane-dev test documentation
2026-06-15 14:46:20 +09:00

455 lines
15 KiB
Go

package http
import (
"context"
"encoding/json"
"errors"
"log/slog"
stdhttp "net/http"
"strconv"
"strings"
"time"
"github.com/go-chi/chi/v5"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/nomadcode/nomadcode-core/internal/db"
"github.com/nomadcode/nomadcode-core/internal/projectsync"
"github.com/nomadcode/nomadcode-core/internal/storage"
"github.com/nomadcode/nomadcode-core/internal/workflow"
"github.com/nomadcode/nomadcode-core/internal/workitem"
"github.com/nomadcode/nomadcode-core/internal/workitempipeline"
)
type WorkItemTaskCreator interface {
CreateTaskFromWorkItem(ctx context.Context, input workitempipeline.CreateTaskInput) (storage.Task, error)
}
type Handler struct {
db *pgxpool.Pool
workflow *workflow.Service
workItemProviders map[workitem.ProviderID]WorkItemTaskCreator
logger *slog.Logger
planeWebhookSecret string
planeDispatch planeWebhookDispatch
}
// PlaneWebhookDispatchConfig carries the non-secret trigger gate identifiers used
// to turn a normalized Plane webhook event into a work item task creation. The
// webhook secret stays separate (SetPlaneWebhookSecret); all values here are
// non-secret Plane identifiers (workspace slug/id, Backlog state id, AGENT
// assignee id, optional self actor id).
type PlaneWebhookDispatchConfig struct {
WorkspaceID string
WorkspaceSlug string
BacklogStateID string
AgentAssigneeID string
SelfActorID string
}
// planeWebhookDispatch is the trimmed, internal form of the dispatch config.
type planeWebhookDispatch struct {
workspaceID string
workspaceSlug string
backlogStateID string
agentAssigneeID string
selfActorID string
}
// ready reports whether the minimum dispatch config is present to build a
// dispatch-ready ref and a meaningful trigger gate. The optional self actor and
// workspace id guard are not required.
func (c planeWebhookDispatch) ready() bool {
return c.workspaceSlug != "" && c.backlogStateID != "" && c.agentAssigneeID != ""
}
type WorkItemProvider struct {
ID workitem.ProviderID
Reader workitem.Reader
}
func NewHandler(pool *pgxpool.Pool, workflowService *workflow.Service, binder workitempipeline.ProjectBinder, logger *slog.Logger, providers ...WorkItemProvider) *Handler {
h := &Handler{
db: pool,
workflow: workflowService,
workItemProviders: make(map[workitem.ProviderID]WorkItemTaskCreator),
logger: logger,
}
if workflowService != nil {
for _, provider := range providers {
if provider.Reader == nil {
continue
}
h.registerWorkItemProvider(provider.ID, workitempipeline.New(provider.Reader, binder, workflowService))
}
}
return h
}
// ProjectSyncStore is the narrow storage dependency the project sync binder needs
// to look up the active setting for a provider/project target and to reserve an
// available workspace slot for it.
type ProjectSyncStore interface {
GetActiveProjectSyncSettingByTarget(ctx context.Context, provider, tenant, project string) (db.ProjectSyncSetting, error)
UpsertWorkspaceSlot(ctx context.Context, args db.UpsertWorkspaceSlotParams) (db.WorkspaceSlot, error)
ReserveWorkspaceSlot(ctx context.Context, projectSyncSettingID int64) (db.WorkspaceSlot, error)
}
type projectSyncBinder struct {
store ProjectSyncStore
provisioner projectsync.WorkspaceProvisioner
}
// NewProjectBinder adapts a project sync store (the storage Store) into the
// workitempipeline.ProjectBinder the work item pipeline consumes. It keys the
// binding strictly on the work item ref's provider/tenant/project so each project
// resolves only its own git/workspace config, and reserves an available slot
// against the resolved setting ID. A nil store yields a nil binder, which makes
// the pipeline fail explicitly as unconfigured.
func NewProjectBinder(store ProjectSyncStore) workitempipeline.ProjectBinder {
if store == nil {
return nil
}
return projectSyncBinder{
store: store,
provisioner: projectsync.FilesystemWorkspaceProvisioner{},
}
}
// NewProjectBinderWithProvisioner allows injecting a custom provisioner (e.g. for testing).
func NewProjectBinderWithProvisioner(store ProjectSyncStore, provisioner projectsync.WorkspaceProvisioner) workitempipeline.ProjectBinder {
if store == nil {
return nil
}
return projectSyncBinder{
store: store,
provisioner: provisioner,
}
}
func (b projectSyncBinder) ResolveProjectBinding(ctx context.Context, ref workitem.Ref) (workitempipeline.ProjectBinding, error) {
setting, err := b.store.GetActiveProjectSyncSettingByTarget(ctx, string(ref.Provider), ref.Tenant, ref.Project)
if err != nil {
return workitempipeline.ProjectBinding{}, err
}
return workitempipeline.ProjectBinding{
SettingID: setting.ID,
Config: projectsync.ConfigFromDBRecord(setting),
}, nil
}
func (b projectSyncBinder) EnsureProjectWorkspace(ctx context.Context, binding workitempipeline.ProjectBinding) error {
config, err := binding.Config.Normalize()
if err != nil {
return err
}
plan, err := projectsync.BuildProvisionPlan(config)
if err != nil {
return err
}
if b.provisioner != nil {
if err := b.provisioner.EnsureProvisioned(ctx, plan); err != nil {
return err
}
}
params, err := projectsync.ToUpsertWorkspaceSlotParams(binding.SettingID, projectsync.DefaultSlotIndex, plan.DefaultSlotPath)
if err != nil {
return err
}
_, err = b.store.UpsertWorkspaceSlot(ctx, params)
return err
}
func (b projectSyncBinder) ReserveWorkspaceSlot(ctx context.Context, projectSyncSettingID int64) (projectsync.WorkspaceSlot, error) {
slot, err := b.store.ReserveWorkspaceSlot(ctx, projectSyncSettingID)
if err != nil {
return projectsync.WorkspaceSlot{}, err
}
return projectsync.SlotFromDBRecord(slot), nil
}
func (h *Handler) registerWorkItemProvider(id workitem.ProviderID, creator WorkItemTaskCreator) {
trimmed := strings.TrimSpace(string(id))
if trimmed == "" || creator == nil {
return
}
h.workItemProviders[workitem.ProviderID(trimmed)] = creator
}
func (h *Handler) Healthz(w stdhttp.ResponseWriter, r *stdhttp.Request) {
writeJSON(w, stdhttp.StatusOK, map[string]string{"status": "ok"})
}
func (h *Handler) Readyz(w stdhttp.ResponseWriter, r *stdhttp.Request) {
ctx, cancel := context.WithTimeout(r.Context(), readyTimeout)
defer cancel()
if err := h.db.Ping(ctx); err != nil {
writeError(w, stdhttp.StatusServiceUnavailable, "database is not ready")
return
}
writeJSON(w, stdhttp.StatusOK, map[string]string{"status": "ready"})
}
func (h *Handler) CreateTask(w stdhttp.ResponseWriter, r *stdhttp.Request) {
var input workflow.CreateTaskInput
if err := json.NewDecoder(r.Body).Decode(&input); err != nil {
writeError(w, stdhttp.StatusBadRequest, "invalid JSON body")
return
}
task, err := h.workflow.CreateTask(r.Context(), input)
if err != nil {
h.writeServiceError(w, err)
return
}
writeJSON(w, stdhttp.StatusCreated, map[string]string{
"id": task.ID,
"status": task.Status,
"external_provider": stringValue(task.ExternalProvider),
"external_id": stringValue(task.ExternalID),
})
}
type legacyWorkItemTaskRequest struct {
WorkspaceSlug string `json:"workspace_slug"`
ProjectID string `json:"project_id"`
WorkItemID string `json:"work_item_id"`
StateID string `json:"state_id"`
ExternalURL string `json:"external_url"`
Comment string `json:"comment"`
TriggerStateID string `json:"trigger_state_id"`
AgentAssigneeID string `json:"agent_assignee_id"`
Actor string `json:"actor"`
SelfActor string `json:"self_actor"`
}
func (h *Handler) CreatePlaneTask(w stdhttp.ResponseWriter, r *stdhttp.Request) {
if _, ok := h.workItemProviders["plane"]; !ok {
writeError(w, stdhttp.StatusServiceUnavailable, "plane client is not configured")
return
}
var input legacyWorkItemTaskRequest
if err := json.NewDecoder(r.Body).Decode(&input); err != nil {
writeError(w, stdhttp.StatusBadRequest, "invalid JSON body")
return
}
_, err := input.workItemRef(workitem.ProviderID("plane"))
if err != nil {
writeError(w, stdhttp.StatusBadRequest, err.Error())
return
}
var body createWorkItemTaskRequest
body.Tenant = strings.TrimSpace(input.WorkspaceSlug)
body.Project = strings.TrimSpace(input.ProjectID)
body.ID = strings.TrimSpace(input.WorkItemID)
body.ExternalURL = strings.TrimSpace(input.ExternalURL)
body.StateID = strings.TrimSpace(input.StateID)
body.Comment = strings.TrimSpace(input.Comment)
body.TriggerStateID = strings.TrimSpace(input.TriggerStateID)
body.AgentAssigneeID = strings.TrimSpace(input.AgentAssigneeID)
body.Actor = strings.TrimSpace(input.Actor)
body.SelfActor = strings.TrimSpace(input.SelfActor)
// Delegate to generic registry path with provider="plane"
h.createWorkItemTask(w, r, "plane", &body)
}
func (h *Handler) GetTask(w stdhttp.ResponseWriter, r *stdhttp.Request) {
task, err := h.workflow.GetTask(r.Context(), chi.URLParam(r, "id"))
if err != nil {
h.writeServiceError(w, err)
return
}
writeJSON(w, stdhttp.StatusOK, task)
}
func (h *Handler) ListTasks(w stdhttp.ResponseWriter, r *stdhttp.Request) {
limit := int32(20)
if rawLimit := r.URL.Query().Get("limit"); rawLimit != "" {
parsed, err := strconv.Atoi(rawLimit)
if err != nil || parsed < 1 {
writeError(w, stdhttp.StatusBadRequest, "limit must be a positive integer")
return
}
limit = int32(parsed)
}
tasks, err := h.workflow.ListTasks(r.Context(), limit)
if err != nil {
h.writeServiceError(w, err)
return
}
writeJSON(w, stdhttp.StatusOK, tasks)
}
func (h *Handler) EnqueueTask(w stdhttp.ResponseWriter, r *stdhttp.Request) {
task, err := h.workflow.EnqueueTask(r.Context(), chi.URLParam(r, "id"))
if err != nil {
h.writeServiceError(w, err)
return
}
writeJSON(w, stdhttp.StatusOK, map[string]string{
"id": task.ID,
"status": task.Status,
})
}
func (input legacyWorkItemTaskRequest) workItemRef(provider workitem.ProviderID) (workitem.Ref, error) {
tenant := strings.TrimSpace(input.WorkspaceSlug)
project := strings.TrimSpace(input.ProjectID)
id := strings.TrimSpace(input.WorkItemID)
if tenant == "" || project == "" || id == "" {
return workitem.Ref{}, errors.New("workspace_slug, project_id, and work_item_id are required")
}
return workitem.NormalizeRef(workitem.Ref{
Provider: provider,
Tenant: tenant,
Project: project,
ID: id,
URL: strings.TrimSpace(input.ExternalURL),
})
}
type createWorkItemTaskRequest struct {
Tenant string `json:"tenant"`
Project string `json:"project"`
ID string `json:"id"`
Provider string `json:"provider,omitempty"`
ExternalURL string `json:"external_url"`
StateID string `json:"state_id"`
Comment string `json:"comment"`
TriggerStateID string `json:"trigger_state_id"`
AgentAssigneeID string `json:"agent_assignee_id"`
Actor string `json:"actor"`
SelfActor string `json:"self_actor"`
}
func (h *Handler) createWorkItemTask(w stdhttp.ResponseWriter, r *stdhttp.Request, provider string, body *createWorkItemTaskRequest) {
prov := workitem.ProviderID(strings.TrimSpace(provider))
if body != nil {
bodyProvider := strings.TrimSpace(body.Provider)
if bodyProvider != "" && bodyProvider != string(prov) {
writeError(w, stdhttp.StatusBadRequest, "provider mismatch: path="+string(prov)+", body="+bodyProvider)
return
}
}
creator, ok := h.workItemProviders[prov]
if !ok {
writeError(w, stdhttp.StatusBadRequest, "unknown or unconfigured provider: "+string(prov))
return
}
t := strings.TrimSpace(body.Tenant)
p := strings.TrimSpace(body.Project)
id := strings.TrimSpace(body.ID)
if t == "" || p == "" || id == "" {
writeError(w, stdhttp.StatusBadRequest, "tenant, project, and id are required")
return
}
ref, err := workitem.NormalizeRef(workitem.Ref{
Provider: prov,
Tenant: t,
Project: p,
ID: id,
URL: strings.TrimSpace(body.ExternalURL),
})
if err != nil {
writeError(w, stdhttp.StatusBadRequest, "invalid ref: "+err.Error())
return
}
task, err := creator.CreateTaskFromWorkItem(r.Context(), workitempipeline.CreateTaskInput{
Ref: ref,
StateID: strings.TrimSpace(body.StateID),
Comment: strings.TrimSpace(body.Comment),
Trigger: workitempipeline.CreationTrigger{
RequiredStateID: strings.TrimSpace(body.TriggerStateID),
RequiredAssigneeID: strings.TrimSpace(body.AgentAssigneeID),
Actor: strings.TrimSpace(body.Actor),
SelfActor: strings.TrimSpace(body.SelfActor),
},
})
if err != nil {
h.writeServiceError(w, err)
return
}
writeJSON(w, stdhttp.StatusCreated, map[string]string{
"id": task.ID,
"status": task.Status,
"external_provider": stringValue(task.ExternalProvider),
"external_id": stringValue(task.ExternalID),
})
}
func (h *Handler) CreateWorkItemTask(w stdhttp.ResponseWriter, r *stdhttp.Request) {
var body createWorkItemTaskRequest
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
writeError(w, stdhttp.StatusBadRequest, "invalid JSON body")
return
}
provider := chi.URLParam(r, "provider")
h.createWorkItemTask(w, r, provider, &body)
}
func stringValue(value *string) string {
if value == nil {
return ""
}
return *value
}
func (h *Handler) writeServiceError(w stdhttp.ResponseWriter, err error) {
switch {
case errors.Is(err, workflow.ErrInvalidTaskInput):
writeError(w, stdhttp.StatusBadRequest, err.Error())
case errors.Is(err, workflow.ErrTaskCannotBeEnqueued):
writeError(w, stdhttp.StatusConflict, err.Error())
case errors.Is(err, workitempipeline.ErrProjectSyncNotConfigured):
writeError(w, stdhttp.StatusServiceUnavailable, "project sync resolver is not configured")
case errors.Is(err, storage.ErrProjectSyncNotFound):
writeError(w, stdhttp.StatusUnprocessableEntity, "no active project sync setting for the work item's project")
case errors.Is(err, projectsync.ErrWorkspaceProvisionNotReady):
writeError(w, stdhttp.StatusConflict, "workspace provision is not ready")
case errors.Is(err, storage.ErrNoAvailableWorkspaceSlot):
writeError(w, stdhttp.StatusConflict, "no available workspace slot")
case errors.Is(err, workitempipeline.ErrCreationTriggerIgnored):
writeJSON(w, stdhttp.StatusAccepted, map[string]string{
"status": "ignored",
"reason": err.Error(),
})
case errors.Is(err, pgx.ErrNoRows):
writeError(w, stdhttp.StatusNotFound, "task not found")
default:
if h.logger != nil {
h.logger.Error("request failed", "error", err)
}
writeError(w, stdhttp.StatusInternalServerError, "internal server error")
}
}
const readyTimeout = 2 * time.Second
func writeJSON(w stdhttp.ResponseWriter, status int, value any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(value)
}
func writeError(w stdhttp.ResponseWriter, status int, message string) {
writeJSON(w, status, map[string]string{"error": message})
}