- Add Plane webhook handler for issue events (created, state, assignees) - Add Plane webhook integration tests with testdata fixtures - Add Gito Protosocket consumer wire readiness milestone - Add Plane work item webhook intake milestone - Add agent-task for plane-work-item-webhook-intake (trigger dispatch, idempotency, live smoke) - Update service config, router, handlers for Plane webhook endpoints - Add SOPS env setup script and secrets configuration - Update agent-ops domain rules and phase roadmap
425 lines
14 KiB
Go
425 lines
14 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
|
|
}
|
|
|
|
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})
|
|
}
|