diff --git a/backend-api/cmd/api/main.go b/backend-api/cmd/api/main.go index 3a5a8a5..9dbb304 100644 --- a/backend-api/cmd/api/main.go +++ b/backend-api/cmd/api/main.go @@ -255,12 +255,22 @@ func buildRouter( if err != nil { return nil, err } + freightImages, err := usecase.NewFreightImageService( + store, + erpSession, + files, + clock, + ) + if err != nil { + return nil, err + } freight, err := usecase.NewFreightService( store, erpSession, clock, ids, 90*time.Second, + usecase.WithFreightImageCache(freightImages), ) if err != nil { return nil, err @@ -354,6 +364,7 @@ func buildRouter( Results: results, Authorizations: authorizations, Freight: freight, + FreightImages: freightImages, Procurement: procurement, }, webHandler, diff --git a/backend-api/internal/domain/freight.go b/backend-api/internal/domain/freight.go index 94d4eb7..38df43f 100644 --- a/backend-api/internal/domain/freight.go +++ b/backend-api/internal/domain/freight.go @@ -18,6 +18,16 @@ const ( FreightSyncFailed FreightSyncStatus = "FAILED" ) +type FreightItemImageStatus string + +const ( + FreightItemImageNone FreightItemImageStatus = "NONE" + FreightItemImagePending FreightItemImageStatus = "PENDING" + FreightItemImageReady FreightItemImageStatus = "READY" + FreightItemImageMissing FreightItemImageStatus = "MISSING" + FreightItemImageFailed FreightItemImageStatus = "FAILED" +) + type FreightSyncRun struct { ID string CreatorSubject string @@ -69,6 +79,8 @@ type FreightOrderItem struct { ProductThumbRef *string OriginalUnitPriceMinor *int64 OriginalCurrency string + ImageStatus FreightItemImageStatus + ImageErrorCode *string PurchaseStatus *string CanonicalSHA256 string Revision int @@ -79,6 +91,26 @@ type FreightOrderItem struct { UpdatedAt time.Time } +type FreightItemImageJob struct { + CreatorSubject string + FreightOrderItemID string + ProductThumbRef string +} + +type FreightItemImage struct { + CreatorSubject string + FreightOrderItemID string + ProductThumbRef string + Status FreightItemImageStatus + MediaType string + SizeBytes int64 + SHA256 string + StorageKey string + ErrorCode *string + AttemptCount int + UpdatedAt time.Time +} + type FreightOrderDetail struct { Order FreightOrder Items []FreightOrderItem diff --git a/backend-api/internal/domain/freight_source_errors.go b/backend-api/internal/domain/freight_source_errors.go index 454380f..f386f10 100644 --- a/backend-api/internal/domain/freight_source_errors.go +++ b/backend-api/internal/domain/freight_source_errors.go @@ -12,4 +12,6 @@ var ( ErrFreightSourceProtocol = errors.New("freight source protocol is invalid") ErrFreightSourceOCRInvalid = errors.New("freight source OCR service is invalid") ErrFreightSourceLoginRejected = errors.New("freight source login was rejected") + ErrFreightImageNotFound = errors.New("freight source image was not found") + ErrFreightImageInvalid = errors.New("freight source image is invalid") ) diff --git a/backend-api/internal/platform/migration/claims_migration_test.go b/backend-api/internal/platform/migration/claims_migration_test.go index b855fad..c7e0d79 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 != 15 { - t.Fatalf("initial Up() applied = %d, want 15", applied) + } else if applied != 16 { + t.Fatalf("initial Up() applied = %d, want 16", applied) + } + if err := runner.Down(ctx); err != nil { + t.Fatalf("initial Down(v16) error = %v", err) } if err := runner.Down(ctx); err != nil { t.Fatalf("initial Down(v15) error = %v", err) @@ -74,9 +77,14 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { seedClaimsHistoricalFixture(t, db) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("Up(v5-v15) over historical data error = %v", err) - } else if applied != 11 { - t.Fatalf("Up(v5-v15) applied = %d, want 11", applied) + 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) + } + assertClaimsHistory(t, db, true) + + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v16) with compatible history error = %v", err) } assertClaimsHistory(t, db, true) @@ -141,9 +149,9 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { assertClaimsHistory(t, db, false) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("final Up(v4-v15) error = %v", err) - } else if applied != 12 { - t.Fatalf("final Up(v4-v15) applied = %d, want 12", applied) + t.Fatalf("final Up(v4-v16) error = %v", err) + } else if applied != 13 { + t.Fatalf("final Up(v4-v16) applied = %d, want 13", applied) } assertClaimsHistory(t, db, true) } @@ -379,6 +387,9 @@ func TestClaimsMigrationDownFailsClosedForNewAuditData(t *testing.T) { t.Fatalf("insert v4 audit event: %v", err) } + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v16) error = %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("Down(v15) error = %v", err) } diff --git a/backend-api/internal/platform/migration/runner_test.go b/backend-api/internal/platform/migration/runner_test.go index 9371778..baef96e 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 != 15 { - t.Fatalf("Up() applied = %d, want 15", applied) + if applied != 16 { + t.Fatalf("Up() applied = %d, want 16", applied) } assertStatuses(t, runner, map[int64]bool{ 1: true, @@ -46,6 +46,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 13: true, 14: true, 15: true, + 16: true, }) applied, err = runner.Up(context.Background()) @@ -74,7 +75,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 12: true, 13: true, 14: true, - 15: false, + 15: true, + 16: false, }) applied, err = runner.Up(context.Background()) @@ -100,6 +102,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 13: true, 14: true, 15: true, + 16: true, }) } diff --git a/backend-api/internal/platform/shunyunbao/image.go b/backend-api/internal/platform/shunyunbao/image.go new file mode 100644 index 0000000..bd288fe --- /dev/null +++ b/backend-api/internal/platform/shunyunbao/image.go @@ -0,0 +1,136 @@ +package shunyunbao + +import ( + "context" + "mime" + "net/http" + "net/url" + "strconv" + "strings" + + "cmroubao/backend-api/internal/domain" + "cmroubao/backend-api/internal/usecase" +) + +const ( + ProductImagePath = "/api/p/file" + maxProductImageBytes = 20 << 20 +) + +func (manager *SessionManager) FetchProductImage( + ctx context.Context, + productThumbRef string, +) (usecase.FreightSourceImage, error) { + productThumbRef = strings.TrimSpace(productThumbRef) + value, err := strconv.ParseUint(productThumbRef, 10, 64) + if err != nil || value == 0 || + strconv.FormatUint(value, 10) != productThumbRef { + return usecase.FreightSourceImage{}, domain.ErrFreightImageInvalid + } + + manager.mu.Lock() + if !manager.configuredLocked() { + manager.mu.Unlock() + return usecase.FreightSourceImage{}, + domain.ErrFreightSourceNotConfigured + } + if !manager.authenticated { + manager.mu.Unlock() + return usecase.FreightSourceImage{}, + domain.ErrFreightSourceSessionNeeded + } + baseURL := manager.baseURL + headers := manager.headers.Clone() + client := manager.http + manager.mu.Unlock() + + query := url.Values{} + query.Set("id", productThumbRef) + request, err := http.NewRequestWithContext( + ctx, + http.MethodGet, + baseURL+ProductImagePath+"?"+query.Encode(), + nil, + ) + if err != nil { + return usecase.FreightSourceImage{}, + domain.ErrFreightSourceUnavailable + } + for name, values := range headers { + request.Header[name] = append([]string(nil), values...) + } + request.Header.Set("Accept", "image/jpeg, image/png, image/webp") + manager.logERPRequest(request) + response, err := client.Do(request) + if err != nil { + manager.logERPTransportFailure(request) + return usecase.FreightSourceImage{}, + domain.ErrFreightSourceUnavailable + } + if response.StatusCode < http.StatusOK || + response.StatusCode >= http.StatusMultipleChoices { + defer response.Body.Close() + manager.logERPBinaryResponse(request, response) + switch response.StatusCode { + case http.StatusUnauthorized, http.StatusForbidden: + manager.mu.Lock() + manager.clearAuthenticatedLocked() + manager.mu.Unlock() + return usecase.FreightSourceImage{}, + domain.ErrFreightSourceSessionNeeded + case http.StatusNotFound: + return usecase.FreightSourceImage{}, + domain.ErrFreightImageNotFound + default: + return usecase.FreightSourceImage{}, + domain.ErrFreightSourceUnavailable + } + } + mediaType, _, err := mime.ParseMediaType( + response.Header.Get("Content-Type"), + ) + if err != nil || !supportedProductImageType(mediaType) || + response.ContentLength > maxProductImageBytes { + response.Body.Close() + manager.logERPBinaryResponse(request, response) + return usecase.FreightSourceImage{}, + domain.ErrFreightImageInvalid + } + manager.logERPBinaryResponse(request, response) + return usecase.FreightSourceImage{ + Content: response.Body, + MediaType: mediaType, + }, nil +} + +func supportedProductImageType(value string) bool { + switch strings.ToLower(value) { + case "image/jpeg", "image/png", "image/webp": + return true + default: + return false + } +} + +func (manager *SessionManager) logERPBinaryResponse( + request *http.Request, + response *http.Response, +) { + if !manager.diagnosticsOn { + return + } + byteCount := "unknown" + if response.ContentLength >= 0 { + byteCount = strconv.FormatInt(response.ContentLength, 10) + } + manager.diagnosticLog( + "erp_response method=" + request.Method + + " path=" + request.URL.EscapedPath() + + " status=" + strconv.Itoa(response.StatusCode) + + " content_type=" + + diagnosticContentType(response.Header.Get("Content-Type")) + + " bytes=" + byteCount + " body=omitted_binary", + ) +} + +var _ usecase.FreightImageSource = (*SessionManager)(nil) diff --git a/backend-api/internal/platform/shunyunbao/image_test.go b/backend-api/internal/platform/shunyunbao/image_test.go new file mode 100644 index 0000000..b15a820 --- /dev/null +++ b/backend-api/internal/platform/shunyunbao/image_test.go @@ -0,0 +1,190 @@ +package shunyunbao + +import ( + "context" + "errors" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "cmroubao/backend-api/internal/domain" +) + +func TestSessionManagerFetchProductImageUsesFixedAuthenticatedEndpoint( + t *testing.T, +) { + var requestPath, rawQuery, accept string + manager, closeServer := imageSessionManager( + t, + func(writer http.ResponseWriter, request *http.Request) { + requestPath = request.URL.Path + rawQuery = request.URL.RawQuery + accept = request.Header.Get("Accept") + if cookie, err := request.Cookie("authenticated"); err != nil || + cookie.Value != "yes" { + t.Fatalf("image cookie = %v, %v", cookie, err) + } + writer.Header().Set("Content-Type", "image/png") + writer.Header().Set("Content-Length", "4") + _, _ = writer.Write([]byte("png!")) + }, + ) + defer closeServer() + + image, err := manager.FetchProductImage(context.Background(), "190") + if err != nil { + t.Fatalf("FetchProductImage() error = %v", err) + } + defer image.Content.Close() + content, err := io.ReadAll(image.Content) + if err != nil { + t.Fatalf("read image: %v", err) + } + if requestPath != ProductImagePath || rawQuery != "id=190" || + image.MediaType != "image/png" || string(content) != "png!" || + !strings.Contains(accept, "image/jpeg") { + t.Fatalf( + "image request/result = %q / %q / %q / %q / %q", + requestPath, + rawQuery, + accept, + image.MediaType, + content, + ) + } +} + +func TestSessionManagerFetchProductImageClassifiesFailures(t *testing.T) { + manager, closeServer := imageSessionManager( + t, + func(writer http.ResponseWriter, request *http.Request) { + switch request.URL.Query().Get("id") { + case "190": + writer.WriteHeader(http.StatusNotFound) + case "191": + writer.Header().Set("Content-Type", "application/json") + _, _ = writer.Write([]byte(`{"status":false}`)) + case "192": + writer.Header().Set("Location", ProductImagePath+"?id=190") + writer.WriteHeader(http.StatusFound) + case "193": + writer.Header().Set( + "Content-Type", + "image/jpeg", + ) + writer.Header().Set( + "Content-Length", + "20971521", + ) + default: + writer.WriteHeader(http.StatusUnauthorized) + } + }, + ) + defer closeServer() + + tests := []struct { + ref string + want error + }{ + {"190", domain.ErrFreightImageNotFound}, + {"191", domain.ErrFreightImageInvalid}, + {"192", domain.ErrFreightSourceUnavailable}, + {"193", domain.ErrFreightImageInvalid}, + {"194", domain.ErrFreightSourceSessionNeeded}, + } + for _, test := range tests { + image, err := manager.FetchProductImage( + context.Background(), + test.ref, + ) + if image.Content != nil { + image.Content.Close() + } + if !errors.Is(err, test.want) { + t.Fatalf("FetchProductImage(%q) error = %v", test.ref, err) + } + } + if _, err := manager.FetchProductImage( + context.Background(), + "195", + ); !errors.Is(err, domain.ErrFreightSourceSessionNeeded) { + t.Fatalf("request after unauthorized error = %v", err) + } +} + +func TestSessionManagerFetchProductImageRejectsUntrustedReferences( + t *testing.T, +) { + manager := testSessionManager(t, "http://127.0.0.1:1", "user", "password") + for _, value := range []string{ + "", + "0", + "0190", + "-1", + "1.5", + "https://invalid.example/image", + } { + if _, err := manager.FetchProductImage( + context.Background(), + value, + ); !errors.Is(err, domain.ErrFreightImageInvalid) { + t.Fatalf("FetchProductImage(%q) error = %v", value, err) + } + } +} + +func imageSessionManager( + t *testing.T, + imageHandler http.HandlerFunc, +) (*SessionManager, func()) { + t.Helper() + server := httptest.NewServer(http.HandlerFunc(func( + writer http.ResponseWriter, + request *http.Request, + ) { + switch request.URL.Path { + case CaptchaPath: + http.SetCookie( + writer, + &http.Cookie{Name: "captcha", Value: "ready", Path: "/"}, + ) + writer.Header().Set("Content-Type", "image/png") + _, _ = writer.Write([]byte("captcha")) + case LoginPath: + http.SetCookie( + writer, + &http.Cookie{ + Name: "authenticated", + Value: "yes", + Path: "/", + }, + ) + _, _ = writer.Write( + []byte( + `{"status":true,"data":{"user":{"id":1,"username":"test-user"}}}`, + ), + ) + case UserPath: + _, _ = writer.Write( + []byte( + `{"status":true,"data":{"id":1,"username":"test-user"}}`, + ), + ) + case ProductImagePath: + imageHandler(writer, request) + default: + writer.WriteHeader(http.StatusNotFound) + } + })) + manager := testSessionManager( + t, + server.URL, + "test-user", + "test-password", + ) + loginForSource(t, manager) + return manager, server.Close +} diff --git a/backend-api/internal/repository/sqlite/auth_repository_test.go b/backend-api/internal/repository/sqlite/auth_repository_test.go index 5a84550..98656c5 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(v16) error = %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("Down(v15) error = %v", err) } @@ -435,9 +438,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-v15) error = %v", err) - } else if applied != 13 { - t.Fatalf("Up(v3-v15) applied = %d, want 13", applied) + t.Fatalf("Up(v3-v16) error = %v", err) + } else if applied != 14 { + t.Fatalf("Up(v3-v16) applied = %d, want 14", applied) } } diff --git a/backend-api/internal/repository/sqlite/freight_image_repository.go b/backend-api/internal/repository/sqlite/freight_image_repository.go new file mode 100644 index 0000000..3205c32 --- /dev/null +++ b/backend-api/internal/repository/sqlite/freight_image_repository.go @@ -0,0 +1,275 @@ +package sqlite + +import ( + "context" + "database/sql" + "errors" + + "cmroubao/backend-api/internal/domain" + "cmroubao/backend-api/internal/usecase" +) + +func (store *Store) ListFreightImageJobs( + ctx context.Context, + creatorSubject, runID string, + limit int, +) ([]domain.FreightItemImageJob, error) { + rows, err := store.db.QueryContext( + ctx, + `SELECT item.id, item.product_thumb_ref + FROM freight_order_items AS item + JOIN freight_orders AS freight + ON freight.id = item.freight_order_id + LEFT JOIN freight_item_images AS image + ON image.freight_order_item_id = item.id + WHERE freight.creator_subject = ? + AND item.last_sync_run_id = ? + AND item.is_present = 1 + AND item.product_thumb_ref IS NOT NULL + AND ( + image.freight_order_item_id IS NULL + OR image.product_thumb_ref != item.product_thumb_ref + OR image.status != 'READY' + ) + ORDER BY item.id + LIMIT ?`, + creatorSubject, + runID, + limit, + ) + if err != nil { + return nil, repositoryFailure(err) + } + defer rows.Close() + jobs := make([]domain.FreightItemImageJob, 0) + for rows.Next() { + var job domain.FreightItemImageJob + job.CreatorSubject = creatorSubject + if err := rows.Scan( + &job.FreightOrderItemID, + &job.ProductThumbRef, + ); err != nil { + return nil, repositoryFailure(err) + } + jobs = append(jobs, job) + } + if err := rows.Err(); err != nil { + return nil, repositoryFailure(err) + } + return jobs, nil +} + +func (store *Store) SaveFreightItemImage( + ctx context.Context, + candidate domain.FreightItemImage, +) (*string, error) { + tx, err := store.db.BeginTx(ctx, nil) + if err != nil { + return nil, repositoryFailure(err) + } + defer tx.Rollback() + + var currentThumb sql.NullString + err = tx.QueryRowContext( + ctx, + `SELECT item.product_thumb_ref + 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`, + candidate.CreatorSubject, + candidate.FreightOrderItemID, + ).Scan(¤tThumb) + if errors.Is(err, sql.ErrNoRows) { + return nil, usecase.ErrRepositoryNotFound + } + if err != nil { + return nil, repositoryFailure(err) + } + if !currentThumb.Valid || + currentThumb.String != candidate.ProductThumbRef { + return nil, usecase.ErrTaskStateConflict + } + + var existingThumb string + var existingStorage sql.NullString + var existingAttempts int + err = tx.QueryRowContext( + ctx, + `SELECT product_thumb_ref, storage_key, attempt_count + FROM freight_item_images + WHERE freight_order_item_id = ?`, + candidate.FreightOrderItemID, + ).Scan(&existingThumb, &existingStorage, &existingAttempts) + switch { + case errors.Is(err, sql.ErrNoRows): + existingAttempts = 0 + case err != nil: + return nil, repositoryFailure(err) + case existingThumb != candidate.ProductThumbRef: + existingAttempts = 0 + } + attemptCount := existingAttempts + 1 + _, err = tx.ExecContext( + ctx, + `INSERT INTO freight_item_images ( + freight_order_item_id, product_thumb_ref, status, media_type, + size_bytes, sha256, storage_key, error_code, attempt_count, + updated_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (freight_order_item_id) + DO UPDATE SET + product_thumb_ref = excluded.product_thumb_ref, + status = excluded.status, + media_type = excluded.media_type, + size_bytes = excluded.size_bytes, + sha256 = excluded.sha256, + storage_key = excluded.storage_key, + error_code = excluded.error_code, + attempt_count = excluded.attempt_count, + updated_at = excluded.updated_at`, + candidate.FreightOrderItemID, + candidate.ProductThumbRef, + candidate.Status, + nullableFreightImageText(candidate.MediaType), + nullableFreightImageSize(candidate), + nullableFreightImageText(candidate.SHA256), + nullableFreightImageText(candidate.StorageKey), + nullableString(candidate.ErrorCode), + attemptCount, + formatTimestamp(candidate.UpdatedAt), + ) + if err != nil { + return nil, repositoryFailure(err) + } + if err := tx.Commit(); err != nil { + return nil, repositoryFailure(err) + } + if existingStorage.Valid && + existingStorage.String != candidate.StorageKey { + return &existingStorage.String, nil + } + return nil, nil +} + +func (store *Store) GetReadyFreightItemImage( + ctx context.Context, + creatorSubject, itemID string, +) (domain.FreightItemImage, error) { + var image domain.FreightItemImage + var updatedAt string + err := store.db.QueryRowContext( + ctx, + `SELECT freight.creator_subject, image.freight_order_item_id, + image.product_thumb_ref, image.status, image.media_type, + image.size_bytes, image.sha256, image.storage_key, + image.attempt_count, image.updated_at + 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'`, + creatorSubject, + itemID, + ).Scan( + &image.CreatorSubject, + &image.FreightOrderItemID, + &image.ProductThumbRef, + &image.Status, + &image.MediaType, + &image.SizeBytes, + &image.SHA256, + &image.StorageKey, + &image.AttemptCount, + &updatedAt, + ) + if errors.Is(err, sql.ErrNoRows) { + return domain.FreightItemImage{}, usecase.ErrRepositoryNotFound + } + if err != nil { + return domain.FreightItemImage{}, repositoryFailure(err) + } + image.UpdatedAt, err = parseTimestamp(updatedAt) + if err != nil { + return domain.FreightItemImage{}, repositoryFailure(err) + } + return image, nil +} + +func (store *Store) applyFreightItemImageStatuses( + ctx context.Context, + orderID string, + items []domain.FreightOrderItem, +) error { + byID := make(map[string]int, len(items)) + for index := range items { + byID[items[index].ID] = index + if items[index].ProductThumbRef == nil { + items[index].ImageStatus = domain.FreightItemImageNone + } else { + items[index].ImageStatus = domain.FreightItemImagePending + } + } + rows, err := store.db.QueryContext( + ctx, + `SELECT item.id, image.product_thumb_ref, image.status, + image.error_code + FROM freight_order_items AS item + JOIN freight_item_images AS image + ON image.freight_order_item_id = item.id + WHERE item.freight_order_id = ? + AND item.is_present = 1`, + orderID, + ) + if err != nil { + return repositoryFailure(err) + } + defer rows.Close() + for rows.Next() { + var itemID, imageThumb string + var status domain.FreightItemImageStatus + var errorCode sql.NullString + if err := rows.Scan( + &itemID, + &imageThumb, + &status, + &errorCode, + ); err != nil { + return repositoryFailure(err) + } + index, exists := byID[itemID] + if !exists || items[index].ProductThumbRef == nil || + *items[index].ProductThumbRef != imageThumb { + continue + } + items[index].ImageStatus = status + items[index].ImageErrorCode = optionalString(errorCode) + } + if err := rows.Err(); err != nil { + return repositoryFailure(err) + } + return nil +} + +func nullableFreightImageText(value string) any { + if value == "" { + return nil + } + return value +} + +func nullableFreightImageSize(image domain.FreightItemImage) any { + if image.Status != domain.FreightItemImageReady { + return nil + } + return image.SizeBytes +} + +var _ usecase.FreightImageRepository = (*Store)(nil) diff --git a/backend-api/internal/repository/sqlite/freight_image_repository_test.go b/backend-api/internal/repository/sqlite/freight_image_repository_test.go new file mode 100644 index 0000000..3ded4ec --- /dev/null +++ b/backend-api/internal/repository/sqlite/freight_image_repository_test.go @@ -0,0 +1,196 @@ +package sqlite_test + +import ( + "context" + "testing" + "time" + + "cmroubao/backend-api/internal/domain" + "cmroubao/backend-api/internal/platform/migration" + repository "cmroubao/backend-api/internal/repository/sqlite" +) + +func TestFreightItemImagesAreCurrentRetryableAndRollbackGuarded( + t *testing.T, +) { + db := openDatabase(t) + store, err := repository.New(db) + if err != nil { + t.Fatalf("repository.New() error = %v", err) + } + ctx := context.Background() + now := time.Date(2026, 7, 29, 3, 0, 0, 0, time.UTC) + userID := uuid(1000) + seedFreightUser(t, db, userID, now) + + first := freightRun(1001, userID, now) + createAndStartFreightRun(t, store, first, "image-run-1", "1") + batch := freightBatch(1010, "a", "b") + firstThumb := "190" + secondThumb := "191" + batch.Orders[0].Items[0].ProductThumbRef = &firstThumb + batch.Orders[0].Items[1].ProductThumbRef = &secondThumb + if err := store.CompleteFreightSync( + ctx, + first, + batch, + now.Add(time.Second), + ); err != nil { + t.Fatalf("CompleteFreightSync() error = %v", err) + } + jobs, err := store.ListFreightImageJobs( + ctx, + "local-admin", + first.ID, + 10, + ) + if err != nil || len(jobs) != 2 { + t.Fatalf("initial image jobs = %+v, %v", jobs, err) + } + byThumb := make(map[string]domain.FreightItemImageJob, len(jobs)) + for _, job := range jobs { + byThumb[job.ProductThumbRef] = job + } + ready := domain.FreightItemImage{ + CreatorSubject: "local-admin", + FreightOrderItemID: byThumb[firstThumb].FreightOrderItemID, + ProductThumbRef: firstThumb, + Status: domain.FreightItemImageReady, + MediaType: "image/jpeg", + SizeBytes: 100, + SHA256: repeatHex("c"), + StorageKey: "aa/first.jpg", + UpdatedAt: now.Add(2 * time.Second), + } + replaced, err := store.SaveFreightItemImage(ctx, ready) + if err != nil || replaced != nil { + t.Fatalf("save ready image = %v, %v", replaced, err) + } + errorCode := "ERP_IMAGE_NOT_FOUND" + failed := domain.FreightItemImage{ + CreatorSubject: "local-admin", + FreightOrderItemID: byThumb[secondThumb].FreightOrderItemID, + ProductThumbRef: secondThumb, + Status: domain.FreightItemImageMissing, + ErrorCode: &errorCode, + UpdatedAt: now.Add(2 * time.Second), + } + if _, err := store.SaveFreightItemImage(ctx, failed); err != nil { + t.Fatalf("save missing image: %v", err) + } + + jobs, err = store.ListFreightImageJobs( + ctx, + "local-admin", + first.ID, + 10, + ) + if err != nil || len(jobs) != 1 || + jobs[0].ProductThumbRef != secondThumb { + t.Fatalf("retry jobs = %+v, %v", jobs, err) + } + orders, _ := store.ListFreightOrders(ctx, "local-admin", 10) + detail, err := store.GetFreightOrder( + ctx, + "local-admin", + orders[0].ID, + ) + if err != nil || len(detail.Items) != 2 { + t.Fatalf("detail = %+v, %v", detail, err) + } + statuses := map[string]domain.FreightItemImageStatus{} + for _, item := range detail.Items { + statuses[*item.ProductThumbRef] = item.ImageStatus + } + if statuses[firstThumb] != domain.FreightItemImageReady || + statuses[secondThumb] != domain.FreightItemImageMissing { + t.Fatalf("image statuses = %+v", statuses) + } + stored, err := store.GetReadyFreightItemImage( + ctx, + "local-admin", + ready.FreightOrderItemID, + ) + if err != nil || stored.StorageKey != ready.StorageKey || + stored.AttemptCount != 1 { + t.Fatalf("ready image = %+v, %v", stored, err) + } + if _, err := store.SaveFreightItemImage( + ctx, + failed, + ); err != nil { + t.Fatalf("retry missing image: %v", err) + } + var attempts int + if err := db.QueryRow( + `SELECT attempt_count FROM freight_item_images + WHERE freight_order_item_id = ?`, + failed.FreightOrderItemID, + ).Scan(&attempts); err != nil || attempts != 2 { + t.Fatalf("missing attempts = %d, %v", attempts, err) + } + + second := freightRun(1002, userID, now.Add(time.Minute)) + createAndStartFreightRun(t, store, second, "image-run-2", "2") + changed := freightBatch(1020, "d", "b") + changedThumb := "192" + changed.Orders[0].Items[0].ProductThumbRef = &changedThumb + changed.Orders[0].Items[1].ProductThumbRef = &secondThumb + if err := store.CompleteFreightSync( + ctx, + second, + changed, + now.Add(time.Minute+time.Second), + ); err != nil { + t.Fatalf("changed CompleteFreightSync() error = %v", err) + } + jobs, err = store.ListFreightImageJobs( + ctx, + "local-admin", + second.ID, + 10, + ) + if err != nil || len(jobs) != 2 { + t.Fatalf("changed image jobs = %+v, %v", jobs, err) + } + var changedJob domain.FreightItemImageJob + for _, job := range jobs { + if job.ProductThumbRef == changedThumb { + changedJob = job + } + } + changedReady := ready + changedReady.FreightOrderItemID = changedJob.FreightOrderItemID + changedReady.ProductThumbRef = changedThumb + changedReady.StorageKey = "bb/changed.jpg" + changedReady.SHA256 = repeatHex("d") + changedReady.UpdatedAt = now.Add(time.Minute + 2*time.Second) + replaced, err = store.SaveFreightItemImage(ctx, changedReady) + if err != nil || replaced == nil || *replaced != ready.StorageKey { + t.Fatalf("replace ready image = %v, %v", replaced, err) + } + stored, err = store.GetReadyFreightItemImage( + ctx, + "local-admin", + ready.FreightOrderItemID, + ) + if err != nil || stored.ProductThumbRef != changedThumb || + stored.AttemptCount != 1 { + t.Fatalf("changed ready image = %+v, %v", stored, 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("image migration down succeeded with retained image metadata") + } + if _, err := store.GetReadyFreightItemImage( + ctx, + "local-admin", + ready.FreightOrderItemID, + ); err != nil { + t.Fatalf("failed rollback lost ready image: %v", err) + } +} diff --git a/backend-api/internal/repository/sqlite/freight_repository.go b/backend-api/internal/repository/sqlite/freight_repository.go index 17d0a5e..360eae6 100644 --- a/backend-api/internal/repository/sqlite/freight_repository.go +++ b/backend-api/internal/repository/sqlite/freight_repository.go @@ -548,6 +548,13 @@ func (store *Store) GetFreightOrder( if err := rows.Err(); err != nil { return domain.FreightOrderDetail{}, repositoryFailure(err) } + if err := store.applyFreightItemImageStatuses( + ctx, + order.ID, + items, + ); err != nil { + return domain.FreightOrderDetail{}, err + } return domain.FreightOrderDetail{Order: order, Items: items}, nil } diff --git a/backend-api/internal/repository/sqlite/freight_repository_test.go b/backend-api/internal/repository/sqlite/freight_repository_test.go index fbc3454..5acb15e 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("image migration down: %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("metadata migration down: %v", err) } @@ -268,6 +271,9 @@ func TestFreightDateSyncAdvancesWatermarkOnlyOnWholeBatchSuccess( } runner, _ := migration.New(db) + if err := runner.Down(ctx); err != nil { + t.Fatalf("image migration down: %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("metadata migration down: %v", err) } diff --git a/backend-api/internal/repository/sqlite/procurement_repository_test.go b/backend-api/internal/repository/sqlite/procurement_repository_test.go index 6ba5cba..05a3bb1 100644 --- a/backend-api/internal/repository/sqlite/procurement_repository_test.go +++ b/backend-api/internal/repository/sqlite/procurement_repository_test.go @@ -238,6 +238,9 @@ func TestProcurementRequestsArePerItemAndTaskSnapshotIsImmutable( if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(ctx); err != nil { + t.Fatalf("image migration down: %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("metadata migration down: %v", err) } diff --git a/backend-api/internal/transport/httpapi/admin_handlers.go b/backend-api/internal/transport/httpapi/admin_handlers.go index 09bd838..eb26117 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers.go +++ b/backend-api/internal/transport/httpapi/admin_handlers.go @@ -30,6 +30,7 @@ type AdminServices struct { Results *usecase.ExecutionResultService Authorizations *usecase.OrderAuthorizationService Freight *usecase.FreightService + FreightImages *usecase.FreightImageService Procurement *usecase.ProcurementService ERP *shunyunbao.SessionManager } @@ -75,6 +76,12 @@ 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.FreightImages != nil { + routes.GET( + "/api/v1/freight-items/:id/image", + handler.freightItemImage, + ) + } if services.Procurement != nil { routes.POST( "/api/v1/freight-items/:id/procurement-request", @@ -92,6 +99,29 @@ func registerAdminAPI(routes gin.IRoutes, services AdminServices) error { return nil } +func (h *adminHandlers) freightItemImage(ctx *gin.Context) { + result, err := h.services.FreightImages.OpenItemImage( + ctx.Request.Context(), + localAdminSubject, + ctx.Param("id"), + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + defer result.Content.Close() + ctx.Header("Cache-Control", "private, no-store") + ctx.Header("Content-Type", result.Image.MediaType) + ctx.Header( + "Content-Length", + strconv.FormatInt(result.Image.SizeBytes, 10), + ) + ctx.Header("ETag", `"`+result.Image.SHA256+`"`) + ctx.Header("X-Content-Type-Options", "nosniff") + ctx.Status(http.StatusOK) + _, _ = io.Copy(ctx.Writer, result.Content) +} + func (h *adminHandlers) evidenceContent(ctx *gin.Context) { result, err := h.services.Results.OpenEvidence( ctx.Request.Context(), diff --git a/backend-api/internal/transport/httpapi/admin_handlers_test.go b/backend-api/internal/transport/httpapi/admin_handlers_test.go index 7ac9afe..380cae0 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers_test.go +++ b/backend-api/internal/transport/httpapi/admin_handlers_test.go @@ -479,6 +479,41 @@ func TestAdminFreightAPIImportsAllItemsWithoutPII(t *testing.T) { detail.Header().Get("Cache-Control") != "no-store" { t.Fatalf("freight detail status/body = %d / %s", detail.Code, detail.Body) } + var detailBody struct { + Items []struct { + ID string `json:"id"` + ImageStatus string `json:"image_status"` + ImageURL string `json:"image_url"` + } `json:"items"` + } + decodeResponse(t, detail, &detailBody) + if len(detailBody.Items) != 2 || + detailBody.Items[0].ImageStatus != "READY" || + detailBody.Items[0].ImageURL == "" { + t.Fatalf("freight detail images = %+v", detailBody.Items) + } + imageResponse := performAdminRequest( + t, + router, + http.MethodGet, + detailBody.Items[0].ImageURL, + "", + nil, + "", + ) + if imageResponse.Code != http.StatusOK || + imageResponse.Header().Get("Content-Type") != "image/jpeg" || + imageResponse.Header().Get("Cache-Control") != "private, no-store" || + imageResponse.Header().Get("X-Content-Type-Options") != "nosniff" || + imageResponse.Header().Get("ETag") == "" || + !bytes.HasPrefix(imageResponse.Body.Bytes(), []byte{0xff, 0xd8}) { + t.Fatalf( + "freight image status/headers/body = %d / %#v / %x", + imageResponse.Code, + imageResponse.Header(), + imageResponse.Body.Bytes(), + ) + } replay := performAdminRequest( t, router, @@ -875,6 +910,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("image migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("metadata migration down: %v", err) } @@ -1211,12 +1249,22 @@ func newAdminIntegrationFixture(t *testing.T) *adminIntegrationFixture { if err != nil { t.Fatalf("usecase.NewOrderAuthorizationService() error = %v", err) } + freightImages, err := usecase.NewFreightImageService( + repositories, + staticFreightSource{}, + files, + clock, + ) + if err != nil { + t.Fatalf("usecase.NewFreightImageService() error = %v", err) + } freight, err := usecase.NewFreightService( repositories, staticFreightSource{}, clock, ids, time.Second, + usecase.WithFreightImageCache(freightImages), ) if err != nil { t.Fatalf("usecase.NewFreightService() error = %v", err) @@ -1236,6 +1284,7 @@ func newAdminIntegrationFixture(t *testing.T) *adminIntegrationFixture { Results: results, Authorizations: authorizations, Freight: freight, + FreightImages: freightImages, Procurement: procurement, }, emptyAdminWeb{}, @@ -1270,6 +1319,8 @@ func (staticFreightSource) QueryOrder( quantityTwo := 2 priceOne := int64(12950) priceTwo := int64(8800) + thumbOne := "190" + thumbTwo := "191" return domain.FreightSourceBatch{ SchemaVersion: 1, Query: domain.FreightSourceQuery{ @@ -1289,6 +1340,7 @@ func (staticFreightSource) QueryOrder( Quantity: &quantityOne, OriginalUnitPriceMinor: &priceOne, OriginalCurrency: domain.FreightCurrencyTWD, + ProductThumbRef: &thumbOne, }, { ExternalItemID: "89", @@ -1298,12 +1350,41 @@ func (staticFreightSource) QueryOrder( Quantity: &quantityTwo, OriginalUnitPriceMinor: &priceTwo, OriginalCurrency: domain.FreightCurrencyTWD, + ProductThumbRef: &thumbTwo, }, }, }}, }, nil } +func (staticFreightSource) FetchProductImage( + context.Context, + string, +) (usecase.FreightSourceImage, error) { + var imageBytes bytes.Buffer + source := image.NewRGBA(image.Rect(0, 0, 8, 6)) + for y := 0; y < 6; y++ { + for x := 0; x < 8; x++ { + source.Set( + x, + y, + color.RGBA{R: uint8(x * 20), G: 80, B: 160, A: 255}, + ) + } + } + if err := jpeg.Encode( + &imageBytes, + source, + &jpeg.Options{Quality: 85}, + ); err != nil { + return usecase.FreightSourceImage{}, err + } + return usecase.FreightSourceImage{ + Content: io.NopCloser(bytes.NewReader(imageBytes.Bytes())), + MediaType: "image/jpeg", + }, nil +} + func (staticFreightSource) QueryCreatedRange( _ context.Context, createdFrom, createdTo string, diff --git a/backend-api/internal/transport/httpapi/device_handlers_test.go b/backend-api/internal/transport/httpapi/device_handlers_test.go index 1dd2a74..87b2a27 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("image migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("metadata migration down: %v", err) } @@ -739,8 +742,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 != 7 { - t.Fatalf("restored migrations = %d, want 7", applied) + } else if applied != 8 { + t.Fatalf("restored migrations = %d, want 8", applied) } completePayload := fmt.Sprintf( @@ -1429,6 +1432,9 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable( if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("image migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("metadata 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 16855b4..e9bdf55 100644 --- a/backend-api/internal/transport/httpapi/freight_handlers.go +++ b/backend-api/internal/transport/httpapi/freight_handlers.go @@ -205,6 +205,10 @@ func (h *adminHandlers) freightOrderDetail(ctx *gin.Context) { } items := make([]gin.H, 0, len(detail.Items)) for _, item := range detail.Items { + var imageURL any + if item.ImageStatus == domain.FreightItemImageReady { + imageURL = "/api/v1/freight-items/" + item.ID + "/image" + } items = append(items, gin.H{ "id": item.ID, "external_item_id": item.ExternalItemID, @@ -215,6 +219,9 @@ func (h *adminHandlers) freightOrderDetail(ctx *gin.Context) { "product_thumb_ref": item.ProductThumbRef, "original_unit_price_minor": item.OriginalUnitPriceMinor, "original_currency": item.OriginalCurrency, + "image_status": item.ImageStatus, + "image_error_code": item.ImageErrorCode, + "image_url": imageURL, "purchase_status": item.PurchaseStatus, "revision": item.Revision, "canonical_sha256": item.CanonicalSHA256, diff --git a/backend-api/internal/transport/webui/types.go b/backend-api/internal/transport/webui/types.go index 0d3801c..f8c5307 100644 --- a/backend-api/internal/transport/webui/types.go +++ b/backend-api/internal/transport/webui/types.go @@ -159,6 +159,9 @@ type FreightOrderItem struct { ProductThumbRef string OriginalUnitPriceMinor *int64 OriginalCurrency string + ImageStatus string + ImageErrorCode string + ImageURL string PurchaseStatus string Revision int } diff --git a/backend-api/internal/transport/webui/usecase_adapter.go b/backend-api/internal/transport/webui/usecase_adapter.go index 9ae571a..1cea79e 100644 --- a/backend-api/internal/transport/webui/usecase_adapter.go +++ b/backend-api/internal/transport/webui/usecase_adapter.go @@ -168,12 +168,17 @@ func (adapter *UsecaseAdapter) GetFreightOrder( ProductThumbRef: stringValue(item.ProductThumbRef), OriginalUnitPriceMinor: item.OriginalUnitPriceMinor, OriginalCurrency: item.OriginalCurrency, + ImageStatus: string(item.ImageStatus), + ImageErrorCode: stringValue(item.ImageErrorCode), PurchaseStatus: stringValue(item.PurchaseStatus), Revision: item.Revision, } if item.Quantity != nil { view.Quantity = *item.Quantity } + if item.ImageStatus == domain.FreightItemImageReady { + view.ImageURL = "/api/v1/freight-items/" + item.ID + "/image" + } review := FreightItemReview{Item: view} if request, exists := requestByItem[item.ID]; exists { requestView := procurementRequestFrom(request) diff --git a/backend-api/internal/usecase/freight_image_ports.go b/backend-api/internal/usecase/freight_image_ports.go new file mode 100644 index 0000000..e67780f --- /dev/null +++ b/backend-api/internal/usecase/freight_image_ports.go @@ -0,0 +1,42 @@ +package usecase + +import ( + "context" + "io" + + "cmroubao/backend-api/internal/domain" +) + +type FreightSourceImage struct { + Content io.ReadCloser + MediaType string +} + +type FreightImageSource interface { + FetchProductImage( + context.Context, + string, + ) (FreightSourceImage, error) +} + +type FreightImageRepository interface { + ListFreightImageJobs( + context.Context, + string, + string, + int, + ) ([]domain.FreightItemImageJob, error) + SaveFreightItemImage( + context.Context, + domain.FreightItemImage, + ) (*string, error) + GetReadyFreightItemImage( + context.Context, + string, + string, + ) (domain.FreightItemImage, error) +} + +type FreightImageCache interface { + CacheRun(context.Context, string, string) error +} diff --git a/backend-api/internal/usecase/freight_image_service.go b/backend-api/internal/usecase/freight_image_service.go new file mode 100644 index 0000000..eab52dc --- /dev/null +++ b/backend-api/internal/usecase/freight_image_service.go @@ -0,0 +1,270 @@ +package usecase + +import ( + "context" + "errors" + "io" + "sync" + "time" + + "cmroubao/backend-api/internal/domain" +) + +const ( + maxFreightImageJobsPerRun = 1000 + freightImageWorkers = 4 + freightImageCleanup = 5 * time.Second +) + +type FreightImageService struct { + repository FreightImageRepository + source FreightImageSource + store ReferenceImageStore + clock Clock +} + +type FreightItemImageContent struct { + Image domain.FreightItemImage + Content io.ReadCloser +} + +func NewFreightImageService( + repository FreightImageRepository, + source FreightImageSource, + store ReferenceImageStore, + clock Clock, +) (*FreightImageService, error) { + if repository == nil || source == nil || store == nil || clock == nil { + return nil, errors.New("freight image service dependencies are required") + } + return &FreightImageService{ + repository: repository, + source: source, + store: store, + clock: clock, + }, nil +} + +func (service *FreightImageService) CacheRun( + ctx context.Context, + creatorSubject, runID string, +) error { + jobs, err := service.repository.ListFreightImageJobs( + ctx, + creatorSubject, + runID, + maxFreightImageJobsPerRun, + ) + if err != nil { + return wrapRepositoryError(err) + } + if len(jobs) == 0 { + return nil + } + workerCount := freightImageWorkers + if len(jobs) < workerCount { + workerCount = len(jobs) + } + queue := make(chan domain.FreightItemImageJob) + var workers sync.WaitGroup + workers.Add(workerCount) + for range workerCount { + go func() { + defer workers.Done() + for job := range queue { + if ctx.Err() != nil { + return + } + service.cacheJob(ctx, job) + } + }() + } + for _, job := range jobs { + select { + case queue <- job: + case <-ctx.Done(): + close(queue) + workers.Wait() + return nil + } + } + close(queue) + workers.Wait() + return nil +} + +func (service *FreightImageService) cacheJob( + ctx context.Context, + job domain.FreightItemImageJob, +) { + sourceImage, err := service.source.FetchProductImage( + ctx, + job.ProductThumbRef, + ) + if err != nil { + if ctx.Err() == nil { + service.saveFailure(ctx, job, freightImageSourceFailure(err)) + } + return + } + defer sourceImage.Content.Close() + normalized, err := service.store.Put( + ctx, + job.FreightOrderItemID, + sourceImage.MediaType, + sourceImage.Content, + ) + if err != nil { + if ctx.Err() == nil { + service.saveFailure(ctx, job, freightImageStoreFailure(err)) + } + return + } + candidate := domain.FreightItemImage{ + CreatorSubject: job.CreatorSubject, + FreightOrderItemID: job.FreightOrderItemID, + ProductThumbRef: job.ProductThumbRef, + Status: domain.FreightItemImageReady, + MediaType: normalized.MediaType, + SizeBytes: normalized.SizeBytes, + SHA256: normalized.SHA256, + StorageKey: normalized.StorageKey, + UpdatedAt: service.clock.Now().UTC(), + } + replaced, err := service.repository.SaveFreightItemImage(ctx, candidate) + if err != nil { + service.deleteStoredImage(normalized.StorageKey) + return + } + if replaced != nil { + service.deleteStoredImage(*replaced) + } +} + +func (service *FreightImageService) saveFailure( + ctx context.Context, + job domain.FreightItemImageJob, + failure freightImageFailure, +) { + errorCode := failure.code + replaced, err := service.repository.SaveFreightItemImage( + ctx, + domain.FreightItemImage{ + CreatorSubject: job.CreatorSubject, + FreightOrderItemID: job.FreightOrderItemID, + ProductThumbRef: job.ProductThumbRef, + Status: failure.status, + ErrorCode: &errorCode, + UpdatedAt: service.clock.Now().UTC(), + }, + ) + if err == nil && replaced != nil { + service.deleteStoredImage(*replaced) + } +} + +func (service *FreightImageService) deleteStoredImage(storageKey string) { + ctx, cancel := context.WithTimeout(context.Background(), freightImageCleanup) + defer cancel() + _ = service.store.Delete(ctx, storageKey) +} + +func (service *FreightImageService) OpenItemImage( + ctx context.Context, + creatorSubject, itemID string, +) (FreightItemImageContent, error) { + if creatorSubject == "" || !isUUID(itemID) { + return FreightItemImageContent{}, freightImageNotFoundError(nil) + } + image, err := service.repository.GetReadyFreightItemImage( + ctx, + creatorSubject, + itemID, + ) + if err != nil { + wrapped := wrapRepositoryError(err) + var typed *Error + if errors.As(wrapped, &typed) && + typed.Kind == ErrorKindNotFound { + return FreightItemImageContent{}, + freightImageNotFoundError(err) + } + return FreightItemImageContent{}, wrapped + } + content, err := service.store.Open(ctx, image.StorageKey) + if err != nil { + var storeError *ImageStoreError + if errors.As(err, &storeError) && + storeError.Kind == ImageStoreErrorNotFound { + return FreightItemImageContent{}, + freightImageNotFoundError(err) + } + return FreightItemImageContent{}, mapImageStoreError(err) + } + return FreightItemImageContent{Image: image, Content: content}, nil +} + +type freightImageFailure struct { + status domain.FreightItemImageStatus + code string +} + +func freightImageSourceFailure(err error) freightImageFailure { + switch { + case errors.Is(err, domain.ErrFreightImageNotFound): + return freightImageFailure{ + status: domain.FreightItemImageMissing, + code: "ERP_IMAGE_NOT_FOUND", + } + case errors.Is(err, domain.ErrFreightImageInvalid): + return freightImageFailure{ + status: domain.FreightItemImageFailed, + code: "ERP_IMAGE_INVALID", + } + case errors.Is(err, domain.ErrFreightSourceSessionNeeded): + return freightImageFailure{ + status: domain.FreightItemImageFailed, + code: "ERP_SESSION_REQUIRED", + } + case errors.Is(err, domain.ErrFreightSourceNotConfigured): + return freightImageFailure{ + status: domain.FreightItemImageFailed, + code: "ERP_NOT_CONFIGURED", + } + default: + return freightImageFailure{ + status: domain.FreightItemImageFailed, + code: "ERP_IMAGE_UNAVAILABLE", + } + } +} + +func freightImageStoreFailure(err error) freightImageFailure { + var storeError *ImageStoreError + if errors.As(err, &storeError) { + switch storeError.Kind { + case ImageStoreErrorTooLarge, + ImageStoreErrorUnsupported, + ImageStoreErrorInvalid: + return freightImageFailure{ + status: domain.FreightItemImageFailed, + code: "ERP_IMAGE_INVALID", + } + } + } + return freightImageFailure{ + status: domain.FreightItemImageFailed, + code: "IMAGE_STORE_UNAVAILABLE", + } +} + +func freightImageNotFoundError(cause error) error { + return newError( + ErrorKindNotFound, + "FREIGHT_IMAGE_NOT_FOUND", + "freight item image not found", + cause, + ) +} + +var _ FreightImageCache = (*FreightImageService)(nil) diff --git a/backend-api/internal/usecase/freight_image_service_test.go b/backend-api/internal/usecase/freight_image_service_test.go new file mode 100644 index 0000000..2644fbc --- /dev/null +++ b/backend-api/internal/usecase/freight_image_service_test.go @@ -0,0 +1,219 @@ +package usecase + +import ( + "bytes" + "context" + "errors" + "io" + "strings" + "sync" + "testing" + + "cmroubao/backend-api/internal/domain" +) + +func TestFreightImageCacheIsBestEffortAndCleansReplacedFiles(t *testing.T) { + repository := &freightImageRepositoryFake{ + jobs: []domain.FreightItemImageJob{ + { + CreatorSubject: "local-admin", + FreightOrderItemID: "00000000-0000-4000-8000-000000000001", + ProductThumbRef: "190", + }, + { + CreatorSubject: "local-admin", + FreightOrderItemID: "00000000-0000-4000-8000-000000000002", + ProductThumbRef: "191", + }, + }, + replaced: "old/image.jpg", + } + source := &freightImageSourceFake{} + store := &freightImageStoreFake{} + service, err := NewFreightImageService( + repository, + source, + store, + fakeClock{}, + ) + if err != nil { + t.Fatalf("NewFreightImageService() error = %v", err) + } + if err := service.CacheRun( + context.Background(), + "local-admin", + "00000000-0000-4000-8000-000000000099", + ); err != nil { + t.Fatalf("CacheRun() error = %v", err) + } + repository.mu.Lock() + defer repository.mu.Unlock() + if len(repository.saved) != 2 { + t.Fatalf("saved images = %+v", repository.saved) + } + statusByRef := map[string]domain.FreightItemImage{} + for _, image := range repository.saved { + statusByRef[image.ProductThumbRef] = image + } + if statusByRef["190"].Status != domain.FreightItemImageReady || + statusByRef["190"].StorageKey != "new/image.jpg" || + statusByRef["191"].Status != domain.FreightItemImageMissing || + statusByRef["191"].ErrorCode == nil || + *statusByRef["191"].ErrorCode != "ERP_IMAGE_NOT_FOUND" { + t.Fatalf("saved image statuses = %+v", statusByRef) + } + store.mu.Lock() + defer store.mu.Unlock() + if len(store.deleted) != 1 || store.deleted[0] != "old/image.jpg" { + t.Fatalf("deleted storage keys = %#v", store.deleted) + } +} + +func TestFreightImageOpenDoesNotRevealUnavailableItems(t *testing.T) { + repository := &freightImageRepositoryFake{ + ready: domain.FreightItemImage{ + CreatorSubject: "local-admin", + FreightOrderItemID: "00000000-0000-4000-8000-000000000001", + Status: domain.FreightItemImageReady, + StorageKey: "ready/image.jpg", + }, + } + store := &freightImageStoreFake{ + openContent: []byte("normalized-jpeg"), + } + service, err := NewFreightImageService( + repository, + &freightImageSourceFake{}, + store, + fakeClock{}, + ) + if err != nil { + t.Fatalf("NewFreightImageService() error = %v", err) + } + result, err := service.OpenItemImage( + context.Background(), + "local-admin", + repository.ready.FreightOrderItemID, + ) + if err != nil { + t.Fatalf("OpenItemImage() error = %v", err) + } + defer result.Content.Close() + content, _ := io.ReadAll(result.Content) + if string(content) != "normalized-jpeg" { + t.Fatalf("image content = %q", content) + } + if _, err := service.OpenItemImage( + context.Background(), + "another-subject", + repository.ready.FreightOrderItemID, + ); err == nil { + t.Fatal("cross-subject OpenItemImage() error = nil") + } else { + var typed *Error + if !errors.As(err, &typed) || + typed.Code != "FREIGHT_IMAGE_NOT_FOUND" { + t.Fatalf("cross-subject error = %v", err) + } + } +} + +type freightImageRepositoryFake struct { + mu sync.Mutex + jobs []domain.FreightItemImageJob + saved []domain.FreightItemImage + replaced string + ready domain.FreightItemImage +} + +func (repository *freightImageRepositoryFake) ListFreightImageJobs( + context.Context, + string, + string, + int, +) ([]domain.FreightItemImageJob, error) { + return append([]domain.FreightItemImageJob(nil), repository.jobs...), nil +} + +func (repository *freightImageRepositoryFake) SaveFreightItemImage( + _ context.Context, + image domain.FreightItemImage, +) (*string, error) { + repository.mu.Lock() + defer repository.mu.Unlock() + repository.saved = append(repository.saved, image) + if image.Status == domain.FreightItemImageReady && + repository.replaced != "" { + value := repository.replaced + repository.replaced = "" + return &value, nil + } + return nil, nil +} + +func (repository *freightImageRepositoryFake) GetReadyFreightItemImage( + _ context.Context, + creatorSubject, itemID string, +) (domain.FreightItemImage, error) { + if repository.ready.CreatorSubject != creatorSubject || + repository.ready.FreightOrderItemID != itemID { + return domain.FreightItemImage{}, ErrRepositoryNotFound + } + return repository.ready, nil +} + +type freightImageSourceFake struct{} + +func (*freightImageSourceFake) FetchProductImage( + _ context.Context, + productThumbRef string, +) (FreightSourceImage, error) { + if productThumbRef == "191" { + return FreightSourceImage{}, domain.ErrFreightImageNotFound + } + return FreightSourceImage{ + Content: io.NopCloser(bytes.NewReader([]byte("source-image"))), + MediaType: "image/png", + }, nil +} + +type freightImageStoreFake struct { + mu sync.Mutex + deleted []string + openContent []byte +} + +func (*freightImageStoreFake) Put( + context.Context, + string, + string, + io.Reader, +) (NormalizedReferenceImage, error) { + return NormalizedReferenceImage{ + StorageKey: "new/image.jpg", + MediaType: "image/jpeg", + SizeBytes: 100, + SHA256: repeatSHA("a"), + }, nil +} + +func (store *freightImageStoreFake) Open( + context.Context, + string, +) (io.ReadCloser, error) { + return io.NopCloser(bytes.NewReader(store.openContent)), nil +} + +func (store *freightImageStoreFake) Delete( + _ context.Context, + storageKey string, +) error { + store.mu.Lock() + defer store.mu.Unlock() + store.deleted = append(store.deleted, storageKey) + return nil +} + +func repeatSHA(value string) string { + return strings.Repeat(value, 64) +} diff --git a/backend-api/internal/usecase/freight_service.go b/backend-api/internal/usecase/freight_service.go index 0811a77..91dcffa 100644 --- a/backend-api/internal/usecase/freight_service.go +++ b/backend-api/internal/usecase/freight_service.go @@ -22,6 +22,7 @@ const ( freightWatermarkOverlap = 10 * time.Minute orderFreightSyncTimeout = 55 * time.Second freightCleanupTimeout = 5 * time.Second + freightImageCacheBudget = 15 * time.Second ) type FreightService struct { @@ -33,6 +34,19 @@ type FreightService struct { orderTimeout time.Duration cleanupTimeout time.Duration orderSyncGate chan struct{} + imageCache FreightImageCache +} + +type FreightServiceOption func(*FreightService) error + +func WithFreightImageCache(cache FreightImageCache) FreightServiceOption { + return func(service *FreightService) error { + if cache == nil { + return errors.New("freight image cache is required") + } + service.imageCache = cache + return nil + } } type FreightSourcePreflight interface { @@ -66,12 +80,13 @@ func NewFreightService( clock Clock, ids IDGenerator, timeout time.Duration, + options ...FreightServiceOption, ) (*FreightService, error) { if repository == nil || source == nil || clock == nil || ids == nil || timeout <= 0 { return nil, errors.New("freight service dependencies are required") } - return &FreightService{ + service := &FreightService{ repository: repository, source: source, clock: clock, @@ -80,7 +95,16 @@ func NewFreightService( orderTimeout: orderFreightSyncTimeout, cleanupTimeout: freightCleanupTimeout, orderSyncGate: make(chan struct{}, 1), - }, nil + } + for _, option := range options { + if option == nil { + return nil, errors.New("freight service option is required") + } + if err := option(service); err != nil { + return nil, err + } + } + return service, nil } func (service *FreightService) CreateOrderSync( @@ -424,6 +448,18 @@ func (service *FreightService) executeRun( run.OrderCount = len(batch.Orders) run.ItemCount = itemCount run.FinishedAt = &finishedAt + if service.imageCache != nil { + cacheCtx, cancel := context.WithTimeout( + ctx, + freightImageCacheBudget, + ) + _ = service.imageCache.CacheRun( + cacheCtx, + run.CreatorSubject, + run.ID, + ) + cancel() + } return run, nil } diff --git a/backend-api/internal/usecase/freight_service_test.go b/backend-api/internal/usecase/freight_service_test.go index 322f9d0..443b889 100644 --- a/backend-api/internal/usecase/freight_service_test.go +++ b/backend-api/internal/usecase/freight_service_test.go @@ -184,12 +184,16 @@ func TestCreateFreightOrderSyncStopsBeforePersistingWhenOCRIsInvalid(t *testing. func TestCreateFreightOrderSyncReturnsCommittedResult(t *testing.T) { repository := &syncTrackingRepository{} + imageCache := &recordingFreightImageCache{ + err: errors.New("best effort image failure"), + } service, err := NewFreightService( repository, &recordingDateSource{}, fakeClock{}, &sequenceIDs{}, time.Minute, + WithFreightImageCache(imageCache), ) if err != nil { t.Fatalf("NewFreightService() error = %v", err) @@ -209,6 +213,11 @@ func TestCreateFreightOrderSyncReturnsCommittedResult(t *testing.T) { result.Run.StartedAt == nil || result.Run.FinishedAt == nil { t.Fatalf("CreateOrderSync() result = %+v", result) } + if imageCache.calls != 1 || + imageCache.creatorSubject != "local-admin" || + imageCache.runID != result.Run.ID { + t.Fatalf("image cache calls = %+v", imageCache) + } repository.mu.Lock() defer repository.mu.Unlock() if repository.startCalls != 1 || repository.completeCalls != 1 || @@ -222,6 +231,23 @@ func TestCreateFreightOrderSyncReturnsCommittedResult(t *testing.T) { } } +type recordingFreightImageCache struct { + calls int + creatorSubject string + runID string + err error +} + +func (cache *recordingFreightImageCache) CacheRun( + _ context.Context, + creatorSubject, runID string, +) error { + cache.calls++ + cache.creatorSubject = creatorSubject + cache.runID = runID + return cache.err +} + func TestCreateFreightOrderSyncTimeoutUsesLiveCleanupContext(t *testing.T) { repository := &syncTrackingRepository{} source := &blockingOrderSource{} diff --git a/backend-api/migrations/00016_freight_item_images.sql b/backend-api/migrations/00016_freight_item_images.sql new file mode 100644 index 0000000..544a8dc --- /dev/null +++ b/backend-api/migrations/00016_freight_item_images.sql @@ -0,0 +1,65 @@ +-- +goose Up +CREATE TABLE freight_item_images ( + freight_order_item_id TEXT PRIMARY KEY NOT NULL + REFERENCES freight_order_items(id) + ON UPDATE RESTRICT ON DELETE CASCADE, + product_thumb_ref TEXT NOT NULL + CHECK ( + product_thumb_ref GLOB '[1-9]*' + AND product_thumb_ref NOT GLOB '*[^0-9]*' + AND length(product_thumb_ref) <= 20 + ), + status TEXT NOT NULL CHECK (status IN ('READY', 'MISSING', 'FAILED')), + media_type TEXT, + size_bytes INTEGER, + sha256 TEXT, + storage_key TEXT UNIQUE, + error_code TEXT + CHECK ( + error_code IS NULL + OR ( + length(trim(error_code)) > 0 + AND length(CAST(error_code AS BLOB)) <= 64 + ) + ), + attempt_count INTEGER NOT NULL CHECK (attempt_count >= 1), + updated_at TEXT NOT NULL, + CHECK ( + ( + status = 'READY' + AND media_type = 'image/jpeg' + AND size_bytes > 0 + AND length(sha256) = 64 + AND sha256 NOT GLOB '*[^0-9a-f]*' + AND length(trim(storage_key)) > 0 + AND length(CAST(storage_key AS BLOB)) <= 512 + AND error_code IS NULL + ) + OR ( + status IN ('MISSING', 'FAILED') + AND media_type IS NULL + AND size_bytes IS NULL + AND sha256 IS NULL + AND storage_key IS NULL + AND error_code IS NOT NULL + ) + ) +); + +CREATE INDEX freight_item_images_status_idx + ON freight_item_images (status, updated_at, freight_order_item_id); + +-- +goose Down +CREATE TEMP TABLE freight_item_images_v16_down_guard ( + allowed INTEGER NOT NULL CHECK (allowed = 1) +); + +INSERT INTO freight_item_images_v16_down_guard (allowed) +SELECT CASE + WHEN EXISTS (SELECT 1 FROM freight_item_images) + THEN 0 + ELSE 1 +END; + +DROP TABLE freight_item_images_v16_down_guard; +DROP TABLE freight_item_images; diff --git a/docs/api.md b/docs/api.md index 39307da..f8f4a9a 100644 --- a/docs/api.md +++ b/docs/api.md @@ -205,6 +205,33 @@ T-203 成功返回 `201`。使用相同 `Idempotency-Key` 和相同图片内容 `null`,币种固定为 `TWD`。`product_thumb_ref` 来自商品 `productThumb`,存在时必须是 规范化正整数,不能使用货运单 ID、商品明细 ID 或完整 URL。 +货运详情中的每个商品还返回: + +```json +{ + "image_status": "READY", + "image_error_code": null, + "image_url": "/api/v1/freight-items/{item_id}/image" +} +``` + +`image_status` 为 `NONE/PENDING/READY/MISSING/FAILED`。只有 `READY` 返回同源 +`image_url`;其他状态返回 `null`,失败时 `image_error_code` 是不包含上游正文的稳定错误码。 +元数据同步先提交,图片在剩余请求预算内 best-effort 缓存,因此图片失败不改变货运同步的 +`SUCCEEDED` 终态。后续新同步会重试非 `READY` 图片,相同引用的 `READY` 图片不会重复下载。 + +### `GET /api/v1/freight-items/{item_id}/image` + +Admin 会话鉴权的本地货运商品图片流。仅当商品仍在当前货运版本中、当前 `productThumb` +与已缓存引用一致且状态为 `READY` 时返回 `200 image/jpeg`;不存在、未完成、引用已变化或 +无权访问统一返回 `404`。响应包含 `Content-Length`、ETag、 +`X-Content-Type-Options: nosniff` 和 `Cache-Control: private, no-store`,不返回 ERP URL、 +查询参数或服务端文件路径。 + +后端只使用配置的 `CMROUBAO_SHUNYUNBAO_URL` origin 和固定 +`/api/p/file?id=` 请求图片,复用 ERP Cookie jar 并禁止重定向。下载内容经过 +既有大小、尺寸、像素和解码校验,统一白底规范化为最长边不超过 2048 px 的 JPEG。 + ### Go ERP 会话(T-230、T-236) `/erp`、`/erp/captcha`、`/erp/login` 和 `/api/v1/erp-session*` 不再暴露。创建货运同步 diff --git a/docs/current-state.md b/docs/current-state.md index df3dcff..89c7811 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -202,6 +202,7 @@ | `docs/tasks/T-237.md` | DONE | 完整单号在 55 秒预算内同步导入并可靠收敛终态 | | `docs/tasks/T-238.md` | DONE | 兼容顺运宝详情分钟精度时间 | | `docs/tasks/T-239.md` | DONE | 冻结货运规格 SKU、TWD 原始单价和图片引用契约 | +| `docs/tasks/T-240.md` | DONE | 后端受控缓存并鉴权提供货运商品图片 | | `docs/design/` | 已确认 | T-202 原型索引、4 个管理页和 7 个 Android 页面 | | `deepseek总结.txt` | 已有 | 历史讨论摘要,不是正式需求权威 | | `android-buyer/` | 已有 | Roubao `main` 固定 commit 的 Android 基线 | @@ -213,11 +214,11 @@ ## 任务摘要 - 已完成:T-001 至 T-004、T-101 至 T-104、T-201 至 T-219。 -- 已完成:另含 T-220 至 T-239 ERP 契约、货运存储、采购需求生成、日期增量同步、Go +- 已完成:另含 T-220 至 T-240 ERP 契约、货运存储、采购需求生成、日期增量同步、Go 直连协议、OCR 会话预检、稳定预检错误、安全诊断日志、直连 `FreightSource`、旧 Connector - 清理、受控本地凭证加载、同步单号导入、分钟时间兼容和货运商品元数据契约。 + 清理、受控本地凭证加载、同步单号导入、分钟时间兼容、货运商品元数据契约和本地图片缓存。 - 进行中:无。 -- 下一步:实现商品 `productThumb` 图片的后端受控下载、本地保存和鉴权访问。 +- 下一步:在货运列表和详情中展示本地缩略图、原始单价和清晰的图片状态。 ## 当前可运行内容 diff --git a/docs/tasks/T-240.md b/docs/tasks/T-240.md index b887ad5..2262697 100644 --- a/docs/tasks/T-240.md +++ b/docs/tasks/T-240.md @@ -4,7 +4,7 @@ title: 缓存并鉴权提供顺运宝货运商品图片 phase: 2 deps: - T-239 -status: TODO +status: DONE created: 2026-07-29 context_ref: 3bd9726 work_branch: null @@ -57,14 +57,14 @@ write_paths: ## 验收要点 -- [ ] 图片请求固定为配置 ERP origin 的 `/api/p/file?id=`,不接受任意 URL。 -- [ ] 请求复用现有 ERP Cookie jar,禁止重定向;401/403、404、错误类型和无效图片稳定分类。 -- [ ] 有效 JPEG/PNG/WebP 规范化保存为本地 JPEG,大小和像素限制不回归。 -- [ ] 货运元数据先成功提交;单图失败不改变同步 `SUCCEEDED`,状态可审计并在后续同步重试。 -- [ ] 相同 READY 引用不重复下载;引用变化替换关联并 best-effort 删除旧文件。 -- [ ] Admin 图片接口只返回当前商品本地 READY 图片,未授权及其他状态不泄露存在性。 -- [ ] 自动测试不访问真实 ERP,覆盖 Cookie、URL、重定向、状态、去重、替换和内容响应头。 -- [ ] 标准 Go 测试、race、vet 和三个入口构建通过。 +- [x] 图片请求固定为配置 ERP origin 的 `/api/p/file?id=`,不接受任意 URL。 +- [x] 请求复用现有 ERP Cookie jar,禁止重定向;401/403、404、错误类型和无效图片稳定分类。 +- [x] 有效 JPEG/PNG/WebP 规范化保存为本地 JPEG,大小和像素限制不回归。 +- [x] 货运元数据先成功提交;单图失败不改变同步 `SUCCEEDED`,状态可审计并在后续同步重试。 +- [x] 相同 READY 引用不重复下载;引用变化替换关联并 best-effort 删除旧文件。 +- [x] Admin 图片接口只返回当前商品本地 READY 图片,未授权及其他状态不泄露存在性。 +- [x] 自动测试不访问真实 ERP,覆盖 Cookie、URL、重定向、状态、去重、替换和内容响应头。 +- [x] 标准 Go 测试、race、vet 和三个入口构建通过。 ## 边界 @@ -77,3 +77,8 @@ write_paths: - 2026-07-29:创建任务。采用货运商品专属图片元数据表并复用安全文件编码管线;先提交货运 元数据,再 best-effort 缓存图片,避免单个外部文件阻塞或回滚完整订单。 +- 2026-07-29:增加 migration 16、固定 ERP 图片客户端、四 worker/最多 1000 项/15 秒 + best-effort 缓存、图片状态仓储和 Admin 本地内容接口。相同 READY 引用跳过,失败引用在 + 新同步重试,引用替换后删除旧文件;同步成功不受单图失败影响。脱敏测试覆盖 Cookie、 + 固定 query、重定向、401/404/无效类型、JPEG 规范化、去重、替换、回滚保护、鉴权与响应头。 + `go test ./...`、`go test -race ./...`、`go vet ./...` 及三个入口构建均通过;未访问真实 ERP。