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/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 } func NewHandler(pool *pgxpool.Pool, workflowService *workflow.Service, plane workitem.Reader, jira workitem.Reader, logger *slog.Logger) *Handler { h := &Handler{ db: pool, workflow: workflowService, workItemProviders: make(map[workitem.ProviderID]WorkItemTaskCreator), logger: logger, } if plane != nil && workflowService != nil { h.registerWorkItemProvider("plane", workitempipeline.New(plane, workflowService)) } if jira != nil && workflowService != nil { h.workItemProviders[workitem.ProviderID("jira")] = workitempipeline.New(jira, workflowService) } return h } 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"` } 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) // 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"` } 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), }) 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, 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}) }