refactor: db layer separation and updates
This commit is contained in:
parent
e0d2b5b616
commit
5e58770c8a
9 changed files with 280 additions and 196 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
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
|
||||
}
|
||||
|
||||
cfg, err := pgxpool.ParseConfig(databaseURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
func New(db DBTX) *Queries {
|
||||
return &Queries{db: db}
|
||||
}
|
||||
|
||||
pool, err := pgxpool.NewWithConfig(ctx, cfg)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := pool.Ping(ctx); err != nil {
|
||||
pool.Close()
|
||||
return nil, err
|
||||
type Queries struct {
|
||||
db DBTX
|
||||
}
|
||||
|
||||
return pool, nil
|
||||
func (q *Queries) WithTx(tx pgx.Tx) *Queries {
|
||||
return &Queries{
|
||||
db: tx,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
22
internal/db/models.go
Normal file
22
internal/db/models.go
Normal file
|
|
@ -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"`
|
||||
}
|
||||
30
internal/db/pool.go
Normal file
30
internal/db/pool.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
|
||||
}
|
||||
|
|
@ -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
|
||||
}
|
||||
187
internal/db/tasks.sql.go
Normal file
187
internal/db/tasks.sql.go
Normal file
|
|
@ -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
|
||||
}
|
||||
|
|
@ -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,
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in a new issue