From 433db451008e3ca5b31124b6889f647df83f91d6 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Tue, 28 Jul 2026 23:51:59 +0800 Subject: [PATCH] feat(t223): generate tasks from freight items --- backend-api/cmd/api/main.go | 6 + backend-api/internal/domain/procurement.go | 58 ++ .../migration/claims_migration_test.go | 27 +- .../platform/migration/runner_test.go | 9 +- .../repository/sqlite/auth_repository_test.go | 9 +- .../sqlite/freight_repository_test.go | 3 + .../sqlite/procurement_repository.go | 632 ++++++++++++++++++ .../sqlite/procurement_repository_test.go | 325 +++++++++ .../transport/httpapi/admin_handlers.go | 15 + .../transport/httpapi/admin_handlers_test.go | 205 ++++++ .../transport/httpapi/device_handlers_test.go | 10 +- .../transport/httpapi/freight_handlers.go | 24 +- .../transport/httpapi/procurement_handlers.go | 168 +++++ .../internal/transport/webui/handler.go | 174 +++++ .../internal/transport/webui/handler_test.go | 111 +++ .../internal/transport/webui/static/admin.css | 29 + .../webui/templates/freight-detail.gohtml | 67 +- backend-api/internal/transport/webui/types.go | 55 +- .../transport/webui/usecase_adapter.go | 158 ++++- .../internal/usecase/procurement_ports.go | 46 ++ .../internal/usecase/procurement_service.go | 412 ++++++++++++ .../migrations/00013_procurement_requests.sql | 132 ++++ docs/api.md | 33 +- docs/current-state.md | 18 +- docs/tasks/T-223.md | 19 +- 25 files changed, 2697 insertions(+), 48 deletions(-) create mode 100644 backend-api/internal/domain/procurement.go create mode 100644 backend-api/internal/repository/sqlite/procurement_repository.go create mode 100644 backend-api/internal/repository/sqlite/procurement_repository_test.go create mode 100644 backend-api/internal/transport/httpapi/procurement_handlers.go create mode 100644 backend-api/internal/usecase/procurement_ports.go create mode 100644 backend-api/internal/usecase/procurement_service.go create mode 100644 backend-api/migrations/00013_procurement_requests.sql diff --git a/backend-api/cmd/api/main.go b/backend-api/cmd/api/main.go index fa5f7c0..dc73cc1 100644 --- a/backend-api/cmd/api/main.go +++ b/backend-api/cmd/api/main.go @@ -243,6 +243,10 @@ func buildRouter( if _, err := freight.RecoverInterrupted(ctx); err != nil { return nil, err } + procurement, err := usecase.NewProcurementService(store, clock, ids) + if err != nil { + return nil, err + } passwords, err := password.NewBcrypt(12) if err != nil { return nil, err @@ -280,6 +284,7 @@ func buildRouter( if err != nil { return nil, err } + webService.SetProcurement(procurement) renderer, err := webui.NewRenderer() if err != nil { return nil, err @@ -322,6 +327,7 @@ func buildRouter( Results: results, Authorizations: authorizations, Freight: freight, + Procurement: procurement, }, webHandler, ) diff --git a/backend-api/internal/domain/procurement.go b/backend-api/internal/domain/procurement.go new file mode 100644 index 0000000..bf64f4e --- /dev/null +++ b/backend-api/internal/domain/procurement.go @@ -0,0 +1,58 @@ +package domain + +import "time" + +type ProcurementRequestStatus string + +const ( + ProcurementBlocked ProcurementRequestStatus = "BLOCKED" + ProcurementNeedsImage ProcurementRequestStatus = "NEEDS_IMAGE" + ProcurementReady ProcurementRequestStatus = "READY" + ProcurementTaskCreated ProcurementRequestStatus = "TASK_CREATED" + ProcurementSourceChanged ProcurementRequestStatus = "SOURCE_CHANGED" +) + +const ( + ProcurementBlockSourceCanceled = "SOURCE_CANCELED" + ProcurementBlockTitleRequired = "TITLE_REQUIRED" + ProcurementBlockSKURequired = "SKU_REQUIRED" + ProcurementBlockQuantityRequired = "QUANTITY_REQUIRED" + ProcurementBlockSourceChanged = "SOURCE_CHANGED" +) + +type ProcurementSourceItem struct { + Item FreightOrderItem + IsCanceled *bool +} + +type ProcurementRequest struct { + ID string + CreatorSubject string + FreightOrderItemID string + SourceRevision int + SourceSHA256 string + Title string + ProductSpec string + SKU string + Quantity *int + SourcePurchaseStatus *string + SourceIsCanceled *bool + ProcurementConfirmedByUserID string + ProcurementConfirmedAt time.Time + Status ProcurementRequestStatus + BlockingCode *string + ReferenceAssetID *string + PurchaseTaskID *string + SourceChanged bool + CreatedAt time.Time + UpdatedAt time.Time +} + +type PurchaseTaskSource struct { + TaskID string + ProcurementRequestID string + FreightOrderItemID string + SourceRevision int + SourceSHA256 string + CreatedAt time.Time +} diff --git a/backend-api/internal/platform/migration/claims_migration_test.go b/backend-api/internal/platform/migration/claims_migration_test.go index 7a1dfb7..37a996a 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 != 12 { - t.Fatalf("initial Up() applied = %d, want 12", applied) + } else if applied != 13 { + t.Fatalf("initial Up() applied = %d, want 13", applied) + } + if err := runner.Down(ctx); err != nil { + t.Fatalf("initial Down(v13) error = %v", err) } if err := runner.Down(ctx); err != nil { t.Fatalf("initial Down(v12) error = %v", err) @@ -65,9 +68,14 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { seedClaimsHistoricalFixture(t, db) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("Up(v5-v12) over historical data error = %v", err) - } else if applied != 8 { - t.Fatalf("Up(v5-v12) applied = %d, want 8", applied) + t.Fatalf("Up(v5-v13) over historical data error = %v", err) + } else if applied != 9 { + t.Fatalf("Up(v5-v13) applied = %d, want 9", applied) + } + assertClaimsHistory(t, db, true) + + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v13) with compatible history error = %v", err) } assertClaimsHistory(t, db, true) @@ -117,9 +125,9 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { assertClaimsHistory(t, db, false) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("final Up(v4-v12) error = %v", err) - } else if applied != 9 { - t.Fatalf("final Up(v4-v12) applied = %d, want 9", applied) + t.Fatalf("final Up(v4-v13) error = %v", err) + } else if applied != 10 { + t.Fatalf("final Up(v4-v13) applied = %d, want 10", applied) } assertClaimsHistory(t, db, true) } @@ -355,6 +363,9 @@ func TestClaimsMigrationDownFailsClosedForNewAuditData(t *testing.T) { t.Fatalf("insert v4 audit event: %v", err) } + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v13) error = %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("Down(v12) error = %v", err) } diff --git a/backend-api/internal/platform/migration/runner_test.go b/backend-api/internal/platform/migration/runner_test.go index c219091..4ad63b8 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 != 12 { - t.Fatalf("Up() applied = %d, want 12", applied) + if applied != 13 { + t.Fatalf("Up() applied = %d, want 13", applied) } assertStatuses(t, runner, map[int64]bool{ 1: true, @@ -43,6 +43,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 10: true, 11: true, 12: true, + 13: true, }) applied, err = runner.Up(context.Background()) @@ -68,7 +69,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 9: true, 10: true, 11: true, - 12: false, + 12: true, + 13: false, }) applied, err = runner.Up(context.Background()) @@ -91,6 +93,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 10: true, 11: true, 12: true, + 13: true, }) } diff --git a/backend-api/internal/repository/sqlite/auth_repository_test.go b/backend-api/internal/repository/sqlite/auth_repository_test.go index 1e9bef1..9830d0f 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(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) } } diff --git a/backend-api/internal/repository/sqlite/freight_repository_test.go b/backend-api/internal/repository/sqlite/freight_repository_test.go index d251bb7..84d4a15 100644 --- a/backend-api/internal/repository/sqlite/freight_repository_test.go +++ b/backend-api/internal/repository/sqlite/freight_repository_test.go @@ -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") } diff --git a/backend-api/internal/repository/sqlite/procurement_repository.go b/backend-api/internal/repository/sqlite/procurement_repository.go new file mode 100644 index 0000000..d2215d5 --- /dev/null +++ b/backend-api/internal/repository/sqlite/procurement_repository.go @@ -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(¤tRevision, ¤tSHA, &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) diff --git a/backend-api/internal/repository/sqlite/procurement_repository_test.go b/backend-api/internal/repository/sqlite/procurement_repository_test.go new file mode 100644 index 0000000..df7bb99 --- /dev/null +++ b/backend-api/internal/repository/sqlite/procurement_repository_test.go @@ -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 +} diff --git a/backend-api/internal/transport/httpapi/admin_handlers.go b/backend-api/internal/transport/httpapi/admin_handlers.go index e24fe4d..1a3d00c 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers.go +++ b/backend-api/internal/transport/httpapi/admin_handlers.go @@ -29,6 +29,7 @@ type AdminServices struct { Results *usecase.ExecutionResultService Authorizations *usecase.OrderAuthorizationService Freight *usecase.FreightService + Procurement *usecase.ProcurementService } func (s AdminServices) validate() error { @@ -68,6 +69,20 @@ func registerAdminAPI(routes gin.IRoutes, services AdminServices) error { routes.GET("/api/v1/freight-orders", handler.listFreightOrders) routes.GET("/api/v1/freight-orders/:id", handler.freightOrderDetail) } + if services.Procurement != nil { + routes.POST( + "/api/v1/freight-items/:id/procurement-request", + handler.createProcurementRequest, + ) + routes.PUT( + "/api/v1/procurement-requests/:id/reference-asset", + handler.bindProcurementReference, + ) + routes.POST( + "/api/v1/procurement-requests/:id/purchase-task", + handler.createProcurementTask, + ) + } return nil } diff --git a/backend-api/internal/transport/httpapi/admin_handlers_test.go b/backend-api/internal/transport/httpapi/admin_handlers_test.go index 443f484..723aa22 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers_test.go +++ b/backend-api/internal/transport/httpapi/admin_handlers_test.go @@ -395,6 +395,199 @@ func TestAdminFreightAPIImportsAllItemsWithoutPII(t *testing.T) { } } +func TestAdminProcurementAPIProducesImmutablePendingTask(t *testing.T) { + fixture := newAdminIntegrationFixture(t) + createSync := performAdminRequest( + t, + fixture.router, + http.MethodPost, + "/api/v1/freight-syncs", + "application/json", + strings.NewReader( + `{"mode":"ORDER_NUMBER","order_number":"SOURCE-12"}`, + ), + "procurement-freight-sync", + ) + var syncBody struct { + Sync struct { + ID string `json:"id"` + } `json:"sync"` + } + decodeResponse(t, createSync, &syncBody) + succeeded := false + for attempt := 0; attempt < 50; attempt++ { + status := performAdminRequest( + t, + fixture.router, + http.MethodGet, + "/api/v1/freight-syncs/"+syncBody.Sync.ID, + "", + nil, + "", + ) + if strings.Contains(status.Body.String(), `"status":"SUCCEEDED"`) { + succeeded = true + break + } + time.Sleep(10 * time.Millisecond) + } + if !succeeded { + t.Fatal("freight sync did not succeed") + } + orders := performAdminRequest( + t, + fixture.router, + http.MethodGet, + "/api/v1/freight-orders", + "", + nil, + "", + ) + var orderBody struct { + Items []struct { + ID string `json:"id"` + } `json:"items"` + } + decodeResponse(t, orders, &orderBody) + detail := performAdminRequest( + t, + fixture.router, + http.MethodGet, + "/api/v1/freight-orders/"+orderBody.Items[0].ID, + "", + nil, + "", + ) + var detailBody struct { + Items []struct { + ID string `json:"id"` + } `json:"items"` + } + decodeResponse(t, detail, &detailBody) + createRequest := performAdminRequest( + t, + fixture.router, + http.MethodPost, + "/api/v1/freight-items/"+detailBody.Items[0].ID+ + "/procurement-request", + "application/json", + strings.NewReader(`{"confirm_procurement_needed":true}`), + "", + ) + if createRequest.Code != http.StatusCreated || + !strings.Contains(createRequest.Body.String(), `"status":"NEEDS_IMAGE"`) { + t.Fatalf( + "create request status/body = %d / %s", + createRequest.Code, + createRequest.Body, + ) + } + var requestBody struct { + Request struct { + ID string `json:"id"` + } `json:"procurement_request"` + } + decodeResponse(t, createRequest, &requestBody) + imageBody, imageContentType := referenceUpload(t, "procurement-image") + upload := performAdminRequest( + t, + fixture.router, + http.MethodPost, + "/api/v1/assets", + imageContentType, + imageBody, + "procurement-image", + ) + var assetBody struct { + ID string `json:"id"` + } + decodeResponse(t, upload, &assetBody) + bind := performAdminRequest( + t, + fixture.router, + http.MethodPut, + "/api/v1/procurement-requests/"+requestBody.Request.ID+ + "/reference-asset", + "application/json", + strings.NewReader( + fmt.Sprintf(`{"image_asset_id":%q}`, assetBody.ID), + ), + "", + ) + if bind.Code != http.StatusOK || + !strings.Contains(bind.Body.String(), `"status":"READY"`) { + t.Fatalf("bind status/body = %d / %s", bind.Code, bind.Body) + } + createTask := performAdminRequest( + t, + fixture.router, + http.MethodPost, + "/api/v1/procurement-requests/"+requestBody.Request.ID+ + "/purchase-task", + "application/json", + strings.NewReader(`{}`), + "procurement-task-create", + ) + if createTask.Code != http.StatusCreated || + !strings.Contains(createTask.Body.String(), `"status":"PENDING"`) { + t.Fatalf( + "create task status/body = %d / %s", + createTask.Code, + createTask.Body, + ) + } + var taskBody struct { + Task struct { + ID string `json:"id"` + } `json:"task"` + } + decodeResponse(t, createTask, &taskBody) + taskDetail := performAdminRequest( + t, + fixture.router, + http.MethodGet, + "/api/v1/tasks/"+taskBody.Task.ID, + "", + nil, + "", + ) + for _, required := range []string{ + `"title":"商品一"`, + `"sku":"BLACK-L"`, + `"quantity":1`, + `"image_asset_id":"` + assetBody.ID + `"`, + `"description":"ERP 规格:黑色,L"`, + } { + if !strings.Contains(taskDetail.Body.String(), required) { + t.Fatalf("task detail missing %q: %s", required, taskDetail.Body) + } + } + replay := performAdminRequest( + t, + fixture.router, + http.MethodPost, + "/api/v1/procurement-requests/"+requestBody.Request.ID+ + "/purchase-task", + "application/json", + strings.NewReader(`{}`), + "procurement-task-replay", + ) + if replay.Code != http.StatusCreated || + !strings.Contains(replay.Body.String(), `"replayed":true`) || + !strings.Contains(replay.Body.String(), taskBody.Task.ID) { + t.Fatalf("task replay status/body = %d / %s", replay.Code, replay.Body) + } + var sourceCount int + if err := fixture.db.QueryRow( + `SELECT count(*) FROM purchase_task_sources + WHERE task_id = ? AND procurement_request_id = ?`, + taskBody.Task.ID, + requestBody.Request.ID, + ).Scan(&sourceCount); err != nil || sourceCount != 1 { + t.Fatalf("task source count = %d, error = %v", sourceCount, err) + } +} + func TestAdminOrderAuthorizationIsIdempotentAndRevisioned(t *testing.T) { fixture := newAdminIntegrationFixture(t) taskID, executionID, taskHash, firstKey, secondKey := @@ -515,6 +708,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("procurement migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("freight migration down: %v", err) } @@ -852,6 +1048,14 @@ func newAdminIntegrationFixture(t *testing.T) *adminIntegrationFixture { if err != nil { t.Fatalf("usecase.NewFreightService() error = %v", err) } + procurement, err := usecase.NewProcurementService( + repositories, + clock, + ids, + ) + if err != nil { + t.Fatalf("usecase.NewProcurementService() error = %v", err) + } registrar, err := NewAdminRouteRegistrar( AdminServices{ Assets: assets, @@ -859,6 +1063,7 @@ func newAdminIntegrationFixture(t *testing.T) *adminIntegrationFixture { Results: results, Authorizations: authorizations, Freight: freight, + Procurement: procurement, }, emptyAdminWeb{}, ) diff --git a/backend-api/internal/transport/httpapi/device_handlers_test.go b/backend-api/internal/transport/httpapi/device_handlers_test.go index 69e9454..a26d336 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("procurement migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("freight migration down: %v", err) } @@ -730,8 +733,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 != 4 { - t.Fatalf("restored migrations = %d, want 4", applied) + } else if applied != 5 { + t.Fatalf("restored migrations = %d, want 5", applied) } completePayload := fmt.Sprintf( @@ -1420,6 +1423,9 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable( if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("procurement migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("freight migration down: %v", err) } diff --git a/backend-api/internal/transport/httpapi/freight_handlers.go b/backend-api/internal/transport/httpapi/freight_handlers.go index d0e3a8e..93972ef 100644 --- a/backend-api/internal/transport/httpapi/freight_handlers.go +++ b/backend-api/internal/transport/httpapi/freight_handlers.go @@ -147,10 +147,30 @@ func (h *adminHandlers) freightOrderDetail(ctx *gin.Context) { }) } ctx.Header("Cache-Control", "no-store") - ctx.JSON(http.StatusOK, gin.H{ + response := gin.H{ "order": freightOrderResponse(detail.Order), "items": items, - }) + } + if h.services.Procurement != nil { + requests, requestErr := h.services.Procurement.ListForOrder( + ctx.Request.Context(), + localAdminSubject, + detail.Order.ID, + ) + if requestErr != nil { + writeUsecaseError(ctx, requestErr) + return + } + requestItems := make([]gin.H, 0, len(requests)) + for _, request := range requests { + requestItems = append( + requestItems, + procurementRequestResponse(request), + ) + } + response["procurement_requests"] = requestItems + } + ctx.JSON(http.StatusOK, response) } func freightSyncResponse(run domain.FreightSyncRun) gin.H { diff --git a/backend-api/internal/transport/httpapi/procurement_handlers.go b/backend-api/internal/transport/httpapi/procurement_handlers.go new file mode 100644 index 0000000..6feed3c --- /dev/null +++ b/backend-api/internal/transport/httpapi/procurement_handlers.go @@ -0,0 +1,168 @@ +package httpapi + +import ( + "net/http" + + "cmroubao/backend-api/internal/domain" + "cmroubao/backend-api/internal/usecase" + + "github.com/gin-gonic/gin" +) + +func (h *adminHandlers) createProcurementRequest(ctx *gin.Context) { + if !hasMediaType(ctx, "application/json") { + writePublicError( + ctx, + http.StatusUnsupportedMediaType, + "UNSUPPORTED_MEDIA_TYPE", + "application/json is required", + false, + gin.H{}, + ) + return + } + var request struct { + ConfirmProcurementNeeded bool `json:"confirm_procurement_needed"` + } + if err := decodeJSON(ctx, &request); err != nil { + writePublicError( + ctx, + http.StatusBadRequest, + "INVALID_JSON", + "request body must be valid JSON", + false, + gin.H{}, + ) + return + } + result, err := h.services.Procurement.CreateRequest( + ctx.Request.Context(), + usecase.CreateProcurementRequestCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: adminActorUserID(ctx), + FreightOrderItemID: ctx.Param("id"), + ConfirmProcurementNeeded: request.ConfirmProcurementNeeded, + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusCreated, gin.H{ + "procurement_request": procurementRequestResponse(result.Request), + "replayed": result.Replayed, + }) +} + +func (h *adminHandlers) bindProcurementReference(ctx *gin.Context) { + if !hasMediaType(ctx, "application/json") { + writePublicError( + ctx, + http.StatusUnsupportedMediaType, + "UNSUPPORTED_MEDIA_TYPE", + "application/json is required", + false, + gin.H{}, + ) + return + } + var request struct { + ImageAssetID string `json:"image_asset_id"` + } + if err := decodeJSON(ctx, &request); err != nil { + writePublicError( + ctx, + http.StatusBadRequest, + "INVALID_JSON", + "request body must be valid JSON", + false, + gin.H{}, + ) + return + } + result, err := h.services.Procurement.BindReference( + ctx.Request.Context(), + usecase.BindProcurementReferenceCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: adminActorUserID(ctx), + RequestID: ctx.Param("id"), + ImageAssetID: request.ImageAssetID, + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusOK, gin.H{ + "procurement_request": procurementRequestResponse(result), + }) +} + +func (h *adminHandlers) createProcurementTask(ctx *gin.Context) { + if !hasMediaType(ctx, "application/json") { + writePublicError( + ctx, + http.StatusUnsupportedMediaType, + "UNSUPPORTED_MEDIA_TYPE", + "application/json is required", + false, + gin.H{}, + ) + return + } + var request struct{} + if err := decodeJSON(ctx, &request); err != nil { + writePublicError( + ctx, + http.StatusBadRequest, + "INVALID_JSON", + "request body must be an empty JSON object", + false, + gin.H{}, + ) + return + } + result, err := h.services.Procurement.CreateTask( + ctx.Request.Context(), + usecase.CreateProcurementTaskCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: adminActorUserID(ctx), + RequestID: ctx.Param("id"), + IdempotencyKey: ctx.GetHeader("Idempotency-Key"), + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusCreated, gin.H{ + "task": taskSummaryResponse(result.Task), + "replayed": result.Replayed, + }) +} + +func procurementRequestResponse( + request domain.ProcurementRequest, +) gin.H { + return gin.H{ + "id": request.ID, + "freight_order_item_id": request.FreightOrderItemID, + "source_revision": request.SourceRevision, + "source_sha256": request.SourceSHA256, + "title": request.Title, + "product_spec": request.ProductSpec, + "sku": request.SKU, + "quantity": request.Quantity, + "source_purchase_status": request.SourcePurchaseStatus, + "status": request.Status, + "blocking_code": request.BlockingCode, + "reference_asset_id": request.ReferenceAssetID, + "purchase_task_id": request.PurchaseTaskID, + "source_changed": request.SourceChanged, + "created_at": formatTime(request.CreatedAt), + "updated_at": formatTime(request.UpdatedAt), + } +} diff --git a/backend-api/internal/transport/webui/handler.go b/backend-api/internal/transport/webui/handler.go index e3c6c3c..f38e2a4 100644 --- a/backend-api/internal/transport/webui/handler.go +++ b/backend-api/internal/transport/webui/handler.go @@ -77,6 +77,23 @@ func (h *Handler) RegisterProtected(routes gin.IRoutes) { routes.POST("/freight/import", SecurityHeaders(), h.CreateFreightImport) routes.GET("/freight/:id", SecurityHeaders(), h.FreightDetail) } + if _, ok := h.service.(ProcurementService); ok { + routes.POST( + "/freight/items/:id/procurement-request", + SecurityHeaders(), + h.CreateFreightProcurementRequest, + ) + routes.POST( + "/freight/procurement-requests/:id/reference", + SecurityHeaders(), + h.BindFreightProcurementReference, + ) + routes.POST( + "/freight/procurement-requests/:id/purchase-task", + SecurityHeaders(), + h.CreateFreightProcurementTask, + ) + } } func (h *Handler) ListFreight(ctx *gin.Context) { @@ -193,6 +210,23 @@ func (h *Handler) FreightDetail(ctx *gin.Context) { return } token, _ := csrfToken(ctx) + for index := range detail.Items { + if detail.Items[index].Request == nil { + continue + } + uploadKey, keyErr := newToken() + if keyErr != nil { + h.renderError(ctx, http.StatusInternalServerError, "页面暂时无法打开", "请稍后重试。") + return + } + taskKey, keyErr := newToken() + if keyErr != nil { + h.renderError(ctx, http.StatusInternalServerError, "页面暂时无法打开", "请稍后重试。") + return + } + detail.Items[index].Request.UploadKey = uploadKey + detail.Items[index].Request.TaskKey = taskKey + } h.render(ctx, http.StatusOK, "freight-detail", freightDetailPage{ Page: pageView{ Title: "货运详情", @@ -200,9 +234,148 @@ func (h *Handler) FreightDetail(ctx *gin.Context) { CSRFToken: token, }, Detail: detail, + Notice: freightNotice(ctx.Query("notice")), }) } +func (h *Handler) CreateFreightProcurementRequest(ctx *gin.Context) { + ctx.Request.Body = http.MaxBytesReader(ctx.Writer, ctx.Request.Body, 16<<10) + if err := ctx.Request.ParseForm(); err != nil || !validCSRF(ctx) { + h.renderError(ctx, http.StatusForbidden, "请求已失效", "请返回货运详情后重新操作。") + return + } + orderID := strings.TrimSpace(ctx.PostForm("order_id")) + if pathEscape(orderID) == "invalid" || + ctx.PostForm("confirm_procurement_needed") != "1" { + h.renderError(ctx, http.StatusUnprocessableEntity, "必须人工确认", "请核对来源商品后确认仍需采购。") + return + } + service := h.service.(ProcurementService) + _, err := service.CreateProcurementRequest( + ctx.Request.Context(), + CreateProcurementRequestInput{ + ActorUserID: actorUserID(ctx.Request.Context()), + FreightOrderItemID: strings.TrimSpace(ctx.Param("id")), + ConfirmProcurementNeeded: true, + }, + ) + if err != nil { + if errors.Is(err, ErrConflict) || errors.Is(err, ErrValidation) { + ctx.Redirect( + http.StatusSeeOther, + "/freight/"+pathEscape(orderID)+"?notice=request-conflict", + ) + return + } + h.renderServiceError(ctx, err, "采购需求创建失败,请稍后重试。") + return + } + ctx.Redirect( + http.StatusSeeOther, + "/freight/"+pathEscape(orderID)+"?notice=request-created", + ) +} + +func (h *Handler) BindFreightProcurementReference(ctx *gin.Context) { + ctx.Request.Body = http.MaxBytesReader( + ctx.Writer, + ctx.Request.Body, + maxRequestBytes, + ) + if err := ctx.Request.ParseMultipartForm(maxRequestBytes); err != nil || + !validCSRF(ctx) { + h.renderError(ctx, http.StatusForbidden, "请求已失效", "请返回货运详情后重新操作。") + return + } + if ctx.Request.MultipartForm != nil { + defer ctx.Request.MultipartForm.RemoveAll() + } + orderID := strings.TrimSpace(ctx.PostForm("order_id")) + uploadKey := strings.TrimSpace(ctx.PostForm("upload_key")) + if pathEscape(orderID) == "invalid" || !validToken(uploadKey) { + h.renderError(ctx, http.StatusForbidden, "请求已失效", "请返回货运详情后重新操作。") + return + } + asset, err := h.uploadReference(ctx, uploadKey) + if err != nil { + ctx.Redirect( + http.StatusSeeOther, + "/freight/"+pathEscape(orderID)+"?notice=image-invalid", + ) + return + } + service := h.service.(ProcurementService) + _, err = service.BindProcurementReference( + ctx.Request.Context(), + BindProcurementReferenceInput{ + ActorUserID: actorUserID(ctx.Request.Context()), + RequestID: strings.TrimSpace(ctx.Param("id")), + ImageAssetID: asset.ID, + }, + ) + if err != nil { + ctx.Redirect( + http.StatusSeeOther, + "/freight/"+pathEscape(orderID)+"?notice=reference-conflict", + ) + return + } + ctx.Redirect( + http.StatusSeeOther, + "/freight/"+pathEscape(orderID)+"?notice=reference-bound", + ) +} + +func (h *Handler) CreateFreightProcurementTask(ctx *gin.Context) { + ctx.Request.Body = http.MaxBytesReader(ctx.Writer, ctx.Request.Body, 16<<10) + if err := ctx.Request.ParseForm(); err != nil || !validCSRF(ctx) { + h.renderError(ctx, http.StatusForbidden, "请求已失效", "请返回货运详情后重新操作。") + return + } + orderID := strings.TrimSpace(ctx.PostForm("order_id")) + taskKey := strings.TrimSpace(ctx.PostForm("task_key")) + if pathEscape(orderID) == "invalid" || !validToken(taskKey) { + h.renderError(ctx, http.StatusForbidden, "请求已失效", "请返回货运详情后重新操作。") + return + } + service := h.service.(ProcurementService) + task, err := service.CreateProcurementTask( + ctx.Request.Context(), + CreateProcurementTaskInput{ + ActorUserID: actorUserID(ctx.Request.Context()), + RequestID: strings.TrimSpace(ctx.Param("id")), + IdempotencyKey: taskKey, + }, + ) + if err != nil { + ctx.Redirect( + http.StatusSeeOther, + "/freight/"+pathEscape(orderID)+"?notice=task-conflict", + ) + return + } + ctx.Redirect(http.StatusSeeOther, "/tasks/"+pathEscape(task.ID)) +} + +func freightNotice(value string) string { + switch value { + case "request-created": + return "采购需求已创建,请补充参考图。" + case "request-conflict": + return "来源已变化或当前商品不能创建采购需求。" + case "image-invalid": + return "参考图片无效,请选择 JPG、PNG 或 WebP 后重试。" + case "reference-conflict": + return "参考图已被使用或来源已变化,请刷新后重试。" + case "reference-bound": + return "参考图已绑定,可以生成采购任务。" + case "task-conflict": + return "需求状态或来源已变化,当前不能生成任务。" + default: + return "" + } +} + func SecurityHeaders() gin.HandlerFunc { return func(ctx *gin.Context) { ctx.Header( @@ -892,6 +1065,7 @@ type freightImportPage struct { type freightDetailPage struct { Page pageView Detail FreightOrderDetail + Notice string } type statusOption struct { diff --git a/backend-api/internal/transport/webui/handler_test.go b/backend-api/internal/transport/webui/handler_test.go index f86a549..90957f1 100644 --- a/backend-api/internal/transport/webui/handler_test.go +++ b/backend-api/internal/transport/webui/handler_test.go @@ -994,6 +994,83 @@ func TestFreightPagesEscapeSourceDataAndCreateAsyncSync(t *testing.T) { } } +func TestFreightDetailCreatesProcurementTaskWithCSRF(t *testing.T) { + const itemID = "00000000-0000-4000-8000-000000000002" + service := &fakeProcurementService{ + fakeFreightService: &fakeFreightService{ + fakeService: &fakeService{}, + orderDetail: FreightOrderDetail{ + Order: FreightOrder{ + ID: testTaskID, + ExternalStockID: "12", + SourceCode: "SOURCE-12", + }, + Items: []FreightItemReview{{ + Item: FreightOrderItem{ + ID: itemID, + ExternalItemID: "88", + Title: "商品", + SKU: "BLACK-L", + Quantity: 2, + }, + Request: &ProcurementRequest{ + ID: itemID, + FreightOrderItemID: itemID, + Status: "READY", + StatusLabel: "可以生成任务", + }, + }}, + }, + }, + createTaskResult: Task{ + ID: testTaskID, + Status: "PENDING", + }, + } + router := newTestRouter(t, service) + detail := performRequest( + t, + router, + http.MethodGet, + "/freight/"+testTaskID, + nil, + "", + ) + if detail.Code != http.StatusOK || + !strings.Contains(detail.Body.String(), "生成采购任务") || + !strings.Contains(detail.Body.String(), "可以生成任务") { + t.Fatalf("detail status/body = %d / %s", detail.Code, detail.Body) + } + cookie := csrfCookie(t, detail) + taskKey := hiddenValue(t, detail.Body.String(), "task_key") + values := url.Values{ + "csrf_token": {cookie.Value}, + "order_id": {testTaskID}, + "task_key": {taskKey}, + } + request := httptest.NewRequest( + http.MethodPost, + "/freight/procurement-requests/"+itemID+"/purchase-task", + strings.NewReader(values.Encode()), + ) + request.Header.Set("Content-Type", "application/x-www-form-urlencoded") + request.AddCookie(cookie) + response := httptest.NewRecorder() + router.ServeHTTP(response, request) + if response.Code != http.StatusSeeOther || + response.Header().Get("Location") != "/tasks/"+testTaskID { + t.Fatalf( + "create task status/location = %d / %q", + response.Code, + response.Header().Get("Location"), + ) + } + if service.createTaskInput.RequestID != itemID || + service.createTaskInput.IdempotencyKey != taskKey { + t.Fatalf("create task input = %+v", service.createTaskInput) + } +} + type fakeService struct { listInput ListTasksInput listResult TaskList @@ -1027,6 +1104,40 @@ type fakeFreightService struct { err error } +type fakeProcurementService struct { + *fakeFreightService + createRequestInput CreateProcurementRequestInput + bindInput BindProcurementReferenceInput + createTaskInput CreateProcurementTaskInput + procurementResult ProcurementRequest + createTaskResult Task + procurementError error +} + +func (service *fakeProcurementService) CreateProcurementRequest( + _ context.Context, + input CreateProcurementRequestInput, +) (ProcurementRequest, error) { + service.createRequestInput = input + return service.procurementResult, service.procurementError +} + +func (service *fakeProcurementService) BindProcurementReference( + _ context.Context, + input BindProcurementReferenceInput, +) (ProcurementRequest, error) { + service.bindInput = input + return service.procurementResult, service.procurementError +} + +func (service *fakeProcurementService) CreateProcurementTask( + _ context.Context, + input CreateProcurementTaskInput, +) (Task, error) { + service.createTaskInput = input + return service.createTaskResult, service.procurementError +} + func (service *fakeFreightService) ListFreightOrders( context.Context, int, diff --git a/backend-api/internal/transport/webui/static/admin.css b/backend-api/internal/transport/webui/static/admin.css index dc78496..4ec1d8b 100644 --- a/backend-api/internal/transport/webui/static/admin.css +++ b/backend-api/internal/transport/webui/static/admin.css @@ -721,6 +721,35 @@ tbody tr:last-child td { font-weight: 700; } +.procurement-actions { + min-width: 220px; +} + +.procurement-actions form { + display: grid; + gap: 9px; + margin-top: 8px; +} + +.procurement-actions input[type="file"] { + min-width: 0; + font-size: 12px; +} + +.confirm-line { + display: flex; + align-items: flex-start; + gap: 8px; + font-size: 13px; +} + +.confirm-line input { + width: 20px; + min-height: 20px; + flex: 0 0 auto; + margin-top: 1px; +} + .task-form { display: grid; gap: 18px; diff --git a/backend-api/internal/transport/webui/templates/freight-detail.gohtml b/backend-api/internal/transport/webui/templates/freight-detail.gohtml index 4de258d..fe6dd0c 100644 --- a/backend-api/internal/transport/webui/templates/freight-detail.gohtml +++ b/backend-api/internal/transport/webui/templates/freight-detail.gohtml @@ -15,6 +15,7 @@ 返回列表 + {{if .Notice}}
{{.Notice}}
{{end}}

来源信息

@@ -35,23 +36,73 @@ 数量 采购状态 来源版本 + 采购处理 {{range .Detail.Items}} - {{if .Title}}{{.Title}}{{else}}缺少标题{{end}} - 明细 ID:{{.ExternalItemID}} - {{if .ProductThumbRef}}图片引用:{{.ProductThumbRef}}{{end}} + {{if .Item.Title}}{{.Item.Title}}{{else}}缺少标题{{end}} + 明细 ID:{{.Item.ExternalItemID}} + {{if .Item.ProductThumbRef}}图片引用:{{.Item.ProductThumbRef}}{{end}} - {{if .ProductSpec}}{{.ProductSpec}}{{else}}未提供规格{{end}} - SKU:{{if .SKU}}{{.SKU}}{{else}}未提供{{end}} + {{if .Item.ProductSpec}}{{.Item.ProductSpec}}{{else}}未提供规格{{end}} + SKU:{{if .Item.SKU}}{{.Item.SKU}}{{else}}未提供{{end}} + + {{if .Item.Quantity}}{{.Item.Quantity}}{{else}}未提供{{end}} + {{if .Item.PurchaseStatus}}{{.Item.PurchaseStatus}}{{else}}未提供{{end}} + {{.Item.Revision}} + + {{if .Request}} + {{.Request.StatusLabel}} + {{if .Request.BlockingLabel}}{{.Request.BlockingLabel}}{{end}} + {{if .Request.SourceChanged}} +
+ + + + +
+ {{else if eq .Request.Status "NEEDS_IMAGE"}} +
+ + + + + +
+ {{else if eq .Request.Status "READY"}} +
+ + + + +
+ {{else if .Request.PurchaseTaskID}} + 查看采购任务 + {{end}} + {{else}} +
+ + + + +
+ {{end}} - {{if .Quantity}}{{.Quantity}}{{else}}未提供{{end}} - {{if .PurchaseStatus}}{{.PurchaseStatus}}{{else}}未提供{{end}} - {{.Revision}} {{end}} diff --git a/backend-api/internal/transport/webui/types.go b/backend-api/internal/transport/webui/types.go index f6c90cd..f64f16a 100644 --- a/backend-api/internal/transport/webui/types.go +++ b/backend-api/internal/transport/webui/types.go @@ -37,6 +37,39 @@ type FreightService interface { ) (FreightSync, error) } +type ProcurementService interface { + CreateProcurementRequest( + context.Context, + CreateProcurementRequestInput, + ) (ProcurementRequest, error) + BindProcurementReference( + context.Context, + BindProcurementReferenceInput, + ) (ProcurementRequest, error) + CreateProcurementTask( + context.Context, + CreateProcurementTaskInput, + ) (Task, error) +} + +type CreateProcurementRequestInput struct { + ActorUserID string + FreightOrderItemID string + ConfirmProcurementNeeded bool +} + +type BindProcurementReferenceInput struct { + ActorUserID string + RequestID string + ImageAssetID string +} + +type CreateProcurementTaskInput struct { + ActorUserID string + RequestID string + IdempotencyKey string +} + type FreightSync struct { ID string Status string @@ -81,7 +114,27 @@ type FreightOrderItem struct { type FreightOrderDetail struct { Order FreightOrder - Items []FreightOrderItem + Items []FreightItemReview +} + +type FreightItemReview struct { + Item FreightOrderItem + Request *ProcurementRequest +} + +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 } type ListTasksInput struct { diff --git a/backend-api/internal/transport/webui/usecase_adapter.go b/backend-api/internal/transport/webui/usecase_adapter.go index d257643..0c404ed 100644 --- a/backend-api/internal/transport/webui/usecase_adapter.go +++ b/backend-api/internal/transport/webui/usecase_adapter.go @@ -18,6 +18,13 @@ type UsecaseAdapter struct { assets *usecase.AssetService authorizations *usecase.OrderAuthorizationService freight *usecase.FreightService + procurement *usecase.ProcurementService +} + +func (adapter *UsecaseAdapter) SetProcurement( + procurement *usecase.ProcurementService, +) { + adapter.procurement = procurement } func NewUsecaseAdapter( @@ -73,7 +80,23 @@ func (adapter *UsecaseAdapter) GetFreightOrder( if err != nil { return FreightOrderDetail{}, mapUsecaseError(err) } - items := make([]FreightOrderItem, 0, len(detail.Items)) + requestByItem := map[string]domain.ProcurementRequest{} + if adapter.procurement != nil { + requests, requestErr := adapter.procurement.ListForOrder( + ctx, + localAdminSubject, + detail.Order.ID, + ) + if requestErr != nil { + return FreightOrderDetail{}, mapUsecaseError(requestErr) + } + for _, request := range requests { + if _, exists := requestByItem[request.FreightOrderItemID]; !exists { + requestByItem[request.FreightOrderItemID] = request + } + } + } + items := make([]FreightItemReview, 0, len(detail.Items)) for _, item := range detail.Items { view := FreightOrderItem{ ID: item.ID, @@ -88,7 +111,12 @@ func (adapter *UsecaseAdapter) GetFreightOrder( if item.Quantity != nil { view.Quantity = *item.Quantity } - items = append(items, view) + review := FreightItemReview{Item: view} + if request, exists := requestByItem[item.ID]; exists { + requestView := procurementRequestFrom(request) + review.Request = &requestView + } + items = append(items, review) } return FreightOrderDetail{ Order: freightOrderFrom(detail.Order), @@ -96,6 +124,131 @@ func (adapter *UsecaseAdapter) GetFreightOrder( }, nil } +func (adapter *UsecaseAdapter) CreateProcurementRequest( + ctx context.Context, + input CreateProcurementRequestInput, +) (ProcurementRequest, error) { + if adapter.procurement == nil { + return ProcurementRequest{}, ErrUnavailable + } + result, err := adapter.procurement.CreateRequest( + ctx, + usecase.CreateProcurementRequestCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: input.ActorUserID, + FreightOrderItemID: input.FreightOrderItemID, + ConfirmProcurementNeeded: input.ConfirmProcurementNeeded, + }, + ) + if err != nil { + return ProcurementRequest{}, mapUsecaseError(err) + } + return procurementRequestFrom(result.Request), nil +} + +func (adapter *UsecaseAdapter) BindProcurementReference( + ctx context.Context, + input BindProcurementReferenceInput, +) (ProcurementRequest, error) { + if adapter.procurement == nil { + return ProcurementRequest{}, ErrUnavailable + } + request, err := adapter.procurement.BindReference( + ctx, + usecase.BindProcurementReferenceCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: input.ActorUserID, + RequestID: input.RequestID, + ImageAssetID: input.ImageAssetID, + }, + ) + if err != nil { + return ProcurementRequest{}, mapUsecaseError(err) + } + return procurementRequestFrom(request), nil +} + +func (adapter *UsecaseAdapter) CreateProcurementTask( + ctx context.Context, + input CreateProcurementTaskInput, +) (Task, error) { + if adapter.procurement == nil { + return Task{}, ErrUnavailable + } + result, err := adapter.procurement.CreateTask( + ctx, + usecase.CreateProcurementTaskCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: input.ActorUserID, + RequestID: input.RequestID, + IdempotencyKey: input.IdempotencyKey, + }, + ) + if err != nil { + return Task{}, mapUsecaseError(err) + } + return taskFromPurchase(result.Task), nil +} + +func procurementRequestFrom( + request domain.ProcurementRequest, +) ProcurementRequest { + result := ProcurementRequest{ + ID: request.ID, + FreightOrderItemID: request.FreightOrderItemID, + SourceRevision: request.SourceRevision, + Status: string(request.Status), + StatusLabel: procurementStatusLabel(request.Status), + BlockingCode: stringValue(request.BlockingCode), + BlockingLabel: procurementBlockingLabel( + stringValue(request.BlockingCode), + ), + ReferenceAssetID: stringValue(request.ReferenceAssetID), + PurchaseTaskID: stringValue(request.PurchaseTaskID), + SourceChanged: request.SourceChanged, + } + if request.SourceChanged { + result.StatusLabel = "来源已变化" + } + return result +} + +func procurementStatusLabel( + status domain.ProcurementRequestStatus, +) string { + switch status { + case domain.ProcurementBlocked: + return "资料阻塞" + case domain.ProcurementNeedsImage: + return "需要参考图" + case domain.ProcurementReady: + return "可以生成任务" + case domain.ProcurementTaskCreated: + return "任务已生成" + case domain.ProcurementSourceChanged: + return "来源已变化" + default: + return "未知状态" + } +} + +func procurementBlockingLabel(code string) string { + switch code { + case domain.ProcurementBlockSourceCanceled: + return "货运单已取消" + case domain.ProcurementBlockTitleRequired: + return "商品标题缺失或超限" + case domain.ProcurementBlockSKURequired: + return "SKU 缺失或超限" + case domain.ProcurementBlockQuantityRequired: + return "数量缺失或无效" + case domain.ProcurementBlockSourceChanged: + return "ERP 来源已更新" + default: + return "" + } +} + func (adapter *UsecaseAdapter) GetFreightSync( ctx context.Context, syncID string, @@ -660,3 +813,4 @@ func (err *adapterError) Unwrap() []error { var _ Service = (*UsecaseAdapter)(nil) var _ FreightService = (*UsecaseAdapter)(nil) +var _ ProcurementService = (*UsecaseAdapter)(nil) diff --git a/backend-api/internal/usecase/procurement_ports.go b/backend-api/internal/usecase/procurement_ports.go new file mode 100644 index 0000000..edc257f --- /dev/null +++ b/backend-api/internal/usecase/procurement_ports.go @@ -0,0 +1,46 @@ +package usecase + +import ( + "context" + "time" + + "cmroubao/backend-api/internal/domain" +) + +type ProcurementRepository interface { + GetProcurementSourceItem( + context.Context, + string, + string, + ) (domain.ProcurementSourceItem, error) + CreateProcurementRequest( + context.Context, + domain.ProcurementRequest, + ) (domain.ProcurementRequest, bool, error) + ListProcurementRequestsForOrder( + context.Context, + string, + string, + ) ([]domain.ProcurementRequest, error) + GetProcurementRequest( + context.Context, + string, + string, + ) (domain.ProcurementRequest, error) + BindProcurementReference( + context.Context, + string, + string, + string, + time.Time, + ) (domain.ProcurementRequest, error) + CreateProcurementTask( + context.Context, + domain.ProcurementRequest, + domain.PurchaseTask, + domain.TaskEvent, + domain.PurchaseTaskSource, + string, + string, + ) (domain.PurchaseTask, bool, error) +} diff --git a/backend-api/internal/usecase/procurement_service.go b/backend-api/internal/usecase/procurement_service.go new file mode 100644 index 0000000..ad0c215 --- /dev/null +++ b/backend-api/internal/usecase/procurement_service.go @@ -0,0 +1,412 @@ +package usecase + +import ( + "context" + "errors" + "strconv" + "strings" + "unicode/utf8" + + "cmroubao/backend-api/internal/domain" +) + +var ( + ErrProcurementStateConflict = errors.New( + "procurement request state conflict", + ) + ErrProcurementSourceChanged = errors.New( + "procurement request source changed", + ) +) + +type ProcurementService struct { + repository ProcurementRepository + clock Clock + ids IDGenerator +} + +type CreateProcurementRequestCommand struct { + CreatorSubject string + ActorUserID string + FreightOrderItemID string + ConfirmProcurementNeeded bool +} + +type CreateProcurementRequestResult struct { + Request domain.ProcurementRequest + Replayed bool +} + +type BindProcurementReferenceCommand struct { + CreatorSubject string + ActorUserID string + RequestID string + ImageAssetID string +} + +type CreateProcurementTaskCommand struct { + CreatorSubject string + ActorUserID string + RequestID string + IdempotencyKey string +} + +type CreateProcurementTaskResult struct { + Task domain.PurchaseTask + Replayed bool +} + +func NewProcurementService( + repository ProcurementRepository, + clock Clock, + ids IDGenerator, +) (*ProcurementService, error) { + if repository == nil || clock == nil || ids == nil { + return nil, errors.New("procurement service dependencies are required") + } + return &ProcurementService{ + repository: repository, + clock: clock, + ids: ids, + }, nil +} + +func (service *ProcurementService) CreateRequest( + ctx context.Context, + command CreateProcurementRequestCommand, +) (CreateProcurementRequestResult, error) { + command.CreatorSubject = strings.TrimSpace(command.CreatorSubject) + command.ActorUserID = strings.TrimSpace(command.ActorUserID) + command.FreightOrderItemID = strings.TrimSpace( + command.FreightOrderItemID, + ) + fields := map[string]string{} + if command.CreatorSubject == "" { + fields["creator_subject"] = "is required" + } + if !isUUID(command.ActorUserID) { + fields["actor_user_id"] = "must be a UUID" + } + if !isUUID(command.FreightOrderItemID) { + fields["freight_order_item_id"] = "must be a UUID" + } + if !command.ConfirmProcurementNeeded { + fields["confirm_procurement_needed"] = + "must be explicitly confirmed" + } + if len(fields) > 0 { + return CreateProcurementRequestResult{}, invalidError( + "PROCUREMENT_REQUEST_INVALID", + "procurement request is invalid", + fields, + ) + } + source, err := service.repository.GetProcurementSourceItem( + ctx, + command.CreatorSubject, + command.FreightOrderItemID, + ) + if err != nil { + return CreateProcurementRequestResult{}, wrapRepositoryError(err) + } + requestID, err := service.ids.NewID() + if err != nil { + return CreateProcurementRequestResult{}, procurementInternal(err) + } + now := service.clock.Now().UTC() + request := domain.ProcurementRequest{ + ID: requestID, + CreatorSubject: command.CreatorSubject, + FreightOrderItemID: source.Item.ID, + SourceRevision: source.Item.Revision, + SourceSHA256: source.Item.CanonicalSHA256, + Title: strings.TrimSpace(source.Item.Title), + ProductSpec: strings.TrimSpace(source.Item.ProductSpec), + SKU: strings.TrimSpace(source.Item.SKU), + Quantity: source.Item.Quantity, + SourcePurchaseStatus: source.Item.PurchaseStatus, + SourceIsCanceled: source.IsCanceled, + ProcurementConfirmedByUserID: command.ActorUserID, + ProcurementConfirmedAt: now, + Status: domain.ProcurementNeedsImage, + CreatedAt: now, + UpdatedAt: now, + } + if code := procurementBlockingCode(request); code != "" { + request.Status = domain.ProcurementBlocked + request.BlockingCode = &code + } + stored, created, err := service.repository.CreateProcurementRequest( + ctx, + request, + ) + if err != nil { + return CreateProcurementRequestResult{}, wrapProcurementError(err) + } + return CreateProcurementRequestResult{ + Request: stored, + Replayed: !created, + }, nil +} + +func (service *ProcurementService) ListForOrder( + ctx context.Context, + creatorSubject, orderID string, +) ([]domain.ProcurementRequest, error) { + requests, err := service.repository.ListProcurementRequestsForOrder( + ctx, + strings.TrimSpace(creatorSubject), + strings.TrimSpace(orderID), + ) + if err != nil { + return nil, wrapProcurementError(err) + } + return requests, nil +} + +func (service *ProcurementService) Get( + ctx context.Context, + creatorSubject, requestID string, +) (domain.ProcurementRequest, error) { + request, err := service.repository.GetProcurementRequest( + ctx, + strings.TrimSpace(creatorSubject), + strings.TrimSpace(requestID), + ) + if err != nil { + return domain.ProcurementRequest{}, wrapProcurementError(err) + } + return request, nil +} + +func (service *ProcurementService) BindReference( + ctx context.Context, + command BindProcurementReferenceCommand, +) (domain.ProcurementRequest, error) { + command.CreatorSubject = strings.TrimSpace(command.CreatorSubject) + command.ActorUserID = strings.TrimSpace(command.ActorUserID) + command.RequestID = strings.TrimSpace(command.RequestID) + command.ImageAssetID = strings.TrimSpace(command.ImageAssetID) + fields := map[string]string{} + if command.CreatorSubject == "" { + fields["creator_subject"] = "is required" + } + if !isUUID(command.ActorUserID) { + fields["actor_user_id"] = "must be a UUID" + } + if !isUUID(command.RequestID) { + fields["procurement_request_id"] = "must be a UUID" + } + if !isUUID(command.ImageAssetID) { + fields["image_asset_id"] = "must be a UUID" + } + if len(fields) > 0 { + return domain.ProcurementRequest{}, invalidError( + "PROCUREMENT_REFERENCE_INVALID", + "procurement reference is invalid", + fields, + ) + } + request, err := service.repository.BindProcurementReference( + ctx, + command.CreatorSubject, + command.RequestID, + command.ImageAssetID, + service.clock.Now().UTC(), + ) + if err != nil { + return domain.ProcurementRequest{}, wrapProcurementError(err) + } + return request, nil +} + +func (service *ProcurementService) CreateTask( + ctx context.Context, + command CreateProcurementTaskCommand, +) (CreateProcurementTaskResult, error) { + if err := validateWriteIdentity( + command.CreatorSubject, + command.IdempotencyKey, + ); err != nil { + return CreateProcurementTaskResult{}, err + } + command.CreatorSubject = strings.TrimSpace(command.CreatorSubject) + command.ActorUserID = strings.TrimSpace(command.ActorUserID) + command.RequestID = strings.TrimSpace(command.RequestID) + if !isUUID(command.ActorUserID) || !isUUID(command.RequestID) { + return CreateProcurementTaskResult{}, invalidError( + "PROCUREMENT_TASK_INVALID", + "procurement task request is invalid", + map[string]string{ + "request": "actor and procurement request must be UUIDs", + }, + ) + } + request, err := service.repository.GetProcurementRequest( + ctx, + command.CreatorSubject, + command.RequestID, + ) + if err != nil { + return CreateProcurementTaskResult{}, wrapProcurementError(err) + } + if request.Status != domain.ProcurementTaskCreated && + request.SourceChanged { + return CreateProcurementTaskResult{}, wrapProcurementError( + ErrProcurementSourceChanged, + ) + } + if request.Status != domain.ProcurementReady && + request.Status != domain.ProcurementTaskCreated { + return CreateProcurementTaskResult{}, wrapProcurementError( + ErrProcurementStateConflict, + ) + } + if request.ReferenceAssetID == nil { + return CreateProcurementTaskResult{}, wrapProcurementError( + ErrProcurementStateConflict, + ) + } + taskID, err := service.ids.NewID() + if err != nil { + return CreateProcurementTaskResult{}, procurementInternal(err) + } + eventID, err := service.ids.NewID() + if err != nil { + return CreateProcurementTaskResult{}, procurementInternal(err) + } + now := service.clock.Now().UTC() + sourceRefValue := "erp-procurement:" + request.ID + + ":r" + strconv.Itoa(request.SourceRevision) + actorUserID := command.ActorUserID + description := "" + if request.ProductSpec != "" { + description = "ERP 规格:" + request.ProductSpec + } + task := domain.PurchaseTask{ + ID: taskID, + CreatorSubject: command.CreatorSubject, + CreatedByUserID: &actorUserID, + SourceRef: &sourceRefValue, + Title: request.Title, + Description: description, + SKU: request.SKU, + ImageAssetID: *request.ReferenceAssetID, + Quantity: *request.Quantity, + Currency: domain.CurrencyCNY, + Status: domain.TaskStatusPending, + Version: 1, + CreatedAt: now, + UpdatedAt: now, + } + if err := domain.ValidateTaskInput( + task.CreatorSubject, + task.SourceRef, + task.Title, + task.Description, + task.SKU, + task.ImageAssetID, + task.Quantity, + ); err != nil { + return CreateProcurementTaskResult{}, wrapProcurementError( + ErrProcurementStateConflict, + ) + } + event := domain.TaskEvent{ + ID: eventID, + TaskID: task.ID, + ActorUserID: &actorUserID, + Type: "TASK_CREATED", + Message: "task created from procurement request", + OccurredAt: now, + } + source := domain.PurchaseTaskSource{ + TaskID: task.ID, + ProcurementRequestID: request.ID, + FreightOrderItemID: request.FreightOrderItemID, + SourceRevision: request.SourceRevision, + SourceSHA256: request.SourceSHA256, + CreatedAt: now, + } + requestHash := hashJSON(struct { + RequestID string `json:"request_id"` + SourceSHA256 string `json:"source_sha256"` + }{ + RequestID: request.ID, + SourceSHA256: request.SourceSHA256, + }) + created, wasCreated, err := service.repository.CreateProcurementTask( + ctx, + request, + task, + event, + source, + strings.TrimSpace(command.IdempotencyKey), + requestHash, + ) + if err != nil { + return CreateProcurementTaskResult{}, wrapProcurementError(err) + } + return CreateProcurementTaskResult{ + Task: created, + Replayed: !wasCreated, + }, nil +} + +func procurementBlockingCode( + request domain.ProcurementRequest, +) string { + if request.SourceIsCanceled != nil && *request.SourceIsCanceled { + return domain.ProcurementBlockSourceCanceled + } + if request.Title == "" || + utf8.RuneCountInString(request.Title) > domain.MaxTitleRunes || + len([]byte(request.Title)) > domain.MaxTitleBytes { + return domain.ProcurementBlockTitleRequired + } + if request.SKU == "" || + len([]byte(request.SKU)) > domain.MaxSKUBytes { + return domain.ProcurementBlockSKURequired + } + if request.Quantity == nil || *request.Quantity <= 0 { + return domain.ProcurementBlockQuantityRequired + } + return "" +} + +func wrapProcurementError(err error) error { + switch { + case errors.Is(err, ErrProcurementSourceChanged): + return newError( + ErrorKindConflict, + "PROCUREMENT_SOURCE_CHANGED", + "freight source changed; create a request for the current revision", + err, + ) + case errors.Is(err, ErrProcurementStateConflict): + return newError( + ErrorKindConflict, + "PROCUREMENT_STATE_CONFLICT", + "procurement request state does not allow this operation", + err, + ) + case errors.Is(err, ErrAssetUnavailable): + return newError( + ErrorKindConflict, + "PROCUREMENT_ASSET_UNAVAILABLE", + "reference image is not available for this procurement request", + err, + ) + default: + return wrapRepositoryError(err) + } +} + +func procurementInternal(err error) error { + return newError( + ErrorKindInternal, + "INTERNAL_ERROR", + "internal server error", + err, + ) +} diff --git a/backend-api/migrations/00013_procurement_requests.sql b/backend-api/migrations/00013_procurement_requests.sql new file mode 100644 index 0000000..73b074b --- /dev/null +++ b/backend-api/migrations/00013_procurement_requests.sql @@ -0,0 +1,132 @@ +-- +goose Up +CREATE TABLE procurement_requests ( + id TEXT PRIMARY KEY NOT NULL CHECK (length(id) = 36), + creator_subject TEXT NOT NULL, + freight_order_item_id TEXT NOT NULL + REFERENCES freight_order_items(id) + ON UPDATE RESTRICT ON DELETE RESTRICT, + source_revision INTEGER NOT NULL CHECK (source_revision >= 1), + source_sha256 TEXT NOT NULL + CHECK ( + length(source_sha256) = 64 + AND source_sha256 NOT GLOB '*[^0-9a-f]*' + ), + title TEXT NOT NULL + CHECK (length(CAST(title AS BLOB)) <= 2048), + product_spec TEXT NOT NULL + CHECK (length(CAST(product_spec AS BLOB)) <= 1024), + sku TEXT NOT NULL + CHECK (length(CAST(sku AS BLOB)) <= 512), + quantity INTEGER CHECK (quantity IS NULL OR quantity > 0), + source_purchase_status TEXT + CHECK ( + source_purchase_status IS NULL + OR length(CAST(source_purchase_status AS BLOB)) <= 128 + ), + source_is_canceled INTEGER + CHECK ( + source_is_canceled IS NULL + OR source_is_canceled IN (0, 1) + ), + procurement_confirmed_by_user_id TEXT NOT NULL + REFERENCES users(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + procurement_confirmed_at TEXT NOT NULL, + status TEXT NOT NULL + CHECK ( + status IN ( + 'BLOCKED', + 'NEEDS_IMAGE', + 'READY', + 'TASK_CREATED', + 'SOURCE_CHANGED' + ) + ), + blocking_code TEXT + CHECK ( + blocking_code IS NULL + OR blocking_code IN ( + 'SOURCE_CANCELED', + 'TITLE_REQUIRED', + 'SKU_REQUIRED', + 'QUANTITY_REQUIRED', + 'SOURCE_CHANGED' + ) + ), + reference_asset_id TEXT UNIQUE + REFERENCES assets(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + purchase_task_id TEXT UNIQUE + REFERENCES purchase_tasks(id) + ON UPDATE RESTRICT ON DELETE RESTRICT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + UNIQUE (freight_order_item_id, source_revision), + CHECK ( + (status = 'BLOCKED' + AND blocking_code IS NOT NULL + AND reference_asset_id IS NULL + AND purchase_task_id IS NULL) + OR (status = 'NEEDS_IMAGE' + AND blocking_code IS NULL + AND reference_asset_id IS NULL + AND purchase_task_id IS NULL) + OR (status = 'READY' + AND blocking_code IS NULL + AND reference_asset_id IS NOT NULL + AND purchase_task_id IS NULL) + OR (status = 'TASK_CREATED' + AND blocking_code IS NULL + AND reference_asset_id IS NOT NULL + AND purchase_task_id IS NOT NULL) + OR (status = 'SOURCE_CHANGED' + AND blocking_code = 'SOURCE_CHANGED' + AND purchase_task_id IS NULL) + ) +); + +CREATE INDEX procurement_requests_creator_updated_idx + ON procurement_requests (creator_subject, updated_at DESC, id DESC); + +CREATE INDEX procurement_requests_item_created_idx + ON procurement_requests ( + freight_order_item_id, + source_revision DESC, + id DESC + ); + +CREATE TABLE purchase_task_sources ( + task_id TEXT PRIMARY KEY NOT NULL + REFERENCES purchase_tasks(id) + ON UPDATE RESTRICT ON DELETE RESTRICT, + procurement_request_id TEXT NOT NULL UNIQUE + REFERENCES procurement_requests(id) + ON UPDATE RESTRICT ON DELETE RESTRICT, + freight_order_item_id TEXT NOT NULL + REFERENCES freight_order_items(id) + ON UPDATE RESTRICT ON DELETE RESTRICT, + source_revision INTEGER NOT NULL CHECK (source_revision >= 1), + source_sha256 TEXT NOT NULL + CHECK ( + length(source_sha256) = 64 + AND source_sha256 NOT GLOB '*[^0-9a-f]*' + ), + created_at TEXT NOT NULL +); + +-- +goose Down +CREATE TEMP TABLE procurement_v13_down_guard ( + allowed INTEGER NOT NULL CHECK (allowed = 1) +); + +INSERT INTO procurement_v13_down_guard (allowed) +SELECT CASE + WHEN EXISTS (SELECT 1 FROM procurement_requests) + OR EXISTS (SELECT 1 FROM purchase_task_sources) + THEN 0 + ELSE 1 +END; + +DROP TABLE procurement_v13_down_guard; +DROP TABLE purchase_task_sources; +DROP INDEX procurement_requests_item_created_idx; +DROP INDEX procurement_requests_creator_updated_idx; +DROP TABLE procurement_requests; diff --git a/docs/api.md b/docs/api.md index be7dbe2..4b92847 100644 --- a/docs/api.md +++ b/docs/api.md @@ -257,18 +257,41 @@ Connector 未登录、找不到货运单、响应协议错误和暂时不可用 ### `POST /api/v1/freight-items/{item_id}/procurement-request` -T-223 创建或重放该商品当前 revision 的采购需求。标题、SKU/规格或数量不合法时返回 -`422 FREIGHT_ITEM_NOT_PURCHASABLE`;不调用 VLM 补字段。 +T-223 创建或重放该商品当前 revision 的采购需求。ERP 数字状态不自动解释,ADMIN +必须显式确认仍需采购: + +```json +{"confirm_procurement_needed":true} +``` + +请求快照保留 title、product spec、SKU、quantity、来源 hash/revision 和当时的取消/ +采购状态。缺标题、SKU、正整数数量或来源明确取消时创建为 `BLOCKED`,不调用 VLM +补字段;字段合格但没有受控参考图时为 `NEEDS_IMAGE`。 ### `PUT /api/v1/procurement-requests/{id}/reference-asset` 绑定当前 ADMIN 已上传的 `TASK_REFERENCE` asset。asset 必须尚未属于其他任务/请求, -内容仍走既有解码、像素、规范化和 hash 校验。 +内容仍走既有解码、像素、规范化和 hash 校验。请求体: + +```json +{"image_asset_id":"2bbc1bf2-30f7-497e-a52b-bcb4c61d57c3"} +``` + +绑定成功后进入 `READY`。来源 revision/hash 或货运取消标记已变化时返回 +`409 PROCUREMENT_SOURCE_CHANGED`。 ### `POST /api/v1/procurement-requests/{id}/purchase-task` -必须带 `Idempotency-Key`。只有 `READY` request 可生成任务;同一 request revision -并发或重试返回同一 task。来源后续更新不修改已生成任务。 +必须带 `Idempotency-Key`,请求体是空 JSON 对象 `{}`。只有 `READY` request 可生成 +任务;task、`purchase_task_sources` 和 request 的 `TASK_CREATED` 状态在同一事务内 +提交。同一 request revision 即使使用不同幂等 key 重试也返回同一 task。来源后续 +更新只显示 `source_changed=true`,不修改 PENDING、CLAIMED、RUNNING 或终态任务。 + +管理 Web 在 `/freight/{id}` 内提供对应的人工确认、参考图上传和生成任务操作;写操作 +分别使用 `/freight/items/{id}/procurement-request`、 +`/freight/procurement-requests/{id}/reference` 和 +`/freight/procurement-requests/{id}/purchase-task`,均受 ADMIN session 与 CSRF +保护。 ## 采购任务 diff --git a/docs/current-state.md b/docs/current-state.md index af5d1dc..7286890 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -5,7 +5,7 @@ ## 当前快照 - 日期:2026-07-28 -- 阶段:T-222 已完成,T-223 待采购需求提取与任务生成已领取 +- 阶段:T-223 已完成,T-224 ERP 日期增量同步已领取 - Git:当前分支为 `main`;T-001 至 T-004、T-101 至 T-104、T-201 至 T-219 均按文档提交、实现提交的顺序纳入历史 - 生产代码:`android-buyer/` 已接入 Roubao Android 源码 @@ -17,14 +17,16 @@ - ERP 货运:Go 后端通过仅限 loopback、服务密钥鉴权的 Python Connector 异步按 完整单号同步;v12 保存同步记录、货运头和全部明细,canonical hash 控制 revision, Admin 已有 `/freight`、`/freight/import`、`/freight/{id}` 与对应 JSON API。 +- 采购需求:v13 按货运商品 revision 保存可复核快照;ADMIN 显式确认采购、上传参考 + 图后生成不可变 PENDING 任务。来源变化不会改写任务,Roubao 可沿用现有 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-222 运行 `go test ./...`、`go test -race ./...`、`go vet ./...`; - 覆盖货运 migration 降级保护、严格 Connector 字段白名单、重复导入/revision、 - 中断恢复、Admin API/SSR 鉴权和 PII 缺失断言 +- 后端测试:T-223 运行 `go test ./...`、`go test -race ./...`、`go vet ./...`; + 覆盖采购需求 migration 降级保护、一商品一需求、阻塞字段、图片唯一绑定、任务 + 并发重放、来源变化后任务不变、Roubao claim 和 Admin API/SSR - 原型:4 个管理 Web 页面和 7 个 Android 页面均可离线独立打开;Playwright 以 1440×900、390×844、360×800 验证 36 个页面/视口组合,无页面横向溢出、 脚本错误或外部请求,Android 可见交互控件均不小于 44px @@ -162,8 +164,8 @@ | `docs/tasks/T-220.md` | DONE | 顺运宝 ERP 字段契约与凭证安全基线 | | `docs/tasks/T-221.md` | DONE | 顺运宝精确单号 Connector | | `docs/tasks/T-222.md` | DONE | 货运信息存储、API 与 Admin 页面 | -| `docs/tasks/T-223.md` | DOING | 待采购需求提取与任务生成 | -| `docs/tasks/T-224.md` | TODO | ERP 日期增量同步 | +| `docs/tasks/T-223.md` | DONE | 待采购需求提取与任务生成 | +| `docs/tasks/T-224.md` | DOING | ERP 日期增量同步 | | `docs/design/` | 已确认 | T-202 原型索引、4 个管理页和 7 个 Android 页面 | | `deepseek总结.txt` | 已有 | 历史讨论摘要,不是正式需求权威 | | `android-buyer/` | 已有 | Roubao `main` 固定 commit 的 Android 基线 | @@ -176,8 +178,8 @@ ## 任务摘要 - 已完成:T-001 至 T-004、T-101 至 T-104、T-201 至 T-219。 -- 已完成:另含 T-220 至 T-222 ERP 契约、Connector、货运存储/API/Admin。 -- 进行中:T-223 待采购需求提取与任务生成。 +- 已完成:另含 T-220 至 T-223 ERP 契约、Connector、货运存储和采购需求生成。 +- 进行中:T-224 ERP 日期增量同步。 - 下一步:完成 T-223 采购需求生成和 T-224 日期增量同步,再进入 T-301 P0 UI 完整交互验收。 diff --git a/docs/tasks/T-223.md b/docs/tasks/T-223.md index 4829bdc..8239e0e 100644 --- a/docs/tasks/T-223.md +++ b/docs/tasks/T-223.md @@ -4,7 +4,7 @@ title: 待采购需求提取与任务生成 phase: 2 deps: - T-222 -status: TODO +status: DONE created: 2026-07-28 context_ref: 78dc595 work_branch: null @@ -50,11 +50,11 @@ Roubao 队列。 ## 验收要点 -- [ ] 一货运单多商品形成多条独立采购需求。 -- [ ] 缺标题/SKU/数量/图片和来源状态不明时不能生成任务。 -- [ ] 同一 revision 并发/重试只生成一个采购任务。 -- [ ] 来源变化不覆写 PENDING/CLAIMED/RUNNING/终态任务。 -- [ ] 生成任务可被现有 Roubao claim,原始字段和参考图 hash 一致。 +- [x] 一货运单多商品形成多条独立采购需求。 +- [x] 缺标题/SKU/数量/图片、明确取消或未人工确认时不能生成任务。 +- [x] 同一 revision 并发/重试只生成一个采购任务。 +- [x] 来源变化不覆写 PENDING/CLAIMED/RUNNING/终态任务。 +- [x] 生成任务可被现有 Roubao claim,原始字段和参考图 hash 一致。 ## 边界 @@ -65,3 +65,10 @@ Roubao 队列。 ## 执行记录 - 2026-07-28:任务合约已冻结,等待 T-222。 +- 2026-07-28:新增 v13 `procurement_requests` 与不可变 + `purchase_task_sources`;来源数字状态不猜测,由 ADMIN 显式确认是否仍需采购。 +- 2026-07-28:参考图绑定、任务和来源快照采用唯一约束与单事务;ERP 更新后旧任务 + 保持原始标题、SKU、数量和 asset,旧请求显示来源变化,当前 revision 可新建需求。 +- 2026-07-28:真实 Gin/SQLite SSR 流程完成确认需求、上传规范化参考图、生成 PENDING + 任务;同任务可由既有 Roubao claim。Playwright 三种视口无横向溢出、短按钮或页头 + 重叠;未调用 VLM 或真实 ERP。