feat(backend): implement task creation and admin web

This commit is contained in:
QiuSW
2026-07-26 14:03:32 +08:00
parent 2b265c92fc
commit c5d3b215ff
58 changed files with 8773 additions and 83 deletions
@@ -0,0 +1,100 @@
package sqlite
import (
"context"
"database/sql"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/usecase"
)
const assetUploadOperation = "UPLOAD_TASK_REFERENCE"
func (s *Store) CreateAssetIdempotent(
ctx context.Context,
candidate domain.Asset,
idempotencyKey string,
requestHash string,
) (domain.Asset, bool, error) {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return domain.Asset{}, false, repositoryFailure(err)
}
defer func() { _ = tx.Rollback() }()
existingHash, resourceID, found, err := lookupIdempotency(
ctx,
tx,
candidate.CreatorSubject,
assetUploadOperation,
idempotencyKey,
)
if err != nil {
return domain.Asset{}, false, err
}
if found {
if existingHash != requestHash {
return domain.Asset{}, false, usecase.ErrIdempotencyConflict
}
existing, err := getAssetByID(
ctx,
tx,
candidate.CreatorSubject,
resourceID,
)
if err != nil {
return domain.Asset{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.Asset{}, false, repositoryFailure(err)
}
return existing, false, nil
}
_, err = tx.ExecContext(
ctx,
`INSERT INTO assets (
id, creator_subject, purpose, media_type, size_bytes,
sha256, storage_key, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
candidate.ID,
candidate.CreatorSubject,
candidate.Purpose,
candidate.MediaType,
candidate.SizeBytes,
candidate.SHA256,
candidate.StorageKey,
formatTimestamp(candidate.CreatedAt),
)
if err != nil {
return domain.Asset{}, false, repositoryFailure(err)
}
if err := insertIdempotency(
ctx,
tx,
candidate.CreatorSubject,
assetUploadOperation,
idempotencyKey,
requestHash,
"ASSET",
candidate.ID,
candidate.CreatedAt,
); err != nil {
return domain.Asset{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.Asset{}, false, repositoryFailure(err)
}
return candidate, true, nil
}
func (s *Store) GetAsset(
ctx context.Context,
creatorSubject string,
assetID string,
) (domain.Asset, error) {
return getAssetByID(ctx, s.db, creatorSubject, assetID)
}
var _ usecase.AssetRepository = (*Store)(nil)
var _ queryRower = (*sql.Tx)(nil)
@@ -0,0 +1,270 @@
package sqlite
import (
"context"
"database/sql"
"errors"
"fmt"
"strings"
"time"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/usecase"
sqlite3 "github.com/mattn/go-sqlite3"
)
const timestampLayout = time.RFC3339Nano
type queryRower interface {
QueryRowContext(context.Context, string, ...any) *sql.Row
}
type rowScanner interface {
Scan(...any) error
}
func scanAsset(scanner rowScanner) (domain.Asset, error) {
var asset domain.Asset
var createdAt string
err := scanner.Scan(
&asset.ID,
&asset.CreatorSubject,
&asset.Purpose,
&asset.MediaType,
&asset.SizeBytes,
&asset.SHA256,
&asset.StorageKey,
&createdAt,
)
if err != nil {
return domain.Asset{}, err
}
asset.CreatedAt, err = parseTimestamp(createdAt)
if err != nil {
return domain.Asset{}, err
}
return asset, nil
}
func scanTask(scanner rowScanner) (domain.PurchaseTask, error) {
var task domain.PurchaseTask
var sourceRef sql.NullString
var maxBudget sql.NullInt64
var cancelReason sql.NullString
var canceledAt sql.NullString
var createdAt string
var updatedAt string
err := scanner.Scan(
&task.ID,
&task.CreatorSubject,
&sourceRef,
&task.Title,
&task.Description,
&task.SKU,
&task.ImageAssetID,
&task.Quantity,
&maxBudget,
&task.Currency,
&task.Status,
&task.Version,
&cancelReason,
&canceledAt,
&createdAt,
&updatedAt,
)
if err != nil {
return domain.PurchaseTask{}, err
}
if sourceRef.Valid {
task.SourceRef = &sourceRef.String
}
if maxBudget.Valid {
task.MaxBudgetCents = &maxBudget.Int64
}
if cancelReason.Valid {
task.CancelReason = &cancelReason.String
}
if canceledAt.Valid {
value, err := parseTimestamp(canceledAt.String)
if err != nil {
return domain.PurchaseTask{}, err
}
task.CanceledAt = &value
}
task.CreatedAt, err = parseTimestamp(createdAt)
if err != nil {
return domain.PurchaseTask{}, err
}
task.UpdatedAt, err = parseTimestamp(updatedAt)
if err != nil {
return domain.PurchaseTask{}, err
}
return task, nil
}
func getAssetByID(
ctx context.Context,
queryer queryRower,
creatorSubject string,
assetID string,
) (domain.Asset, error) {
asset, err := scanAsset(queryer.QueryRowContext(
ctx,
`SELECT
id, creator_subject, purpose, media_type, size_bytes,
sha256, storage_key, created_at
FROM assets
WHERE creator_subject = ? AND id = ?`,
creatorSubject,
assetID,
))
if errors.Is(err, sql.ErrNoRows) {
return domain.Asset{}, usecase.ErrRepositoryNotFound
}
if err != nil {
return domain.Asset{}, repositoryFailure(err)
}
return asset, nil
}
func getTaskByID(
ctx context.Context,
queryer queryRower,
creatorSubject string,
taskID string,
) (domain.PurchaseTask, error) {
task, err := scanTask(queryer.QueryRowContext(
ctx,
`SELECT
id, creator_subject, source_ref, title, description, sku,
image_asset_id, quantity, max_budget_cents, currency, status,
version, cancel_reason, canceled_at, created_at, updated_at
FROM purchase_tasks
WHERE creator_subject = ? AND id = ?`,
creatorSubject,
taskID,
))
if errors.Is(err, sql.ErrNoRows) {
return domain.PurchaseTask{}, usecase.ErrRepositoryNotFound
}
if err != nil {
return domain.PurchaseTask{}, repositoryFailure(err)
}
return task, nil
}
func lookupIdempotency(
ctx context.Context,
tx *sql.Tx,
creatorSubject string,
operation string,
idempotencyKey string,
) (requestHash string, resourceID string, found bool, err error) {
err = tx.QueryRowContext(
ctx,
`SELECT request_sha256, resource_id
FROM idempotency_records
WHERE creator_subject = ?
AND operation = ?
AND idempotency_key = ?`,
creatorSubject,
operation,
idempotencyKey,
).Scan(&requestHash, &resourceID)
if errors.Is(err, sql.ErrNoRows) {
return "", "", false, nil
}
if err != nil {
return "", "", false, repositoryFailure(err)
}
return requestHash, resourceID, true, nil
}
func insertIdempotency(
ctx context.Context,
tx *sql.Tx,
creatorSubject string,
operation string,
idempotencyKey string,
requestHash string,
resourceType string,
resourceID string,
createdAt time.Time,
) error {
_, err := tx.ExecContext(
ctx,
`INSERT INTO idempotency_records (
creator_subject, operation, idempotency_key, request_sha256,
resource_type, resource_id, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?)`,
creatorSubject,
operation,
idempotencyKey,
requestHash,
resourceType,
resourceID,
formatTimestamp(createdAt),
)
if err != nil {
return repositoryFailure(err)
}
return nil
}
func formatTimestamp(value time.Time) string {
return value.UTC().Format(timestampLayout)
}
func parseTimestamp(value string) (time.Time, error) {
parsed, err := time.Parse(timestampLayout, value)
if err != nil {
return time.Time{}, fmt.Errorf(
"%w: invalid stored timestamp",
usecase.ErrRepositoryInvariant,
)
}
return parsed.UTC(), nil
}
func nullableString(value *string) any {
if value == nil {
return nil
}
return *value
}
func nullableInt64(value *int64) any {
if value == nil {
return nil
}
return *value
}
func repositoryFailure(err error) error {
if err == nil {
return nil
}
var sqliteError sqlite3.Error
if errors.As(err, &sqliteError) {
switch sqliteError.Code {
case sqlite3.ErrBusy, sqlite3.ErrLocked, sqlite3.ErrIoErr,
sqlite3.ErrCantOpen, sqlite3.ErrFull:
return fmt.Errorf("%w", usecase.ErrRepositoryUnavailable)
}
}
if errors.Is(err, context.Canceled) ||
errors.Is(err, context.DeadlineExceeded) {
return fmt.Errorf("%w", usecase.ErrRepositoryUnavailable)
}
return fmt.Errorf("%w", usecase.ErrRepositoryInvariant)
}
func isUniqueConstraint(err error, fragment string) bool {
var sqliteError sqlite3.Error
if !errors.As(err, &sqliteError) ||
sqliteError.ExtendedCode != sqlite3.ErrConstraintUnique {
return false
}
return strings.Contains(err.Error(), fragment)
}
@@ -0,0 +1,17 @@
package sqlite
import (
"database/sql"
"errors"
)
type Store struct {
db *sql.DB
}
func New(db *sql.DB) (*Store, error) {
if db == nil {
return nil, errors.New("SQLite database is required")
}
return &Store{db: db}, nil
}
@@ -0,0 +1,436 @@
package sqlite_test
import (
"context"
"database/sql"
"errors"
"path/filepath"
"strings"
"testing"
"time"
"cmroubao/backend-api/internal/domain"
platformdatabase "cmroubao/backend-api/internal/platform/database"
"cmroubao/backend-api/internal/platform/migration"
repository "cmroubao/backend-api/internal/repository/sqlite"
"cmroubao/backend-api/internal/usecase"
)
func TestMigrationRejectsTaskFieldLimitViolations(t *testing.T) {
db := openDatabase(t)
ctx := context.Background()
asset := testAsset(1, time.Now().UTC())
_, err := db.ExecContext(
ctx,
`INSERT INTO assets (
id, creator_subject, purpose, media_type, size_bytes,
sha256, storage_key, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
asset.ID,
asset.CreatorSubject,
asset.Purpose,
asset.MediaType,
asset.SizeBytes,
asset.SHA256,
asset.StorageKey,
asset.CreatedAt.Format(time.RFC3339Nano),
)
if err != nil {
t.Fatalf("insert asset: %v", err)
}
tests := []struct {
name string
title string
description string
sku string
sourceRef string
}{
{
name: "title characters",
title: repeatText("a", domain.MaxTitleRunes+1),
description: "",
sku: "sku",
sourceRef: "source-1",
},
{
name: "SKU bytes",
title: "title",
description: "",
sku: repeatText("a", domain.MaxSKUBytes+1),
sourceRef: "source-1",
},
{
name: "description bytes",
title: "title",
description: repeatText("a", domain.MaxDescriptionBytes+1),
sku: "sku",
sourceRef: "source-1",
},
{
name: "source reference bytes",
title: "title",
description: "",
sku: "sku",
sourceRef: repeatText("a", domain.MaxSourceRefBytes+1),
},
}
for index, test := range tests {
t.Run(test.name, func(t *testing.T) {
_, err := db.ExecContext(
ctx,
`INSERT INTO purchase_tasks (
id, creator_subject, source_ref, title, description, sku,
image_asset_id, quantity, max_budget_cents, currency,
status, version, created_at, updated_at
) VALUES (?, 'local-admin', ?, ?, ?, ?, ?, 1, NULL, 'CNY',
'PENDING', 1, ?, ?)`,
uuid(500+index),
test.sourceRef,
test.title,
test.description,
test.sku,
asset.ID,
asset.CreatedAt.Format(time.RFC3339Nano),
asset.CreatedAt.Format(time.RFC3339Nano),
)
if err == nil {
t.Fatal("constraint violation error = nil")
}
})
}
}
func TestStoreAssetAndTaskLifecycleIsTransactionalAndIdempotent(
t *testing.T,
) {
store := openStore(t)
ctx := context.Background()
now := time.Date(2026, 7, 26, 2, 3, 4, 5, time.UTC)
asset := testAsset(1, now)
createdAsset, created, err := store.CreateAssetIdempotent(
ctx,
asset,
"asset-key-1",
repeatHex("1"),
)
if err != nil || !created || createdAsset.ID != asset.ID {
t.Fatalf(
"CreateAssetIdempotent() = %+v, %t, %v",
createdAsset,
created,
err,
)
}
replayedAsset, created, err := store.CreateAssetIdempotent(
ctx,
testAsset(2, now.Add(time.Second)),
"asset-key-1",
repeatHex("1"),
)
if err != nil || created || replayedAsset.ID != asset.ID {
t.Fatalf(
"asset replay = %+v, %t, %v",
replayedAsset,
created,
err,
)
}
_, _, err = store.CreateAssetIdempotent(
ctx,
testAsset(3, now.Add(2*time.Second)),
"asset-key-1",
repeatHex("2"),
)
if !errors.Is(err, usecase.ErrIdempotencyConflict) {
t.Fatalf("different asset replay error = %v", err)
}
task := testTask(1, asset.ID, "source-1", now)
event := testEvent(1, task.ID, "TASK_CREATED", now)
createdTask, created, err := store.CreateTaskIdempotent(
ctx,
task,
event,
"task-key-1",
repeatHex("3"),
)
if err != nil || !created || createdTask.Status != domain.TaskStatusPending {
t.Fatalf(
"CreateTaskIdempotent() = %+v, %t, %v",
createdTask,
created,
err,
)
}
replayedTask, created, err := store.CreateTaskIdempotent(
ctx,
testTask(2, asset.ID, "different", now.Add(time.Second)),
testEvent(2, uuid(2), "TASK_CREATED", now.Add(time.Second)),
"task-key-1",
repeatHex("3"),
)
if err != nil || created || replayedTask.ID != task.ID {
t.Fatalf(
"task replay = %+v, %t, %v",
replayedTask,
created,
err,
)
}
detail, err := store.GetTaskDetail(ctx, "local-admin", task.ID)
if err != nil {
t.Fatalf("GetTaskDetail() error = %v", err)
}
if detail.Asset.ID != asset.ID ||
len(detail.Events) != 1 ||
detail.Events[0].Type != "TASK_CREATED" {
t.Fatalf("detail = %+v", detail)
}
canceled, err := store.CancelPendingTask(
ctx,
"local-admin",
task.ID,
"no longer needed",
now.Add(time.Minute),
testEvent(3, task.ID, "TASK_CANCELED", now.Add(time.Minute)),
)
if err != nil {
t.Fatalf("CancelPendingTask() error = %v", err)
}
if canceled.Status != domain.TaskStatusCanceled ||
canceled.Version != 2 ||
canceled.CancelReason == nil ||
*canceled.CancelReason != "no longer needed" {
t.Fatalf("canceled task = %+v", canceled)
}
_, err = store.CancelPendingTask(
ctx,
"local-admin",
task.ID,
"again",
now.Add(2*time.Minute),
testEvent(4, task.ID, "TASK_CANCELED", now.Add(2*time.Minute)),
)
if !errors.Is(err, usecase.ErrTaskStateConflict) {
t.Fatalf("repeat cancel error = %v", err)
}
detail, err = store.GetTaskDetail(ctx, "local-admin", task.ID)
if err != nil || len(detail.Events) != 2 {
t.Fatalf("canceled detail events = %d, error = %v", len(detail.Events), err)
}
}
func TestStoreEnforcesAssetOwnershipSourceReferenceAndStableCursor(
t *testing.T,
) {
store := openStore(t)
ctx := context.Background()
now := time.Date(2026, 7, 26, 3, 4, 5, 0, time.UTC)
for index := 1; index <= 4; index++ {
asset := testAsset(index, now)
if _, _, err := store.CreateAssetIdempotent(
ctx,
asset,
"asset-key-"+string(rune('0'+index)),
repeatHex(string(rune('0'+index))),
); err != nil {
t.Fatalf("create asset %d: %v", index, err)
}
if index <= 3 {
task := testTask(index, asset.ID, "source-"+string(rune('0'+index)), now)
if _, _, err := store.CreateTaskIdempotent(
ctx,
task,
testEvent(index, task.ID, "TASK_CREATED", now),
"task-key-"+string(rune('0'+index)),
repeatHex(string(rune('4'+index))),
); err != nil {
t.Fatalf("create task %d: %v", index, err)
}
}
}
first, err := store.ListTasks(ctx, usecase.TaskListFilter{
CreatorSubject: "local-admin",
Query: "title",
Limit: 2,
})
if err != nil {
t.Fatalf("first ListTasks() error = %v", err)
}
if len(first) != 2 || first[0].ID <= first[1].ID {
t.Fatalf("first page order = %+v", first)
}
second, err := store.ListTasks(ctx, usecase.TaskListFilter{
CreatorSubject: "local-admin",
Limit: 2,
After: &usecase.TaskCursor{
CreatedAt: first[1].CreatedAt,
ID: first[1].ID,
},
})
if err != nil {
t.Fatalf("second ListTasks() error = %v", err)
}
if len(second) != 1 || second[0].ID == first[0].ID ||
second[0].ID == first[1].ID {
t.Fatalf("second page = %+v", second)
}
conflicting := testTask(4, testAsset(4, now).ID, "source-1", now)
_, _, err = store.CreateTaskIdempotent(
ctx,
conflicting,
testEvent(4, conflicting.ID, "TASK_CREATED", now),
"task-key-conflict",
repeatHex("a"),
)
if !errors.Is(err, usecase.ErrSourceReferenceConflict) {
t.Fatalf("source conflict error = %v", err)
}
otherAsset := testAsset(8, now)
otherAsset.CreatorSubject = "other-admin"
if _, _, err := store.CreateAssetIdempotent(
ctx,
otherAsset,
"other-asset",
repeatHex("b"),
); err != nil {
t.Fatalf("create other asset: %v", err)
}
foreignTask := testTask(8, otherAsset.ID, "foreign-source", now)
foreignTask.CreatorSubject = "local-admin"
_, _, err = store.CreateTaskIdempotent(
ctx,
foreignTask,
testEvent(8, foreignTask.ID, "TASK_CREATED", now),
"foreign-task",
repeatHex("c"),
)
if !errors.Is(err, usecase.ErrAssetUnavailable) {
t.Fatalf("foreign asset error = %v", err)
}
}
func openStore(t *testing.T) *repository.Store {
t.Helper()
db := openDatabase(t)
store, err := repository.New(db)
if err != nil {
t.Fatalf("repository.New() error = %v", err)
}
return store
}
func openDatabase(t *testing.T) *sql.DB {
t.Helper()
db, err := platformdatabase.Open(
context.Background(),
filepath.Join(t.TempDir(), "store.db"),
)
if err != nil {
t.Fatalf("database.Open() error = %v", err)
}
t.Cleanup(func() { _ = db.Close() })
runner, err := migration.New(db)
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if _, err := runner.Up(context.Background()); err != nil {
t.Fatalf("migration Up() error = %v", err)
}
return db
}
func testAsset(index int, createdAt time.Time) domain.Asset {
return domain.Asset{
ID: uuid(100 + index),
CreatorSubject: "local-admin",
Purpose: domain.AssetPurposeTaskReference,
MediaType: domain.NormalizedImageMediaType,
SizeBytes: int64(1000 + index),
SHA256: repeatHex("d"),
StorageKey: "aa/" + uuid(200+index) + ".jpg",
CreatedAt: createdAt,
}
}
func testTask(
index int,
assetID string,
source string,
createdAt time.Time,
) domain.PurchaseTask {
budget := int64(2000 + index)
return domain.PurchaseTask{
ID: uuid(index),
CreatorSubject: "local-admin",
SourceRef: &source,
Title: "title " + string(rune('0'+index)),
Description: "description",
SKU: "SKU-" + string(rune('0'+index)),
ImageAssetID: assetID,
Quantity: index,
MaxBudgetCents: &budget,
Currency: domain.CurrencyCNY,
Status: domain.TaskStatusPending,
Version: 1,
CreatedAt: createdAt,
UpdatedAt: createdAt,
}
}
func testEvent(
index int,
taskID string,
eventType string,
occurredAt time.Time,
) domain.TaskEvent {
return domain.TaskEvent{
ID: uuid(300 + index),
TaskID: taskID,
Type: eventType,
Message: "event",
OccurredAt: occurredAt,
}
}
func uuid(index int) string {
return "00000000-0000-4000-8000-" + twelveDigits(index)
}
func twelveDigits(value int) string {
result := "000000000000"
digits := ""
for value > 0 {
digits = string(rune('0'+value%10)) + digits
value /= 10
}
if digits == "" {
digits = "0"
}
return result[:len(result)-len(digits)] + digits
}
func repeatHex(value string) string {
result := ""
for len(result) < 64 {
result += value
}
return result[:64]
}
func repeatText(value string, count int) string {
var result strings.Builder
result.Grow(len(value) * count)
for index := 0; index < count; index++ {
result.WriteString(value)
}
return result.String()
}
@@ -0,0 +1,383 @@
package sqlite
import (
"context"
"database/sql"
"strings"
"time"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/usecase"
)
const createTaskOperation = "CREATE_PURCHASE_TASK"
func (s *Store) CreateTaskIdempotent(
ctx context.Context,
candidate domain.PurchaseTask,
event domain.TaskEvent,
idempotencyKey string,
requestHash string,
) (domain.PurchaseTask, bool, error) {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
defer func() { _ = tx.Rollback() }()
existingHash, resourceID, found, err := lookupIdempotency(
ctx,
tx,
candidate.CreatorSubject,
createTaskOperation,
idempotencyKey,
)
if err != nil {
return domain.PurchaseTask{}, false, err
}
if found {
if existingHash != requestHash {
return domain.PurchaseTask{}, false, usecase.ErrIdempotencyConflict
}
existing, err := getTaskByID(
ctx,
tx,
candidate.CreatorSubject,
resourceID,
)
if err != nil {
return domain.PurchaseTask{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
return existing, false, nil
}
var available int
err = tx.QueryRowContext(
ctx,
`SELECT EXISTS (
SELECT 1
FROM assets AS asset
WHERE asset.id = ?
AND asset.creator_subject = ?
AND asset.purpose = 'TASK_REFERENCE'
AND NOT EXISTS (
SELECT 1
FROM purchase_tasks AS task
WHERE task.image_asset_id = asset.id
)
)`,
candidate.ImageAssetID,
candidate.CreatorSubject,
).Scan(&available)
if err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
if available != 1 {
return domain.PurchaseTask{}, false, usecase.ErrAssetUnavailable
}
_, err = tx.ExecContext(
ctx,
`INSERT INTO purchase_tasks (
id, creator_subject, source_ref, title, description, sku,
image_asset_id, quantity, max_budget_cents, currency, status,
version, cancel_reason, canceled_at, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, NULL, ?, ?)`,
candidate.ID,
candidate.CreatorSubject,
nullableString(candidate.SourceRef),
candidate.Title,
candidate.Description,
candidate.SKU,
candidate.ImageAssetID,
candidate.Quantity,
nullableInt64(candidate.MaxBudgetCents),
candidate.Currency,
candidate.Status,
candidate.Version,
formatTimestamp(candidate.CreatedAt),
formatTimestamp(candidate.UpdatedAt),
)
if err != nil {
switch {
case isUniqueConstraint(
err,
"purchase_tasks.creator_subject, purchase_tasks.source_ref",
):
return domain.PurchaseTask{}, false, usecase.ErrSourceReferenceConflict
case isUniqueConstraint(err, "purchase_tasks.image_asset_id"):
return domain.PurchaseTask{}, false, usecase.ErrAssetUnavailable
default:
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
}
if err := insertTaskEvent(ctx, tx, event); err != nil {
return domain.PurchaseTask{}, false, err
}
if err := insertIdempotency(
ctx,
tx,
candidate.CreatorSubject,
createTaskOperation,
idempotencyKey,
requestHash,
"PURCHASE_TASK",
candidate.ID,
candidate.CreatedAt,
); err != nil {
return domain.PurchaseTask{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
return candidate, true, nil
}
func (s *Store) ListTasks(
ctx context.Context,
filter usecase.TaskListFilter,
) ([]domain.PurchaseTask, error) {
var query strings.Builder
query.WriteString(`SELECT
id, creator_subject, source_ref, title, description, sku,
image_asset_id, quantity, max_budget_cents, currency, status,
version, cancel_reason, canceled_at, created_at, updated_at
FROM purchase_tasks
WHERE creator_subject = ?`)
arguments := []any{filter.CreatorSubject}
if filter.Status != nil {
query.WriteString(" AND status = ?")
arguments = append(arguments, *filter.Status)
}
if filter.Query != "" {
query.WriteString(` AND (
id LIKE ? ESCAPE '\'
OR COALESCE(source_ref, '') LIKE ? ESCAPE '\'
OR title LIKE ? ESCAPE '\'
OR sku LIKE ? ESCAPE '\'
)`)
pattern := "%" + escapeLike(filter.Query) + "%"
arguments = append(
arguments,
pattern,
pattern,
pattern,
pattern,
)
}
if filter.CreatedFrom != nil {
query.WriteString(" AND created_at >= ?")
arguments = append(
arguments,
formatTimestamp(*filter.CreatedFrom),
)
}
if filter.CreatedTo != nil {
query.WriteString(" AND created_at <= ?")
arguments = append(
arguments,
formatTimestamp(*filter.CreatedTo),
)
}
if filter.After != nil {
query.WriteString(
" AND (created_at < ? OR (created_at = ? AND id < ?))",
)
createdAt := formatTimestamp(filter.After.CreatedAt)
arguments = append(
arguments,
createdAt,
createdAt,
filter.After.ID,
)
}
query.WriteString(" ORDER BY created_at DESC, id DESC LIMIT ?")
arguments = append(arguments, filter.Limit)
rows, err := s.db.QueryContext(ctx, query.String(), arguments...)
if err != nil {
return nil, repositoryFailure(err)
}
defer rows.Close()
tasks := make([]domain.PurchaseTask, 0)
for rows.Next() {
task, err := scanTask(rows)
if err != nil {
return nil, repositoryFailure(err)
}
tasks = append(tasks, task)
}
if err := rows.Err(); err != nil {
return nil, repositoryFailure(err)
}
return tasks, nil
}
func (s *Store) GetTaskDetail(
ctx context.Context,
creatorSubject string,
taskID string,
) (domain.TaskDetail, error) {
tx, err := s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
if err != nil {
return domain.TaskDetail{}, repositoryFailure(err)
}
defer func() { _ = tx.Rollback() }()
task, err := getTaskByID(ctx, tx, creatorSubject, taskID)
if err != nil {
return domain.TaskDetail{}, err
}
asset, err := getAssetByID(
ctx,
tx,
creatorSubject,
task.ImageAssetID,
)
if err != nil {
return domain.TaskDetail{}, err
}
rows, err := tx.QueryContext(
ctx,
`SELECT id, task_id, event_type, message, occurred_at
FROM task_events
WHERE task_id = ?
ORDER BY occurred_at ASC, id ASC`,
taskID,
)
if err != nil {
return domain.TaskDetail{}, repositoryFailure(err)
}
defer rows.Close()
events := make([]domain.TaskEvent, 0)
for rows.Next() {
var event domain.TaskEvent
var occurredAt string
if err := rows.Scan(
&event.ID,
&event.TaskID,
&event.Type,
&event.Message,
&occurredAt,
); err != nil {
return domain.TaskDetail{}, repositoryFailure(err)
}
event.OccurredAt, err = parseTimestamp(occurredAt)
if err != nil {
return domain.TaskDetail{}, err
}
events = append(events, event)
}
if err := rows.Err(); err != nil {
return domain.TaskDetail{}, repositoryFailure(err)
}
if err := rows.Close(); err != nil {
return domain.TaskDetail{}, repositoryFailure(err)
}
detail := domain.TaskDetail{
Task: task,
Asset: asset,
Events: events,
}
if err := tx.Commit(); err != nil {
return domain.TaskDetail{}, repositoryFailure(err)
}
return detail, nil
}
func (s *Store) CancelPendingTask(
ctx context.Context,
creatorSubject string,
taskID string,
reason string,
canceledAt time.Time,
event domain.TaskEvent,
) (domain.PurchaseTask, error) {
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return domain.PurchaseTask{}, repositoryFailure(err)
}
defer func() { _ = tx.Rollback() }()
task, err := getTaskByID(ctx, tx, creatorSubject, taskID)
if err != nil {
return domain.PurchaseTask{}, err
}
if !domain.CanCancel(task.Status) {
return domain.PurchaseTask{}, usecase.ErrTaskStateConflict
}
result, err := tx.ExecContext(
ctx,
`UPDATE purchase_tasks
SET status = 'CANCELED',
version = version + 1,
cancel_reason = NULLIF(?, ''),
canceled_at = ?,
updated_at = ?
WHERE id = ?
AND creator_subject = ?
AND status = 'PENDING'
AND version = ?`,
reason,
formatTimestamp(canceledAt),
formatTimestamp(canceledAt),
taskID,
creatorSubject,
task.Version,
)
if err != nil {
return domain.PurchaseTask{}, repositoryFailure(err)
}
affected, err := result.RowsAffected()
if err != nil {
return domain.PurchaseTask{}, repositoryFailure(err)
}
if affected != 1 {
return domain.PurchaseTask{}, usecase.ErrTaskStateConflict
}
if err := insertTaskEvent(ctx, tx, event); err != nil {
return domain.PurchaseTask{}, err
}
updated, err := getTaskByID(ctx, tx, creatorSubject, taskID)
if err != nil {
return domain.PurchaseTask{}, err
}
if err := tx.Commit(); err != nil {
return domain.PurchaseTask{}, repositoryFailure(err)
}
return updated, nil
}
func insertTaskEvent(
ctx context.Context,
tx *sql.Tx,
event domain.TaskEvent,
) error {
_, err := tx.ExecContext(
ctx,
`INSERT INTO task_events (
id, task_id, event_type, message, occurred_at
) VALUES (?, ?, ?, ?, ?)`,
event.ID,
event.TaskID,
event.Type,
event.Message,
formatTimestamp(event.OccurredAt),
)
if err != nil {
return repositoryFailure(err)
}
return nil
}
func escapeLike(value string) string {
replacer := strings.NewReplacer(
`\`, `\\`,
`%`, `\%`,
`_`, `\_`,
)
return replacer.Replace(value)
}
var _ usecase.TaskRepository = (*Store)(nil)