diff --git a/README.md b/README.md index fd653c2..8eb10f7 100644 --- a/README.md +++ b/README.md @@ -11,7 +11,7 @@ NomadCode Core는 사용자 요청을 작업 단위로 받고, 작업 상태를 - task 생성 / 조회 / 목록 / enqueue API - PostgreSQL 연결 - goose migration -- sqlc 기반 query 구조 +- sqlc 기반 DB query 생성 구조 - River dummy job - Plane Adapter stub - Mattermost Adapter stub @@ -100,6 +100,8 @@ sqlc 실행: ./bin/sqlc ``` +DB 접근 코드는 `migrations/`의 스키마와 `queries/`의 SQL을 기준으로 `internal/db/`에 생성합니다. `internal/db/db.go`, `internal/db/models.go`, `internal/db/tasks.sql.go`는 생성 파일이므로 직접 수정하지 않습니다. + Docker Compose 실행은 통합 확인이 필요할 때 선택적으로 사용합니다. ```bash diff --git a/internal/db/db.go b/internal/db/db.go index 1a90c27..468d1fa 100644 --- a/internal/db/db.go +++ b/internal/db/db.go @@ -1,30 +1,32 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.31.1 + package db import ( "context" - "errors" - "github.com/jackc/pgx/v5/pgxpool" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" ) -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 +type DBTX interface { + Exec(context.Context, string, ...interface{}) (pgconn.CommandTag, error) + Query(context.Context, string, ...interface{}) (pgx.Rows, error) + QueryRow(context.Context, string, ...interface{}) pgx.Row +} + +func New(db DBTX) *Queries { + return &Queries{db: db} +} + +type Queries struct { + db DBTX +} + +func (q *Queries) WithTx(tx pgx.Tx) *Queries { + return &Queries{ + db: tx, + } } diff --git a/internal/db/models.go b/internal/db/models.go new file mode 100644 index 0000000..777c6aa --- /dev/null +++ b/internal/db/models.go @@ -0,0 +1,22 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.31.1 + +package db + +import ( + "encoding/json" + "time" +) + +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"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} diff --git a/internal/db/pool.go b/internal/db/pool.go new file mode 100644 index 0000000..1a90c27 --- /dev/null +++ b/internal/db/pool.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 deleted file mode 100644 index 0f4b954..0000000 --- a/internal/db/queries.go +++ /dev/null @@ -1,164 +0,0 @@ -// 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/db/tasks.sql.go b/internal/db/tasks.sql.go new file mode 100644 index 0000000..4c8728d --- /dev/null +++ b/internal/db/tasks.sql.go @@ -0,0 +1,187 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.31.1 +// source: tasks.sql + +package db + +import ( + "context" + "encoding/json" +) + +const completeTask = `-- name: CompleteTask :one +UPDATE tasks +SET status = 'completed', result = $2, error = NULL, updated_at = now() +WHERE id::text = $1 +RETURNING id, title, source, status, payload, result, error, created_at, updated_at +` + +type CompleteTaskParams struct { + ID string `json:"id"` + Result json.RawMessage `json:"result"` +} + +func (q *Queries) CompleteTask(ctx context.Context, arg CompleteTaskParams) (Task, error) { + row := q.db.QueryRow(ctx, completeTask, arg.ID, arg.Result) + var i Task + err := row.Scan( + &i.ID, + &i.Title, + &i.Source, + &i.Status, + &i.Payload, + &i.Result, + &i.Error, + &i.CreatedAt, + &i.UpdatedAt, + ) + return i, err +} + +const createTask = `-- name: CreateTask :one +INSERT INTO tasks (title, source, payload) +VALUES ($1, $2, $3) +RETURNING id, title, source, status, payload, result, error, created_at, updated_at +` + +type CreateTaskParams struct { + Title string `json:"title"` + Source string `json:"source"` + Payload json.RawMessage `json:"payload"` +} + +func (q *Queries) CreateTask(ctx context.Context, arg CreateTaskParams) (Task, error) { + row := q.db.QueryRow(ctx, createTask, arg.Title, arg.Source, arg.Payload) + var i Task + err := row.Scan( + &i.ID, + &i.Title, + &i.Source, + &i.Status, + &i.Payload, + &i.Result, + &i.Error, + &i.CreatedAt, + &i.UpdatedAt, + ) + return i, err +} + +const failTask = `-- name: FailTask :one +UPDATE tasks +SET status = 'failed', error = $2, updated_at = now() +WHERE id::text = $1 +RETURNING id, title, source, status, payload, result, error, created_at, updated_at +` + +type FailTaskParams struct { + ID string `json:"id"` + Error *string `json:"error"` +} + +func (q *Queries) FailTask(ctx context.Context, arg FailTaskParams) (Task, error) { + row := q.db.QueryRow(ctx, failTask, arg.ID, arg.Error) + var i Task + err := row.Scan( + &i.ID, + &i.Title, + &i.Source, + &i.Status, + &i.Payload, + &i.Result, + &i.Error, + &i.CreatedAt, + &i.UpdatedAt, + ) + return i, err +} + +const getTask = `-- name: GetTask :one +SELECT id, title, source, status, payload, result, error, created_at, updated_at +FROM tasks +WHERE id::text = $1 +` + +func (q *Queries) GetTask(ctx context.Context, id string) (Task, error) { + row := q.db.QueryRow(ctx, getTask, id) + var i Task + err := row.Scan( + &i.ID, + &i.Title, + &i.Source, + &i.Status, + &i.Payload, + &i.Result, + &i.Error, + &i.CreatedAt, + &i.UpdatedAt, + ) + return i, err +} + +const listTasks = `-- name: ListTasks :many +SELECT id, title, source, status, payload, result, error, created_at, updated_at +FROM tasks +ORDER BY created_at DESC +LIMIT $1 +` + +func (q *Queries) ListTasks(ctx context.Context, limit int32) ([]Task, error) { + rows, err := q.db.Query(ctx, listTasks, limit) + if err != nil { + return nil, err + } + defer rows.Close() + var items []Task + for rows.Next() { + var i Task + if err := rows.Scan( + &i.ID, + &i.Title, + &i.Source, + &i.Status, + &i.Payload, + &i.Result, + &i.Error, + &i.CreatedAt, + &i.UpdatedAt, + ); err != nil { + return nil, err + } + items = append(items, i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + +const updateTaskStatus = `-- name: UpdateTaskStatus :one +UPDATE tasks +SET status = $2, updated_at = now() +WHERE id::text = $1 +RETURNING id, title, source, status, payload, result, error, created_at, updated_at +` + +type UpdateTaskStatusParams struct { + ID string `json:"id"` + Status string `json:"status"` +} + +func (q *Queries) UpdateTaskStatus(ctx context.Context, arg UpdateTaskStatusParams) (Task, error) { + row := q.db.QueryRow(ctx, updateTaskStatus, arg.ID, arg.Status) + var i Task + err := row.Scan( + &i.ID, + &i.Title, + &i.Source, + &i.Status, + &i.Payload, + &i.Result, + &i.Error, + &i.CreatedAt, + &i.UpdatedAt, + ) + return i, err +} diff --git a/internal/storage/store.go b/internal/storage/store.go index 65453c6..857edc3 100644 --- a/internal/storage/store.go +++ b/internal/storage/store.go @@ -51,6 +51,6 @@ func (s *Store) CompleteTask(ctx context.Context, id string, result json.RawMess func (s *Store) FailTask(ctx context.Context, id, message string) (db.Task, error) { return s.queries.FailTask(ctx, db.FailTaskParams{ ID: id, - Error: message, + Error: &message, }) } diff --git a/queries/tasks.sql b/queries/tasks.sql index 1265254..2ef3f11 100644 --- a/queries/tasks.sql +++ b/queries/tasks.sql @@ -1,15 +1,15 @@ -- name: CreateTask :one INSERT INTO tasks (title, source, payload) VALUES ($1, $2, $3) -RETURNING *; +RETURNING id, title, source, status, payload, result, error, created_at, updated_at; -- name: GetTask :one -SELECT * +SELECT id, title, source, status, payload, result, error, created_at, updated_at FROM tasks WHERE id::text = $1; -- name: ListTasks :many -SELECT * +SELECT id, title, source, status, payload, result, error, created_at, updated_at FROM tasks ORDER BY created_at DESC LIMIT $1; @@ -18,16 +18,16 @@ LIMIT $1; UPDATE tasks SET status = $2, updated_at = now() WHERE id::text = $1 -RETURNING *; +RETURNING id, title, source, status, payload, result, error, created_at, updated_at; -- name: CompleteTask :one UPDATE tasks SET status = 'completed', result = $2, error = NULL, updated_at = now() WHERE id::text = $1 -RETURNING *; +RETURNING id, title, source, status, payload, result, error, created_at, updated_at; -- name: FailTask :one UPDATE tasks SET status = 'failed', error = $2, updated_at = now() WHERE id::text = $1 -RETURNING *; +RETURNING id, title, source, status, payload, result, error, created_at, updated_at; diff --git a/sqlc.yaml b/sqlc.yaml index 8c2acf0..72bffa9 100644 --- a/sqlc.yaml +++ b/sqlc.yaml @@ -5,10 +5,11 @@ sql: queries: queries gen: go: - package: generated - out: internal/db/generated + package: db + out: internal/db sql_package: pgx/v5 emit_json_tags: true + emit_pointers_for_null_types: true overrides: - db_type: uuid go_type: string @@ -16,3 +17,7 @@ sql: go_type: import: encoding/json type: RawMessage + - db_type: timestamptz + go_type: + import: time + type: Time