- Update agent-ops project rules and roadmap files - Add proto-socket infrastructure communication rail milestone - Update Flutter pubspec.lock and contracts notes - Enhance core service: config, HTTP middleware, router - Add notification module improvements - Add protosocket internal package
182 lines
4.7 KiB
Go
182 lines
4.7 KiB
Go
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
|
|
}
|