feat(t245): use ERP images as procurement references

This commit is contained in:
QiuSW
2026-07-29 17:29:20 +08:00
parent a20590b4b7
commit 15074692a3
22 changed files with 1057 additions and 117 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(v17) error = %v", err)
}
if err := runner.Down(context.Background()); err != nil {
t.Fatalf("Down(v16) error = %v", err)
}
@@ -439,8 +442,8 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks(
}
if applied, err := runner.Up(context.Background()); err != nil {
t.Fatalf("Up(v3-v16) error = %v", err)
} else if applied != 14 {
t.Fatalf("Up(v3-v16) applied = %d, want 14", applied)
} else if applied != 15 {
t.Fatalf("Up(v3-v17) applied = %d, want 15", applied)
}
}
@@ -189,6 +189,9 @@ func TestFreightItemImagesAreCurrentRetryableAndRollbackGuarded(
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("reference source migration down: %v", err)
}
if err := runner.Down(ctx); err == nil {
t.Fatal("image migration down succeeded with retained image metadata")
}
@@ -131,6 +131,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("reference source migration down: %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("image migration down: %v", err)
}
@@ -271,6 +274,9 @@ func TestFreightDateSyncAdvancesWatermarkOnlyOnWholeBatchSuccess(
}
runner, _ := migration.New(db)
if err := runner.Down(ctx); err != nil {
t.Fatalf("reference source migration down: %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("image migration down: %v", err)
}
@@ -0,0 +1,281 @@
package sqlite_test
import (
"bytes"
"context"
"image"
"image/jpeg"
"testing"
"time"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/platform/assetstore"
repository "cmroubao/backend-api/internal/repository/sqlite"
"cmroubao/backend-api/internal/usecase"
)
func TestProcurementAutomaticallyCopiesAndCanReplaceERPReference(
t *testing.T,
) {
db := openDatabase(t)
repositories, _ := repository.New(db)
ctx := context.Background()
now := time.Date(2026, 7, 29, 8, 0, 0, 0, time.UTC)
userID := uuid(1700)
seedFreightUser(t, db, userID, now)
run := freightRun(1701, userID, now)
createAndStartFreightRun(t, repositories, run, "auto-reference", "1")
batch := freightBatch(1710, "a", "b")
firstThumb := "190100"
secondThumb := "190101"
batch.Orders[0].Items[0].ProductThumbRef = &firstThumb
batch.Orders[0].Items[1].ProductThumbRef = &secondThumb
if err := repositories.CompleteFreightSync(
ctx,
run,
batch,
now.Add(time.Second),
); err != nil {
t.Fatalf("CompleteFreightSync() error = %v", err)
}
orders, _ := repositories.ListFreightOrders(ctx, "local-admin", 10)
detail, _ := repositories.GetFreightOrder(
ctx,
"local-admin",
orders[0].ID,
)
files, err := assetstore.New(t.TempDir())
if err != nil {
t.Fatalf("assetstore.New() error = %v", err)
}
firstImage := saveReadyFreightImage(
t,
ctx,
repositories,
files,
detail.Items[0],
firstThumb,
now.Add(2*time.Second),
)
ids := &procurementIDs{next: 1800}
service, err := usecase.NewProcurementService(
repositories,
procurementClock{now: now.Add(3 * time.Second)},
ids,
usecase.WithProcurementReferenceStore(files),
)
if err != nil {
t.Fatalf("NewProcurementService() error = %v", err)
}
command := usecase.CreateProcurementRequestCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
FreightOrderItemID: detail.Items[0].ID,
ConfirmProcurementNeeded: true,
}
created, err := service.CreateRequest(ctx, command)
if err != nil ||
created.Request.Status != domain.ProcurementReady ||
created.Request.ReferenceAssetID == nil ||
created.Request.ReferenceOrigin == nil ||
*created.Request.ReferenceOrigin != domain.ProcurementReferenceERP ||
created.Request.ReferenceSourceImageSHA256 == nil ||
*created.Request.ReferenceSourceImageSHA256 != firstImage.SHA256 {
t.Fatalf("automatic request = %+v, %v", created, err)
}
autoAsset, err := repositories.GetAsset(
ctx,
"local-admin",
*created.Request.ReferenceAssetID,
)
if err != nil || autoAsset.StorageKey == firstImage.StorageKey {
t.Fatalf("automatic asset = %+v, %v", autoAsset, err)
}
assertStoredImageExists(t, ctx, files, firstImage.StorageKey)
assertStoredImageExists(t, ctx, files, autoAsset.StorageKey)
replayed, err := service.CreateRequest(ctx, command)
if err != nil || !replayed.Replayed ||
replayed.Request.ID != created.Request.ID ||
replayed.Request.ReferenceAssetID == nil ||
*replayed.Request.ReferenceAssetID != autoAsset.ID {
t.Fatalf("automatic replay = %+v, %v", replayed, err)
}
var assetCount int
if err := db.QueryRow(
`SELECT count(*) FROM assets WHERE creator_subject = ?`,
"local-admin",
).Scan(&assetCount); err != nil || assetCount != 1 {
t.Fatalf("asset count after replay = %d, %v", assetCount, err)
}
manualID := uuid(1900)
manualNormalized := putTestImage(t, ctx, files, manualID)
manual := domain.Asset{
ID: manualID,
CreatorSubject: "local-admin",
Purpose: domain.AssetPurposeTaskReference,
MediaType: manualNormalized.MediaType,
SizeBytes: manualNormalized.SizeBytes,
SHA256: manualNormalized.SHA256,
StorageKey: manualNormalized.StorageKey,
CreatedAt: now.Add(4 * time.Second),
}
if _, _, err := repositories.CreateAssetIdempotent(
ctx,
manual,
"manual-reference",
repeatHex("a"),
); err != nil {
t.Fatalf("CreateAssetIdempotent() error = %v", err)
}
replaced, err := service.BindReference(
ctx,
usecase.BindProcurementReferenceCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
RequestID: created.Request.ID,
ImageAssetID: manual.ID,
},
)
if err != nil ||
replaced.ReferenceOrigin == nil ||
*replaced.ReferenceOrigin != domain.ProcurementReferenceManual ||
replaced.ReferenceAssetID == nil ||
*replaced.ReferenceAssetID != manual.ID {
t.Fatalf("manual replacement = %+v, %v", replaced, err)
}
if _, err := repositories.GetAsset(
ctx,
"local-admin",
autoAsset.ID,
); err == nil {
t.Fatal("replaced automatic asset remains in database")
}
assertStoredImageMissing(t, ctx, files, autoAsset.StorageKey)
assertStoredImageExists(t, ctx, files, firstImage.StorageKey)
assertStoredImageExists(t, ctx, files, manual.StorageKey)
legacy, _ := usecase.NewProcurementService(
repositories,
procurementClock{now: now.Add(5 * time.Second)},
ids,
)
pending, err := legacy.CreateRequest(
ctx,
usecase.CreateProcurementRequestCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
FreightOrderItemID: detail.Items[1].ID,
ConfirmProcurementNeeded: true,
},
)
if err != nil || pending.Request.Status != domain.ProcurementNeedsImage {
t.Fatalf("pending request = %+v, %v", pending, err)
}
saveReadyFreightImage(
t,
ctx,
repositories,
files,
detail.Items[1],
secondThumb,
now.Add(6*time.Second),
)
upgraded, err := service.CreateRequest(
ctx,
usecase.CreateProcurementRequestCommand{
CreatorSubject: "local-admin",
ActorUserID: userID,
FreightOrderItemID: detail.Items[1].ID,
ConfirmProcurementNeeded: true,
},
)
if err != nil || !upgraded.Replayed ||
upgraded.Request.ID != pending.Request.ID ||
upgraded.Request.Status != domain.ProcurementReady {
t.Fatalf("upgraded request = %+v, %v", upgraded, err)
}
}
func saveReadyFreightImage(
t *testing.T,
ctx context.Context,
repositories *repository.Store,
files *assetstore.Store,
item domain.FreightOrderItem,
thumb string,
now time.Time,
) domain.FreightItemImage {
t.Helper()
normalized := putTestImage(t, ctx, files, item.ID)
candidate := domain.FreightItemImage{
CreatorSubject: "local-admin",
FreightOrderItemID: item.ID,
ProductThumbRef: thumb,
Status: domain.FreightItemImageReady,
MediaType: normalized.MediaType,
SizeBytes: normalized.SizeBytes,
SHA256: normalized.SHA256,
StorageKey: normalized.StorageKey,
UpdatedAt: now,
}
if _, err := repositories.SaveFreightItemImage(ctx, candidate); err != nil {
t.Fatalf("SaveFreightItemImage() error = %v", err)
}
return candidate
}
func putTestImage(
t *testing.T,
ctx context.Context,
files *assetstore.Store,
id string,
) usecase.NormalizedReferenceImage {
t.Helper()
var content bytes.Buffer
if err := jpeg.Encode(
&content,
image.NewRGBA(image.Rect(0, 0, 8, 6)),
nil,
); err != nil {
t.Fatalf("jpeg.Encode() error = %v", err)
}
normalized, err := files.Put(
ctx,
id,
"image/jpeg",
bytes.NewReader(content.Bytes()),
)
if err != nil {
t.Fatalf("files.Put() error = %v", err)
}
return normalized
}
func assertStoredImageExists(
t *testing.T,
ctx context.Context,
files *assetstore.Store,
key string,
) {
t.Helper()
content, err := files.Open(ctx, key)
if err != nil {
t.Fatalf("files.Open(%q) error = %v", key, err)
}
_ = content.Close()
}
func assertStoredImageMissing(
t *testing.T,
ctx context.Context,
files *assetstore.Store,
key string,
) {
t.Helper()
if content, err := files.Open(ctx, key); err == nil {
_ = content.Close()
t.Fatalf("files.Open(%q) succeeded after cleanup", key)
}
}
@@ -66,10 +66,12 @@ func (store *Store) GetProcurementSourceItem(
func (store *Store) CreateProcurementRequest(
ctx context.Context,
candidate domain.ProcurementRequest,
) (domain.ProcurementRequest, bool, error) {
auto *usecase.AutoProcurementReference,
) (domain.ProcurementRequest, bool, bool, error) {
tx, err := store.db.BeginTx(ctx, nil)
if err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
return domain.ProcurementRequest{}, false, false,
repositoryFailure(err)
}
defer tx.Rollback()
existing, err := getProcurementRequestBySource(
@@ -80,36 +82,60 @@ func (store *Store) CreateProcurementRequest(
candidate.SourceRevision,
)
if err == nil {
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
if auto != nil &&
existing.Status == domain.ProcurementNeedsImage {
updated, bindErr := bindAutomaticProcurementReference(
ctx,
tx,
existing,
*auto,
candidate.UpdatedAt,
)
if bindErr != nil {
return domain.ProcurementRequest{}, false, false, bindErr
}
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, false, false,
repositoryFailure(err)
}
return updated, false, true, nil
}
return existing, false, nil
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, false, false,
repositoryFailure(err)
}
return existing, false, false, nil
}
if !errors.Is(err, usecase.ErrRepositoryNotFound) {
return domain.ProcurementRequest{}, false, err
return domain.ProcurementRequest{}, false, 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
if err := validateProcurementSource(ctx, tx, candidate); err != nil {
return domain.ProcurementRequest{}, false, false, err
}
referenceUsed := false
if auto != nil {
if candidate.Status != domain.ProcurementNeedsImage {
return domain.ProcurementRequest{}, false, false,
repositoryFailure(usecase.ErrRepositoryInvariant)
}
return domain.ProcurementRequest{}, false, repositoryFailure(err)
}
if !present || currentRevision != candidate.SourceRevision ||
currentSHA != candidate.SourceSHA256 {
return domain.ProcurementRequest{}, false,
usecase.ErrProcurementSourceChanged
if err := validateAutomaticProcurementReference(
ctx,
tx,
candidate,
*auto,
); err != nil {
return domain.ProcurementRequest{}, false, false, err
}
if err := insertProcurementReferenceAsset(
ctx,
tx,
auto.Asset,
); err != nil {
return domain.ProcurementRequest{}, false, false, err
}
candidate.ReferenceAssetID = &auto.Asset.ID
candidate.Status = domain.ProcurementReady
referenceUsed = true
}
_, err = tx.ExecContext(
ctx,
@@ -120,7 +146,7 @@ func (store *Store) CreateProcurementRequest(
procurement_confirmed_by_user_id, procurement_confirmed_at,
status, blocking_code, reference_asset_id, purchase_task_id,
created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, NULL, ?, ?)`,
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, ?, ?)`,
candidate.ID,
candidate.CreatorSubject,
candidate.FreightOrderItemID,
@@ -136,16 +162,43 @@ func (store *Store) CreateProcurementRequest(
formatTimestamp(candidate.ProcurementConfirmedAt),
candidate.Status,
nullableString(candidate.BlockingCode),
nullableString(candidate.ReferenceAssetID),
formatTimestamp(candidate.CreatedAt),
formatTimestamp(candidate.UpdatedAt),
)
if err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
return domain.ProcurementRequest{}, false, false,
repositoryFailure(err)
}
if auto != nil {
if err := saveProcurementReferenceSource(
ctx,
tx,
candidate.ID,
domain.ProcurementReferenceSource{
Origin: domain.ProcurementReferenceERP,
ProductThumbRef: &auto.ProductThumbRef,
SourceImageSHA256: &auto.SourceImageSHA256,
BoundAt: candidate.UpdatedAt,
},
); err != nil {
return domain.ProcurementRequest{}, false, false, err
}
}
stored, err := getProcurementRequest(
ctx,
tx,
candidate.CreatorSubject,
candidate.ID,
)
if err != nil {
return domain.ProcurementRequest{}, false, false, err
}
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, false, repositoryFailure(err)
return domain.ProcurementRequest{}, false, false,
repositoryFailure(err)
}
return candidate, true, nil
return stored, true, referenceUsed, nil
}
func (store *Store) ListProcurementRequestsForOrder(
@@ -190,10 +243,10 @@ func (store *Store) BindProcurementReference(
ctx context.Context,
creatorSubject, requestID, assetID string,
now time.Time,
) (domain.ProcurementRequest, error) {
) (domain.ProcurementRequest, *string, error) {
tx, err := store.db.BeginTx(ctx, nil)
if err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
return domain.ProcurementRequest{}, nil, repositoryFailure(err)
}
defer tx.Rollback()
request, err := getProcurementRequest(
@@ -203,30 +256,33 @@ func (store *Store) BindProcurementReference(
requestID,
)
if err != nil {
return domain.ProcurementRequest{}, err
return domain.ProcurementRequest{}, nil, err
}
if request.SourceChanged {
if request.PurchaseTaskID == nil {
if err := markProcurementSourceChanged(ctx, tx, request.ID, now); err != nil {
return domain.ProcurementRequest{}, err
return domain.ProcurementRequest{}, nil, err
}
}
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
return domain.ProcurementRequest{}, nil, repositoryFailure(err)
}
return domain.ProcurementRequest{}, usecase.ErrProcurementSourceChanged
return domain.ProcurementRequest{}, nil,
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 domain.ProcurementRequest{}, nil, repositoryFailure(err)
}
return request, nil
return request, nil, nil
}
if request.Status != domain.ProcurementNeedsImage {
return domain.ProcurementRequest{}, usecase.ErrProcurementStateConflict
if request.Status != domain.ProcurementNeedsImage &&
request.Status != domain.ProcurementReady {
return domain.ProcurementRequest{}, nil,
usecase.ErrProcurementStateConflict
}
var available int
if err := tx.QueryRowContext(
@@ -250,30 +306,68 @@ func (store *Store) BindProcurementReference(
creatorSubject,
requestID,
).Scan(&available); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
return domain.ProcurementRequest{}, nil, repositoryFailure(err)
}
if available != 1 {
return domain.ProcurementRequest{}, usecase.ErrAssetUnavailable
return domain.ProcurementRequest{}, nil, usecase.ErrAssetUnavailable
}
var replacedStorageKey *string
var replacedAssetID *string
if request.ReferenceAssetID != nil &&
request.ReferenceOrigin != nil &&
*request.ReferenceOrigin == domain.ProcurementReferenceERP {
replaced, getErr := getAssetByID(
ctx,
tx,
creatorSubject,
*request.ReferenceAssetID,
)
if getErr != nil {
return domain.ProcurementRequest{}, nil, getErr
}
replacedStorageKey = &replaced.StorageKey
replacedAssetID = &replaced.ID
}
if _, err := tx.ExecContext(
ctx,
`UPDATE procurement_requests
SET reference_asset_id = ?, status = 'READY', updated_at = ?
WHERE id = ? AND status = 'NEEDS_IMAGE'`,
SET reference_asset_id = ?, status = 'READY', blocking_code = NULL,
updated_at = ?
WHERE id = ? AND status IN ('NEEDS_IMAGE', 'READY')`,
assetID,
formatTimestamp(now),
requestID,
); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
return domain.ProcurementRequest{}, nil, repositoryFailure(err)
}
if err := saveProcurementReferenceSource(
ctx,
tx,
requestID,
domain.ProcurementReferenceSource{
Origin: domain.ProcurementReferenceManual,
BoundAt: now,
},
); err != nil {
return domain.ProcurementRequest{}, nil, err
}
if replacedAssetID != nil {
if _, err := tx.ExecContext(
ctx,
`DELETE FROM assets WHERE id = ?`,
*replacedAssetID,
); err != nil {
return domain.ProcurementRequest{}, nil, repositoryFailure(err)
}
}
updated, err := getProcurementRequest(ctx, tx, creatorSubject, requestID)
if err != nil {
return domain.ProcurementRequest{}, err
return domain.ProcurementRequest{}, nil, err
}
if err := tx.Commit(); err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
return domain.ProcurementRequest{}, nil, repositoryFailure(err)
}
return updated, nil
return updated, replacedStorageKey, nil
}
func (store *Store) CreateProcurementTask(
@@ -488,6 +582,8 @@ const procurementRequestSelect = `SELECT
request.procurement_confirmed_by_user_id,
request.procurement_confirmed_at, request.status,
request.blocking_code, request.reference_asset_id,
reference_source.origin, reference_source.product_thumb_ref,
reference_source.source_image_sha256, reference_source.bound_at,
request.purchase_task_id, request.created_at, request.updated_at,
CASE
WHEN item.revision != request.source_revision
@@ -502,6 +598,8 @@ const procurementRequestSelect = `SELECT
ON item.id = request.freight_order_item_id
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
LEFT JOIN procurement_reference_sources AS reference_source
ON reference_source.procurement_request_id = request.id
`
func getProcurementRequest(
@@ -557,7 +655,9 @@ func scanProcurementRequest(
var quantity sql.NullInt64
var purchaseStatus sql.NullString
var isCanceled sql.NullBool
var blockingCode, assetID, taskID sql.NullString
var blockingCode, assetID, referenceOrigin sql.NullString
var referenceThumb, referenceSHA, referenceBoundAt sql.NullString
var taskID sql.NullString
var confirmedAt, createdAt, updatedAt string
var sourceChanged bool
err := scanner.Scan(
@@ -577,6 +677,10 @@ func scanProcurementRequest(
&request.Status,
&blockingCode,
&assetID,
&referenceOrigin,
&referenceThumb,
&referenceSHA,
&referenceBoundAt,
&taskID,
&createdAt,
&updatedAt,
@@ -595,6 +699,19 @@ func scanProcurementRequest(
}
request.BlockingCode = optionalString(blockingCode)
request.ReferenceAssetID = optionalString(assetID)
if referenceOrigin.Valid {
origin := domain.ProcurementReferenceOrigin(referenceOrigin.String)
request.ReferenceOrigin = &origin
}
request.ReferenceProductThumbRef = optionalString(referenceThumb)
request.ReferenceSourceImageSHA256 = optionalString(referenceSHA)
if referenceBoundAt.Valid {
boundAt, parseErr := parseTimestamp(referenceBoundAt.String)
if parseErr != nil {
return domain.ProcurementRequest{}, parseErr
}
request.ReferenceBoundAt = &boundAt
}
request.PurchaseTaskID = optionalString(taskID)
request.SourceChanged = sourceChanged
request.ProcurementConfirmedAt, err = parseTimestamp(confirmedAt)
@@ -609,6 +726,204 @@ func scanProcurementRequest(
return request, err
}
func validateProcurementSource(
ctx context.Context,
tx *sql.Tx,
candidate domain.ProcurementRequest,
) error {
var revision int
var sourceSHA string
var present bool
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(&revision, &sourceSHA, &present)
if errors.Is(err, sql.ErrNoRows) {
return usecase.ErrRepositoryNotFound
}
if err != nil {
return repositoryFailure(err)
}
if !present || revision != candidate.SourceRevision ||
sourceSHA != candidate.SourceSHA256 {
return usecase.ErrProcurementSourceChanged
}
return nil
}
func validateAutomaticProcurementReference(
ctx context.Context,
tx *sql.Tx,
candidate domain.ProcurementRequest,
auto usecase.AutoProcurementReference,
) error {
if auto.Asset.CreatorSubject != candidate.CreatorSubject ||
auto.Asset.Purpose != domain.AssetPurposeTaskReference ||
auto.Asset.MediaType != domain.NormalizedImageMediaType ||
auto.Asset.ID == "" || auto.Asset.StorageKey == "" ||
auto.Asset.SizeBytes <= 0 || len(auto.Asset.SHA256) != 64 {
return usecase.ErrRepositoryInvariant
}
var thumbRef, sourceSHA, storageKey string
err := tx.QueryRowContext(
ctx,
`SELECT image.product_thumb_ref, image.sha256, image.storage_key
FROM freight_item_images AS image
JOIN freight_order_items AS item
ON item.id = image.freight_order_item_id
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
WHERE freight.creator_subject = ?
AND item.id = ?
AND item.is_present = 1
AND item.product_thumb_ref = image.product_thumb_ref
AND image.status = 'READY'`,
candidate.CreatorSubject,
candidate.FreightOrderItemID,
).Scan(&thumbRef, &sourceSHA, &storageKey)
if errors.Is(err, sql.ErrNoRows) {
return usecase.ErrProcurementSourceChanged
}
if err != nil {
return repositoryFailure(err)
}
if thumbRef != auto.ProductThumbRef ||
sourceSHA != auto.SourceImageSHA256 ||
storageKey != auto.SourceImageStorageKey {
return usecase.ErrProcurementSourceChanged
}
return nil
}
func insertProcurementReferenceAsset(
ctx context.Context,
tx *sql.Tx,
asset domain.Asset,
) error {
_, err := tx.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,
formatTimestamp(asset.CreatedAt),
)
if err != nil {
return repositoryFailure(err)
}
return nil
}
func saveProcurementReferenceSource(
ctx context.Context,
tx *sql.Tx,
requestID string,
source domain.ProcurementReferenceSource,
) error {
_, err := tx.ExecContext(
ctx,
`INSERT INTO procurement_reference_sources (
procurement_request_id, origin, product_thumb_ref,
source_image_sha256, bound_at
) VALUES (?, ?, ?, ?, ?)
ON CONFLICT(procurement_request_id) DO UPDATE SET
origin = excluded.origin,
product_thumb_ref = excluded.product_thumb_ref,
source_image_sha256 = excluded.source_image_sha256,
bound_at = excluded.bound_at`,
requestID,
source.Origin,
nullableString(source.ProductThumbRef),
nullableString(source.SourceImageSHA256),
formatTimestamp(source.BoundAt),
)
if err != nil {
return repositoryFailure(err)
}
return nil
}
func bindAutomaticProcurementReference(
ctx context.Context,
tx *sql.Tx,
request domain.ProcurementRequest,
auto usecase.AutoProcurementReference,
now time.Time,
) (domain.ProcurementRequest, error) {
if request.SourceChanged {
return domain.ProcurementRequest{},
usecase.ErrProcurementSourceChanged
}
if err := validateProcurementSource(ctx, tx, request); err != nil {
return domain.ProcurementRequest{}, err
}
if err := validateAutomaticProcurementReference(
ctx,
tx,
request,
auto,
); err != nil {
return domain.ProcurementRequest{}, err
}
if err := insertProcurementReferenceAsset(ctx, tx, auto.Asset); err != nil {
return domain.ProcurementRequest{}, err
}
result, err := tx.ExecContext(
ctx,
`UPDATE procurement_requests
SET reference_asset_id = ?, status = 'READY',
blocking_code = NULL, updated_at = ?
WHERE id = ? AND status = 'NEEDS_IMAGE'
AND reference_asset_id IS NULL`,
auto.Asset.ID,
formatTimestamp(now),
request.ID,
)
if err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
affected, err := result.RowsAffected()
if err != nil {
return domain.ProcurementRequest{}, repositoryFailure(err)
}
if affected != 1 {
return domain.ProcurementRequest{},
usecase.ErrProcurementStateConflict
}
if err := saveProcurementReferenceSource(
ctx,
tx,
request.ID,
domain.ProcurementReferenceSource{
Origin: domain.ProcurementReferenceERP,
ProductThumbRef: &auto.ProductThumbRef,
SourceImageSHA256: &auto.SourceImageSHA256,
BoundAt: now,
},
); err != nil {
return domain.ProcurementRequest{}, err
}
return getProcurementRequest(
ctx,
tx,
request.CreatorSubject,
request.ID,
)
}
func markProcurementSourceChanged(
ctx context.Context,
tx *sql.Tx,
@@ -238,6 +238,17 @@ func TestProcurementRequestsArePerItemAndTaskSnapshotIsImmutable(
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(ctx); err == nil {
t.Fatal("reference source migration down succeeded with retained provenance")
}
if _, err := db.Exec(
`DELETE FROM procurement_reference_sources`,
); err != nil {
t.Fatalf("delete reference provenance: %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("reference source migration down: %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("image migration down: %v", err)
}