commit e0d2b5b616fdcf7705bc872f9f9d5f9a6f8411c0 Author: toki Date: Tue May 19 07:32:31 2026 +0900 Initial nomadcode-core project diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..0394978 --- /dev/null +++ b/.gitignore @@ -0,0 +1,3 @@ +/.build/ +/bin/nomadcode-core +coverage.out diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..77631b6 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,22 @@ +# syntax=docker/dockerfile:1 + +FROM golang:1.25-alpine AS build +WORKDIR /src + +RUN apk add --no-cache ca-certificates + +COPY go.mod go.sum ./ +RUN go mod download + +COPY . . +RUN CGO_ENABLED=0 GOOS=linux go build -o /bin/nomadcode-core ./cmd/server + +FROM alpine:3.21 +RUN adduser -D -H appuser +WORKDIR /app + +COPY --from=build /bin/nomadcode-core /usr/local/bin/nomadcode-core + +USER appuser +EXPOSE 8080 +ENTRYPOINT ["nomadcode-core"] diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..196e438 --- /dev/null +++ b/Makefile @@ -0,0 +1,29 @@ +export APP_ENV ?= local +export HTTP_ADDR ?= :8080 +export DATABASE_URL ?= postgres://nomadcode:nomadcode@localhost:5432/nomadcode?sslmode=disable +export GOOSE_DRIVER ?= postgres +export GOOSE_DBSTRING ?= $(DATABASE_URL) +export OUTPUT ?= .build/nomadcode-core + +.PHONY: run test build docker-up docker-down migrate-up sqlc + +run: + ./bin/run + +test: + ./bin/test + +build: + ./bin/build + +docker-up: + ./bin/docker-up + +docker-down: + ./bin/docker-down + +migrate-up: + ./bin/migrate-up + +sqlc: + ./bin/sqlc diff --git a/README.md b/README.md new file mode 100644 index 0000000..fd653c2 --- /dev/null +++ b/README.md @@ -0,0 +1,135 @@ +# NomadCode Core + +NomadCode Core는 사용자 요청을 작업 단위로 받고, 작업 상태를 저장하며, 비동기 Agent 작업 흐름을 관리하기 위한 서버입니다. + +초기 목표는 작업 생성과 조회, PostgreSQL 기반 작업 상태 저장, River 기반 비동기 작업 실행, Plane/Mattermost Adapter stub 구성, 이후 IOP / Agent Integrator와 연결 가능한 구조 확보입니다. + +## 현재 구현 범위 + +- Go HTTP Server +- `GET /healthz`, `GET /readyz` +- task 생성 / 조회 / 목록 / enqueue API +- PostgreSQL 연결 +- goose migration +- sqlc 기반 query 구조 +- River dummy job +- Plane Adapter stub +- Mattermost Adapter stub +- 선택적 Docker Compose 실행 환경 + +## 현재 구현하지 않은 범위 + +- 실제 Plane API 호출 +- 실제 Mattermost 메시지 발송 +- IOP 연동 +- Agent Integrator 연동 +- Outline / Forgejo / Nextcloud 연동 +- MCP 서버 +- Web Agent UI +- Flutter 앱 +- 복잡한 권한 정책 +- 복잡한 workflow DSL + +## 단계별 다음 작업 + +### Phase 1. Server Skeleton + +현재 스캐폴드 단계입니다. + +목표: +- 서버 실행 골격 구성 +- task 저장 구조 구성 +- dummy 비동기 job 실행 +- Adapter stub 생성 + +### Phase 2. Workflow Core + +목표: +- task 상태 전이 정리 +- enqueue / running / completed / failed 흐름 안정화 +- retry / timeout 기본 구조 추가 +- notification event 구조 정리 + +### Phase 3. External Integration + +목표: +- Plane issue 생성 / comment / status update 구현 +- Mattermost 메시지 발송 구현 +- Agent Integrator 호출 구조 추가 +- IOP 호출 구조 추가 + +## 실행 방법 + +로컬 실행은 호스트에 Go와 PostgreSQL이 설치되어 있다는 전제로 진행합니다. 기본 `DATABASE_URL`은 `postgres://nomadcode:nomadcode@localhost:5432/nomadcode?sslmode=disable` 입니다. + +PostgreSQL 사용자와 DB 생성 예시: + +```bash +psql postgres -c "CREATE USER nomadcode WITH PASSWORD 'nomadcode';" +psql postgres -c "CREATE DATABASE nomadcode OWNER nomadcode;" +``` + +Migration 실행: + +```bash +./bin/migrate-up +``` + +서버 실행: + +```bash +./bin/run +``` + +다른 DB 주소를 사용할 경우: + +```bash +DATABASE_URL="postgres://user:password@localhost:5432/dbname?sslmode=disable" ./bin/migrate-up +DATABASE_URL="postgres://user:password@localhost:5432/dbname?sslmode=disable" ./bin/run +``` + +테스트 실행: + +```bash +./bin/test +``` + +sqlc 실행: + +```bash +./bin/sqlc +``` + +Docker Compose 실행은 통합 확인이 필요할 때 선택적으로 사용합니다. + +```bash +./bin/docker-up +``` + +Makefile은 같은 명령을 감싸는 얇은 alias입니다. + +```bash +make migrate-up +make run +make test +make sqlc +``` + +API 테스트: + +```bash +curl localhost:8080/healthz +curl localhost:8080/readyz +``` + +```bash +curl -X POST localhost:8080/api/tasks \ + -H 'Content-Type: application/json' \ + -d '{"title":"README 수정 작업","source":"manual","payload":{"message":"README 초안을 정리해줘"}}' +``` + +```bash +curl localhost:8080/api/tasks +curl localhost:8080/api/tasks/{id} +curl -X POST localhost:8080/api/tasks/{id}/enqueue +``` diff --git a/bin/build b/bin/build new file mode 100755 index 0000000..6346104 --- /dev/null +++ b/bin/build @@ -0,0 +1,10 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT_DIR" + +OUTPUT="${OUTPUT:-.build/nomadcode-core}" +mkdir -p "$(dirname "$OUTPUT")" + +exec go build -o "$OUTPUT" ./cmd/server diff --git a/bin/docker-down b/bin/docker-down new file mode 100755 index 0000000..2ce9574 --- /dev/null +++ b/bin/docker-down @@ -0,0 +1,7 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT_DIR" + +exec docker compose down "$@" diff --git a/bin/docker-up b/bin/docker-up new file mode 100755 index 0000000..180fa17 --- /dev/null +++ b/bin/docker-up @@ -0,0 +1,7 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT_DIR" + +exec docker compose up --build "$@" diff --git a/bin/migrate-up b/bin/migrate-up new file mode 100755 index 0000000..a7a6ad8 --- /dev/null +++ b/bin/migrate-up @@ -0,0 +1,17 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT_DIR" + +GOOSE_DRIVER="${GOOSE_DRIVER:-postgres}" +DATABASE_URL="${DATABASE_URL:-postgres://nomadcode:nomadcode@localhost:5432/nomadcode?sslmode=disable}" +GOOSE_DBSTRING="${GOOSE_DBSTRING:-$DATABASE_URL}" + +if [[ -n "${GOOSE_BIN:-}" ]]; then + goose_cmd=("$GOOSE_BIN") +else + goose_cmd=(go run github.com/pressly/goose/v3/cmd/goose@v3.27.1) +fi + +exec "${goose_cmd[@]}" -dir migrations "$GOOSE_DRIVER" "$GOOSE_DBSTRING" up diff --git a/bin/run b/bin/run new file mode 100755 index 0000000..d024278 --- /dev/null +++ b/bin/run @@ -0,0 +1,11 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT_DIR" + +export APP_ENV="${APP_ENV:-local}" +export HTTP_ADDR="${HTTP_ADDR:-:8080}" +export DATABASE_URL="${DATABASE_URL:-postgres://nomadcode:nomadcode@localhost:5432/nomadcode?sslmode=disable}" + +exec go run ./cmd/server "$@" diff --git a/bin/sqlc b/bin/sqlc new file mode 100755 index 0000000..6c62921 --- /dev/null +++ b/bin/sqlc @@ -0,0 +1,13 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT_DIR" + +if [[ -n "${SQLC_BIN:-}" ]]; then + sqlc_cmd=("$SQLC_BIN") +else + sqlc_cmd=(go run github.com/sqlc-dev/sqlc/cmd/sqlc@v1.31.1) +fi + +exec "${sqlc_cmd[@]}" generate "$@" diff --git a/bin/test b/bin/test new file mode 100755 index 0000000..38a515f --- /dev/null +++ b/bin/test @@ -0,0 +1,7 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +cd "$ROOT_DIR" + +exec go test ./... "$@" diff --git a/cmd/server/main.go b/cmd/server/main.go new file mode 100644 index 0000000..8c102f1 --- /dev/null +++ b/cmd/server/main.go @@ -0,0 +1,98 @@ +package main + +import ( + "context" + "errors" + "log/slog" + "net/http" + "os" + "os/signal" + "syscall" + "time" + + "github.com/nomadcode/nomadcode-core/internal/adapters/mattermost" + "github.com/nomadcode/nomadcode-core/internal/config" + "github.com/nomadcode/nomadcode-core/internal/db" + apphttp "github.com/nomadcode/nomadcode-core/internal/http" + "github.com/nomadcode/nomadcode-core/internal/notification" + "github.com/nomadcode/nomadcode-core/internal/scheduler" + "github.com/nomadcode/nomadcode-core/internal/storage" + "github.com/nomadcode/nomadcode-core/internal/workflow" +) + +func main() { + logger := slog.New(slog.NewJSONHandler(os.Stdout, nil)) + if err := run(logger); err != nil { + logger.Error("server stopped", "error", err) + os.Exit(1) + } +} + +func run(logger *slog.Logger) error { + cfg := config.Load() + + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + + pool, err := db.NewPool(ctx, cfg.DatabaseURL) + if err != nil { + return err + } + defer pool.Close() + + store := storage.NewStore(pool) + mattermostClient := mattermost.NewClient(mattermost.Config{ + BaseURL: cfg.MattermostBaseURL, + Token: cfg.MattermostToken, + }, logger) + notificationService := notification.NewService(mattermostClient, logger) + + taskScheduler, err := scheduler.New(pool, store, notificationService, logger) + if err != nil { + return err + } + if err := taskScheduler.Migrate(ctx); err != nil { + return err + } + if err := taskScheduler.Start(ctx); err != nil { + return err + } + + workflowService := workflow.NewService(store, taskScheduler, logger) + handler := apphttp.NewHandler(pool, workflowService, logger) + server := &http.Server{ + Addr: cfg.HTTPAddr, + Handler: apphttp.NewRouter(handler, logger), + ReadHeaderTimeout: 5 * time.Second, + } + + serverErr := make(chan error, 1) + go func() { + logger.Info("server listening", "addr", cfg.HTTPAddr, "env", cfg.AppEnv) + if err := server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { + serverErr <- err + return + } + serverErr <- nil + }() + + select { + case <-ctx.Done(): + case err := <-serverErr: + if err != nil { + return err + } + } + + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + if err := server.Shutdown(shutdownCtx); err != nil { + return err + } + if err := taskScheduler.Stop(shutdownCtx); err != nil { + return err + } + + return nil +} diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..094239b --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,35 @@ +services: + postgres: + image: postgres:17-alpine + environment: + POSTGRES_DB: nomadcode + POSTGRES_USER: nomadcode + POSTGRES_PASSWORD: nomadcode + ports: + - "5432:5432" + healthcheck: + test: ["CMD-SHELL", "pg_isready -U nomadcode -d nomadcode"] + interval: 5s + timeout: 5s + retries: 10 + volumes: + - postgres-data:/var/lib/postgresql/data + + nomadcode-core: + build: . + environment: + APP_ENV: local + HTTP_ADDR: :8080 + DATABASE_URL: postgres://nomadcode:nomadcode@postgres:5432/nomadcode?sslmode=disable + MATTERMOST_BASE_URL: "" + MATTERMOST_TOKEN: "" + PLANE_BASE_URL: "" + PLANE_TOKEN: "" + depends_on: + postgres: + condition: service_healthy + ports: + - "8080:8080" + +volumes: + postgres-data: diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..20f7431 --- /dev/null +++ b/go.mod @@ -0,0 +1,30 @@ +module github.com/nomadcode/nomadcode-core + +go 1.25.0 + +require ( + github.com/go-chi/chi/v5 v5.2.5 + github.com/jackc/pgx/v5 v5.9.2 + github.com/riverqueue/river v0.37.1 + github.com/riverqueue/river/riverdriver/riverpgxv5 v0.37.1 +) + +require ( + github.com/davecgh/go-spew v1.1.1 // indirect + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + github.com/riverqueue/river/riverdriver v0.37.1 // indirect + github.com/riverqueue/river/rivershared v0.37.1 // indirect + github.com/riverqueue/river/rivertype v0.37.1 // indirect + github.com/stretchr/testify v1.11.1 // indirect + github.com/tidwall/gjson v1.19.0 // indirect + github.com/tidwall/match v1.2.0 // indirect + github.com/tidwall/pretty v1.2.1 // indirect + github.com/tidwall/sjson v1.2.5 // indirect + go.uber.org/goleak v1.3.0 // indirect + golang.org/x/sync v0.20.0 // indirect + golang.org/x/text v0.37.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..9d7375f --- /dev/null +++ b/go.sum @@ -0,0 +1,63 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/go-chi/chi/v5 v5.2.5 h1:Eg4myHZBjyvJmAFjFvWgrqDTXFyOzjj7YIm3L3mu6Ug= +github.com/go-chi/chi/v5 v5.2.5/go.mod h1:X7Gx4mteadT3eDOMTsXzmI4/rwUpOwBHLpAfupzFJP0= +github.com/jackc/pgerrcode v0.0.0-20240316143900-6e2875d9b438 h1:Dj0L5fhJ9F82ZJyVOmBx6msDp/kfd1t9GRfny/mfJA0= +github.com/jackc/pgerrcode v0.0.0-20240316143900-6e2875d9b438/go.mod h1:a/s9Lp5W7n/DD0VrVoyJ00FbP2ytTPDVOivvn2bMlds= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.9.2 h1:3ZhOzMWnR4yJ+RW1XImIPsD1aNSz4T4fyP7zlQb56hw= +github.com/jackc/pgx/v5 v5.9.2/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= +github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/riverqueue/river v0.37.1 h1:lZgXooqtGulMrItWFzhEDLeSBWoiJ3ZmKofKkHgBOwo= +github.com/riverqueue/river v0.37.1/go.mod h1:+5tefuBbaBRWiIEWOIyQexwQwhPVX23tsWgomIAcP88= +github.com/riverqueue/river/riverdriver v0.37.1 h1:IXV2xdS+LvsBvrJK+/GOy9s/a8P+9JW41RYIV/Vdya0= +github.com/riverqueue/river/riverdriver v0.37.1/go.mod h1:e7x1Q9gUFpbtKJ6/K4qpGxPT1ayc9NGlknZbwl0Re14= +github.com/riverqueue/river/riverdriver/riverpgxv5 v0.37.1 h1:dja9aBIEZTRk+gX+38t6/h+hLdWex7TmJdxft0Y0EXU= +github.com/riverqueue/river/riverdriver/riverpgxv5 v0.37.1/go.mod h1:48ij+FKRcarxG2fDscCreFLUiDXs+u6Wxjx5/TntKKU= +github.com/riverqueue/river/rivershared v0.37.1 h1:lvsxy9RAU+f1gRQ8+B4lVylyb184z4/82ilyU0DrsnY= +github.com/riverqueue/river/rivershared v0.37.1/go.mod h1:NAJdSXrjUjqdudLGCEm5heD2YRccClI9I9Yq+F3fyQM= +github.com/riverqueue/river/rivertype v0.37.1 h1:XZlhRR+c4RIgTTR//mJUZGAR26FNAHlMaNERpgUjb0c= +github.com/riverqueue/river/rivertype v0.37.1/go.mod h1:D1Ad+EaZiaXbQbJcJcfeicXJMBKno0n6UcfKI5Q7DIQ= +github.com/robfig/cron/v3 v3.0.1 h1:WdRxkvbJztn8LMz/QEvLN5sBU+xKpSqwwUO1Pjr4qDs= +github.com/robfig/cron/v3 v3.0.1/go.mod h1:eQICP3HwyT7UooqI/z+Ov+PtYAWygg1TEWWzGIFLtro= +github.com/rogpeppe/go-internal v1.12.0 h1:exVL4IDcn6na9z1rAb56Vxr+CgyK3nn3O+epU5NdKM8= +github.com/rogpeppe/go-internal v1.12.0/go.mod h1:E+RYuTGaKKdloAfM02xzb0FW3Paa99yedzYV+kq4uf4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/tidwall/gjson v1.14.2/go.mod h1:/wbyibRr2FHMks5tjHJ5F8dMZh3AcwJEMf5vlfC0lxk= +github.com/tidwall/gjson v1.19.0 h1:xwxm7n691Uf3u5OFjzngavjGTh55KX5q/9w9xHW88JU= +github.com/tidwall/gjson v1.19.0/go.mod h1:V37/opeE/JbLUOfH0QTXiNez2l0RUjYUhpT4szFQAfc= +github.com/tidwall/match v1.1.1/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM= +github.com/tidwall/match v1.2.0 h1:0pt8FlkOwjN2fPt4bIl4BoNxb98gGHN2ObFEDkrfZnM= +github.com/tidwall/match v1.2.0/go.mod h1:eRSPERbgtNPcGhD8UCthc6PmLEQXEWd3PRB5JTxsfmM= +github.com/tidwall/pretty v1.2.0/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= +github.com/tidwall/pretty v1.2.1 h1:qjsOFOWWQl+N3RsoF5/ssm1pHmJJwhjlSbZ51I6wMl4= +github.com/tidwall/pretty v1.2.1/go.mod h1:ITEVvHYasfjBbM0u2Pg8T2nJnzm8xPwvNhhsoaGGjNU= +github.com/tidwall/sjson v1.2.5 h1:kLy8mja+1c9jlljvWTlSazM7cKDRfJuR/bOJhcY5NcY= +github.com/tidwall/sjson v1.2.5/go.mod h1:Fvgq9kS/6ociJEDnK0Fk1cpYF4FIW6ZF7LAe+6jwd28= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= +golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= +golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/goose.env.example b/goose.env.example new file mode 100644 index 0000000..b55e3db --- /dev/null +++ b/goose.env.example @@ -0,0 +1,2 @@ +GOOSE_DRIVER=postgres +GOOSE_DBSTRING=postgres://nomadcode:nomadcode@localhost:5432/nomadcode?sslmode=disable diff --git a/internal/adapters/mattermost/client.go b/internal/adapters/mattermost/client.go new file mode 100644 index 0000000..7005a1b --- /dev/null +++ b/internal/adapters/mattermost/client.go @@ -0,0 +1,54 @@ +package mattermost + +import ( + "context" + "log/slog" +) + +type Config struct { + BaseURL string + Token string +} + +type Client struct { + cfg Config + logger *slog.Logger +} + +type MessageInput struct { + ChannelID string + Text string +} + +type TaskNotificationInput struct { + TaskID string + Title string + Status string + Message string +} + +func NewClient(cfg Config, logger *slog.Logger) *Client { + return &Client{cfg: cfg, logger: logger} +} + +func (c *Client) SendMessage(ctx context.Context, input MessageInput) error { + _ = ctx + if c.logger != nil { + c.logger.Info("mattermost send message skipped", "channel_id", input.ChannelID, "text", input.Text) + } + return nil +} + +func (c *Client) SendTaskNotification(ctx context.Context, input TaskNotificationInput) error { + _ = ctx + if c.logger != nil { + c.logger.Info( + "mattermost task notification skipped", + "task_id", input.TaskID, + "title", input.Title, + "status", input.Status, + "message", input.Message, + ) + } + return nil +} diff --git a/internal/adapters/plane/client.go b/internal/adapters/plane/client.go new file mode 100644 index 0000000..cb66721 --- /dev/null +++ b/internal/adapters/plane/client.go @@ -0,0 +1,60 @@ +package plane + +import ( + "context" + "errors" + "log/slog" +) + +var ErrNotImplemented = errors.New("plane adapter not implemented") + +type Config struct { + BaseURL string + Token string +} + +type Client struct { + cfg Config + logger *slog.Logger +} + +type CreateIssueInput struct { + Title string + Description string +} + +type AddCommentInput struct { + IssueID string + Body string +} + +type UpdateIssueStatusInput struct { + IssueID string + Status string +} + +func NewClient(cfg Config, logger *slog.Logger) *Client { + return &Client{cfg: cfg, logger: logger} +} + +func (c *Client) CreateIssue(ctx context.Context, input CreateIssueInput) error { + c.log(ctx, "plane create issue skipped", "title", input.Title) + return ErrNotImplemented +} + +func (c *Client) AddComment(ctx context.Context, input AddCommentInput) error { + c.log(ctx, "plane add comment skipped", "issue_id", input.IssueID) + return ErrNotImplemented +} + +func (c *Client) UpdateIssueStatus(ctx context.Context, input UpdateIssueStatusInput) error { + c.log(ctx, "plane update issue status skipped", "issue_id", input.IssueID, "status", input.Status) + return ErrNotImplemented +} + +func (c *Client) log(ctx context.Context, message string, args ...any) { + _ = ctx + if c.logger != nil { + c.logger.Info(message, args...) + } +} diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..5961d76 --- /dev/null +++ b/internal/config/config.go @@ -0,0 +1,33 @@ +package config + +import "os" + +type Config struct { + AppEnv string + HTTPAddr string + DatabaseURL string + MattermostBaseURL string + MattermostToken string + PlaneBaseURL string + PlaneToken string +} + +func Load() Config { + return Config{ + AppEnv: getEnv("APP_ENV", "local"), + HTTPAddr: getEnv("HTTP_ADDR", ":8080"), + DatabaseURL: os.Getenv("DATABASE_URL"), + MattermostBaseURL: os.Getenv("MATTERMOST_BASE_URL"), + MattermostToken: os.Getenv("MATTERMOST_TOKEN"), + PlaneBaseURL: os.Getenv("PLANE_BASE_URL"), + PlaneToken: os.Getenv("PLANE_TOKEN"), + } +} + +func getEnv(key, fallback string) string { + value := os.Getenv(key) + if value == "" { + return fallback + } + return value +} diff --git a/internal/db/db.go b/internal/db/db.go new file mode 100644 index 0000000..1a90c27 --- /dev/null +++ b/internal/db/db.go @@ -0,0 +1,30 @@ +package db + +import ( + "context" + "errors" + + "github.com/jackc/pgx/v5/pgxpool" +) + +func NewPool(ctx context.Context, databaseURL string) (*pgxpool.Pool, error) { + if databaseURL == "" { + return nil, errors.New("DATABASE_URL is required") + } + + cfg, err := pgxpool.ParseConfig(databaseURL) + if err != nil { + return nil, err + } + + pool, err := pgxpool.NewWithConfig(ctx, cfg) + if err != nil { + return nil, err + } + if err := pool.Ping(ctx); err != nil { + pool.Close() + return nil, err + } + + return pool, nil +} diff --git a/internal/db/queries.go b/internal/db/queries.go new file mode 100644 index 0000000..0f4b954 --- /dev/null +++ b/internal/db/queries.go @@ -0,0 +1,164 @@ +// Code generated placeholder for sqlc-style queries. DO NOT EDIT lightly. +package db + +import ( + "context" + "encoding/json" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" +) + +type DBTX interface { + Exec(context.Context, string, ...any) (pgconn.CommandTag, error) + Query(context.Context, string, ...any) (pgx.Rows, error) + QueryRow(context.Context, string, ...any) pgx.Row +} + +type Queries struct { + db DBTX +} + +func New(db DBTX) *Queries { + return &Queries{db: db} +} + +type Task struct { + ID string `json:"id"` + Title string `json:"title"` + Source string `json:"source"` + Status string `json:"status"` + Payload json.RawMessage `json:"payload"` + Result json.RawMessage `json:"result"` + Error *string `json:"error,omitempty"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +type CreateTaskParams struct { + Title string `json:"title"` + Source string `json:"source"` + Payload json.RawMessage `json:"payload"` +} + +type UpdateTaskStatusParams struct { + ID string `json:"id"` + Status string `json:"status"` +} + +type CompleteTaskParams struct { + ID string `json:"id"` + Result json.RawMessage `json:"result"` +} + +type FailTaskParams struct { + ID string `json:"id"` + Error string `json:"error"` +} + +const taskColumns = ` + id::text, + title, + source, + status, + payload, + result, + error, + created_at, + updated_at +` + +func (q *Queries) CreateTask(ctx context.Context, arg CreateTaskParams) (Task, error) { + row := q.db.QueryRow(ctx, ` + INSERT INTO tasks (title, source, payload) + VALUES ($1, $2, $3) + RETURNING `+taskColumns, arg.Title, arg.Source, arg.Payload) + return scanTask(row) +} + +func (q *Queries) GetTask(ctx context.Context, id string) (Task, error) { + row := q.db.QueryRow(ctx, ` + SELECT `+taskColumns+` + FROM tasks + WHERE id::text = $1 + `, id) + return scanTask(row) +} + +func (q *Queries) ListTasks(ctx context.Context, limit int32) ([]Task, error) { + rows, err := q.db.Query(ctx, ` + SELECT `+taskColumns+` + FROM tasks + ORDER BY created_at DESC + LIMIT $1 + `, limit) + if err != nil { + return nil, err + } + defer rows.Close() + + var tasks []Task + for rows.Next() { + task, err := scanTask(rows) + if err != nil { + return nil, err + } + tasks = append(tasks, task) + } + if err := rows.Err(); err != nil { + return nil, err + } + + return tasks, nil +} + +func (q *Queries) UpdateTaskStatus(ctx context.Context, arg UpdateTaskStatusParams) (Task, error) { + row := q.db.QueryRow(ctx, ` + UPDATE tasks + SET status = $2, updated_at = now() + WHERE id::text = $1 + RETURNING `+taskColumns, arg.ID, arg.Status) + return scanTask(row) +} + +func (q *Queries) CompleteTask(ctx context.Context, arg CompleteTaskParams) (Task, error) { + row := q.db.QueryRow(ctx, ` + UPDATE tasks + SET status = 'completed', result = $2, error = NULL, updated_at = now() + WHERE id::text = $1 + RETURNING `+taskColumns, arg.ID, arg.Result) + return scanTask(row) +} + +func (q *Queries) FailTask(ctx context.Context, arg FailTaskParams) (Task, error) { + row := q.db.QueryRow(ctx, ` + UPDATE tasks + SET status = 'failed', error = $2, updated_at = now() + WHERE id::text = $1 + RETURNING `+taskColumns, arg.ID, arg.Error) + return scanTask(row) +} + +type taskScanner interface { + Scan(dest ...any) error +} + +func scanTask(row taskScanner) (Task, error) { + var task Task + err := row.Scan( + &task.ID, + &task.Title, + &task.Source, + &task.Status, + &task.Payload, + &task.Result, + &task.Error, + &task.CreatedAt, + &task.UpdatedAt, + ) + if err != nil { + return Task{}, err + } + return task, nil +} diff --git a/internal/http/handlers.go b/internal/http/handlers.go new file mode 100644 index 0000000..d098fbc --- /dev/null +++ b/internal/http/handlers.go @@ -0,0 +1,137 @@ +package http + +import ( + "context" + "encoding/json" + "errors" + "log/slog" + stdhttp "net/http" + "strconv" + "time" + + "github.com/go-chi/chi/v5" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/nomadcode/nomadcode-core/internal/workflow" +) + +type Handler struct { + db *pgxpool.Pool + workflow *workflow.Service + logger *slog.Logger +} + +func NewHandler(pool *pgxpool.Pool, workflowService *workflow.Service, logger *slog.Logger) *Handler { + return &Handler{ + db: pool, + workflow: workflowService, + 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, + }) +} + +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 (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}) +} diff --git a/internal/http/middleware.go b/internal/http/middleware.go new file mode 100644 index 0000000..d08d153 --- /dev/null +++ b/internal/http/middleware.go @@ -0,0 +1,37 @@ +package http + +import ( + "log/slog" + stdhttp "net/http" + "time" +) + +type statusRecorder struct { + stdhttp.ResponseWriter + status int +} + +func (r *statusRecorder) WriteHeader(status int) { + r.status = status + r.ResponseWriter.WriteHeader(status) +} + +func loggingMiddleware(logger *slog.Logger) func(stdhttp.Handler) stdhttp.Handler { + return func(next stdhttp.Handler) stdhttp.Handler { + return stdhttp.HandlerFunc(func(w stdhttp.ResponseWriter, r *stdhttp.Request) { + start := time.Now() + recorder := &statusRecorder{ResponseWriter: w, status: stdhttp.StatusOK} + next.ServeHTTP(recorder, r) + + if logger != nil { + logger.Info( + "http request", + "method", r.Method, + "path", r.URL.Path, + "status", recorder.status, + "duration", time.Since(start).String(), + ) + } + }) + } +} diff --git a/internal/http/router.go b/internal/http/router.go new file mode 100644 index 0000000..8de899c --- /dev/null +++ b/internal/http/router.go @@ -0,0 +1,29 @@ +package http + +import ( + "log/slog" + stdhttp "net/http" + + "github.com/go-chi/chi/v5" + chimiddleware "github.com/go-chi/chi/v5/middleware" +) + +func NewRouter(handler *Handler, logger *slog.Logger) stdhttp.Handler { + r := chi.NewRouter() + r.Use(chimiddleware.RequestID) + r.Use(chimiddleware.RealIP) + r.Use(loggingMiddleware(logger)) + r.Use(chimiddleware.Recoverer) + + r.Get("/healthz", handler.Healthz) + r.Get("/readyz", handler.Readyz) + + r.Route("/api", func(r chi.Router) { + r.Post("/tasks", handler.CreateTask) + r.Get("/tasks", handler.ListTasks) + r.Get("/tasks/{id}", handler.GetTask) + r.Post("/tasks/{id}/enqueue", handler.EnqueueTask) + }) + + return r +} diff --git a/internal/notification/model.go b/internal/notification/model.go new file mode 100644 index 0000000..13051e7 --- /dev/null +++ b/internal/notification/model.go @@ -0,0 +1,8 @@ +package notification + +type TaskNotification struct { + TaskID string + Title string + Status string + Message string +} diff --git a/internal/notification/service.go b/internal/notification/service.go new file mode 100644 index 0000000..690d39a --- /dev/null +++ b/internal/notification/service.go @@ -0,0 +1,36 @@ +package notification + +import ( + "context" + "log/slog" + + "github.com/nomadcode/nomadcode-core/internal/adapters/mattermost" +) + +type Service struct { + mattermost *mattermost.Client + logger *slog.Logger +} + +func NewService(mattermostClient *mattermost.Client, logger *slog.Logger) *Service { + return &Service{ + mattermost: mattermostClient, + logger: logger, + } +} + +func (s *Service) NotifyTaskCompleted(ctx context.Context, input TaskNotification) error { + if s.logger != nil { + s.logger.Info("task completed notification requested", "task_id", input.TaskID) + } + if s.mattermost == nil { + return nil + } + + return s.mattermost.SendTaskNotification(ctx, mattermost.TaskNotificationInput{ + TaskID: input.TaskID, + Title: input.Title, + Status: input.Status, + Message: input.Message, + }) +} diff --git a/internal/scheduler/jobs.go b/internal/scheduler/jobs.go new file mode 100644 index 0000000..747be9c --- /dev/null +++ b/internal/scheduler/jobs.go @@ -0,0 +1,86 @@ +package scheduler + +import ( + "context" + "encoding/json" + "log/slog" + "time" + + "github.com/riverqueue/river" + + "github.com/nomadcode/nomadcode-core/internal/notification" + "github.com/nomadcode/nomadcode-core/internal/storage" + "github.com/nomadcode/nomadcode-core/internal/workflow" +) + +type TaskJobArgs struct { + TaskID string `json:"task_id"` +} + +func (TaskJobArgs) Kind() string { + return "task_run" +} + +type TaskWorker struct { + river.WorkerDefaults[TaskJobArgs] + + Store *storage.Store + Notifications *notification.Service + Logger *slog.Logger +} + +func (w *TaskWorker) Work(ctx context.Context, job *river.Job[TaskJobArgs]) error { + taskID := job.Args.TaskID + if w.Logger != nil { + w.Logger.Info("task job started", "task_id", taskID) + } + + task, err := w.Store.UpdateStatus(ctx, taskID, string(workflow.StatusRunning)) + if err != nil { + return err + } + + select { + case <-time.After(500 * time.Millisecond): + case <-ctx.Done(): + w.markFailed(taskID, ctx.Err()) + return ctx.Err() + } + + result := json.RawMessage(`{"message":"dummy task completed"}`) + task, err = w.Store.CompleteTask(ctx, taskID, result) + if err != nil { + w.markFailed(taskID, err) + return err + } + + if w.Notifications != nil { + err = w.Notifications.NotifyTaskCompleted(ctx, notification.TaskNotification{ + TaskID: task.ID, + Title: task.Title, + Status: task.Status, + Message: "dummy task completed", + }) + if err != nil && w.Logger != nil { + w.Logger.Warn("task notification failed", "task_id", taskID, "error", err) + } + } + + if w.Logger != nil { + w.Logger.Info("task job completed", "task_id", taskID) + } + return nil +} + +func (w *TaskWorker) markFailed(taskID string, err error) { + if err == nil || w.Store == nil { + return + } + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + if _, failErr := w.Store.FailTask(ctx, taskID, err.Error()); failErr != nil && w.Logger != nil { + w.Logger.Error("failed to mark task failed", "task_id", taskID, "error", failErr) + } +} diff --git a/internal/scheduler/river.go b/internal/scheduler/river.go new file mode 100644 index 0000000..c2d5d83 --- /dev/null +++ b/internal/scheduler/river.go @@ -0,0 +1,81 @@ +package scheduler + +import ( + "context" + "log/slog" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/riverqueue/river" + "github.com/riverqueue/river/riverdriver/riverpgxv5" + "github.com/riverqueue/river/rivermigrate" + + "github.com/nomadcode/nomadcode-core/internal/notification" + "github.com/nomadcode/nomadcode-core/internal/storage" +) + +type Client struct { + river *river.Client[pgx.Tx] + driver *riverpgxv5.Driver + logger *slog.Logger +} + +func New(pool *pgxpool.Pool, store *storage.Store, notifications *notification.Service, logger *slog.Logger) (*Client, error) { + workers := river.NewWorkers() + river.AddWorker(workers, &TaskWorker{ + Store: store, + Notifications: notifications, + Logger: logger, + }) + + driver := riverpgxv5.New(pool) + riverClient, err := river.NewClient[pgx.Tx](driver, &river.Config{ + Logger: logger, + Queues: map[string]river.QueueConfig{ + river.QueueDefault: {MaxWorkers: 2}, + }, + Workers: workers, + MaxAttempts: 3, + }) + if err != nil { + return nil, err + } + + return &Client{ + river: riverClient, + driver: driver, + logger: logger, + }, nil +} + +func (c *Client) Migrate(ctx context.Context) error { + migrator, err := rivermigrate.New[pgx.Tx](c.driver, &rivermigrate.Config{ + Logger: c.logger, + }) + if err != nil { + return err + } + + result, err := migrator.Migrate(ctx, rivermigrate.DirectionUp, nil) + if err != nil { + return err + } + if c.logger != nil && len(result.Versions) > 0 { + c.logger.Info("river migrations applied", "count", len(result.Versions)) + } + + return nil +} + +func (c *Client) Start(ctx context.Context) error { + return c.river.Start(ctx) +} + +func (c *Client) Stop(ctx context.Context) error { + return c.river.Stop(ctx) +} + +func (c *Client) EnqueueTask(ctx context.Context, taskID string) error { + _, err := c.river.Insert(ctx, TaskJobArgs{TaskID: taskID}, nil) + return err +} diff --git a/internal/storage/store.go b/internal/storage/store.go new file mode 100644 index 0000000..65453c6 --- /dev/null +++ b/internal/storage/store.go @@ -0,0 +1,56 @@ +package storage + +import ( + "context" + "encoding/json" + + "github.com/jackc/pgx/v5/pgxpool" + "github.com/nomadcode/nomadcode-core/internal/db" +) + +type Task = db.Task + +type Store struct { + queries *db.Queries +} + +func NewStore(pool *pgxpool.Pool) *Store { + return &Store{queries: db.New(pool)} +} + +func (s *Store) CreateTask(ctx context.Context, title, source string, payload json.RawMessage) (db.Task, error) { + return s.queries.CreateTask(ctx, db.CreateTaskParams{ + Title: title, + Source: source, + Payload: payload, + }) +} + +func (s *Store) GetTask(ctx context.Context, id string) (db.Task, error) { + return s.queries.GetTask(ctx, id) +} + +func (s *Store) ListTasks(ctx context.Context, limit int32) ([]db.Task, error) { + return s.queries.ListTasks(ctx, limit) +} + +func (s *Store) UpdateStatus(ctx context.Context, id, status string) (db.Task, error) { + return s.queries.UpdateTaskStatus(ctx, db.UpdateTaskStatusParams{ + ID: id, + Status: status, + }) +} + +func (s *Store) CompleteTask(ctx context.Context, id string, result json.RawMessage) (db.Task, error) { + return s.queries.CompleteTask(ctx, db.CompleteTaskParams{ + ID: id, + Result: result, + }) +} + +func (s *Store) FailTask(ctx context.Context, id, message string) (db.Task, error) { + return s.queries.FailTask(ctx, db.FailTaskParams{ + ID: id, + Error: message, + }) +} diff --git a/internal/workflow/model.go b/internal/workflow/model.go new file mode 100644 index 0000000..759f4f1 --- /dev/null +++ b/internal/workflow/model.go @@ -0,0 +1,20 @@ +package workflow + +import "encoding/json" + +type TaskStatus string + +const ( + StatusPending TaskStatus = "pending" + StatusQueued TaskStatus = "queued" + StatusRunning TaskStatus = "running" + StatusCompleted TaskStatus = "completed" + StatusFailed TaskStatus = "failed" + StatusCanceled TaskStatus = "canceled" +) + +type CreateTaskInput struct { + Title string `json:"title"` + Source string `json:"source"` + Payload json.RawMessage `json:"payload"` +} diff --git a/internal/workflow/service.go b/internal/workflow/service.go new file mode 100644 index 0000000..5ae3133 --- /dev/null +++ b/internal/workflow/service.go @@ -0,0 +1,102 @@ +package workflow + +import ( + "context" + "encoding/json" + "errors" + "log/slog" + "strings" + + "github.com/nomadcode/nomadcode-core/internal/storage" +) + +var ( + ErrInvalidTaskInput = errors.New("invalid task input") + ErrTaskCannotBeEnqueued = errors.New("task cannot be enqueued in current status") +) + +type TaskEnqueuer interface { + EnqueueTask(ctx context.Context, taskID string) error +} + +type Service struct { + store *storage.Store + enqueuer TaskEnqueuer + logger *slog.Logger +} + +func NewService(store *storage.Store, enqueuer TaskEnqueuer, logger *slog.Logger) *Service { + return &Service{ + store: store, + enqueuer: enqueuer, + logger: logger, + } +} + +func (s *Service) CreateTask(ctx context.Context, input CreateTaskInput) (storage.Task, error) { + title := strings.TrimSpace(input.Title) + source := strings.TrimSpace(input.Source) + if title == "" || source == "" { + return storage.Task{}, ErrInvalidTaskInput + } + + payload := input.Payload + if len(payload) == 0 || string(payload) == "null" { + payload = json.RawMessage(`{}`) + } + if !json.Valid(payload) { + return storage.Task{}, ErrInvalidTaskInput + } + + return s.store.CreateTask(ctx, title, source, payload) +} + +func (s *Service) GetTask(ctx context.Context, id string) (storage.Task, error) { + return s.store.GetTask(ctx, id) +} + +func (s *Service) ListTasks(ctx context.Context, limit int32) ([]storage.Task, error) { + if limit <= 0 { + limit = 20 + } + if limit > 100 { + limit = 100 + } + return s.store.ListTasks(ctx, limit) +} + +func (s *Service) EnqueueTask(ctx context.Context, id string) (storage.Task, error) { + task, err := s.store.GetTask(ctx, id) + if err != nil { + return storage.Task{}, err + } + if !canEnqueue(task.Status) { + return storage.Task{}, ErrTaskCannotBeEnqueued + } + + queuedTask, err := s.store.UpdateStatus(ctx, id, string(StatusQueued)) + if err != nil { + return storage.Task{}, err + } + + if s.enqueuer == nil { + return queuedTask, errors.New("task enqueuer is not configured") + } + if err := s.enqueuer.EnqueueTask(ctx, id); err != nil { + if _, failErr := s.store.FailTask(ctx, id, err.Error()); failErr != nil && s.logger != nil { + s.logger.Error("failed to mark task enqueue error", "task_id", id, "error", failErr) + } + return queuedTask, err + } + + return queuedTask, nil +} + +func canEnqueue(status string) bool { + switch TaskStatus(status) { + case StatusPending, StatusFailed: + return true + default: + return false + } +} diff --git a/migrations/00001_create_tasks.sql b/migrations/00001_create_tasks.sql new file mode 100644 index 0000000..4bd9dbb --- /dev/null +++ b/migrations/00001_create_tasks.sql @@ -0,0 +1,23 @@ +-- +goose Up +CREATE EXTENSION IF NOT EXISTS pgcrypto; + +CREATE TABLE IF NOT EXISTS tasks ( + id uuid PRIMARY KEY DEFAULT gen_random_uuid(), + title text NOT NULL, + source text NOT NULL, + status text NOT NULL DEFAULT 'pending', + payload jsonb NOT NULL DEFAULT '{}', + result jsonb NOT NULL DEFAULT '{}', + error text NULL, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + CONSTRAINT tasks_status_check CHECK ( + status IN ('pending', 'queued', 'running', 'completed', 'failed', 'canceled') + ) +); + +CREATE INDEX IF NOT EXISTS tasks_status_idx ON tasks (status); +CREATE INDEX IF NOT EXISTS tasks_created_at_idx ON tasks (created_at DESC); + +-- +goose Down +DROP TABLE IF EXISTS tasks; diff --git a/queries/tasks.sql b/queries/tasks.sql new file mode 100644 index 0000000..1265254 --- /dev/null +++ b/queries/tasks.sql @@ -0,0 +1,33 @@ +-- name: CreateTask :one +INSERT INTO tasks (title, source, payload) +VALUES ($1, $2, $3) +RETURNING *; + +-- name: GetTask :one +SELECT * +FROM tasks +WHERE id::text = $1; + +-- name: ListTasks :many +SELECT * +FROM tasks +ORDER BY created_at DESC +LIMIT $1; + +-- name: UpdateTaskStatus :one +UPDATE tasks +SET status = $2, updated_at = now() +WHERE id::text = $1 +RETURNING *; + +-- name: CompleteTask :one +UPDATE tasks +SET status = 'completed', result = $2, error = NULL, updated_at = now() +WHERE id::text = $1 +RETURNING *; + +-- name: FailTask :one +UPDATE tasks +SET status = 'failed', error = $2, updated_at = now() +WHERE id::text = $1 +RETURNING *; diff --git a/sqlc.yaml b/sqlc.yaml new file mode 100644 index 0000000..8c2acf0 --- /dev/null +++ b/sqlc.yaml @@ -0,0 +1,18 @@ +version: "2" +sql: + - engine: postgresql + schema: migrations + queries: queries + gen: + go: + package: generated + out: internal/db/generated + sql_package: pgx/v5 + emit_json_tags: true + overrides: + - db_type: uuid + go_type: string + - db_type: jsonb + go_type: + import: encoding/json + type: RawMessage