277 lines
7.2 KiB
Go
277 lines
7.2 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/adapters/plane"
|
|
"github.com/nomadcode/nomadcode-core/internal/workflow"
|
|
)
|
|
|
|
type PlaneWorkItemClient interface {
|
|
GetWorkItem(ctx context.Context, ref plane.WorkItemRef) (plane.WorkItem, error)
|
|
}
|
|
|
|
type Handler struct {
|
|
db *pgxpool.Pool
|
|
workflow *workflow.Service
|
|
plane PlaneWorkItemClient
|
|
logger *slog.Logger
|
|
}
|
|
|
|
func NewHandler(pool *pgxpool.Pool, workflowService *workflow.Service, planeClient PlaneWorkItemClient, logger *slog.Logger) *Handler {
|
|
return &Handler{
|
|
db: pool,
|
|
workflow: workflowService,
|
|
plane: planeClient,
|
|
logger: logger,
|
|
}
|
|
}
|
|
|
|
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.plane == 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
|
|
}
|
|
|
|
workItem, err := h.plane.GetWorkItem(r.Context(), ref)
|
|
if err != nil {
|
|
h.writeServiceError(w, err)
|
|
return
|
|
}
|
|
|
|
createInput, err := buildPlaneCreateTaskInput(input, workItem)
|
|
if err != nil {
|
|
writeError(w, stdhttp.StatusBadRequest, err.Error())
|
|
return
|
|
}
|
|
|
|
task, err := h.workflow.CreateTask(r.Context(), createInput)
|
|
if err != nil {
|
|
h.writeServiceError(w, err)
|
|
return
|
|
}
|
|
|
|
status := stdhttp.StatusCreated
|
|
|
|
writeJSON(w, status, 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() (plane.WorkItemRef, error) {
|
|
ref := plane.WorkItemRef{
|
|
WorkspaceSlug: strings.TrimSpace(input.WorkspaceSlug),
|
|
ProjectID: strings.TrimSpace(input.ProjectID),
|
|
WorkItemID: strings.TrimSpace(input.WorkItemID),
|
|
}
|
|
if ref.WorkspaceSlug == "" || ref.ProjectID == "" || ref.WorkItemID == "" {
|
|
return plane.WorkItemRef{}, errors.New("workspace_slug, project_id, and work_item_id are required")
|
|
}
|
|
return ref, nil
|
|
}
|
|
|
|
func buildPlaneCreateTaskInput(input createPlaneTaskRequest, workItem plane.WorkItem) (workflow.CreateTaskInput, error) {
|
|
ref, err := input.workItemRef()
|
|
if err != nil {
|
|
return workflow.CreateTaskInput{}, err
|
|
}
|
|
|
|
title := strings.TrimSpace(workItem.Name)
|
|
if title == "" {
|
|
title = ref.WorkItemID
|
|
}
|
|
|
|
stateID := strings.TrimSpace(input.StateID)
|
|
message := firstNonEmpty(input.Comment, workItem.DescriptionStripped, workItem.Description, workItem.DescriptionHTML, title)
|
|
metadata := map[string]string{
|
|
"workspace_slug": ref.WorkspaceSlug,
|
|
"project_id": ref.ProjectID,
|
|
"work_item_id": ref.WorkItemID,
|
|
"state_id": stateID,
|
|
"external_url": strings.TrimSpace(input.ExternalURL),
|
|
}
|
|
payload := map[string]any{
|
|
"message": message,
|
|
"plane": metadata,
|
|
}
|
|
|
|
rawPayload, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return workflow.CreateTaskInput{}, err
|
|
}
|
|
rawMetadata, err := json.Marshal(metadata)
|
|
if err != nil {
|
|
return workflow.CreateTaskInput{}, err
|
|
}
|
|
|
|
return workflow.CreateTaskInput{
|
|
Title: title,
|
|
Source: "plane",
|
|
Payload: rawPayload,
|
|
External: &workflow.ExternalRefInput{
|
|
Provider: "plane",
|
|
ID: ref.WorkItemID,
|
|
URL: strings.TrimSpace(input.ExternalURL),
|
|
Metadata: rawMetadata,
|
|
},
|
|
}, nil
|
|
}
|
|
|
|
func firstNonEmpty(values ...string) string {
|
|
for _, value := range values {
|
|
if trimmed := strings.TrimSpace(value); trimmed != "" {
|
|
return trimmed
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
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})
|
|
}
|