Initial nomadcode-core project
This commit is contained in:
commit
e0d2b5b616
34 changed files with 1496 additions and 0 deletions
3
.gitignore
vendored
Normal file
3
.gitignore
vendored
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
/.build/
|
||||
/bin/nomadcode-core
|
||||
coverage.out
|
||||
22
Dockerfile
Normal file
22
Dockerfile
Normal file
|
|
@ -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"]
|
||||
29
Makefile
Normal file
29
Makefile
Normal file
|
|
@ -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
|
||||
135
README.md
Normal file
135
README.md
Normal file
|
|
@ -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
|
||||
```
|
||||
10
bin/build
Executable file
10
bin/build
Executable file
|
|
@ -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
|
||||
7
bin/docker-down
Executable file
7
bin/docker-down
Executable file
|
|
@ -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 "$@"
|
||||
7
bin/docker-up
Executable file
7
bin/docker-up
Executable file
|
|
@ -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 "$@"
|
||||
17
bin/migrate-up
Executable file
17
bin/migrate-up
Executable file
|
|
@ -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
|
||||
11
bin/run
Executable file
11
bin/run
Executable file
|
|
@ -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 "$@"
|
||||
13
bin/sqlc
Executable file
13
bin/sqlc
Executable file
|
|
@ -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 "$@"
|
||||
7
bin/test
Executable file
7
bin/test
Executable file
|
|
@ -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 ./... "$@"
|
||||
98
cmd/server/main.go
Normal file
98
cmd/server/main.go
Normal file
|
|
@ -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
|
||||
}
|
||||
35
docker-compose.yml
Normal file
35
docker-compose.yml
Normal file
|
|
@ -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:
|
||||
30
go.mod
Normal file
30
go.mod
Normal file
|
|
@ -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
|
||||
)
|
||||
63
go.sum
Normal file
63
go.sum
Normal file
|
|
@ -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=
|
||||
2
goose.env.example
Normal file
2
goose.env.example
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
GOOSE_DRIVER=postgres
|
||||
GOOSE_DBSTRING=postgres://nomadcode:nomadcode@localhost:5432/nomadcode?sslmode=disable
|
||||
54
internal/adapters/mattermost/client.go
Normal file
54
internal/adapters/mattermost/client.go
Normal file
|
|
@ -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
|
||||
}
|
||||
60
internal/adapters/plane/client.go
Normal file
60
internal/adapters/plane/client.go
Normal file
|
|
@ -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...)
|
||||
}
|
||||
}
|
||||
33
internal/config/config.go
Normal file
33
internal/config/config.go
Normal file
|
|
@ -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
|
||||
}
|
||||
30
internal/db/db.go
Normal file
30
internal/db/db.go
Normal file
|
|
@ -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
|
||||
}
|
||||
164
internal/db/queries.go
Normal file
164
internal/db/queries.go
Normal file
|
|
@ -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
|
||||
}
|
||||
137
internal/http/handlers.go
Normal file
137
internal/http/handlers.go
Normal file
|
|
@ -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})
|
||||
}
|
||||
37
internal/http/middleware.go
Normal file
37
internal/http/middleware.go
Normal file
|
|
@ -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(),
|
||||
)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
29
internal/http/router.go
Normal file
29
internal/http/router.go
Normal file
|
|
@ -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
|
||||
}
|
||||
8
internal/notification/model.go
Normal file
8
internal/notification/model.go
Normal file
|
|
@ -0,0 +1,8 @@
|
|||
package notification
|
||||
|
||||
type TaskNotification struct {
|
||||
TaskID string
|
||||
Title string
|
||||
Status string
|
||||
Message string
|
||||
}
|
||||
36
internal/notification/service.go
Normal file
36
internal/notification/service.go
Normal file
|
|
@ -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,
|
||||
})
|
||||
}
|
||||
86
internal/scheduler/jobs.go
Normal file
86
internal/scheduler/jobs.go
Normal file
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
81
internal/scheduler/river.go
Normal file
81
internal/scheduler/river.go
Normal file
|
|
@ -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
|
||||
}
|
||||
56
internal/storage/store.go
Normal file
56
internal/storage/store.go
Normal file
|
|
@ -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,
|
||||
})
|
||||
}
|
||||
20
internal/workflow/model.go
Normal file
20
internal/workflow/model.go
Normal file
|
|
@ -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"`
|
||||
}
|
||||
102
internal/workflow/service.go
Normal file
102
internal/workflow/service.go
Normal file
|
|
@ -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
|
||||
}
|
||||
}
|
||||
23
migrations/00001_create_tasks.sql
Normal file
23
migrations/00001_create_tasks.sql
Normal file
|
|
@ -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;
|
||||
33
queries/tasks.sql
Normal file
33
queries/tasks.sql
Normal file
|
|
@ -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 *;
|
||||
18
sqlc.yaml
Normal file
18
sqlc.yaml
Normal file
|
|
@ -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
|
||||
Loading…
Reference in a new issue