package protosocket import ( "context" "encoding/json" "errors" "github.com/jackc/pgx/v5" "github.com/nomadcode/nomadcode-core/internal/storage" "github.com/nomadcode/nomadcode-core/internal/workflow" ) // TaskService is the subset of the workflow service required to map REST task // semantics onto proto-socket actions. type TaskService interface { CreateTask(context.Context, workflow.CreateTaskInput) (storage.Task, error) ListTasks(context.Context, int32) ([]storage.Task, error) GetTask(context.Context, string) (storage.Task, error) EnqueueTask(context.Context, string) (storage.Task, error) } // TaskChannels registers task.* request handlers on a dispatcher. type TaskChannels struct { tasks TaskService } func NewTaskChannels(tasks TaskService) *TaskChannels { return &TaskChannels{tasks: tasks} } func (c *TaskChannels) Register(dispatcher *Dispatcher) { dispatcher.Register("task.create", c.handleCreate) dispatcher.Register("task.list", c.handleList) dispatcher.Register("task.get", c.handleGet) dispatcher.Register("task.enqueue", c.handleEnqueue) } type listRequest struct { Limit int32 `json:"limit"` } type idRequest struct { ID string `json:"id"` } func (c *TaskChannels) handleCreate(ctx context.Context, req Envelope) Envelope { input, err := payloadAs[workflow.CreateTaskInput](req.Payload) if err != nil { return invalidPayloadResponse(req, err) } task, err := c.tasks.CreateTask(ctx, input) if err != nil { return mapTaskError(req, err) } taskMap, err := taskToMap(task) if err != nil { return ErrorResponse(req, "internal.error", "internal server error", true) } return SuccessResponse(req, map[string]any{ "id": task.ID, "status": task.Status, "external_provider": stringOrEmpty(task.ExternalProvider), "external_id": stringOrEmpty(task.ExternalID), "task": taskMap, }) } func (c *TaskChannels) handleList(ctx context.Context, req Envelope) Envelope { input, err := payloadAs[listRequest](req.Payload) if err != nil { return invalidPayloadResponse(req, err) } tasks, err := c.tasks.ListTasks(ctx, input.Limit) if err != nil { return mapTaskError(req, err) } encoded := make([]any, 0, len(tasks)) for _, task := range tasks { taskMap, err := taskToMap(task) if err != nil { return ErrorResponse(req, "internal.error", "internal server error", true) } encoded = append(encoded, taskMap) } return SuccessResponse(req, map[string]any{"tasks": encoded}) } func (c *TaskChannels) handleGet(ctx context.Context, req Envelope) Envelope { input, err := payloadAs[idRequest](req.Payload) if err != nil { return invalidPayloadResponse(req, err) } task, err := c.tasks.GetTask(ctx, input.ID) if err != nil { return mapTaskError(req, err) } taskMap, err := taskToMap(task) if err != nil { return ErrorResponse(req, "internal.error", "internal server error", true) } return SuccessResponse(req, map[string]any{"task": taskMap}) } func (c *TaskChannels) handleEnqueue(ctx context.Context, req Envelope) Envelope { input, err := payloadAs[idRequest](req.Payload) if err != nil { return invalidPayloadResponse(req, err) } task, err := c.tasks.EnqueueTask(ctx, input.ID) if err != nil { return mapTaskError(req, err) } taskMap, err := taskToMap(task) if err != nil { return ErrorResponse(req, "internal.error", "internal server error", true) } return SuccessResponse(req, map[string]any{ "id": task.ID, "status": task.Status, "task": taskMap, }) } func mapTaskError(req Envelope, err error) Envelope { switch { case errors.Is(err, workflow.ErrInvalidTaskInput): return ErrorResponse(req, "task.invalid_input", err.Error(), false) case errors.Is(err, workflow.ErrTaskCannotBeEnqueued): return ErrorResponse(req, "task.conflict", err.Error(), false) case errors.Is(err, pgx.ErrNoRows): return ErrorResponse(req, "task.not_found", "task not found", false) default: return ErrorResponse(req, "internal.error", "internal server error", true) } } func invalidPayloadResponse(req Envelope, err error) Envelope { return ErrorResponse(req, "task.invalid_payload", "invalid payload: "+err.Error(), false) } // payloadAs decodes an envelope payload into a typed request by round-tripping // through JSON instead of ad hoc map indexing. func payloadAs[T any](payload map[string]any) (T, error) { var out T b, err := json.Marshal(payload) if err != nil { return out, err } return out, json.Unmarshal(b, &out) } func taskToMap(task storage.Task) (map[string]any, error) { b, err := json.Marshal(task) if err != nil { return nil, err } var out map[string]any if err := json.Unmarshal(b, &out); err != nil { return nil, err } return out, nil } func stringOrEmpty(value *string) string { if value == nil { return "" } return *value }