nomadcode/services/core/internal/http/handlers.go

220 lines
6 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/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
workItemTasks WorkItemTaskCreator
logger *slog.Logger
}
func NewHandler(pool *pgxpool.Pool, workflowService *workflow.Service, workItemReader workitem.Reader, logger *slog.Logger) *Handler {
h := &Handler{
db: pool,
workflow: workflowService,
logger: logger,
}
if workItemReader != nil && workflowService != nil {
h.workItemTasks = workitempipeline.New(workItemReader, workflowService)
}
return h
}
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 createPlaneTaskRequest 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 h.workItemTasks == nil {
writeError(w, stdhttp.StatusServiceUnavailable, "plane client is not configured")
return
}
var input createPlaneTaskRequest
if err := json.NewDecoder(r.Body).Decode(&input); err != nil {
writeError(w, stdhttp.StatusBadRequest, "invalid JSON body")
return
}
ref, err := input.workItemRef()
if err != nil {
writeError(w, stdhttp.StatusBadRequest, err.Error())
return
}
task, err := h.workItemTasks.CreateTaskFromWorkItem(r.Context(), workitempipeline.CreateTaskInput{
Ref: ref,
StateID: input.StateID,
Comment: input.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) 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 createPlaneTaskRequest) workItemRef() (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: workitem.ProviderID("plane"),
Tenant: tenant,
Project: project,
ID: id,
URL: strings.TrimSpace(input.ExternalURL),
})
}
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})
}