feat(t223): generate tasks from freight items

This commit is contained in:
QiuSW
2026-07-28 23:51:59 +08:00
parent 6d052b0462
commit 433db45100
25 changed files with 2697 additions and 48 deletions
@@ -383,6 +383,9 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks(
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("Down(v13) error = %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("Down(v12) error = %v", err)
}
@@ -426,9 +429,9 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks(
t.Fatal("purchase_tasks was lost during auth migration rollback")
}
if applied, err := runner.Up(context.Background()); err != nil {
t.Fatalf("Up(v3-v12) error = %v", err)
} else if applied != 10 {
t.Fatalf("Up(v3-v12) applied = %d, want 10", applied)
t.Fatalf("Up(v3-v13) error = %v", err)
} else if applied != 11 {
t.Fatalf("Up(v3-v13) applied = %d, want 11", applied)
}
}
@@ -117,6 +117,9 @@ func TestFreightImportIsAtomicIdempotentAndRevisioned(t *testing.T) {
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("procurement migration down: %v", err)
}
if err := runner.Down(ctx); err == nil {
t.Fatal("freight migration down succeeded with retained source data")
}
@@ -0,0 +1,632 @@
package sqlite
import (
"context"
"database/sql"
"errors"
"time"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/usecase"
)
const createProcurementTaskOperation = "CREATE_PROCUREMENT_TASK"
func (store *Store) GetProcurementSourceItem(
ctx context.Context,
creatorSubject, itemID string,
) (domain.ProcurementSourceItem, error) {
item, err := scanFreightOrderItem(store.db.QueryRowContext(
ctx,
`SELECT
item.id, item.freight_order_id, item.external_item_id,
item.title, item.product_spec, item.sku, item.quantity,
item.product_thumb_ref, item.purchase_status,
item.canonical_sha256, item.revision, item.is_present,
item.first_sync_run_id, item.last_sync_run_id,
item.created_at, item.updated_at
FROM freight_order_items AS item
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
WHERE freight.creator_subject = ? AND item.id = ?
AND item.is_present = 1`,
creatorSubject,
itemID,
))
if errors.Is(err, sql.ErrNoRows) {
return domain.ProcurementSourceItem{}, usecase.ErrRepositoryNotFound
}
if err != nil {
return domain.ProcurementSourceItem{}, repositoryFailure(err)
}
var canceled sql.NullBool
if err := store.db.QueryRowContext(
ctx,
`SELECT freight.is_canceled
FROM freight_orders AS freight
JOIN freight_order_items AS item
ON item.freight_order_id = freight.id
WHERE freight.creator_subject = ? AND item.id = ?`,
creatorSubject,
itemID,
).Scan(&canceled); err != nil {
return domain.ProcurementSourceItem{}, repositoryFailure(err)
}
var isCanceled *bool
if canceled.Valid {
isCanceled = &canceled.Bool
}
return domain.ProcurementSourceItem{
Item: item,
IsCanceled: isCanceled,
}, nil
}
func (store *Store) CreateProcurementRequest(
ctx context.Context,
candidate domain.ProcurementRequest,
) (domain.ProcurementRequest, bool, error) {
tx, err := store.db.BeginTx(ctx, nil)
if err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
}
defer tx.Rollback()
existing, err := getProcurementRequestBySource(
ctx,
tx,
candidate.CreatorSubject,
candidate.FreightOrderItemID,
candidate.SourceRevision,
)
if err == nil {
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
}
return existing, false, nil
}
if !errors.Is(err, usecase.ErrRepositoryNotFound) {
return domain.ProcurementRequest{}, false, err
}
var currentRevision int
var currentSHA string
var present bool
if err := tx.QueryRowContext(
ctx,
`SELECT item.revision, item.canonical_sha256, item.is_present
FROM freight_order_items AS item
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
WHERE freight.creator_subject = ? AND item.id = ?`,
candidate.CreatorSubject,
candidate.FreightOrderItemID,
).Scan(&currentRevision, &currentSHA, &present); err != nil {
if errors.Is(err, sql.ErrNoRows) {
return domain.ProcurementRequest{}, false, usecase.ErrRepositoryNotFound
}
return domain.ProcurementRequest{}, false, repositoryFailure(err)
}
if !present || currentRevision != candidate.SourceRevision ||
currentSHA != candidate.SourceSHA256 {
return domain.ProcurementRequest{}, false,
usecase.ErrProcurementSourceChanged
}
_, err = tx.ExecContext(
ctx,
`INSERT INTO procurement_requests (
id, creator_subject, freight_order_item_id, source_revision,
source_sha256, title, product_spec, sku, quantity,
source_purchase_status, source_is_canceled,
procurement_confirmed_by_user_id, procurement_confirmed_at,
status, blocking_code, reference_asset_id, purchase_task_id,
created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, NULL, ?, ?)`,
candidate.ID,
candidate.CreatorSubject,
candidate.FreightOrderItemID,
candidate.SourceRevision,
candidate.SourceSHA256,
candidate.Title,
candidate.ProductSpec,
candidate.SKU,
nullableFreightQuantity(candidate.Quantity),
nullableString(candidate.SourcePurchaseStatus),
nullableBool(candidate.SourceIsCanceled),
candidate.ProcurementConfirmedByUserID,
formatTimestamp(candidate.ProcurementConfirmedAt),
candidate.Status,
nullableString(candidate.BlockingCode),
formatTimestamp(candidate.CreatedAt),
formatTimestamp(candidate.UpdatedAt),
)
if err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
}
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
}
return candidate, true, nil
}
func (store *Store) ListProcurementRequestsForOrder(
ctx context.Context,
creatorSubject, orderID string,
) ([]domain.ProcurementRequest, error) {
rows, err := store.db.QueryContext(
ctx,
procurementRequestSelect+`
WHERE request.creator_subject = ?
AND item.freight_order_id = ?
ORDER BY item.id, request.source_revision DESC, request.id DESC`,
creatorSubject,
orderID,
)
if err != nil {
return nil, repositoryFailure(err)
}
defer rows.Close()
requests := make([]domain.ProcurementRequest, 0)
for rows.Next() {
request, err := scanProcurementRequest(rows)
if err != nil {
return nil, repositoryFailure(err)
}
requests = append(requests, request)
}
if err := rows.Err(); err != nil {
return nil, repositoryFailure(err)
}
return requests, nil
}
func (store *Store) GetProcurementRequest(
ctx context.Context,
creatorSubject, requestID string,
) (domain.ProcurementRequest, error) {
return getProcurementRequest(ctx, store.db, creatorSubject, requestID)
}
func (store *Store) BindProcurementReference(
ctx context.Context,
creatorSubject, requestID, assetID string,
now time.Time,
) (domain.ProcurementRequest, error) {
tx, err := store.db.BeginTx(ctx, nil)
if err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
defer tx.Rollback()
request, err := getProcurementRequest(
ctx,
tx,
creatorSubject,
requestID,
)
if err != nil {
return domain.ProcurementRequest{}, err
}
if request.SourceChanged {
if request.PurchaseTaskID == nil {
if err := markProcurementSourceChanged(ctx, tx, request.ID, now); err != nil {
return domain.ProcurementRequest{}, err
}
}
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
return domain.ProcurementRequest{}, usecase.ErrProcurementSourceChanged
}
if request.ReferenceAssetID != nil &&
*request.ReferenceAssetID == assetID &&
(request.Status == domain.ProcurementReady ||
request.Status == domain.ProcurementTaskCreated) {
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
return request, nil
}
if request.Status != domain.ProcurementNeedsImage {
return domain.ProcurementRequest{}, usecase.ErrProcurementStateConflict
}
var available int
if 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
)
AND NOT EXISTS (
SELECT 1 FROM procurement_requests AS other
WHERE other.reference_asset_id = asset.id
AND other.id != ?
)
)`,
assetID,
creatorSubject,
requestID,
).Scan(&available); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
if available != 1 {
return domain.ProcurementRequest{}, usecase.ErrAssetUnavailable
}
if _, err := tx.ExecContext(
ctx,
`UPDATE procurement_requests
SET reference_asset_id = ?, status = 'READY', updated_at = ?
WHERE id = ? AND status = 'NEEDS_IMAGE'`,
assetID,
formatTimestamp(now),
requestID,
); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
updated, err := getProcurementRequest(ctx, tx, creatorSubject, requestID)
if err != nil {
return domain.ProcurementRequest{}, err
}
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
return updated, nil
}
func (store *Store) CreateProcurementTask(
ctx context.Context,
expected domain.ProcurementRequest,
task domain.PurchaseTask,
event domain.TaskEvent,
source domain.PurchaseTaskSource,
idempotencyKey, requestSHA256 string,
) (domain.PurchaseTask, bool, error) {
tx, err := store.db.BeginTx(ctx, nil)
if err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
defer tx.Rollback()
existingHash, resourceID, found, err := lookupIdempotency(
ctx,
tx,
task.CreatorSubject,
createProcurementTaskOperation,
idempotencyKey,
)
if err != nil {
return domain.PurchaseTask{}, false, err
}
if found {
if existingHash != requestSHA256 {
return domain.PurchaseTask{}, false, usecase.ErrIdempotencyConflict
}
existing, err := getTaskByID(
ctx,
tx,
task.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
}
request, err := getProcurementRequest(
ctx,
tx,
task.CreatorSubject,
expected.ID,
)
if err != nil {
return domain.PurchaseTask{}, false, err
}
if request.Status == domain.ProcurementTaskCreated &&
request.PurchaseTaskID != nil {
existing, err := getTaskByID(
ctx,
tx,
task.CreatorSubject,
*request.PurchaseTaskID,
)
if err != nil {
return domain.PurchaseTask{}, false, err
}
if err := insertIdempotency(
ctx,
tx,
task.CreatorSubject,
createProcurementTaskOperation,
idempotencyKey,
requestSHA256,
"PURCHASE_TASK",
existing.ID,
task.CreatedAt,
); err != nil {
return domain.PurchaseTask{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
return existing, false, nil
}
if request.SourceChanged {
if request.PurchaseTaskID == nil {
if err := markProcurementSourceChanged(
ctx,
tx,
request.ID,
task.CreatedAt,
); err != nil {
return domain.PurchaseTask{}, false, err
}
}
if err := tx.Commit(); err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
return domain.PurchaseTask{}, false,
usecase.ErrProcurementSourceChanged
}
if request.Status != domain.ProcurementReady ||
request.ReferenceAssetID == nil ||
request.SourceSHA256 != expected.SourceSHA256 {
return domain.PurchaseTask{}, false,
usecase.ErrProcurementStateConflict
}
var assetAvailable int
if 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 other
WHERE other.image_asset_id = asset.id
)
)`,
*request.ReferenceAssetID,
task.CreatorSubject,
).Scan(&assetAvailable); err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
if assetAvailable != 1 {
return domain.PurchaseTask{}, false, usecase.ErrAssetUnavailable
}
_, err = tx.ExecContext(
ctx,
`INSERT INTO purchase_tasks (
id, creator_subject, created_by_user_id, 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, 'CNY', 'PENDING', 1,
NULL, NULL, ?, ?)`,
task.ID,
task.CreatorSubject,
nullableString(task.CreatedByUserID),
nullableString(task.SourceRef),
task.Title,
task.Description,
task.SKU,
task.ImageAssetID,
task.Quantity,
formatTimestamp(task.CreatedAt),
formatTimestamp(task.UpdatedAt),
)
if err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
if err := insertTaskEvent(ctx, tx, event); err != nil {
return domain.PurchaseTask{}, false, err
}
_, err = tx.ExecContext(
ctx,
`INSERT INTO purchase_task_sources (
task_id, procurement_request_id, freight_order_item_id,
source_revision, source_sha256, created_at
) VALUES (?, ?, ?, ?, ?, ?)`,
source.TaskID,
source.ProcurementRequestID,
source.FreightOrderItemID,
source.SourceRevision,
source.SourceSHA256,
formatTimestamp(source.CreatedAt),
)
if err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
result, err := tx.ExecContext(
ctx,
`UPDATE procurement_requests
SET status = 'TASK_CREATED', purchase_task_id = ?, updated_at = ?
WHERE id = ? AND status = 'READY'`,
task.ID,
formatTimestamp(task.CreatedAt),
request.ID,
)
if err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
changed, err := result.RowsAffected()
if err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
if changed != 1 {
return domain.PurchaseTask{}, false,
usecase.ErrProcurementStateConflict
}
if err := insertIdempotency(
ctx,
tx,
task.CreatorSubject,
createProcurementTaskOperation,
idempotencyKey,
requestSHA256,
"PURCHASE_TASK",
task.ID,
task.CreatedAt,
); err != nil {
return domain.PurchaseTask{}, false, err
}
if err := tx.Commit(); err != nil {
return domain.PurchaseTask{}, false, repositoryFailure(err)
}
return task, true, nil
}
const procurementRequestSelect = `SELECT
request.id, request.creator_subject, request.freight_order_item_id,
request.source_revision, request.source_sha256, request.title,
request.product_spec, request.sku, request.quantity,
request.source_purchase_status, request.source_is_canceled,
request.procurement_confirmed_by_user_id,
request.procurement_confirmed_at, request.status,
request.blocking_code, request.reference_asset_id,
request.purchase_task_id, request.created_at, request.updated_at,
CASE
WHEN item.revision != request.source_revision
OR item.canonical_sha256 != request.source_sha256
OR item.is_present != 1
OR COALESCE(freight.is_canceled, -1)
!= COALESCE(request.source_is_canceled, -1)
THEN 1 ELSE 0
END
FROM procurement_requests AS request
JOIN freight_order_items AS item
ON item.id = request.freight_order_item_id
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
`
func getProcurementRequest(
ctx context.Context,
query queryRower,
creatorSubject, requestID string,
) (domain.ProcurementRequest, error) {
request, err := scanProcurementRequest(query.QueryRowContext(
ctx,
procurementRequestSelect+`
WHERE request.creator_subject = ? AND request.id = ?`,
creatorSubject,
requestID,
))
if errors.Is(err, sql.ErrNoRows) {
return domain.ProcurementRequest{}, usecase.ErrRepositoryNotFound
}
if err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
return request, nil
}
func getProcurementRequestBySource(
ctx context.Context,
query queryRower,
creatorSubject, itemID string,
revision int,
) (domain.ProcurementRequest, error) {
request, err := scanProcurementRequest(query.QueryRowContext(
ctx,
procurementRequestSelect+`
WHERE request.creator_subject = ?
AND request.freight_order_item_id = ?
AND request.source_revision = ?`,
creatorSubject,
itemID,
revision,
))
if errors.Is(err, sql.ErrNoRows) {
return domain.ProcurementRequest{}, usecase.ErrRepositoryNotFound
}
if err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
return request, nil
}
func scanProcurementRequest(
scanner rowScanner,
) (domain.ProcurementRequest, error) {
var request domain.ProcurementRequest
var quantity sql.NullInt64
var purchaseStatus sql.NullString
var isCanceled sql.NullBool
var blockingCode, assetID, taskID sql.NullString
var confirmedAt, createdAt, updatedAt string
var sourceChanged bool
err := scanner.Scan(
&request.ID,
&request.CreatorSubject,
&request.FreightOrderItemID,
&request.SourceRevision,
&request.SourceSHA256,
&request.Title,
&request.ProductSpec,
&request.SKU,
&quantity,
&purchaseStatus,
&isCanceled,
&request.ProcurementConfirmedByUserID,
&confirmedAt,
&request.Status,
&blockingCode,
&assetID,
&taskID,
&createdAt,
&updatedAt,
&sourceChanged,
)
if err != nil {
return domain.ProcurementRequest{}, err
}
if quantity.Valid {
value := int(quantity.Int64)
request.Quantity = &value
}
request.SourcePurchaseStatus = optionalString(purchaseStatus)
if isCanceled.Valid {
request.SourceIsCanceled = &isCanceled.Bool
}
request.BlockingCode = optionalString(blockingCode)
request.ReferenceAssetID = optionalString(assetID)
request.PurchaseTaskID = optionalString(taskID)
request.SourceChanged = sourceChanged
request.ProcurementConfirmedAt, err = parseTimestamp(confirmedAt)
if err != nil {
return domain.ProcurementRequest{}, err
}
request.CreatedAt, err = parseTimestamp(createdAt)
if err != nil {
return domain.ProcurementRequest{}, err
}
request.UpdatedAt, err = parseTimestamp(updatedAt)
return request, err
}
func markProcurementSourceChanged(
ctx context.Context,
tx *sql.Tx,
requestID string,
now time.Time,
) error {
_, err := tx.ExecContext(
ctx,
`UPDATE procurement_requests
SET status = 'SOURCE_CHANGED', blocking_code = 'SOURCE_CHANGED',
updated_at = ?
WHERE id = ? AND purchase_task_id IS NULL`,
formatTimestamp(now),
requestID,
)
if err != nil {
return repositoryFailure(err)
}
return nil
}
var _ usecase.ProcurementRepository = (*Store)(nil)
@@ -0,0 +1,325 @@
package sqlite_test
import (
"context"
"errors"
"testing"
"time"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/platform/migration"
repository "cmroubao/backend-api/internal/repository/sqlite"
"cmroubao/backend-api/internal/usecase"
)
func TestProcurementRequestsArePerItemAndTaskSnapshotIsImmutable(
t *testing.T,
) {
db := openDatabase(t)
store, _ := repository.New(db)
ctx := context.Background()
now := time.Date(2026, 7, 28, 5, 6, 7, 0, time.UTC)
userID := uuid(1100)
seedFreightUser(t, db, userID, now)
run := freightRun(1101, userID, now)
createAndStartFreightRun(t, store, run, "procurement-source-1", "1")
if err := store.CompleteFreightSync(
ctx,
run,
freightBatch(1110, "a", "b"),
now.Add(time.Second),
); err != nil {
t.Fatalf("CompleteFreightSync() error = %v", err)
}
orders, _ := store.ListFreightOrders(ctx, "local-admin", 10)
detail, _ := store.GetFreightOrder(ctx, "local-admin", orders[0].ID)
service, err := usecase.NewProcurementService(
store,
procurementClock{now: now.Add(2 * time.Second)},
&procurementIDs{next: 1200},
)
if err != nil {
t.Fatalf("NewProcurementService() error = %v", err)
}
requests := make([]domain.ProcurementRequest, 0, len(detail.Items))
for _, item := range detail.Items {
result, err := service.CreateRequest(
ctx,
usecase.CreateProcurementRequestCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
FreightOrderItemID: item.ID,
ConfirmProcurementNeeded: true,
},
)
if err != nil || result.Request.Status !=
domain.ProcurementNeedsImage {
t.Fatalf("CreateRequest() = %+v, %v", result, err)
}
requests = append(requests, result.Request)
}
if len(requests) != 2 || requests[0].ID == requests[1].ID {
t.Fatalf("requests = %+v", requests)
}
replay, err := service.CreateRequest(
ctx,
usecase.CreateProcurementRequestCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
FreightOrderItemID: detail.Items[0].ID,
ConfirmProcurementNeeded: true,
},
)
if err != nil || !replay.Replayed ||
replay.Request.ID != requests[0].ID {
t.Fatalf("request replay = %+v, %v", replay, err)
}
asset := domain.Asset{
ID: uuid(1300),
CreatorSubject: "local-admin",
Purpose: domain.AssetPurposeTaskReference,
MediaType: domain.NormalizedImageMediaType,
SizeBytes: 100,
SHA256: repeatHex("f"),
StorageKey: "procurement/reference.jpg",
CreatedAt: now,
}
if _, _, err := store.CreateAssetIdempotent(
ctx,
asset,
"procurement-asset",
repeatHex("e"),
); err != nil {
t.Fatalf("CreateAssetIdempotent() error = %v", err)
}
bound, err := service.BindReference(
ctx,
usecase.BindProcurementReferenceCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
RequestID: requests[0].ID,
ImageAssetID: asset.ID,
},
)
if err != nil || bound.Status != domain.ProcurementReady {
t.Fatalf("BindReference() = %+v, %v", bound, err)
}
created, err := service.CreateTask(
ctx,
usecase.CreateProcurementTaskCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
RequestID: requests[0].ID,
IdempotencyKey: "procurement-task-1",
},
)
if err != nil || created.Task.Status != domain.TaskStatusPending ||
created.Task.Title != detail.Items[0].Title ||
created.Task.SKU != detail.Items[0].SKU ||
created.Task.ImageAssetID != asset.ID {
t.Fatalf("CreateTask() = %+v, %v", created, err)
}
replayedTask, err := service.CreateTask(
ctx,
usecase.CreateProcurementTaskCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
RequestID: requests[0].ID,
IdempotencyKey: "procurement-task-2",
},
)
if err != nil || !replayedTask.Replayed ||
replayedTask.Task.ID != created.Task.ID {
t.Fatalf("task replay = %+v, %v", replayedTask, err)
}
buyerID := uuid(1600)
deviceID := uuid(1601)
seedLifecycleUser(
t,
db,
buyerID,
"procurement-buyer",
domain.UserRoleBuyer,
now,
)
seedLifecycleDevice(
t,
db,
deviceID,
"procurement-device",
buyerID,
repeatHex("6"),
now,
)
lifecycle := &lifecycleFixture{
store: store,
db: db,
buyerOneID: buyerID,
deviceOneID: deviceID,
}
readyAt := now.Add(3 * time.Second)
recordReadyHeartbeat(t, lifecycle, buyerID, deviceID, readyAt)
claimed, err := store.ClaimNext(
ctx,
lifecycleClaimRequest(
lifecycle,
buyerID,
deviceID,
1602,
readyAt.Add(time.Second),
readyAt.Add(10*time.Minute),
readyAt.Add(-time.Minute),
),
)
if err != nil || claimed.Task == nil ||
claimed.Task.ID != created.Task.ID ||
claimed.Task.Status != domain.TaskStatusClaimed {
t.Fatalf("Roubao claim = %+v, %v", claimed, err)
}
updateRun := freightRun(1102, userID, now.Add(time.Minute))
createAndStartFreightRun(
t,
store,
updateRun,
"procurement-source-2",
"2",
)
changed := freightBatch(1120, "c", "b")
changed.Orders[0].CanonicalSHA256 = repeatHex("d")
changed.Orders[0].Items[0].Title = "ERP 更新后的标题"
changed.Orders[0].Items[0].CanonicalSHA256 = repeatHex("c")
if err := store.CompleteFreightSync(
ctx,
updateRun,
changed,
now.Add(time.Minute+time.Second),
); err != nil {
t.Fatalf("updated freight import error = %v", err)
}
storedRequest, err := service.Get(
ctx,
"local-admin",
requests[0].ID,
)
if err != nil || !storedRequest.SourceChanged ||
storedRequest.PurchaseTaskID == nil ||
*storedRequest.PurchaseTaskID != created.Task.ID {
t.Fatalf("source changed request = %+v, %v", storedRequest, err)
}
replayAfterUpdate, err := service.CreateTask(
ctx,
usecase.CreateProcurementTaskCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
RequestID: requests[0].ID,
IdempotencyKey: "procurement-task-after-source-update",
},
)
if err != nil || !replayAfterUpdate.Replayed ||
replayAfterUpdate.Task.ID != created.Task.ID {
t.Fatalf(
"task replay after source update = %+v, %v",
replayAfterUpdate,
err,
)
}
storedTask, err := store.GetTaskDetail(
ctx,
"local-admin",
created.Task.ID,
)
if err != nil || storedTask.Task.Title != detail.Items[0].Title ||
storedTask.Task.Title == "ERP 更新后的标题" {
t.Fatalf("immutable task = %+v, %v", storedTask.Task, err)
}
runner, err := migration.New(db)
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(ctx); err == nil {
t.Fatal("procurement migration down succeeded with retained requests")
}
var retained int
if err := db.QueryRow(
`SELECT count(*) FROM procurement_requests`,
).Scan(&retained); err != nil || retained != 2 {
t.Fatalf("retained requests = %d, error = %v", retained, err)
}
}
func TestProcurementBlocksInvalidSourceAndChangedRequest(t *testing.T) {
db := openDatabase(t)
store, _ := repository.New(db)
ctx := context.Background()
now := time.Date(2026, 7, 28, 6, 7, 8, 0, time.UTC)
userID := uuid(1400)
seedFreightUser(t, db, userID, now)
run := freightRun(1401, userID, now)
createAndStartFreightRun(t, store, run, "blocked-source-1", "3")
batch := freightBatch(1410, "a", "b")
batch.Orders[0].Items[0].SKU = ""
batch.Orders[0].Items[0].CanonicalSHA256 = repeatHex("7")
batch.Orders[0].CanonicalSHA256 = repeatHex("8")
if err := store.CompleteFreightSync(
ctx,
run,
batch,
now.Add(time.Second),
); err != nil {
t.Fatalf("CompleteFreightSync() error = %v", err)
}
orders, _ := store.ListFreightOrders(ctx, "local-admin", 10)
detail, _ := store.GetFreightOrder(ctx, "local-admin", orders[0].ID)
service, _ := usecase.NewProcurementService(
store,
procurementClock{now: now.Add(2 * time.Second)},
&procurementIDs{next: 1500},
)
blocked, err := service.CreateRequest(
ctx,
usecase.CreateProcurementRequestCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
FreightOrderItemID: detail.Items[0].ID,
ConfirmProcurementNeeded: true,
},
)
if err != nil || blocked.Request.Status != domain.ProcurementBlocked ||
blocked.Request.BlockingCode == nil ||
*blocked.Request.BlockingCode != domain.ProcurementBlockSKURequired {
t.Fatalf("blocked request = %+v, %v", blocked, err)
}
_, err = service.CreateTask(
ctx,
usecase.CreateProcurementTaskCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
RequestID: blocked.Request.ID,
IdempotencyKey: "blocked-task",
},
)
var typed *usecase.Error
if !errors.As(err, &typed) ||
typed.Code != "PROCUREMENT_STATE_CONFLICT" {
t.Fatalf("blocked CreateTask() error = %v", err)
}
}
type procurementClock struct {
now time.Time
}
func (clock procurementClock) Now() time.Time {
return clock.now
}
type procurementIDs struct {
next int
}
func (ids *procurementIDs) NewID() (string, error) {
ids.next++
return uuid(ids.next), nil
}