diff --git a/backend-api/cmd/api/main.go b/backend-api/cmd/api/main.go index 9dbb304..bf2e25e 100644 --- a/backend-api/cmd/api/main.go +++ b/backend-api/cmd/api/main.go @@ -278,7 +278,12 @@ func buildRouter( if _, err := freight.RecoverInterrupted(ctx); err != nil { return nil, err } - procurement, err := usecase.NewProcurementService(store, clock, ids) + procurement, err := usecase.NewProcurementService( + store, + clock, + ids, + usecase.WithProcurementReferenceStore(files), + ) if err != nil { return nil, err } diff --git a/backend-api/internal/domain/procurement.go b/backend-api/internal/domain/procurement.go index bf64f4e..48f5729 100644 --- a/backend-api/internal/domain/procurement.go +++ b/backend-api/internal/domain/procurement.go @@ -3,6 +3,7 @@ package domain import "time" type ProcurementRequestStatus string +type ProcurementReferenceOrigin string const ( ProcurementBlocked ProcurementRequestStatus = "BLOCKED" @@ -20,6 +21,11 @@ const ( ProcurementBlockSourceChanged = "SOURCE_CHANGED" ) +const ( + ProcurementReferenceERP ProcurementReferenceOrigin = "ERP_FREIGHT_IMAGE" + ProcurementReferenceManual ProcurementReferenceOrigin = "MANUAL_UPLOAD" +) + type ProcurementSourceItem struct { Item FreightOrderItem IsCanceled *bool @@ -42,12 +48,23 @@ type ProcurementRequest struct { Status ProcurementRequestStatus BlockingCode *string ReferenceAssetID *string + ReferenceOrigin *ProcurementReferenceOrigin + ReferenceProductThumbRef *string + ReferenceSourceImageSHA256 *string + ReferenceBoundAt *time.Time PurchaseTaskID *string SourceChanged bool CreatedAt time.Time UpdatedAt time.Time } +type ProcurementReferenceSource struct { + Origin ProcurementReferenceOrigin + ProductThumbRef *string + SourceImageSHA256 *string + BoundAt time.Time +} + type PurchaseTaskSource struct { TaskID string ProcurementRequestID string diff --git a/backend-api/internal/platform/migration/claims_migration_test.go b/backend-api/internal/platform/migration/claims_migration_test.go index c7e0d79..563a553 100644 --- a/backend-api/internal/platform/migration/claims_migration_test.go +++ b/backend-api/internal/platform/migration/claims_migration_test.go @@ -34,8 +34,11 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { if applied, err := runner.Up(ctx); err != nil { t.Fatalf("initial Up() error = %v", err) - } else if applied != 16 { - t.Fatalf("initial Up() applied = %d, want 16", applied) + } else if applied != 17 { + t.Fatalf("initial Up() applied = %d, want 17", applied) + } + if err := runner.Down(ctx); err != nil { + t.Fatalf("initial Down(v17) error = %v", err) } if err := runner.Down(ctx); err != nil { t.Fatalf("initial Down(v16) error = %v", err) @@ -77,9 +80,14 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { seedClaimsHistoricalFixture(t, db) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("Up(v5-v16) over historical data error = %v", err) - } else if applied != 12 { - t.Fatalf("Up(v5-v16) applied = %d, want 12", applied) + t.Fatalf("Up(v5-v17) over historical data error = %v", err) + } else if applied != 13 { + t.Fatalf("Up(v5-v17) applied = %d, want 13", applied) + } + assertClaimsHistory(t, db, true) + + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v17) with compatible history error = %v", err) } assertClaimsHistory(t, db, true) @@ -149,9 +157,9 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { assertClaimsHistory(t, db, false) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("final Up(v4-v16) error = %v", err) - } else if applied != 13 { - t.Fatalf("final Up(v4-v16) applied = %d, want 13", applied) + t.Fatalf("final Up(v4-v17) error = %v", err) + } else if applied != 14 { + t.Fatalf("final Up(v4-v17) applied = %d, want 14", applied) } assertClaimsHistory(t, db, true) } @@ -387,6 +395,9 @@ func TestClaimsMigrationDownFailsClosedForNewAuditData(t *testing.T) { t.Fatalf("insert v4 audit event: %v", err) } + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v17) error = %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("Down(v16) error = %v", err) } diff --git a/backend-api/internal/platform/migration/runner_test.go b/backend-api/internal/platform/migration/runner_test.go index baef96e..8748a56 100644 --- a/backend-api/internal/platform/migration/runner_test.go +++ b/backend-api/internal/platform/migration/runner_test.go @@ -27,8 +27,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { if err != nil { t.Fatalf("Up() error = %v", err) } - if applied != 16 { - t.Fatalf("Up() applied = %d, want 16", applied) + if applied != 17 { + t.Fatalf("Up() applied = %d, want 17", applied) } assertStatuses(t, runner, map[int64]bool{ 1: true, @@ -47,6 +47,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 14: true, 15: true, 16: true, + 17: true, }) applied, err = runner.Up(context.Background()) @@ -76,7 +77,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 13: true, 14: true, 15: true, - 16: false, + 16: true, + 17: false, }) applied, err = runner.Up(context.Background()) @@ -103,6 +105,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 14: true, 15: true, 16: true, + 17: true, }) } diff --git a/backend-api/internal/repository/sqlite/auth_repository_test.go b/backend-api/internal/repository/sqlite/auth_repository_test.go index 98656c5..afc1980 100644 --- a/backend-api/internal/repository/sqlite/auth_repository_test.go +++ b/backend-api/internal/repository/sqlite/auth_repository_test.go @@ -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) } } diff --git a/backend-api/internal/repository/sqlite/freight_image_repository_test.go b/backend-api/internal/repository/sqlite/freight_image_repository_test.go index 1260d69..c5b4ad1 100644 --- a/backend-api/internal/repository/sqlite/freight_image_repository_test.go +++ b/backend-api/internal/repository/sqlite/freight_image_repository_test.go @@ -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") } diff --git a/backend-api/internal/repository/sqlite/freight_repository_test.go b/backend-api/internal/repository/sqlite/freight_repository_test.go index 5acb15e..f72aab3 100644 --- a/backend-api/internal/repository/sqlite/freight_repository_test.go +++ b/backend-api/internal/repository/sqlite/freight_repository_test.go @@ -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) } diff --git a/backend-api/internal/repository/sqlite/procurement_auto_reference_test.go b/backend-api/internal/repository/sqlite/procurement_auto_reference_test.go new file mode 100644 index 0000000..6dc855c --- /dev/null +++ b/backend-api/internal/repository/sqlite/procurement_auto_reference_test.go @@ -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) + } +} diff --git a/backend-api/internal/repository/sqlite/procurement_repository.go b/backend-api/internal/repository/sqlite/procurement_repository.go index 1b4f8c5..7bb93c6 100644 --- a/backend-api/internal/repository/sqlite/procurement_repository.go +++ b/backend-api/internal/repository/sqlite/procurement_repository.go @@ -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(¤tRevision, ¤tSHA, &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, diff --git a/backend-api/internal/repository/sqlite/procurement_repository_test.go b/backend-api/internal/repository/sqlite/procurement_repository_test.go index 05a3bb1..3cd5683 100644 --- a/backend-api/internal/repository/sqlite/procurement_repository_test.go +++ b/backend-api/internal/repository/sqlite/procurement_repository_test.go @@ -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) } diff --git a/backend-api/internal/transport/httpapi/admin_handlers_test.go b/backend-api/internal/transport/httpapi/admin_handlers_test.go index 8f15f9c..dcf5342 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers_test.go +++ b/backend-api/internal/transport/httpapi/admin_handlers_test.go @@ -766,7 +766,11 @@ func TestAdminProcurementAPIProducesImmutablePendingTask(t *testing.T) { "", ) if createRequest.Code != http.StatusCreated || - !strings.Contains(createRequest.Body.String(), `"status":"NEEDS_IMAGE"`) { + !strings.Contains(createRequest.Body.String(), `"status":"READY"`) || + !strings.Contains( + createRequest.Body.String(), + `"reference_origin":"ERP_FREIGHT_IMAGE"`, + ) { t.Fatalf( "create request status/body = %d / %s", createRequest.Code, @@ -806,7 +810,11 @@ func TestAdminProcurementAPIProducesImmutablePendingTask(t *testing.T) { "", ) if bind.Code != http.StatusOK || - !strings.Contains(bind.Body.String(), `"status":"READY"`) { + !strings.Contains(bind.Body.String(), `"status":"READY"`) || + !strings.Contains( + bind.Body.String(), + `"reference_origin":"MANUAL_UPLOAD"`, + ) { t.Fatalf("bind status/body = %d / %s", bind.Code, bind.Body) } deleteWithReference := performAdminRequest( @@ -920,7 +928,7 @@ func TestAdminProcurementAPIProducesImmutablePendingTask(t *testing.T) { } func TestAdminFreightDeleteRemovesNeedsImageDraft(t *testing.T) { - fixture := newAdminIntegrationFixture(t) + fixture := newAdminIntegrationFixtureWithReferences(t, false) createSync := performAdminRequest( t, fixture.router, @@ -1123,6 +1131,9 @@ func TestAdminOrderAuthorizationIsIdempotentAndRevisioned(t *testing.T) { if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("reference source migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("image migration down: %v", err) } @@ -1406,6 +1417,13 @@ type adminIntegrationFixture struct { } func newAdminIntegrationFixture(t *testing.T) *adminIntegrationFixture { + return newAdminIntegrationFixtureWithReferences(t, true) +} + +func newAdminIntegrationFixtureWithReferences( + t *testing.T, + autoReferences bool, +) *adminIntegrationFixture { t.Helper() ctx := context.Background() db, err := database.Open(ctx, filepath.Join(t.TempDir(), "admin.db")) @@ -1482,10 +1500,18 @@ func newAdminIntegrationFixture(t *testing.T) *adminIntegrationFixture { if err != nil { t.Fatalf("usecase.NewFreightService() error = %v", err) } + procurementOptions := make([]usecase.ProcurementServiceOption, 0, 1) + if autoReferences { + procurementOptions = append( + procurementOptions, + usecase.WithProcurementReferenceStore(files), + ) + } procurement, err := usecase.NewProcurementService( repositories, clock, ids, + procurementOptions..., ) if err != nil { t.Fatalf("usecase.NewProcurementService() error = %v", err) diff --git a/backend-api/internal/transport/httpapi/device_handlers_test.go b/backend-api/internal/transport/httpapi/device_handlers_test.go index 87b2a27..cd20e15 100644 --- a/backend-api/internal/transport/httpapi/device_handlers_test.go +++ b/backend-api/internal/transport/httpapi/device_handlers_test.go @@ -713,6 +713,9 @@ func TestDeviceExecutionResultsAreIdempotentAndAuditable(t *testing.T) { if err != nil { t.Fatalf("migration.New() after review error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("reference source migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("image migration down: %v", err) } @@ -742,8 +745,8 @@ func TestDeviceExecutionResultsAreIdempotentAndAuditable(t *testing.T) { } if applied, err := runner.Up(context.Background()); err != nil { t.Fatalf("restore device command migration: %v", err) - } else if applied != 8 { - t.Fatalf("restored migrations = %d, want 8", applied) + } else if applied != 9 { + t.Fatalf("restored migrations = %d, want 9", applied) } completePayload := fmt.Sprintf( @@ -1432,6 +1435,9 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable( if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("reference source migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("image migration down: %v", err) } diff --git a/backend-api/internal/transport/httpapi/procurement_handlers.go b/backend-api/internal/transport/httpapi/procurement_handlers.go index 6feed3c..b42fa38 100644 --- a/backend-api/internal/transport/httpapi/procurement_handlers.go +++ b/backend-api/internal/transport/httpapi/procurement_handlers.go @@ -160,6 +160,7 @@ func procurementRequestResponse( "status": request.Status, "blocking_code": request.BlockingCode, "reference_asset_id": request.ReferenceAssetID, + "reference_origin": request.ReferenceOrigin, "purchase_task_id": request.PurchaseTaskID, "source_changed": request.SourceChanged, "created_at": formatTime(request.CreatedAt), diff --git a/backend-api/internal/transport/webui/handler_test.go b/backend-api/internal/transport/webui/handler_test.go index 55a2ecf..3921bb8 100644 --- a/backend-api/internal/transport/webui/handler_test.go +++ b/backend-api/internal/transport/webui/handler_test.go @@ -1239,10 +1239,12 @@ func TestFreightDetailCreatesProcurementTaskWithCSRF(t *testing.T) { ImageURL: "/api/v1/freight-items/" + itemID + "/image", }, Request: &ProcurementRequest{ - ID: itemID, - FreightOrderItemID: itemID, - Status: "READY", - StatusLabel: "可以生成任务", + ID: itemID, + FreightOrderItemID: itemID, + Status: "READY", + StatusLabel: "可以生成任务", + ReferenceOrigin: "ERP_FREIGHT_IMAGE", + ReferenceOriginLabel: "参考图来自 ERP 商品图", }, }}, }, @@ -1263,6 +1265,8 @@ func TestFreightDetailCreatesProcurementTaskWithCSRF(t *testing.T) { ) if detail.Code != http.StatusOK || !strings.Contains(detail.Body.String(), "生成采购任务") || + !strings.Contains(detail.Body.String(), "更换参考图") || + !strings.Contains(detail.Body.String(), "参考图来自 ERP 商品图") || !strings.Contains(detail.Body.String(), "可以生成任务") || !strings.Contains(detail.Body.String(), "TWD 129.50") || !strings.Contains(detail.Body.String(), `width="84" height="84"`) || @@ -1273,6 +1277,13 @@ func TestFreightDetailCreatesProcurementTaskWithCSRF(t *testing.T) { } cookie := csrfCookie(t, detail) taskKey := hiddenValue(t, detail.Body.String(), "task_key") + if uploadKey := hiddenValue( + t, + detail.Body.String(), + "upload_key", + ); !validToken(uploadKey) { + t.Fatalf("upload key = %q", uploadKey) + } values := url.Values{ "csrf_token": {cookie.Value}, "order_id": {testTaskID}, diff --git a/backend-api/internal/transport/webui/templates/freight-detail.gohtml b/backend-api/internal/transport/webui/templates/freight-detail.gohtml index d55d5a7..67dea40 100644 --- a/backend-api/internal/transport/webui/templates/freight-detail.gohtml +++ b/backend-api/internal/transport/webui/templates/freight-detail.gohtml @@ -103,6 +103,7 @@ {{else if eq .Request.Status "READY"}} + {{if .Request.ReferenceOriginLabel}}{{.Request.ReferenceOriginLabel}}{{end}}
@@ -110,6 +111,17 @@
+
+ + + + + +
{{else if .Request.PurchaseTaskID}} 查看采购任务 {{end}} diff --git a/backend-api/internal/transport/webui/types.go b/backend-api/internal/transport/webui/types.go index 24f8bac..8811ea3 100644 --- a/backend-api/internal/transport/webui/types.go +++ b/backend-api/internal/transport/webui/types.go @@ -192,18 +192,20 @@ type FreightItemReview struct { } type ProcurementRequest struct { - ID string - FreightOrderItemID string - SourceRevision int - Status string - StatusLabel string - BlockingCode string - BlockingLabel string - ReferenceAssetID string - PurchaseTaskID string - SourceChanged bool - UploadKey string - TaskKey string + ID string + FreightOrderItemID string + SourceRevision int + Status string + StatusLabel string + BlockingCode string + BlockingLabel string + ReferenceAssetID string + ReferenceOrigin string + ReferenceOriginLabel string + PurchaseTaskID string + SourceChanged bool + UploadKey string + TaskKey string } type ListTasksInput struct { diff --git a/backend-api/internal/transport/webui/usecase_adapter.go b/backend-api/internal/transport/webui/usecase_adapter.go index 977a411..88e231b 100644 --- a/backend-api/internal/transport/webui/usecase_adapter.go +++ b/backend-api/internal/transport/webui/usecase_adapter.go @@ -292,9 +292,11 @@ func procurementRequestFrom( BlockingLabel: procurementBlockingLabel( stringValue(request.BlockingCode), ), - ReferenceAssetID: stringValue(request.ReferenceAssetID), - PurchaseTaskID: stringValue(request.PurchaseTaskID), - SourceChanged: request.SourceChanged, + ReferenceAssetID: stringValue(request.ReferenceAssetID), + ReferenceOrigin: procurementReferenceOrigin(request), + ReferenceOriginLabel: procurementReferenceOriginLabel(request), + PurchaseTaskID: stringValue(request.PurchaseTaskID), + SourceChanged: request.SourceChanged, } if request.SourceChanged { result.StatusLabel = "来源已变化" @@ -302,6 +304,28 @@ func procurementRequestFrom( return result } +func procurementReferenceOrigin( + request domain.ProcurementRequest, +) string { + if request.ReferenceOrigin == nil { + return "" + } + return string(*request.ReferenceOrigin) +} + +func procurementReferenceOriginLabel( + request domain.ProcurementRequest, +) string { + switch procurementReferenceOrigin(request) { + case string(domain.ProcurementReferenceERP): + return "参考图来自 ERP 商品图" + case string(domain.ProcurementReferenceManual): + return "参考图由人工上传" + default: + return "" + } +} + func procurementStatusLabel( status domain.ProcurementRequestStatus, ) string { diff --git a/backend-api/internal/usecase/procurement_ports.go b/backend-api/internal/usecase/procurement_ports.go index edc257f..b0f21c4 100644 --- a/backend-api/internal/usecase/procurement_ports.go +++ b/backend-api/internal/usecase/procurement_ports.go @@ -7,16 +7,29 @@ import ( "cmroubao/backend-api/internal/domain" ) +type AutoProcurementReference struct { + Asset domain.Asset + ProductThumbRef string + SourceImageSHA256 string + SourceImageStorageKey string +} + type ProcurementRepository interface { GetProcurementSourceItem( context.Context, string, string, ) (domain.ProcurementSourceItem, error) + GetReadyFreightItemImage( + context.Context, + string, + string, + ) (domain.FreightItemImage, error) CreateProcurementRequest( context.Context, domain.ProcurementRequest, - ) (domain.ProcurementRequest, bool, error) + *AutoProcurementReference, + ) (domain.ProcurementRequest, bool, bool, error) ListProcurementRequestsForOrder( context.Context, string, @@ -33,7 +46,7 @@ type ProcurementRepository interface { string, string, time.Time, - ) (domain.ProcurementRequest, error) + ) (domain.ProcurementRequest, *string, error) CreateProcurementTask( context.Context, domain.ProcurementRequest, diff --git a/backend-api/internal/usecase/procurement_service.go b/backend-api/internal/usecase/procurement_service.go index ad0c215..a13188f 100644 --- a/backend-api/internal/usecase/procurement_service.go +++ b/backend-api/internal/usecase/procurement_service.go @@ -5,6 +5,7 @@ import ( "errors" "strconv" "strings" + "time" "unicode/utf8" "cmroubao/backend-api/internal/domain" @@ -21,10 +22,25 @@ var ( type ProcurementService struct { repository ProcurementRepository + store ReferenceImageStore clock Clock ids IDGenerator } +type ProcurementServiceOption func(*ProcurementService) error + +func WithProcurementReferenceStore( + store ReferenceImageStore, +) ProcurementServiceOption { + return func(service *ProcurementService) error { + if store == nil { + return errors.New("procurement reference store is required") + } + service.store = store + return nil + } +} + type CreateProcurementRequestCommand struct { CreatorSubject string ActorUserID string @@ -60,15 +76,25 @@ func NewProcurementService( repository ProcurementRepository, clock Clock, ids IDGenerator, + options ...ProcurementServiceOption, ) (*ProcurementService, error) { if repository == nil || clock == nil || ids == nil { return nil, errors.New("procurement service dependencies are required") } - return &ProcurementService{ + service := &ProcurementService{ repository: repository, clock: clock, ids: ids, - }, nil + } + for _, option := range options { + if option == nil { + return nil, errors.New("procurement service option is required") + } + if err := option(service); err != nil { + return nil, err + } + } + return service, nil } func (service *ProcurementService) CreateRequest( @@ -136,13 +162,27 @@ func (service *ProcurementService) CreateRequest( request.Status = domain.ProcurementBlocked request.BlockingCode = &code } - stored, created, err := service.repository.CreateProcurementRequest( - ctx, - request, - ) + var auto *AutoProcurementReference + if request.Status == domain.ProcurementNeedsImage && + service.store != nil { + auto, err = service.copyFreightImageReference(ctx, request, now) + if err != nil { + return CreateProcurementRequestResult{}, err + } + } + stored, created, referenceUsed, err := + service.repository.CreateProcurementRequest( + ctx, + request, + auto, + ) if err != nil { + service.cleanupAutomaticReference(auto) return CreateProcurementRequestResult{}, wrapProcurementError(err) } + if !referenceUsed { + service.cleanupAutomaticReference(auto) + } return CreateProcurementRequestResult{ Request: stored, Replayed: !created, @@ -207,19 +247,95 @@ func (service *ProcurementService) BindReference( fields, ) } - request, err := service.repository.BindProcurementReference( - ctx, - command.CreatorSubject, - command.RequestID, - command.ImageAssetID, - service.clock.Now().UTC(), - ) + request, replacedStorageKey, err := + service.repository.BindProcurementReference( + ctx, + command.CreatorSubject, + command.RequestID, + command.ImageAssetID, + service.clock.Now().UTC(), + ) if err != nil { return domain.ProcurementRequest{}, wrapProcurementError(err) } + if replacedStorageKey != nil { + service.deleteStoredReference(*replacedStorageKey) + } return request, nil } +func (service *ProcurementService) copyFreightImageReference( + ctx context.Context, + request domain.ProcurementRequest, + now time.Time, +) (*AutoProcurementReference, error) { + image, err := service.repository.GetReadyFreightItemImage( + ctx, + request.CreatorSubject, + request.FreightOrderItemID, + ) + if errors.Is(err, ErrRepositoryNotFound) { + return nil, nil + } + if err != nil { + return nil, wrapRepositoryError(err) + } + assetID, err := service.ids.NewID() + if err != nil { + return nil, procurementInternal(err) + } + content, err := service.store.Open(ctx, image.StorageKey) + if err != nil { + return nil, mapImageStoreError(err) + } + normalized, putErr := service.store.Put( + ctx, + assetID, + image.MediaType, + content, + ) + closeErr := content.Close() + if putErr != nil { + return nil, mapImageStoreError(putErr) + } + if closeErr != nil { + service.deleteStoredReference(normalized.StorageKey) + return nil, procurementInternal(closeErr) + } + return &AutoProcurementReference{ + Asset: domain.Asset{ + ID: assetID, + CreatorSubject: request.CreatorSubject, + Purpose: domain.AssetPurposeTaskReference, + MediaType: normalized.MediaType, + SizeBytes: normalized.SizeBytes, + SHA256: normalized.SHA256, + StorageKey: normalized.StorageKey, + CreatedAt: now, + }, + ProductThumbRef: image.ProductThumbRef, + SourceImageSHA256: image.SHA256, + SourceImageStorageKey: image.StorageKey, + }, nil +} + +func (service *ProcurementService) cleanupAutomaticReference( + auto *AutoProcurementReference, +) { + if auto != nil { + service.deleteStoredReference(auto.Asset.StorageKey) + } +} + +func (service *ProcurementService) deleteStoredReference(storageKey string) { + if service.store == nil || storageKey == "" { + return + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + _ = service.store.Delete(ctx, storageKey) +} + func (service *ProcurementService) CreateTask( ctx context.Context, command CreateProcurementTaskCommand, diff --git a/backend-api/migrations/00017_procurement_reference_sources.sql b/backend-api/migrations/00017_procurement_reference_sources.sql new file mode 100644 index 0000000..fb2b0ae --- /dev/null +++ b/backend-api/migrations/00017_procurement_reference_sources.sql @@ -0,0 +1,61 @@ +-- +goose Up +CREATE TABLE procurement_reference_sources ( + procurement_request_id TEXT PRIMARY KEY NOT NULL + REFERENCES procurement_requests(id) + ON UPDATE RESTRICT ON DELETE RESTRICT, + origin TEXT NOT NULL + CHECK (origin IN ('ERP_FREIGHT_IMAGE', 'MANUAL_UPLOAD')), + product_thumb_ref TEXT + CHECK ( + product_thumb_ref IS NULL + OR ( + product_thumb_ref GLOB '[1-9]*' + AND product_thumb_ref NOT GLOB '*[^0-9]*' + AND length(product_thumb_ref) <= 20 + ) + ), + source_image_sha256 TEXT + CHECK ( + source_image_sha256 IS NULL + OR ( + length(source_image_sha256) = 64 + AND source_image_sha256 NOT GLOB '*[^0-9a-f]*' + ) + ), + bound_at TEXT NOT NULL, + CHECK ( + (origin = 'ERP_FREIGHT_IMAGE' + AND product_thumb_ref IS NOT NULL + AND source_image_sha256 IS NOT NULL) + OR (origin = 'MANUAL_UPLOAD' + AND product_thumb_ref IS NULL + AND source_image_sha256 IS NULL) + ) +); + +INSERT INTO procurement_reference_sources ( + procurement_request_id, origin, product_thumb_ref, + source_image_sha256, bound_at +) +SELECT id, 'MANUAL_UPLOAD', NULL, NULL, updated_at +FROM procurement_requests +WHERE reference_asset_id IS NOT NULL; + +CREATE INDEX procurement_reference_sources_origin_idx + ON procurement_reference_sources (origin, bound_at, procurement_request_id); + +-- +goose Down +CREATE TEMP TABLE procurement_reference_sources_v17_down_guard ( + allowed INTEGER NOT NULL CHECK (allowed = 1) +); + +INSERT INTO procurement_reference_sources_v17_down_guard (allowed) +SELECT CASE + WHEN EXISTS (SELECT 1 FROM procurement_reference_sources) + THEN 0 + ELSE 1 +END; + +DROP TABLE procurement_reference_sources_v17_down_guard; +DROP INDEX procurement_reference_sources_origin_idx; +DROP TABLE procurement_reference_sources; diff --git a/docs/current-state.md b/docs/current-state.md index b502603..d32ad30 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -5,7 +5,7 @@ ## 当前快照 - 日期:2026-07-29 -- 阶段:T-244 已完成 ERP 图片缺失媒体类型兼容 +- 阶段:T-245 已完成 ERP 商品图自动作为采购参考图 - Git:当前分支为 `main`;T-001 至 T-004、T-101 至 T-104、T-201 至 T-219 均按文档提交、实现提交的顺序纳入历史 - 生产代码:`android-buyer/` 已接入 Roubao Android 源码 @@ -41,13 +41,17 @@ 货运落库、同步成功和水位推进同事务完成;失败与较旧范围成功不推进水位。Admin 可查看查询范围、计数、稳定错误码和最近成功水位。 - 采购需求:v13 按货运商品 revision 保存可复核快照;ADMIN 显式确认采购、上传参考 - 图后生成不可变 PENDING 任务。来源变化不会改写任务,Roubao 可沿用现有 claim 流程。 + 图后生成不可变 PENDING 任务。T-245 在当前商品有 `READY` 本地 ERP 缓存图时,将其 + 重新校验并复制为独立 `TASK_REFERENCE`,需求直接进入 `READY`;v17 记录 + `ERP_FREIGHT_IMAGE`/`MANUAL_UPLOAD` 来源。无图仍进入 `NEEDS_IMAGE`,既有需求可在 + 图片恢复后显式升级;任务生成前可人工替换并清理旧自动副本。来源变化不会改写任务, + Roubao 继续使用原 `image_asset_id` 和 claim 流程。 - 本机 Android 工具:JDK 17.0.13、Command-line Tools 22.0、SDK 34、 Build Tools 34.0.0、Platform Tools/ADB 37.0.0;用户级 SDK 环境变量已设置 - Android Studio:未安装;`winget` 静默安装卡住后已终止,不阻塞命令行构建 - 测试:T-219 Android Debug/Release 单元测试与构建和根 `init.ps1` 通过; Debug APK `1.4.16 (21)` 已覆盖安装到 PKG110 -- 后端测试:T-226 至 T-244 已运行 `go test ./...`、`go test -race ./...`、`go vet ./...` +- 后端测试:T-226 至 T-245 已运行 `go test ./...`、`go test -race ./...`、`go vet ./...` 和三个 Go 入口构建;T-227 增加 Go source 的伪 ERP 会话预检、完整单号、日期分页去重、 详情 allowlist 和稳定错误码覆盖;根 `init.ps1` 的 Android 测试/Debug APK 与 Go 标准 验证也通过,均未访问真实 ERP; diff --git a/docs/tasks/T-245.md b/docs/tasks/T-245.md index 329f649..0537891 100644 --- a/docs/tasks/T-245.md +++ b/docs/tasks/T-245.md @@ -5,7 +5,7 @@ phase: 2 deps: - T-243 - T-244 -status: TODO +status: DONE created: 2026-07-29 context_ref: 9195019 work_branch: null @@ -69,17 +69,17 @@ write_paths: ## 验收要点 -- [ ] 当前商品有 `READY` ERP 图时,创建需求返回 `READY`、独立 asset 和 +- [x] 当前商品有 `READY` ERP 图时,创建需求返回 `READY`、独立 asset 和 `ERP_FREIGHT_IMAGE` 来源。 -- [ ] 自动 asset 使用不同 storage key;货运缓存图替换不会破坏采购参考图。 -- [ ] 无可用 ERP 图时返回 `NEEDS_IMAGE`,人工上传路径不回归。 -- [ ] 资料阻塞时不复制文件、不创建 asset,状态保持 `BLOCKED`。 -- [ ] 重复创建不产生重复需求/asset;既有 `NEEDS_IMAGE` 可在图片恢复后升级。 -- [ ] 商品或图片并发变化时无采购 asset、来源记录或孤立复制文件。 -- [ ] `READY` 需求可人工替换;旧 ERP 自动 asset/文件被清理,任务创建后禁止替换。 -- [ ] API/SSR 展示参考图来源,生成的采购任务和 Roubao 图片契约不变化。 -- [ ] 已自动绑定参考图的货运单继续受 T-243 删除保护。 -- [ ] v17 上下迁移、标准 Go 测试、race、vet 和三个入口构建通过。 +- [x] 自动 asset 使用不同 storage key;货运缓存图替换不会破坏采购参考图。 +- [x] 无可用 ERP 图时返回 `NEEDS_IMAGE`,人工上传路径不回归。 +- [x] 资料阻塞时不复制文件、不创建 asset,状态保持 `BLOCKED`。 +- [x] 重复创建不产生重复需求/asset;既有 `NEEDS_IMAGE` 可在图片恢复后升级。 +- [x] 商品或图片并发变化时无采购 asset、来源记录或孤立复制文件。 +- [x] `READY` 需求可人工替换;旧 ERP 自动 asset/文件被清理,任务创建后禁止替换。 +- [x] API/SSR 展示参考图来源,生成的采购任务和 Roubao 图片契约不变化。 +- [x] 已自动绑定参考图的货运单继续受 T-243 删除保护。 +- [x] v17 上下迁移、标准 Go 测试、race、vet 和三个入口构建通过。 ## 边界 @@ -96,3 +96,12 @@ write_paths: - 2026-07-29:审计确认当前 `CreateRequest` 对非阻塞商品固定写入 `NEEDS_IMAGE`;ERP 缓存图只有 `freight_item_images` 记录,而采购任务要求 `assets` 外键。确定采用独立 文件/asset 快照,不共享可变货运 storage key;Roubao 协议无需变化。 +- 2026-07-29:v17 已增加参考图来源记录;生产组合启用本地 ERP 缓存图复制。SQLite + 事务复核商品 revision/hash 和图片引用/hash/key 后,原子写 asset、需求及来源;重复 + 创建可升级既有 `NEEDS_IMAGE`,未使用或失败的复制文件由用例补偿清理。 +- 2026-07-29:人工绑定已支持任务生成前替换自动参考图,提交后清理旧自动 asset/文件; + API 返回 `reference_origin`,货运详情显示来源并提供“更换参考图”。任务仍使用既有 + `image_asset_id`,未改变 Roubao 协议。 +- 2026-07-29:`go test ./...`、`go test -race ./...`、`go vet ./...` 通过; + `go build ./cmd/api`、`go build ./cmd/authctl`、`go build ./cmd/migrate` 通过。 + 测试仅使用运行时生成图片和脱敏 ID,未访问真实 ERP。