diff --git a/backend-api/internal/domain/freight.go b/backend-api/internal/domain/freight.go index 409a70a..374984b 100644 --- a/backend-api/internal/domain/freight.go +++ b/backend-api/internal/domain/freight.go @@ -5,6 +5,7 @@ import "time" const ( FreightSourceShunyunbao = "SHUNYUNBAO" FreightSyncOrderNumber = "ORDER_NUMBER" + FreightSyncCreatedRange = "CREATED_RANGE" ) type FreightSyncStatus string @@ -17,19 +18,22 @@ const ( ) type FreightSyncRun struct { - ID string - CreatorSubject string - CreatedByUserID string - Mode string - OrderNumber string - QuerySHA256 string - Status FreightSyncStatus - ErrorCode *string - OrderCount int - ItemCount int - CreatedAt time.Time - StartedAt *time.Time - FinishedAt *time.Time + ID string + CreatorSubject string + CreatedByUserID string + Mode string + OrderNumber string + CreatedFrom string + CreatedTo string + WatermarkThrough *time.Time + QuerySHA256 string + Status FreightSyncStatus + ErrorCode *string + OrderCount int + ItemCount int + CreatedAt time.Time + StartedAt *time.Time + FinishedAt *time.Time } type FreightOrder struct { @@ -84,7 +88,17 @@ type FreightSourceBatch struct { } type FreightSourceQuery struct { - Mode string `json:"mode"` + Mode string `json:"mode"` + CreatedFrom *string `json:"created_from"` + CreatedTo *string `json:"created_to"` +} + +type FreightSyncWatermark struct { + CreatorSubject string + SourceSystem string + LastSuccessfulTo time.Time + LastSuccessfulRunID string + UpdatedAt time.Time } type FreightSourceOrder struct { diff --git a/backend-api/internal/platform/erpconnector/client.go b/backend-api/internal/platform/erpconnector/client.go index 13c6b37..5249f97 100644 --- a/backend-api/internal/platform/erpconnector/client.go +++ b/backend-api/internal/platform/erpconnector/client.go @@ -57,11 +57,33 @@ func New(baseURL, apiKey string, timeout time.Duration) (*Client, error) { func (client *Client) QueryOrder( ctx context.Context, orderNumber string, +) (domain.FreightSourceBatch, error) { + return client.query(ctx, map[string]string{ + "mode": domain.FreightSyncOrderNumber, + "order_number": orderNumber, + }, domain.FreightSyncOrderNumber, "", "") +} + +func (client *Client) QueryCreatedRange( + ctx context.Context, + createdFrom, createdTo string, +) (domain.FreightSourceBatch, error) { + return client.query(ctx, map[string]string{ + "mode": domain.FreightSyncCreatedRange, + "created_from": createdFrom, + "created_to": createdTo, + }, domain.FreightSyncCreatedRange, createdFrom, createdTo) +} + +func (client *Client) query( + ctx context.Context, + payload map[string]string, + expectedMode, expectedFrom, expectedTo string, ) (domain.FreightSourceBatch, error) { if client.apiKey == "" { return domain.FreightSourceBatch{}, ErrNotConfigured } - body, err := json.Marshal(map[string]string{"order_number": orderNumber}) + body, err := json.Marshal(payload) if err != nil { return domain.FreightSourceBatch{}, ErrProtocol } @@ -106,9 +128,19 @@ func (client *Client) QueryOrder( if err := decoder.Decode(&extra); !errors.Is(err, io.EOF) { return domain.FreightSourceBatch{}, ErrProtocol } - if result.SchemaVersion != 1 || result.Query.Mode != "ORDER_NUMBER" || + if result.SchemaVersion != 1 || result.Query.Mode != expectedMode || result.Orders == nil { return domain.FreightSourceBatch{}, ErrProtocol } + if expectedMode == domain.FreightSyncOrderNumber { + if result.Query.CreatedFrom != nil || result.Query.CreatedTo != nil { + return domain.FreightSourceBatch{}, ErrProtocol + } + } else if result.Query.CreatedFrom == nil || + result.Query.CreatedTo == nil || + *result.Query.CreatedFrom != expectedFrom || + *result.Query.CreatedTo != expectedTo { + return domain.FreightSourceBatch{}, ErrProtocol + } return result, nil } diff --git a/backend-api/internal/platform/erpconnector/client_test.go b/backend-api/internal/platform/erpconnector/client_test.go index a20e02d..6154432 100644 --- a/backend-api/internal/platform/erpconnector/client_test.go +++ b/backend-api/internal/platform/erpconnector/client_test.go @@ -2,6 +2,7 @@ package erpconnector import ( "context" + "encoding/json" "errors" "net/http" "net/http/httptest" @@ -61,6 +62,85 @@ func TestQueryOrderAcceptsAllowlistResponse(t *testing.T) { } } +func TestQueryCreatedRangeUsesStrictContract(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func( + writer http.ResponseWriter, + request *http.Request, + ) { + var body map[string]string + if err := json.NewDecoder(request.Body).Decode(&body); err != nil { + t.Fatalf("decode request: %v", err) + } + if body["mode"] != "CREATED_RANGE" || + body["created_from"] != "2026-07-22" || + body["created_to"] != "2026-07-28" || + len(body) != 3 { + t.Fatalf("request body = %#v", body) + } + writer.Header().Set("Content-Type", "application/json") + _, _ = writer.Write([]byte(`{ + "schema_version":1, + "query":{ + "mode":"CREATED_RANGE", + "created_from":"2026-07-22", + "created_to":"2026-07-28" + }, + "orders":[] + }`)) + })) + defer server.Close() + client, _ := New( + server.URL, + "12345678901234567890123456789012", + time.Second, + ) + + result, err := client.QueryCreatedRange( + context.Background(), + "2026-07-22", + "2026-07-28", + ) + + if err != nil || result.Query.CreatedFrom == nil || + *result.Query.CreatedFrom != "2026-07-22" { + t.Fatalf("QueryCreatedRange() = %+v, %v", result, err) + } +} + +func TestQueryCreatedRangeRejectsMismatchedResponseWindow(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func( + writer http.ResponseWriter, + _ *http.Request, + ) { + writer.Header().Set("Content-Type", "application/json") + _, _ = writer.Write([]byte(`{ + "schema_version":1, + "query":{ + "mode":"CREATED_RANGE", + "created_from":"2026-07-21", + "created_to":"2026-07-28" + }, + "orders":[] + }`)) + })) + defer server.Close() + client, _ := New( + server.URL, + "12345678901234567890123456789012", + time.Second, + ) + + _, err := client.QueryCreatedRange( + context.Background(), + "2026-07-22", + "2026-07-28", + ) + + if !errors.Is(err, ErrProtocol) { + t.Fatalf("QueryCreatedRange() error = %v", err) + } +} + func TestQueryOrderRejectsUnexpectedPIIField(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func( writer http.ResponseWriter, diff --git a/backend-api/internal/platform/migration/claims_migration_test.go b/backend-api/internal/platform/migration/claims_migration_test.go index 37a996a..5cb8703 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 != 13 { - t.Fatalf("initial Up() applied = %d, want 13", applied) + } else if applied != 14 { + t.Fatalf("initial Up() applied = %d, want 14", applied) + } + if err := runner.Down(ctx); err != nil { + t.Fatalf("initial Down(v14) error = %v", err) } if err := runner.Down(ctx); err != nil { t.Fatalf("initial Down(v13) error = %v", err) @@ -68,9 +71,14 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { seedClaimsHistoricalFixture(t, db) if applied, err := runner.Up(ctx); err != nil { - 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) + t.Fatalf("Up(v5-v14) over historical data error = %v", err) + } else if applied != 10 { + t.Fatalf("Up(v5-v14) applied = %d, want 10", applied) + } + assertClaimsHistory(t, db, true) + + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v14) with compatible history error = %v", err) } assertClaimsHistory(t, db, true) @@ -125,9 +133,9 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { assertClaimsHistory(t, db, false) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("final Up(v4-v13) error = %v", err) - } else if applied != 10 { - t.Fatalf("final Up(v4-v13) applied = %d, want 10", applied) + t.Fatalf("final Up(v4-v14) error = %v", err) + } else if applied != 11 { + t.Fatalf("final Up(v4-v14) applied = %d, want 11", applied) } assertClaimsHistory(t, db, true) } @@ -363,6 +371,9 @@ func TestClaimsMigrationDownFailsClosedForNewAuditData(t *testing.T) { t.Fatalf("insert v4 audit event: %v", err) } + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v14) error = %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("Down(v13) error = %v", err) } diff --git a/backend-api/internal/platform/migration/runner_test.go b/backend-api/internal/platform/migration/runner_test.go index 4ad63b8..91f4d90 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 != 13 { - t.Fatalf("Up() applied = %d, want 13", applied) + if applied != 14 { + t.Fatalf("Up() applied = %d, want 14", applied) } assertStatuses(t, runner, map[int64]bool{ 1: true, @@ -44,6 +44,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 11: true, 12: true, 13: true, + 14: true, }) applied, err = runner.Up(context.Background()) @@ -70,7 +71,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 10: true, 11: true, 12: true, - 13: false, + 13: true, + 14: false, }) applied, err = runner.Up(context.Background()) @@ -94,6 +96,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 11: true, 12: true, 13: true, + 14: true, }) } diff --git a/backend-api/internal/repository/sqlite/auth_repository_test.go b/backend-api/internal/repository/sqlite/auth_repository_test.go index 9830d0f..8354026 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(v14) error = %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("Down(v13) error = %v", err) } @@ -429,9 +432,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-v13) error = %v", err) - } else if applied != 11 { - t.Fatalf("Up(v3-v13) applied = %d, want 11", applied) + t.Fatalf("Up(v3-v14) error = %v", err) + } else if applied != 12 { + t.Fatalf("Up(v3-v14) applied = %d, want 12", applied) } } diff --git a/backend-api/internal/repository/sqlite/freight_repository.go b/backend-api/internal/repository/sqlite/freight_repository.go index 089b9ec..2e5cf94 100644 --- a/backend-api/internal/repository/sqlite/freight_repository.go +++ b/backend-api/internal/repository/sqlite/freight_repository.go @@ -52,14 +52,18 @@ func (store *Store) CreateFreightSync( ctx, `INSERT INTO erp_sync_runs ( id, creator_subject, created_by_user_id, mode, order_number, - query_sha256, idempotency_key, request_sha256, status, - order_count, item_count, created_at - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 0, 0, ?)`, + created_from, created_to, watermark_through, query_sha256, + idempotency_key, request_sha256, status, order_count, item_count, + created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, 0, ?)`, run.ID, run.CreatorSubject, run.CreatedByUserID, run.Mode, - run.OrderNumber, + nullableFreightSyncValue(run.OrderNumber), + nullableFreightSyncValue(run.CreatedFrom), + nullableFreightSyncValue(run.CreatedTo), + nullableTimestamp(run.WatermarkThrough), run.QuerySHA256, idempotencyKey, requestSHA256, @@ -106,6 +110,32 @@ func (store *Store) CompleteFreightSync( run domain.FreightSyncRun, batch domain.FreightImportBatch, finishedAt time.Time, +) error { + return store.completeFreightSync(ctx, run, batch, nil, finishedAt) +} + +func (store *Store) CompleteFreightDateSync( + ctx context.Context, + run domain.FreightSyncRun, + batch domain.FreightImportBatch, + watermarkThrough time.Time, + finishedAt time.Time, +) error { + return store.completeFreightSync( + ctx, + run, + batch, + &watermarkThrough, + finishedAt, + ) +} + +func (store *Store) completeFreightSync( + ctx context.Context, + run domain.FreightSyncRun, + batch domain.FreightImportBatch, + watermarkThrough *time.Time, + finishedAt time.Time, ) error { tx, err := store.db.BeginTx(ctx, nil) if err != nil { @@ -184,6 +214,28 @@ func (store *Store) CompleteFreightSync( if changed != 1 { return usecase.ErrTaskStateConflict } + if watermarkThrough != nil { + if _, err := tx.ExecContext( + ctx, + `INSERT INTO erp_sync_watermarks ( + creator_subject, source_system, last_successful_to, + last_successful_run_id, updated_at + ) VALUES (?, 'SHUNYUNBAO', ?, ?, ?) + ON CONFLICT (creator_subject, source_system) + DO UPDATE SET + last_successful_to = excluded.last_successful_to, + last_successful_run_id = excluded.last_successful_run_id, + updated_at = excluded.updated_at + WHERE julianday(erp_sync_watermarks.last_successful_to) + < julianday(excluded.last_successful_to)`, + run.CreatorSubject, + formatTimestamp(*watermarkThrough), + run.ID, + formatTimestamp(finishedAt), + ); err != nil { + return repositoryFailure(err) + } + } if err := tx.Commit(); err != nil { return repositoryFailure(err) } @@ -377,6 +429,43 @@ func (store *Store) GetFreightSync( return run, nil } +func (store *Store) GetFreightSyncWatermark( + ctx context.Context, + creatorSubject string, +) (*domain.FreightSyncWatermark, error) { + var watermark domain.FreightSyncWatermark + var lastSuccessfulTo, updatedAt string + err := store.db.QueryRowContext( + ctx, + `SELECT creator_subject, source_system, last_successful_to, + last_successful_run_id, updated_at + FROM erp_sync_watermarks + WHERE creator_subject = ? AND source_system = 'SHUNYUNBAO'`, + creatorSubject, + ).Scan( + &watermark.CreatorSubject, + &watermark.SourceSystem, + &lastSuccessfulTo, + &watermark.LastSuccessfulRunID, + &updatedAt, + ) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, repositoryFailure(err) + } + watermark.LastSuccessfulTo, err = parseTimestamp(lastSuccessfulTo) + if err != nil { + return nil, repositoryFailure(err) + } + watermark.UpdatedAt, err = parseTimestamp(updatedAt) + if err != nil { + return nil, repositoryFailure(err) + } + return &watermark, nil +} + func (store *Store) ListFreightOrders( ctx context.Context, creatorSubject string, @@ -458,13 +547,15 @@ func (store *Store) GetFreightOrder( const freightSyncSelect = `SELECT id, creator_subject, created_by_user_id, mode, order_number, - query_sha256, status, error_code, order_count, item_count, - created_at, started_at, finished_at + created_from, created_to, watermark_through, query_sha256, status, + error_code, order_count, item_count, created_at, started_at, finished_at FROM erp_sync_runs ` func scanFreightSync(scanner rowScanner) (domain.FreightSyncRun, error) { var run domain.FreightSyncRun + var orderNumber, createdFrom, createdTo sql.NullString + var watermarkThrough sql.NullString var errorCode sql.NullString var createdAt string var startedAt sql.NullString @@ -474,7 +565,10 @@ func scanFreightSync(scanner rowScanner) (domain.FreightSyncRun, error) { &run.CreatorSubject, &run.CreatedByUserID, &run.Mode, - &run.OrderNumber, + &orderNumber, + &createdFrom, + &createdTo, + &watermarkThrough, &run.QuerySHA256, &run.Status, &errorCode, @@ -490,6 +584,19 @@ func scanFreightSync(scanner rowScanner) (domain.FreightSyncRun, error) { if errorCode.Valid { run.ErrorCode = &errorCode.String } + if orderNumber.Valid { + run.OrderNumber = orderNumber.String + } + if createdFrom.Valid { + run.CreatedFrom = createdFrom.String + } + if createdTo.Valid { + run.CreatedTo = createdTo.String + } + run.WatermarkThrough, err = parseNullableTimestamp(watermarkThrough) + if err != nil { + return domain.FreightSyncRun{}, err + } run.CreatedAt, err = parseTimestamp(createdAt) if err != nil { return domain.FreightSyncRun{}, err @@ -622,6 +729,13 @@ func nullableFreightQuantity(value *int) any { return *value } +func nullableFreightSyncValue(value string) any { + if value == "" { + return nil + } + return value +} + func optionalString(value sql.NullString) *string { if !value.Valid { return nil diff --git a/backend-api/internal/repository/sqlite/freight_repository_test.go b/backend-api/internal/repository/sqlite/freight_repository_test.go index 84d4a15..9eec2c5 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("date sync migration down: %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("procurement migration down: %v", err) } @@ -163,6 +166,106 @@ func TestFreightSyncRecoveryAndFailedBatchDoNotPersistOrders(t *testing.T) { } } +func TestFreightDateSyncAdvancesWatermarkOnlyOnWholeBatchSuccess( + t *testing.T, +) { + db := openDatabase(t) + store, _ := repository.New(db) + ctx := context.Background() + now := time.Date(2026, 7, 28, 8, 0, 0, 0, time.UTC) + userID := uuid(960) + seedFreightUser(t, db, userID, now) + + firstThrough := now.Add(30 * time.Minute) + first := freightDateRun( + 961, + userID, + now, + "2026-07-27", + "2026-07-28", + firstThrough, + ) + createAndStartFreightRun(t, store, first, "date-success", "6") + if err := store.CompleteFreightDateSync( + ctx, + first, + domain.FreightImportBatch{}, + firstThrough, + now.Add(time.Minute), + ); err != nil { + t.Fatalf("CompleteFreightDateSync() error = %v", err) + } + watermark, err := store.GetFreightSyncWatermark(ctx, "local-admin") + if err != nil || watermark == nil || + !watermark.LastSuccessfulTo.Equal(firstThrough) || + watermark.LastSuccessfulRunID != first.ID { + t.Fatalf("first watermark = %+v, %v", watermark, err) + } + + failed := freightDateRun( + 962, + userID, + now.Add(time.Hour), + "2026-07-28", + "2026-07-28", + now.Add(2*time.Hour), + ) + createAndStartFreightRun(t, store, failed, "date-failed", "7") + if err := store.FailFreightSync( + ctx, + failed.ID, + "ERP_CONNECTOR_UNAVAILABLE", + now.Add(time.Hour+time.Minute), + ); err != nil { + t.Fatalf("FailFreightSync() error = %v", err) + } + watermark, _ = store.GetFreightSyncWatermark(ctx, "local-admin") + if !watermark.LastSuccessfulTo.Equal(firstThrough) || + watermark.LastSuccessfulRunID != first.ID { + t.Fatalf("failed run advanced watermark = %+v", watermark) + } + + olderThrough := now.Add(-time.Hour) + older := freightDateRun( + 963, + userID, + now.Add(2*time.Hour), + "2026-07-26", + "2026-07-26", + olderThrough, + ) + createAndStartFreightRun(t, store, older, "date-older", "8") + if err := store.CompleteFreightDateSync( + ctx, + older, + domain.FreightImportBatch{}, + olderThrough, + now.Add(2*time.Hour+time.Minute), + ); err != nil { + t.Fatalf("older CompleteFreightDateSync() error = %v", err) + } + watermark, _ = store.GetFreightSyncWatermark(ctx, "local-admin") + if !watermark.LastSuccessfulTo.Equal(firstThrough) || + watermark.LastSuccessfulRunID != first.ID { + t.Fatalf("older run retreated watermark = %+v", watermark) + } + + runner, _ := migration.New(db) + if err := runner.Down(ctx); err == nil { + t.Fatal("date sync migration down succeeded with retained watermark") + } + var foreignKeysEnabled int + if err := db.QueryRow(`PRAGMA foreign_keys`).Scan( + &foreignKeysEnabled, + ); err != nil || foreignKeysEnabled != 1 { + t.Fatalf( + "failed down disabled foreign keys = %d, %v", + foreignKeysEnabled, + err, + ) + } +} + func seedFreightUser( t *testing.T, db *sql.DB, @@ -200,6 +303,27 @@ func freightRun( } } +func freightDateRun( + index int, + userID string, + now time.Time, + createdFrom, createdTo string, + watermarkThrough time.Time, +) domain.FreightSyncRun { + return domain.FreightSyncRun{ + ID: uuid(index), + CreatorSubject: "local-admin", + CreatedByUserID: userID, + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: createdFrom, + CreatedTo: createdTo, + WatermarkThrough: &watermarkThrough, + QuerySHA256: repeatHex("b"), + Status: domain.FreightSyncPending, + CreatedAt: now, + } +} + func freightBatch(index int, firstHash, secondHash string) domain.FreightImportBatch { quantityOne := 1 quantityTwo := 2 diff --git a/backend-api/internal/repository/sqlite/procurement_repository_test.go b/backend-api/internal/repository/sqlite/procurement_repository_test.go index df7bb99..6203a69 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("date sync migration down: %v", err) + } if err := runner.Down(ctx); err == nil { t.Fatal("procurement migration down succeeded with retained requests") } diff --git a/backend-api/internal/transport/httpapi/admin_handlers.go b/backend-api/internal/transport/httpapi/admin_handlers.go index 1a3d00c..68c2514 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers.go +++ b/backend-api/internal/transport/httpapi/admin_handlers.go @@ -66,6 +66,10 @@ func registerAdminAPI(routes gin.IRoutes, services AdminServices) error { if services.Freight != nil { routes.POST("/api/v1/freight-syncs", handler.createFreightSync) routes.GET("/api/v1/freight-syncs/:id", handler.freightSyncDetail) + routes.GET( + "/api/v1/freight-sync-watermark", + handler.freightSyncWatermark, + ) routes.GET("/api/v1/freight-orders", handler.listFreightOrders) routes.GET("/api/v1/freight-orders/:id", handler.freightOrderDetail) } diff --git a/backend-api/internal/transport/httpapi/admin_handlers_test.go b/backend-api/internal/transport/httpapi/admin_handlers_test.go index 723aa22..68640a1 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers_test.go +++ b/backend-api/internal/transport/httpapi/admin_handlers_test.go @@ -395,6 +395,89 @@ func TestAdminFreightAPIImportsAllItemsWithoutPII(t *testing.T) { } } +func TestAdminFreightDateSyncAdvancesInspectableWatermark(t *testing.T) { + fixture := newAdminIntegrationFixture(t) + location, err := time.LoadLocation("Asia/Shanghai") + if err != nil { + t.Fatalf("LoadLocation() error = %v", err) + } + today := time.Now().In(location).Format(time.DateOnly) + body := fmt.Sprintf( + `{"mode":"CREATED_RANGE","created_from":%q,"created_to":%q}`, + today, + today, + ) + create := performAdminRequest( + t, + fixture.router, + http.MethodPost, + "/api/v1/freight-syncs", + "application/json", + strings.NewReader(body), + "freight-date-sync-1", + ) + requireAdminStatus(t, create, http.StatusAccepted) + var created struct { + Sync struct { + ID string `json:"id"` + } `json:"sync"` + } + decodeResponse(t, create, &created) + var status *httptest.ResponseRecorder + for attempt := 0; attempt < 50; attempt++ { + status = performAdminRequest( + t, + fixture.router, + http.MethodGet, + "/api/v1/freight-syncs/"+created.Sync.ID, + "", + nil, + "", + ) + if strings.Contains(status.Body.String(), `"status":"SUCCEEDED"`) { + break + } + time.Sleep(10 * time.Millisecond) + } + requireAdminStatus(t, status, http.StatusOK) + if !strings.Contains(status.Body.String(), `"mode":"CREATED_RANGE"`) || + !strings.Contains(status.Body.String(), `"created_from":"`+today+`"`) || + !strings.Contains(status.Body.String(), `"order_count":0`) { + t.Fatalf("date sync response = %s", status.Body) + } + + watermark := performAdminRequest( + t, + fixture.router, + http.MethodGet, + "/api/v1/freight-sync-watermark", + "", + nil, + "", + ) + requireAdminStatus(t, watermark, http.StatusOK) + if !strings.Contains( + watermark.Body.String(), + `"last_successful_run_id":"`+created.Sync.ID+`"`, + ) || strings.Contains(watermark.Body.String(), "receiver") { + t.Fatalf("watermark response = %s", watermark.Body) + } + + mixed := performAdminRequest( + t, + fixture.router, + http.MethodPost, + "/api/v1/freight-syncs", + "application/json", + strings.NewReader( + `{"mode":"CREATED_RANGE","order_number":"must-not-be-ignored",`+ + `"created_from":"`+today+`","created_to":"`+today+`"}`, + ), + "freight-date-mixed", + ) + requireAdminStatus(t, mixed, http.StatusUnprocessableEntity) +} + func TestAdminProcurementAPIProducesImmutablePendingTask(t *testing.T) { fixture := newAdminIntegrationFixture(t) createSync := performAdminRequest( @@ -708,6 +791,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("date sync migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("procurement migration down: %v", err) } @@ -1125,6 +1211,21 @@ func (staticFreightSource) QueryOrder( }, nil } +func (staticFreightSource) QueryCreatedRange( + _ context.Context, + createdFrom, createdTo string, +) (domain.FreightSourceBatch, error) { + return domain.FreightSourceBatch{ + SchemaVersion: 1, + Query: domain.FreightSourceQuery{ + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: &createdFrom, + CreatedTo: &createdTo, + }, + Orders: []domain.FreightSourceOrder{}, + }, nil +} + func mustDecodeAny( t *testing.T, response *httptest.ResponseRecorder, diff --git a/backend-api/internal/transport/httpapi/device_handlers_test.go b/backend-api/internal/transport/httpapi/device_handlers_test.go index a26d336..b06f106 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("date sync migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("procurement migration down: %v", err) } @@ -733,8 +736,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 != 5 { - t.Fatalf("restored migrations = %d, want 5", applied) + } else if applied != 6 { + t.Fatalf("restored migrations = %d, want 6", applied) } completePayload := fmt.Sprintf( @@ -1423,6 +1426,9 @@ func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable( if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("date sync migration down: %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("procurement 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 93972ef..584be19 100644 --- a/backend-api/internal/transport/httpapi/freight_handlers.go +++ b/backend-api/internal/transport/httpapi/freight_handlers.go @@ -26,6 +26,9 @@ func (h *adminHandlers) createFreightSync(ctx *gin.Context) { var request struct { Mode string `json:"mode"` OrderNumber string `json:"order_number"` + CreatedFrom string `json:"created_from"` + CreatedTo string `json:"created_to"` + SyncToNow bool `json:"sync_to_now"` } if err := decodeJSON(ctx, &request); err != nil { writePublicError( @@ -38,26 +41,66 @@ func (h *adminHandlers) createFreightSync(ctx *gin.Context) { ) return } - if request.Mode != domain.FreightSyncOrderNumber { + var result usecase.CreateFreightSyncResult + var err error + switch request.Mode { + case domain.FreightSyncOrderNumber: + if strings.TrimSpace(request.CreatedFrom) != "" || + strings.TrimSpace(request.CreatedTo) != "" || + request.SyncToNow { + writePublicError( + ctx, + http.StatusUnprocessableEntity, + "FREIGHT_SYNC_INVALID", + "freight sync request is invalid", + false, + fieldDetails("mode", "ORDER_NUMBER cannot include date parameters"), + ) + return + } + result, err = h.services.Freight.CreateOrderSync( + ctx.Request.Context(), + usecase.CreateFreightSyncCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: adminActorUserID(ctx), + IdempotencyKey: ctx.GetHeader("Idempotency-Key"), + OrderNumber: request.OrderNumber, + }, + ) + case domain.FreightSyncCreatedRange: + if strings.TrimSpace(request.OrderNumber) != "" { + writePublicError( + ctx, + http.StatusUnprocessableEntity, + "FREIGHT_SYNC_INVALID", + "freight sync request is invalid", + false, + fieldDetails("mode", "CREATED_RANGE cannot include order_number"), + ) + return + } + result, err = h.services.Freight.CreateDateSync( + ctx.Request.Context(), + usecase.CreateFreightDateSyncCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: adminActorUserID(ctx), + IdempotencyKey: ctx.GetHeader("Idempotency-Key"), + CreatedFrom: request.CreatedFrom, + CreatedTo: request.CreatedTo, + SyncToNow: request.SyncToNow, + }, + ) + default: writePublicError( ctx, http.StatusUnprocessableEntity, "FREIGHT_SYNC_MODE_INVALID", "freight sync mode is not supported", false, - fieldDetails("mode", "must be ORDER_NUMBER"), + fieldDetails("mode", "must be ORDER_NUMBER or CREATED_RANGE"), ) return } - result, err := h.services.Freight.CreateOrderSync( - ctx.Request.Context(), - usecase.CreateFreightSyncCommand{ - CreatorSubject: localAdminSubject, - ActorUserID: adminActorUserID(ctx), - IdempotencyKey: ctx.GetHeader("Idempotency-Key"), - OrderNumber: request.OrderNumber, - }, - ) if err != nil { writeUsecaseError(ctx, err) return @@ -69,6 +112,28 @@ func (h *adminHandlers) createFreightSync(ctx *gin.Context) { }) } +func (h *adminHandlers) freightSyncWatermark(ctx *gin.Context) { + watermark, err := h.services.Freight.GetWatermark( + ctx.Request.Context(), + localAdminSubject, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + if watermark == nil { + ctx.JSON(http.StatusOK, gin.H{"watermark": nil}) + return + } + ctx.JSON(http.StatusOK, gin.H{"watermark": gin.H{ + "source_system": watermark.SourceSystem, + "last_successful_to": formatTime(watermark.LastSuccessfulTo), + "last_successful_run_id": watermark.LastSuccessfulRunID, + "updated_at": formatTime(watermark.UpdatedAt), + }}) +} + func (h *adminHandlers) freightSyncDetail(ctx *gin.Context) { run, err := h.services.Freight.GetSync( ctx.Request.Context(), @@ -175,19 +240,29 @@ func (h *adminHandlers) freightOrderDetail(ctx *gin.Context) { func freightSyncResponse(run domain.FreightSyncRun) gin.H { return gin.H{ - "id": run.ID, - "mode": run.Mode, - "query_sha256": run.QuerySHA256, - "status": run.Status, - "error_code": run.ErrorCode, - "order_count": run.OrderCount, - "item_count": run.ItemCount, - "created_at": formatTime(run.CreatedAt), - "started_at": formatOptionalTime(run.StartedAt), - "finished_at": formatOptionalTime(run.FinishedAt), + "id": run.ID, + "mode": run.Mode, + "created_from": nullableResponseString(run.CreatedFrom), + "created_to": nullableResponseString(run.CreatedTo), + "watermark_through": formatOptionalTime(run.WatermarkThrough), + "query_sha256": run.QuerySHA256, + "status": run.Status, + "error_code": run.ErrorCode, + "order_count": run.OrderCount, + "item_count": run.ItemCount, + "created_at": formatTime(run.CreatedAt), + "started_at": formatOptionalTime(run.StartedAt), + "finished_at": formatOptionalTime(run.FinishedAt), } } +func nullableResponseString(value string) any { + if value == "" { + return nil + } + return value +} + func freightOrderResponse(order domain.FreightOrder) gin.H { return gin.H{ "id": order.ID, diff --git a/backend-api/internal/transport/webui/handler.go b/backend-api/internal/transport/webui/handler.go index f38e2a4..59ce058 100644 --- a/backend-api/internal/transport/webui/handler.go +++ b/backend-api/internal/transport/webui/handler.go @@ -140,6 +140,17 @@ func (h *Handler) ImportFreight(ctx *gin.Context) { CSRFToken: token, }, IdempotencyKey: key, + Mode: "ORDER_NUMBER", + } + if location, locationErr := time.LoadLocation("Asia/Shanghai"); locationErr == nil { + today := time.Now().In(location).Format(time.DateOnly) + page.CreatedFrom = today + page.CreatedTo = today + } + if watermark, watermarkErr := service.GetFreightWatermark( + ctx.Request.Context(), + ); watermarkErr == nil { + page.Watermark = watermark } if syncID := strings.TrimSpace(ctx.Query("sync")); syncID != "" { run, getErr := service.GetFreightSync(ctx.Request.Context(), syncID) @@ -157,9 +168,29 @@ func (h *Handler) CreateFreightImport(ctx *gin.Context) { return } orderNumber := strings.TrimSpace(ctx.PostForm("order_number")) + mode := strings.TrimSpace(ctx.PostForm("mode")) + if mode == "" { + mode = "ORDER_NUMBER" + } + createdFrom := strings.TrimSpace(ctx.PostForm("created_from")) + createdTo := strings.TrimSpace(ctx.PostForm("created_to")) + syncToNow := ctx.PostForm("sync_to_now") == "true" + if syncToNow { + createdFrom = "" + createdTo = "" + } key := strings.TrimSpace(ctx.PostForm("idempotency_key")) - if orderNumber == "" || len([]byte(orderNumber)) > 128 || - !validToken(key) { + validInput := validToken(key) + if mode == "ORDER_NUMBER" { + validInput = validInput && orderNumber != "" && + len([]byte(orderNumber)) <= 128 + } else if mode == "CREATED_RANGE" { + validInput = validInput && + (syncToNow || (createdFrom != "" && createdTo != "")) + } else { + validInput = false + } + if !validInput { token, _ := csrfToken(ctx) h.render(ctx, http.StatusUnprocessableEntity, "freight-import", freightImportPage{ Page: pageView{ @@ -168,8 +199,11 @@ func (h *Handler) CreateFreightImport(ctx *gin.Context) { CSRFToken: token, }, OrderNumber: orderNumber, + Mode: mode, + CreatedFrom: createdFrom, + CreatedTo: createdTo, IdempotencyKey: key, - Error: "请输入完整单号后重试。", + Error: "请检查同步方式和查询条件后重试。", }) return } @@ -179,7 +213,11 @@ func (h *Handler) CreateFreightImport(ctx *gin.Context) { CreateFreightSyncInput{ ActorUserID: actorUserID(ctx.Request.Context()), IdempotencyKey: key, + Mode: mode, OrderNumber: orderNumber, + CreatedFrom: createdFrom, + CreatedTo: createdTo, + SyncToNow: syncToNow, }, ) if err != nil { @@ -191,6 +229,9 @@ func (h *Handler) CreateFreightImport(ctx *gin.Context) { CSRFToken: token, }, OrderNumber: orderNumber, + Mode: mode, + CreatedFrom: createdFrom, + CreatedTo: createdTo, IdempotencyKey: key, Error: "同步任务创建失败,请稍后使用相同提交标识重试。", }) @@ -1056,10 +1097,14 @@ type freightPage struct { type freightImportPage struct { Page pageView + Mode string OrderNumber string + CreatedFrom string + CreatedTo string IdempotencyKey string Error string Sync *FreightSync + Watermark *FreightWatermark } type freightDetailPage struct { diff --git a/backend-api/internal/transport/webui/handler_test.go b/backend-api/internal/transport/webui/handler_test.go index 90957f1..344c7bb 100644 --- a/backend-api/internal/transport/webui/handler_test.go +++ b/backend-api/internal/transport/webui/handler_test.go @@ -959,6 +959,11 @@ func TestFreightPagesEscapeSourceDataAndCreateAsyncSync(t *testing.T) { nil, "", ) + if !strings.Contains(form.Body.String(), "同步至现在") || + !strings.Contains(form.Body.String(), `name="created_from"`) || + !strings.Contains(form.Body.String(), "尚无日期同步水位") { + t.Fatalf("freight import form = %s", form.Body) + } cookie := csrfCookie(t, form) idempotencyKey := hiddenValue( t, @@ -994,6 +999,102 @@ func TestFreightPagesEscapeSourceDataAndCreateAsyncSync(t *testing.T) { } } +func TestFreightDateFormSubmitsManualRange(t *testing.T) { + now := time.Date(2026, 7, 28, 3, 4, 5, 0, time.UTC) + service := &fakeFreightService{ + fakeService: &fakeService{}, + watermark: &FreightWatermark{ + LastSuccessfulTo: now, + LastSuccessfulRunID: testTaskID, + UpdatedAt: now, + }, + createResult: FreightSync{ + ID: testTaskID, + Mode: "CREATED_RANGE", + Status: "PENDING", + CreatedAt: now, + }, + } + router := newTestRouter(t, service) + form := performRequest(t, router, http.MethodGet, "/freight/import", nil, "") + if !strings.Contains(form.Body.String(), "增量同步水位") || + !strings.Contains(form.Body.String(), testTaskID) { + t.Fatalf("watermark form = %s", form.Body) + } + cookie := csrfCookie(t, form) + key := hiddenValue(t, form.Body.String(), "idempotency_key") + values := url.Values{ + "csrf_token": {cookie.Value}, + "idempotency_key": {key}, + "mode": {"CREATED_RANGE"}, + "created_from": {"2026-07-22"}, + "created_to": {"2026-07-28"}, + } + request := httptest.NewRequest( + http.MethodPost, + "/freight/import", + 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 || + service.createInput.Mode != "CREATED_RANGE" || + service.createInput.CreatedFrom != "2026-07-22" || + service.createInput.CreatedTo != "2026-07-28" { + t.Fatalf( + "date form response/input = %d / %+v", + response.Code, + service.createInput, + ) + } +} + +func TestFreightSyncToNowIgnoresPrefilledManualDates(t *testing.T) { + service := &fakeFreightService{ + fakeService: &fakeService{}, + createResult: FreightSync{ + ID: testTaskID, + Mode: "CREATED_RANGE", + Status: "PENDING", + }, + } + router := newTestRouter(t, service) + form := performRequest(t, router, http.MethodGet, "/freight/import", nil, "") + cookie := csrfCookie(t, form) + key := hiddenValue(t, form.Body.String(), "idempotency_key") + values := url.Values{ + "csrf_token": {cookie.Value}, + "idempotency_key": {key}, + "mode": {"CREATED_RANGE"}, + "created_from": {"2026-07-22"}, + "created_to": {"2026-07-28"}, + "sync_to_now": {"true"}, + } + request := httptest.NewRequest( + http.MethodPost, + "/freight/import", + 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 || + !service.createInput.SyncToNow || + service.createInput.CreatedFrom != "" || + service.createInput.CreatedTo != "" { + t.Fatalf( + "sync-to-now response/input = %d / %+v", + response.Code, + service.createInput, + ) + } +} + func TestFreightDetailCreatesProcurementTaskWithCSRF(t *testing.T) { const itemID = "00000000-0000-4000-8000-000000000002" service := &fakeProcurementService{ @@ -1101,6 +1202,7 @@ type fakeFreightService struct { sync FreightSync createInput CreateFreightSyncInput createResult FreightSync + watermark *FreightWatermark err error } @@ -1159,6 +1261,12 @@ func (service *fakeFreightService) GetFreightSync( return service.sync, service.err } +func (service *fakeFreightService) GetFreightWatermark( + context.Context, +) (*FreightWatermark, error) { + return service.watermark, service.err +} + func (service *fakeFreightService) CreateFreightSync( _ context.Context, input CreateFreightSyncInput, diff --git a/backend-api/internal/transport/webui/static/admin.css b/backend-api/internal/transport/webui/static/admin.css index 4ec1d8b..4ed46a7 100644 --- a/backend-api/internal/transport/webui/static/admin.css +++ b/backend-api/internal/transport/webui/static/admin.css @@ -688,6 +688,27 @@ tbody tr:last-child td { background: var(--surface); } +.form-panel + .form-panel { + margin-top: 18px; +} + +.form-panel h2 { + margin: 0; + font-size: 18px; +} + +.date-range-fields { + display: grid; + grid-template-columns: repeat(2, minmax(0, 1fr)); + gap: 14px; +} + +.button-row { + display: flex; + flex-wrap: wrap; + gap: 10px; +} + .detail-section { margin-bottom: 22px; } @@ -1273,7 +1294,8 @@ tbody tr:last-child td { } .filters, - .form-grid { + .form-grid, + .date-range-fields { grid-template-columns: 1fr; } diff --git a/backend-api/internal/transport/webui/templates/freight-import.gohtml b/backend-api/internal/transport/webui/templates/freight-import.gohtml index b25485e..3ae2bff 100644 --- a/backend-api/internal/transport/webui/templates/freight-import.gohtml +++ b/backend-api/internal/transport/webui/templates/freight-import.gohtml @@ -11,7 +11,7 @@

导入 ERP 货运

-

使用 ERP 页面“全部单号”中的完整单号

+

按完整单号导入,或按创建日期发现新增货运

返回列表
@@ -20,6 +20,10 @@

同步状态

+
方式
{{if eq .Sync.Mode "CREATED_RANGE"}}创建日期{{else}}完整单号{{end}}
+ {{if eq .Sync.Mode "CREATED_RANGE"}} +
查询范围
{{.Sync.CreatedFrom}} 至 {{.Sync.CreatedTo}}
+ {{end}}
状态
{{.Sync.Status}}
货运单
{{.Sync.OrderCount}}
商品明细
{{.Sync.ItemCount}}
@@ -30,15 +34,49 @@ {{end}}
{{end}} + {{if .Watermark}} +
+

增量同步水位

+
+
已成功同步至
+
最近成功任务
{{.Watermark.LastSuccessfulRunID}}
+
+
+ {{else}} +
尚无日期同步水位,“同步至现在”将从今天开始。
+ {{end}}
+ +

按完整单号

- + +
+
+ + + +

按创建日期

+
+
+ + +
+
+ + +
+
+
+ + +
diff --git a/backend-api/internal/transport/webui/types.go b/backend-api/internal/transport/webui/types.go index f64f16a..4ab6f06 100644 --- a/backend-api/internal/transport/webui/types.go +++ b/backend-api/internal/transport/webui/types.go @@ -31,6 +31,7 @@ type FreightService interface { ListFreightOrders(context.Context, int) ([]FreightOrder, error) GetFreightOrder(context.Context, string) (FreightOrderDetail, error) GetFreightSync(context.Context, string) (FreightSync, error) + GetFreightWatermark(context.Context) (*FreightWatermark, error) CreateFreightSync( context.Context, CreateFreightSyncInput, @@ -71,19 +72,33 @@ type CreateProcurementTaskInput struct { } type FreightSync struct { - ID string - Status string - ErrorCode string - OrderCount int - ItemCount int - CreatedAt time.Time - FinishedAt time.Time + ID string + Mode string + CreatedFrom string + CreatedTo string + WatermarkThrough time.Time + Status string + ErrorCode string + OrderCount int + ItemCount int + CreatedAt time.Time + FinishedAt time.Time +} + +type FreightWatermark struct { + LastSuccessfulTo time.Time + LastSuccessfulRunID string + UpdatedAt time.Time } type CreateFreightSyncInput struct { ActorUserID string IdempotencyKey string + Mode string OrderNumber string + CreatedFrom string + CreatedTo string + SyncToNow bool } type FreightOrder struct { diff --git a/backend-api/internal/transport/webui/usecase_adapter.go b/backend-api/internal/transport/webui/usecase_adapter.go index 0c404ed..fd6f685 100644 --- a/backend-api/internal/transport/webui/usecase_adapter.go +++ b/backend-api/internal/transport/webui/usecase_adapter.go @@ -270,21 +270,57 @@ func (adapter *UsecaseAdapter) CreateFreightSync( if adapter.freight == nil { return FreightSync{}, ErrUnavailable } - result, err := adapter.freight.CreateOrderSync( - ctx, - usecase.CreateFreightSyncCommand{ - CreatorSubject: localAdminSubject, - ActorUserID: input.ActorUserID, - IdempotencyKey: input.IdempotencyKey, - OrderNumber: input.OrderNumber, - }, - ) + var result usecase.CreateFreightSyncResult + var err error + if input.Mode == domain.FreightSyncCreatedRange { + result, err = adapter.freight.CreateDateSync( + ctx, + usecase.CreateFreightDateSyncCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: input.ActorUserID, + IdempotencyKey: input.IdempotencyKey, + CreatedFrom: input.CreatedFrom, + CreatedTo: input.CreatedTo, + SyncToNow: input.SyncToNow, + }, + ) + } else { + result, err = adapter.freight.CreateOrderSync( + ctx, + usecase.CreateFreightSyncCommand{ + CreatorSubject: localAdminSubject, + ActorUserID: input.ActorUserID, + IdempotencyKey: input.IdempotencyKey, + OrderNumber: input.OrderNumber, + }, + ) + } if err != nil { return FreightSync{}, mapUsecaseError(err) } return freightSyncFrom(result.Run), nil } +func (adapter *UsecaseAdapter) GetFreightWatermark( + ctx context.Context, +) (*FreightWatermark, error) { + if adapter.freight == nil { + return nil, ErrUnavailable + } + watermark, err := adapter.freight.GetWatermark(ctx, localAdminSubject) + if err != nil { + return nil, mapUsecaseError(err) + } + if watermark == nil { + return nil, nil + } + return &FreightWatermark{ + LastSuccessfulTo: watermark.LastSuccessfulTo, + LastSuccessfulRunID: watermark.LastSuccessfulRunID, + UpdatedAt: watermark.UpdatedAt, + }, nil +} + func freightOrderFrom(order domain.FreightOrder) FreightOrder { result := FreightOrder{ ID: order.ID, @@ -308,12 +344,18 @@ func freightOrderFrom(order domain.FreightOrder) FreightOrder { func freightSyncFrom(run domain.FreightSyncRun) FreightSync { result := FreightSync{ - ID: run.ID, - Status: string(run.Status), - ErrorCode: stringValue(run.ErrorCode), - OrderCount: run.OrderCount, - ItemCount: run.ItemCount, - CreatedAt: run.CreatedAt, + ID: run.ID, + Mode: run.Mode, + CreatedFrom: run.CreatedFrom, + CreatedTo: run.CreatedTo, + Status: string(run.Status), + ErrorCode: stringValue(run.ErrorCode), + OrderCount: run.OrderCount, + ItemCount: run.ItemCount, + CreatedAt: run.CreatedAt, + } + if run.WatermarkThrough != nil { + result.WatermarkThrough = *run.WatermarkThrough } if run.FinishedAt != nil { result.FinishedAt = *run.FinishedAt diff --git a/backend-api/internal/usecase/freight_ports.go b/backend-api/internal/usecase/freight_ports.go index ce311a4..3946bd0 100644 --- a/backend-api/internal/usecase/freight_ports.go +++ b/backend-api/internal/usecase/freight_ports.go @@ -9,6 +9,11 @@ import ( type FreightSource interface { QueryOrder(context.Context, string) (domain.FreightSourceBatch, error) + QueryCreatedRange( + context.Context, + string, + string, + ) (domain.FreightSourceBatch, error) } type FreightRepository interface { @@ -25,6 +30,13 @@ type FreightRepository interface { domain.FreightImportBatch, time.Time, ) error + CompleteFreightDateSync( + context.Context, + domain.FreightSyncRun, + domain.FreightImportBatch, + time.Time, + time.Time, + ) error FailFreightSync(context.Context, string, string, time.Time) error RecoverFreightSyncs(context.Context, time.Time) (int64, error) GetFreightSync( @@ -32,6 +44,10 @@ type FreightRepository interface { string, string, ) (domain.FreightSyncRun, error) + GetFreightSyncWatermark( + context.Context, + string, + ) (*domain.FreightSyncWatermark, error) ListFreightOrders( context.Context, string, diff --git a/backend-api/internal/usecase/freight_service.go b/backend-api/internal/usecase/freight_service.go index 865af8c..09f8d97 100644 --- a/backend-api/internal/usecase/freight_service.go +++ b/backend-api/internal/usecase/freight_service.go @@ -19,6 +19,8 @@ import ( const ( maxFreightOrdersPerSync = 100 maxFreightItemsPerOrder = 1000 + maxFreightWindowDays = 7 + freightWatermarkOverlap = 10 * time.Minute ) type FreightService struct { @@ -36,6 +38,15 @@ type CreateFreightSyncCommand struct { OrderNumber string } +type CreateFreightDateSyncCommand struct { + CreatorSubject string + ActorUserID string + IdempotencyKey string + CreatedFrom string + CreatedTo string + SyncToNow bool +} + type CreateFreightSyncResult struct { Run domain.FreightSyncRun Replayed bool @@ -127,6 +138,147 @@ func (service *FreightService) CreateOrderSync( return CreateFreightSyncResult{Run: run, Replayed: !created}, nil } +func (service *FreightService) CreateDateSync( + ctx context.Context, + command CreateFreightDateSyncCommand, +) (CreateFreightSyncResult, error) { + command.CreatorSubject = strings.TrimSpace(command.CreatorSubject) + command.ActorUserID = strings.TrimSpace(command.ActorUserID) + command.IdempotencyKey = strings.TrimSpace(command.IdempotencyKey) + command.CreatedFrom = strings.TrimSpace(command.CreatedFrom) + command.CreatedTo = strings.TrimSpace(command.CreatedTo) + fields := validateFreightActor( + command.CreatorSubject, + command.ActorUserID, + command.IdempotencyKey, + ) + if command.SyncToNow && + (command.CreatedFrom != "" || command.CreatedTo != "") { + fields["sync_to_now"] = "cannot be combined with a manual range" + } + if !command.SyncToNow && + (command.CreatedFrom == "" || command.CreatedTo == "") { + fields["created_range"] = "created_from and created_to are required" + } + if len(fields) > 0 { + return CreateFreightSyncResult{}, invalidError( + "FREIGHT_SYNC_INVALID", + "freight sync request is invalid", + fields, + ) + } + + now := service.clock.Now().UTC() + location, err := time.LoadLocation("Asia/Shanghai") + if err != nil { + return CreateFreightSyncResult{}, wrapRepositoryError(err) + } + var startDate, endDate time.Time + var watermarkThrough time.Time + if command.SyncToNow { + watermark, err := service.repository.GetFreightSyncWatermark( + ctx, + command.CreatorSubject, + ) + if err != nil { + return CreateFreightSyncResult{}, wrapRepositoryError(err) + } + localNow := now.In(location) + endDate = startOfLocalDay(localNow, location) + startDate = endDate + if watermark != nil { + overlap := watermark.LastSuccessfulTo.Add( + -freightWatermarkOverlap, + ).In(location) + startDate = startOfLocalDay(overlap, location) + if startDate.After(endDate) { + startDate = endDate + } + } + watermarkThrough = now + } else { + startDate, err = parseFreightDate(command.CreatedFrom, location) + if err != nil { + fields["created_from"] = "must be YYYY-MM-DD" + } + endDate, err = parseFreightDate(command.CreatedTo, location) + if err != nil { + fields["created_to"] = "must be YYYY-MM-DD" + } + if len(fields) == 0 { + if endDate.Before(startDate) { + fields["created_to"] = "must not be before created_from" + } else if daysInclusive(startDate, endDate) > + maxFreightWindowDays { + fields["created_range"] = "must not exceed 7 inclusive days" + } + if endDate.After(startOfLocalDay(now.In(location), location)) { + fields["created_to"] = "must not be in the future" + } + } + if len(fields) > 0 { + return CreateFreightSyncResult{}, invalidError( + "FREIGHT_SYNC_INVALID", + "freight sync request is invalid", + fields, + ) + } + watermarkThrough = endDate.AddDate(0, 0, 1).Add(-time.Nanosecond).UTC() + if watermarkThrough.After(now) { + watermarkThrough = now + } + } + + createdFrom := startDate.Format(time.DateOnly) + createdTo := endDate.Format(time.DateOnly) + runID, err := service.ids.NewID() + if err != nil { + return CreateFreightSyncResult{}, wrapRepositoryError(err) + } + queryHash := hashJSON(struct { + Mode string `json:"mode"` + CreatedFrom string `json:"created_from"` + CreatedTo string `json:"created_to"` + }{ + domain.FreightSyncCreatedRange, + createdFrom, + createdTo, + }) + requestHash := hashJSON(struct { + SyncToNow bool `json:"sync_to_now"` + CreatedFrom string `json:"created_from"` + CreatedTo string `json:"created_to"` + }{ + command.SyncToNow, + command.CreatedFrom, + command.CreatedTo, + }) + run, created, err := service.repository.CreateFreightSync( + ctx, + domain.FreightSyncRun{ + ID: runID, + CreatorSubject: command.CreatorSubject, + CreatedByUserID: command.ActorUserID, + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: createdFrom, + CreatedTo: createdTo, + WatermarkThrough: &watermarkThrough, + QuerySHA256: queryHash, + Status: domain.FreightSyncPending, + CreatedAt: now, + }, + command.IdempotencyKey, + requestHash, + ) + if err != nil { + return CreateFreightSyncResult{}, wrapRepositoryError(err) + } + if created { + go service.execute(run) + } + return CreateFreightSyncResult{Run: run, Replayed: !created}, nil +} + func (service *FreightService) execute(run domain.FreightSyncRun) { ctx, cancel := context.WithTimeout(context.Background(), service.timeout) defer cancel() @@ -134,7 +286,13 @@ func (service *FreightService) execute(run domain.FreightSyncRun) { if err := service.repository.StartFreightSync(ctx, run.ID, now); err != nil { return } - source, err := service.source.QueryOrder(ctx, run.OrderNumber) + var source domain.FreightSourceBatch + var err error + if run.Mode == domain.FreightSyncCreatedRange { + source, err = service.queryCreatedRange(ctx, run) + } else { + source, err = service.source.QueryOrder(ctx, run.OrderNumber) + } if err != nil { _ = service.repository.FailFreightSync( ctx, @@ -154,12 +312,25 @@ func (service *FreightService) execute(run domain.FreightSyncRun) { ) return } - if err := service.repository.CompleteFreightSync( - ctx, - run, - batch, - service.clock.Now().UTC(), - ); err != nil { + finishedAt := service.clock.Now().UTC() + if run.Mode == domain.FreightSyncCreatedRange && + run.WatermarkThrough != nil { + err = service.repository.CompleteFreightDateSync( + ctx, + run, + batch, + *run.WatermarkThrough, + finishedAt, + ) + } else { + err = service.repository.CompleteFreightSync( + ctx, + run, + batch, + finishedAt, + ) + } + if err != nil { _ = service.repository.FailFreightSync( ctx, run.ID, @@ -169,6 +340,76 @@ func (service *FreightService) execute(run domain.FreightSyncRun) { } } +func (service *FreightService) queryCreatedRange( + ctx context.Context, + run domain.FreightSyncRun, +) (domain.FreightSourceBatch, error) { + location, err := time.LoadLocation("Asia/Shanghai") + if err != nil { + return domain.FreightSourceBatch{}, err + } + start, err := parseFreightDate(run.CreatedFrom, location) + if err != nil { + return domain.FreightSourceBatch{}, erpconnector.ErrProtocol + } + end, err := parseFreightDate(run.CreatedTo, location) + if err != nil || end.Before(start) { + return domain.FreightSourceBatch{}, erpconnector.ErrProtocol + } + orders := make([]domain.FreightSourceOrder, 0) + seen := make(map[string]domain.FreightSourceOrder) + for windowStart := start; !windowStart.After(end); { + windowEnd := windowStart.AddDate(0, 0, maxFreightWindowDays-1) + if windowEnd.After(end) { + windowEnd = end + } + fromValue := windowStart.Format(time.DateOnly) + toValue := windowEnd.Format(time.DateOnly) + batch, err := service.source.QueryCreatedRange( + ctx, + fromValue, + toValue, + ) + if err != nil { + return domain.FreightSourceBatch{}, err + } + if batch.SchemaVersion != 1 || + batch.Query.Mode != domain.FreightSyncCreatedRange || + batch.Query.CreatedFrom == nil || + batch.Query.CreatedTo == nil || + *batch.Query.CreatedFrom != fromValue || + *batch.Query.CreatedTo != toValue { + return domain.FreightSourceBatch{}, erpconnector.ErrProtocol + } + for _, order := range batch.Orders { + existing, exists := seen[order.ExternalStockID] + if exists { + if hashJSON(existing) != hashJSON(order) { + return domain.FreightSourceBatch{}, erpconnector.ErrProtocol + } + continue + } + seen[order.ExternalStockID] = order + orders = append(orders, order) + if len(orders) > maxFreightOrdersPerSync { + return domain.FreightSourceBatch{}, erpconnector.ErrProtocol + } + } + windowStart = windowEnd.AddDate(0, 0, 1) + } + createdFrom := run.CreatedFrom + createdTo := run.CreatedTo + return domain.FreightSourceBatch{ + SchemaVersion: 1, + Query: domain.FreightSourceQuery{ + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: &createdFrom, + CreatedTo: &createdTo, + }, + Orders: orders, + }, nil +} + func (service *FreightService) RecoverInterrupted( ctx context.Context, ) (int64, error) { @@ -197,6 +438,20 @@ func (service *FreightService) GetSync( return run, nil } +func (service *FreightService) GetWatermark( + ctx context.Context, + creatorSubject string, +) (*domain.FreightSyncWatermark, error) { + watermark, err := service.repository.GetFreightSyncWatermark( + ctx, + strings.TrimSpace(creatorSubject), + ) + if err != nil { + return nil, wrapRepositoryError(err) + } + return watermark, nil +} + func (service *FreightService) ListOrders( ctx context.Context, creatorSubject string, @@ -242,7 +497,8 @@ func (service *FreightService) normalize( source domain.FreightSourceBatch, ) (domain.FreightImportBatch, error) { if source.SchemaVersion != 1 || - source.Query.Mode != domain.FreightSyncOrderNumber || + (source.Query.Mode != domain.FreightSyncOrderNumber && + source.Query.Mode != domain.FreightSyncCreatedRange) || len(source.Orders) > maxFreightOrdersPerSync { return domain.FreightImportBatch{}, errors.New("invalid source envelope") } @@ -373,6 +629,47 @@ func (service *FreightService) normalize( return result, nil } +func validateFreightActor( + creatorSubject, actorUserID, idempotencyKey string, +) map[string]string { + fields := map[string]string{} + if creatorSubject == "" { + fields["creator_subject"] = "is required" + } + if actorUserID == "" { + fields["actor_user_id"] = "is required" + } + if idempotencyKey == "" || len([]byte(idempotencyKey)) > 128 { + fields["idempotency_key"] = "must contain 1 to 128 bytes" + } + return fields +} + +func parseFreightDate(value string, location *time.Location) (time.Time, error) { + parsed, err := time.ParseInLocation(time.DateOnly, value, location) + if err != nil || parsed.Format(time.DateOnly) != value { + return time.Time{}, errors.New("invalid freight date") + } + return parsed, nil +} + +func startOfLocalDay(value time.Time, location *time.Location) time.Time { + return time.Date( + value.Year(), + value.Month(), + value.Day(), + 0, + 0, + 0, + 0, + location, + ) +} + +func daysInclusive(start, end time.Time) int { + return int(end.Sub(start)/(24*time.Hour)) + 1 +} + func freightSourceErrorCode(err error) string { switch { case errors.Is(err, erpconnector.ErrNotConfigured): diff --git a/backend-api/internal/usecase/freight_service_test.go b/backend-api/internal/usecase/freight_service_test.go index c0204b8..ceba7b7 100644 --- a/backend-api/internal/usecase/freight_service_test.go +++ b/backend-api/internal/usecase/freight_service_test.go @@ -1,7 +1,10 @@ package usecase import ( + "context" + "errors" "testing" + "time" "cmroubao/backend-api/internal/domain" ) @@ -52,6 +55,125 @@ func TestFreightNormalizationRejectsConflictingIdentityAndInvalidTime( } } +func TestFreightDateQuerySplitsIntoSevenDayWindows(t *testing.T) { + source := &recordingDateSource{} + service := &FreightService{source: source} + + result, err := service.queryCreatedRange( + context.Background(), + domain.FreightSyncRun{ + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: "2026-07-01", + CreatedTo: "2026-07-15", + }, + ) + + if err != nil { + t.Fatalf("queryCreatedRange() error = %v", err) + } + want := [][2]string{ + {"2026-07-01", "2026-07-07"}, + {"2026-07-08", "2026-07-14"}, + {"2026-07-15", "2026-07-15"}, + } + if len(source.calls) != len(want) { + t.Fatalf("calls = %#v", source.calls) + } + for index := range want { + if source.calls[index] != want[index] { + t.Fatalf("call %d = %#v, want %#v", index, source.calls[index], want[index]) + } + } + if result.Query.CreatedFrom == nil || + *result.Query.CreatedFrom != "2026-07-01" || + result.Query.CreatedTo == nil || + *result.Query.CreatedTo != "2026-07-15" { + t.Fatalf("merged query = %+v", result.Query) + } +} + +func TestFreightDateQueryStopsOnMiddleWindowFailure(t *testing.T) { + source := &recordingDateSource{failOnCall: 2} + service := &FreightService{source: source} + + _, err := service.queryCreatedRange( + context.Background(), + domain.FreightSyncRun{ + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: "2026-07-01", + CreatedTo: "2026-07-15", + }, + ) + + if !errors.Is(err, errDateSourceFailure) || len(source.calls) != 2 { + t.Fatalf("queryCreatedRange() error/calls = %v / %#v", err, source.calls) + } +} + +func TestCreateFreightSyncToNowUsesWatermarkOverlapDay(t *testing.T) { + watermarkTime := time.Date(2026, 7, 25, 16, 5, 0, 0, time.UTC) + repository := &dateCaptureRepository{ + watermark: &domain.FreightSyncWatermark{ + CreatorSubject: "local-admin", + SourceSystem: domain.FreightSourceShunyunbao, + LastSuccessfulTo: watermarkTime, + }, + } + service, err := NewFreightService( + repository, + &recordingDateSource{}, + fakeClock{}, + &sequenceIDs{}, + time.Minute, + ) + if err != nil { + t.Fatalf("NewFreightService() error = %v", err) + } + + result, err := service.CreateDateSync( + context.Background(), + CreateFreightDateSyncCommand{ + CreatorSubject: "local-admin", + ActorUserID: "00000000-0000-4000-8000-000000000099", + IdempotencyKey: "sync-to-now", + SyncToNow: true, + }, + ) + + if err != nil { + t.Fatalf("CreateDateSync() error = %v", err) + } + if result.Run.CreatedFrom != "2026-07-25" || + result.Run.CreatedTo != "2026-07-26" || + result.Run.WatermarkThrough == nil || + !result.Run.WatermarkThrough.Equal(fakeClock{}.Now()) { + t.Fatalf("sync-to-now run = %+v", result.Run) + } +} + +func TestCreateFreightManualRangeRejectsMoreThanSevenDays(t *testing.T) { + service, _ := NewFreightService( + &dateCaptureRepository{}, + &recordingDateSource{}, + fakeClock{}, + &sequenceIDs{}, + time.Minute, + ) + + _, err := service.CreateDateSync( + context.Background(), + CreateFreightDateSyncCommand{ + CreatorSubject: "local-admin", + ActorUserID: "00000000-0000-4000-8000-000000000099", + IdempotencyKey: "too-wide", + CreatedFrom: "2026-07-19", + CreatedTo: "2026-07-26", + }, + ) + + assertUsecaseError(t, err, ErrorKindInvalid, "FREIGHT_SYNC_INVALID") +} + func validFreightSource() domain.FreightSourceBatch { shop := "测试店铺" sourceCreatedAt := "2026-07-28 08:00:00" @@ -86,3 +208,122 @@ func validFreightSource() domain.FreightSourceBatch { }}, } } + +var errDateSourceFailure = errors.New("date source failed") + +type recordingDateSource struct { + calls [][2]string + failOnCall int +} + +func (source *recordingDateSource) QueryOrder( + context.Context, + string, +) (domain.FreightSourceBatch, error) { + return validFreightSource(), nil +} + +func (source *recordingDateSource) QueryCreatedRange( + _ context.Context, + createdFrom, createdTo string, +) (domain.FreightSourceBatch, error) { + source.calls = append(source.calls, [2]string{createdFrom, createdTo}) + if source.failOnCall == len(source.calls) { + return domain.FreightSourceBatch{}, errDateSourceFailure + } + return domain.FreightSourceBatch{ + SchemaVersion: 1, + Query: domain.FreightSourceQuery{ + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: &createdFrom, + CreatedTo: &createdTo, + }, + Orders: []domain.FreightSourceOrder{}, + }, nil +} + +type dateCaptureRepository struct { + watermark *domain.FreightSyncWatermark +} + +func (repository *dateCaptureRepository) CreateFreightSync( + _ context.Context, + run domain.FreightSyncRun, + _, _ string, +) (domain.FreightSyncRun, bool, error) { + return run, false, nil +} + +func (*dateCaptureRepository) StartFreightSync( + context.Context, + string, + time.Time, +) error { + return nil +} + +func (*dateCaptureRepository) CompleteFreightSync( + context.Context, + domain.FreightSyncRun, + domain.FreightImportBatch, + time.Time, +) error { + return nil +} + +func (*dateCaptureRepository) CompleteFreightDateSync( + context.Context, + domain.FreightSyncRun, + domain.FreightImportBatch, + time.Time, + time.Time, +) error { + return nil +} + +func (*dateCaptureRepository) FailFreightSync( + context.Context, + string, + string, + time.Time, +) error { + return nil +} + +func (*dateCaptureRepository) RecoverFreightSyncs( + context.Context, + time.Time, +) (int64, error) { + return 0, nil +} + +func (*dateCaptureRepository) GetFreightSync( + context.Context, + string, + string, +) (domain.FreightSyncRun, error) { + return domain.FreightSyncRun{}, nil +} + +func (repository *dateCaptureRepository) GetFreightSyncWatermark( + context.Context, + string, +) (*domain.FreightSyncWatermark, error) { + return repository.watermark, nil +} + +func (*dateCaptureRepository) ListFreightOrders( + context.Context, + string, + int, +) ([]domain.FreightOrder, error) { + return nil, nil +} + +func (*dateCaptureRepository) GetFreightOrder( + context.Context, + string, + string, +) (domain.FreightOrderDetail, error) { + return domain.FreightOrderDetail{}, nil +} diff --git a/backend-api/migrations/00014_erp_date_sync.sql b/backend-api/migrations/00014_erp_date_sync.sql new file mode 100644 index 0000000..1a1f053 --- /dev/null +++ b/backend-api/migrations/00014_erp_date_sync.sql @@ -0,0 +1,234 @@ +-- +goose NO TRANSACTION +-- +goose Up +PRAGMA foreign_keys = OFF; +PRAGMA legacy_alter_table = ON; +BEGIN IMMEDIATE; + +ALTER TABLE erp_sync_runs RENAME TO erp_sync_runs_v13; + +CREATE TABLE erp_sync_runs ( + id TEXT PRIMARY KEY NOT NULL CHECK (length(id) = 36), + creator_subject TEXT NOT NULL, + created_by_user_id TEXT NOT NULL + REFERENCES users(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + mode TEXT NOT NULL CHECK (mode IN ('ORDER_NUMBER', 'CREATED_RANGE')), + order_number TEXT + CHECK ( + order_number IS NULL + OR ( + length(trim(order_number)) > 0 + AND length(CAST(order_number AS BLOB)) <= 128 + ) + ), + created_from TEXT + CHECK ( + created_from IS NULL + OR ( + length(created_from) = 10 + AND date(created_from) = created_from + ) + ), + created_to TEXT + CHECK ( + created_to IS NULL + OR ( + length(created_to) = 10 + AND date(created_to) = created_to + ) + ), + watermark_through TEXT, + query_sha256 TEXT NOT NULL + CHECK ( + length(query_sha256) = 64 + AND query_sha256 NOT GLOB '*[^0-9a-f]*' + ), + idempotency_key TEXT NOT NULL + CHECK ( + length(trim(idempotency_key)) > 0 + AND length(CAST(idempotency_key AS BLOB)) <= 128 + ), + request_sha256 TEXT NOT NULL + CHECK ( + length(request_sha256) = 64 + AND request_sha256 NOT GLOB '*[^0-9a-f]*' + ), + status TEXT NOT NULL + CHECK (status IN ('PENDING', 'RUNNING', 'SUCCEEDED', 'FAILED')), + error_code TEXT + CHECK ( + error_code IS NULL + OR length(CAST(error_code AS BLOB)) <= 64 + ), + order_count INTEGER NOT NULL DEFAULT 0 CHECK (order_count >= 0), + item_count INTEGER NOT NULL DEFAULT 0 CHECK (item_count >= 0), + created_at TEXT NOT NULL, + started_at TEXT, + finished_at TEXT, + CHECK ( + ( + mode = 'ORDER_NUMBER' + AND order_number IS NOT NULL + AND created_from IS NULL + AND created_to IS NULL + AND watermark_through IS NULL + ) + OR ( + mode = 'CREATED_RANGE' + AND order_number IS NULL + AND created_from IS NOT NULL + AND created_to IS NOT NULL + AND created_from <= created_to + AND watermark_through IS NOT NULL + ) + ), + CHECK ( + (status = 'PENDING' + AND started_at IS NULL + AND finished_at IS NULL + AND error_code IS NULL) + OR (status = 'RUNNING' + AND started_at IS NOT NULL + AND finished_at IS NULL + AND error_code IS NULL) + OR (status = 'SUCCEEDED' + AND started_at IS NOT NULL + AND finished_at IS NOT NULL + AND error_code IS NULL) + OR (status = 'FAILED' + AND finished_at IS NOT NULL + AND error_code IS NOT NULL) + ), + UNIQUE (creator_subject, idempotency_key) +); + +INSERT INTO erp_sync_runs ( + id, creator_subject, created_by_user_id, mode, order_number, + created_from, created_to, watermark_through, query_sha256, + idempotency_key, request_sha256, status, error_code, order_count, + item_count, created_at, started_at, finished_at +) +SELECT + id, creator_subject, created_by_user_id, mode, order_number, + NULL, NULL, NULL, query_sha256, idempotency_key, request_sha256, + status, error_code, order_count, item_count, created_at, started_at, + finished_at +FROM erp_sync_runs_v13; + +DROP TABLE erp_sync_runs_v13; + +CREATE INDEX erp_sync_runs_creator_created_idx + ON erp_sync_runs (creator_subject, created_at DESC, id DESC); + +CREATE TABLE erp_sync_watermarks ( + creator_subject TEXT NOT NULL, + source_system TEXT NOT NULL CHECK (source_system = 'SHUNYUNBAO'), + last_successful_to TEXT NOT NULL, + last_successful_run_id TEXT NOT NULL + REFERENCES erp_sync_runs(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + updated_at TEXT NOT NULL, + PRIMARY KEY (creator_subject, source_system) +); + +COMMIT; +PRAGMA legacy_alter_table = OFF; +PRAGMA foreign_keys = ON; + +-- +goose Down +DROP TABLE IF EXISTS erp_date_sync_v14_down_guard; +CREATE TEMP TABLE erp_date_sync_v14_down_guard ( + allowed INTEGER NOT NULL CHECK (allowed = 1) +); + +INSERT INTO erp_date_sync_v14_down_guard (allowed) +SELECT CASE + WHEN EXISTS ( + SELECT 1 FROM erp_sync_runs WHERE mode = 'CREATED_RANGE' + ) OR EXISTS (SELECT 1 FROM erp_sync_watermarks) + THEN 0 + ELSE 1 +END; + +DROP TABLE erp_date_sync_v14_down_guard; +PRAGMA foreign_keys = OFF; +PRAGMA legacy_alter_table = ON; +BEGIN IMMEDIATE; +DROP TABLE erp_sync_watermarks; +ALTER TABLE erp_sync_runs RENAME TO erp_sync_runs_v14; + +CREATE TABLE erp_sync_runs ( + id TEXT PRIMARY KEY NOT NULL CHECK (length(id) = 36), + creator_subject TEXT NOT NULL, + created_by_user_id TEXT NOT NULL + REFERENCES users(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + mode TEXT NOT NULL CHECK (mode = 'ORDER_NUMBER'), + order_number TEXT NOT NULL + CHECK ( + length(trim(order_number)) > 0 + AND length(CAST(order_number AS BLOB)) <= 128 + ), + query_sha256 TEXT NOT NULL + CHECK ( + length(query_sha256) = 64 + AND query_sha256 NOT GLOB '*[^0-9a-f]*' + ), + idempotency_key TEXT NOT NULL + CHECK ( + length(trim(idempotency_key)) > 0 + AND length(CAST(idempotency_key AS BLOB)) <= 128 + ), + request_sha256 TEXT NOT NULL + CHECK ( + length(request_sha256) = 64 + AND request_sha256 NOT GLOB '*[^0-9a-f]*' + ), + status TEXT NOT NULL + CHECK (status IN ('PENDING', 'RUNNING', 'SUCCEEDED', 'FAILED')), + error_code TEXT + CHECK ( + error_code IS NULL + OR length(CAST(error_code AS BLOB)) <= 64 + ), + order_count INTEGER NOT NULL DEFAULT 0 CHECK (order_count >= 0), + item_count INTEGER NOT NULL DEFAULT 0 CHECK (item_count >= 0), + created_at TEXT NOT NULL, + started_at TEXT, + finished_at TEXT, + CHECK ( + (status = 'PENDING' + AND started_at IS NULL + AND finished_at IS NULL + AND error_code IS NULL) + OR (status = 'RUNNING' + AND started_at IS NOT NULL + AND finished_at IS NULL + AND error_code IS NULL) + OR (status = 'SUCCEEDED' + AND started_at IS NOT NULL + AND finished_at IS NOT NULL + AND error_code IS NULL) + OR (status = 'FAILED' + AND finished_at IS NOT NULL + AND error_code IS NOT NULL) + ), + UNIQUE (creator_subject, idempotency_key) +); + +INSERT INTO erp_sync_runs ( + id, creator_subject, created_by_user_id, mode, order_number, + query_sha256, idempotency_key, request_sha256, status, error_code, + order_count, item_count, created_at, started_at, finished_at +) +SELECT + id, creator_subject, created_by_user_id, mode, order_number, + query_sha256, idempotency_key, request_sha256, status, error_code, + order_count, item_count, created_at, started_at, finished_at +FROM erp_sync_runs_v14; + +DROP TABLE erp_sync_runs_v14; + +CREATE INDEX erp_sync_runs_creator_created_idx + ON erp_sync_runs (creator_subject, created_at DESC, id DESC); + +COMMIT; +PRAGMA legacy_alter_table = OFF; +PRAGMA foreign_keys = ON; diff --git a/docs/api.md b/docs/api.md index 4b92847..b807e3e 100644 --- a/docs/api.md +++ b/docs/api.md @@ -172,7 +172,17 @@ Connector 是 loopback 内部服务,不复用 ADMIN/BUYER 凭证。除 `/healt ### Connector `POST /v1/freight/query` ```json -{"order_number":"完整单号"} +{"mode":"ORDER_NUMBER","order_number":"完整单号"} +``` + +或按 Asia/Shanghai 自然日闭区间查询,单次最多 7 天: + +```json +{ + "mode":"CREATED_RANGE", + "created_from":"2026-07-22", + "created_to":"2026-07-28" +} ``` 成功返回最小化规范结构: @@ -219,11 +229,15 @@ T-224 增加: ```json { "mode":"CREATED_RANGE", - "created_from":"2026-07-27T16:00:00Z", - "created_to":"2026-07-28T16:00:00Z" + "created_from":"2026-07-22", + "created_to":"2026-07-28" } ``` +手动范围最多 7 个自然日;也可以只提交 +`{"mode":"CREATED_RANGE","sync_to_now":true}`。后者在没有水位时从当天开始,有水位时 +从成功水位前回看 10 分钟对应的自然日开始,后端再按最多 7 天切窗。 + 响应 `202`,返回 sync id/status。订单号不进入 URL、事件 message 或访问日志;数据库 只保存规范值及用于审计/检索的受控字段,不保存 Connector 密钥。 @@ -232,11 +246,16 @@ T-224 增加: 返回 `PENDING/RUNNING/SUCCEEDED/FAILED`、查询模式、匿名查询摘要、开始/结束时间、 订单/商品计数和稳定错误码。失败不返回 ERP 原始 body 或个人信息。 +### `GET /api/v1/freight-sync-watermark` + +返回当前 ADMIN creator 的顺运宝最近成功日期同步水位;尚未成功同步时返回 +`{"watermark":null}`。只有完整批次落库成功才在同一事务推进水位;失败、空响应以外 +的协议错误和较旧范围成功均不会覆盖较新的水位。 + ### `GET /api/v1/freight-orders` -T-222 首版返回当前 ADMIN creator 最近更新的货运单并支持 `limit`(1..100)。 -T-224 增加 `q`、`sync_status`、`created_from/to` 和 `cursor`;`q` 只匹配受控来源 -单号、店铺或内部 UUID,不匹配电话/地址。响应不包含收件人字段或 Connector 原文。 +返回当前 ADMIN creator 最近更新的货运单并支持 `limit`(1..100)。响应不包含 +收件人字段或 Connector 原文;列表筛选和 cursor 后置,不属于本次日期同步闭环。 ### `GET /api/v1/freight-orders/{id}` diff --git a/docs/current-state.md b/docs/current-state.md index 7286890..7cea319 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -5,7 +5,7 @@ ## 当前快照 - 日期:2026-07-28 -- 阶段:T-223 已完成,T-224 ERP 日期增量同步已领取 +- 阶段:T-220 至 T-224 ERP 货运接入闭环已完成 - Git:当前分支为 `main`;T-001 至 T-004、T-101 至 T-104、T-201 至 T-219 均按文档提交、实现提交的顺序纳入历史 - 生产代码:`android-buyer/` 已接入 Roubao Android 源码 @@ -17,6 +17,10 @@ - ERP 货运:Go 后端通过仅限 loopback、服务密钥鉴权的 Python Connector 异步按 完整单号同步;v12 保存同步记录、货运头和全部明细,canonical hash 控制 revision, Admin 已有 `/freight`、`/freight/import`、`/freight/{id}` 与对应 JSON API。 +- ERP 增量同步:v14 支持 Asia/Shanghai 创建日期闭区间和“同步至现在”,Connector + 单窗最多 7 天,后端对较长水位范围切窗并从成功水位前 10 分钟所在自然日回看。 + 货运落库、同步成功和水位推进同事务完成;失败与较旧范围成功不推进水位。Admin + 可查看查询范围、计数、稳定错误码和最近成功水位。 - 采购需求:v13 按货运商品 revision 保存可复核快照;ADMIN 显式确认采购、上传参考 图后生成不可变 PENDING 任务。来源变化不会改写任务,Roubao 可沿用现有 claim 流程。 - 本机 Android 工具:JDK 17.0.13、Command-line Tools 22.0、SDK 34、 @@ -24,9 +28,10 @@ - Android Studio:未安装;`winget` 静默安装卡住后已终止,不阻塞命令行构建 - 测试:T-219 Android Debug/Release 单元测试与构建和根 `init.ps1` 通过; Debug APK `1.4.16 (21)` 已覆盖安装到 PKG110 -- 后端测试:T-223 运行 `go test ./...`、`go test -race ./...`、`go vet ./...`; - 覆盖采购需求 migration 降级保护、一商品一需求、阻塞字段、图片唯一绑定、任务 - 并发重放、来源变化后任务不变、Roubao claim 和 Admin API/SSR +- 后端测试:T-224 运行 `go test ./...`、`go test -race ./...`、`go vet ./...`; + 覆盖 v14 上下迁移、7 天切窗、水位重叠、中途失败、空窗口、重复页、来源 revision、 + 水位事务/不回退、Admin API/SSR 和 Connector 严格响应窗口;Python 22 项伪响应 + 测试通过,未访问真实 ERP - 原型:4 个管理 Web 页面和 7 个 Android 页面均可离线独立打开;Playwright 以 1440×900、390×844、360×800 验证 36 个页面/视口组合,无页面横向溢出、 脚本错误或外部请求,Android 可见交互控件均不小于 44px @@ -165,7 +170,7 @@ | `docs/tasks/T-221.md` | DONE | 顺运宝精确单号 Connector | | `docs/tasks/T-222.md` | DONE | 货运信息存储、API 与 Admin 页面 | | `docs/tasks/T-223.md` | DONE | 待采购需求提取与任务生成 | -| `docs/tasks/T-224.md` | DOING | ERP 日期增量同步 | +| `docs/tasks/T-224.md` | DONE | ERP 日期增量同步 | | `docs/design/` | 已确认 | T-202 原型索引、4 个管理页和 7 个 Android 页面 | | `deepseek总结.txt` | 已有 | 历史讨论摘要,不是正式需求权威 | | `android-buyer/` | 已有 | Roubao `main` 固定 commit 的 Android 基线 | @@ -178,10 +183,11 @@ ## 任务摘要 - 已完成:T-001 至 T-004、T-101 至 T-104、T-201 至 T-219。 -- 已完成:另含 T-220 至 T-223 ERP 契约、Connector、货运存储和采购需求生成。 -- 进行中:T-224 ERP 日期增量同步。 -- 下一步:完成 T-223 采购需求生成和 T-224 日期增量同步,再进入 T-301 P0 UI - 完整交互验收。 +- 已完成:另含 T-220 至 T-224 ERP 契约、Connector、货运存储、采购需求生成和 + 日期增量同步。 +- 进行中:无。 +- 下一步:进入 T-301 P0 UI 完整交互验收;真实 ERP 上线前仍需确认开放 API、 + 数据使用权限并由人员完成验证码登录。 ## 当前可运行内容 diff --git a/docs/routes.md b/docs/routes.md index b1f96b6..ca3312a 100644 --- a/docs/routes.md +++ b/docs/routes.md @@ -11,8 +11,8 @@ T-202 的离线 P0 页面入口见[原型索引](design/index.html)。原型仅 | `/tasks` | 任务列表 | 查看状态、筛选并进入详情 | US-002 | IX-003 | | `/tasks/new` | 新建任务 | 提交图片和采购约束 | US-001 | IX-002 | | `/tasks/{id}` | 任务详情 | 查看输入/候选/证据、创建下单授权,并展示对账中/人工对账/待人工付款状态 | US-002、US-006、US-010 | IX-003、IX-008、IX-011 | -| `/freight` | ERP 货运列表 | 按来源单号、同步状态和时间查看已导入货运单 | US-011 | IX-012 | -| `/freight/import` | ERP 货运导入 | 创建精确单号或日期范围同步 | US-011 | IX-012 | +| `/freight` | ERP 货运列表 | 查看已导入货运单并进入全部商品明细 | US-011 | IX-012 | +| `/freight/import` | ERP 货运导入 | 创建精确单号、最多 7 天日期范围或“同步至现在”任务,并查看成功水位 | US-011 | IX-012 | | `/freight/{id}` | ERP 货运详情 | 查看全部商品明细、来源变化、补图和采购任务生成状态 | US-011、US-012 | IX-012、IX-013 | MVP 登录后默认进入 `/tasks`。未登录访问受保护页面时跳转 `/login` 并携带安全的 @@ -89,7 +89,7 @@ T-207 的 events/candidates/complete/fail 接受原设备对授权内已产生 | `CandidateReview` | Android/Web | 展示候选与硬约束校验,不执行自动化动作 | | `ExecutionController` | Android | UI 命令入口,委托 workflow,不直接调用 Accessibility | | `PendingPaymentSummary` | 管理 Web/App | 只读展示对账订单与人工付款提醒,不提供付款或重提动作 | -| `FreightSyncForm` | 管理 Web | 创建精确单号/日期同步记录,不持有 ERP 凭证 | +| `FreightSyncForm` | 管理 Web | 创建精确单号/日期同步记录、查看成功水位,不持有 ERP 凭证 | | `FreightItemReview` | 管理 Web | 显示规范化来源字段、缺失项和采购任务生成操作 | 具体状态反馈以[交互清单](08-interaction-checklist.md)为准,组件命名可在接入真实框架后 diff --git a/docs/tasks/T-224.md b/docs/tasks/T-224.md index 1a4a368..6078557 100644 --- a/docs/tasks/T-224.md +++ b/docs/tasks/T-224.md @@ -4,7 +4,7 @@ title: ERP 日期增量同步 phase: 2 deps: - T-223 -status: TODO +status: DONE created: 2026-07-28 context_ref: 78dc595 work_branch: null @@ -49,11 +49,11 @@ write_paths: ## 验收要点 -- [ ] 日期 query 与脱敏 HAR 合约一致,时区和闭区间规则有测试。 -- [ ] 重叠窗口、重复页、来源更新、空窗口和中途失败不丢单/不重复。 -- [ ] 失败不推进 watermark,成功重放得到相同货运与采购需求集合。 -- [ ] Admin 可查看范围、状态、计数和匿名错误,不显示 PII/秘密。 -- [ ] Python、Go 全量测试、race/vet、根验证和响应式 SSR 验收通过。 +- [x] 日期 query 与脱敏 HAR 合约一致,时区和闭区间规则有测试。 +- [x] 重叠窗口、重复页、来源更新、空窗口和中途失败不丢单/不重复。 +- [x] 失败不推进 watermark,成功重放得到相同货运与采购需求集合。 +- [x] Admin 可查看范围、状态、计数和匿名错误,不显示 PII/秘密。 +- [x] Python、Go 全量测试、race/vet、根验证和响应式 SSR 验收通过。 ## 边界 @@ -64,3 +64,16 @@ write_paths: ## 执行记录 - 2026-07-28:任务合约已冻结,等待 T-223。 +- 2026-07-28:Connector 增加 `CREATED_RANGE`,严格复现 + `t_stock.created/op=0/type=3/optType=0` 和最多 7 个自然日闭区间;分页按 stock id + 去重,空范围不请求详情,输出继续使用 allowlist schema。 +- 2026-07-28:v14 将同步运行扩展为 `ORDER_NUMBER/CREATED_RANGE` 并新增 creator + 级成功水位。后端“同步至现在”从水位前 10 分钟所在自然日回看,超过 7 天自动切窗; + 货运 upsert、同步成功和水位推进处于同一 SQLite 事务,失败和较旧成功任务不推进。 +- 2026-07-28:Admin JSON/SSR 支持手动日期范围、“同步至现在”、范围/计数/稳定错误码 + 和水位查看。SSR 实测发现并修复预填日期与 `sync_to_now` 同时提交的 422,API 仍严格 + 拒绝混合语义。 +- 2026-07-28:Python 22 项伪响应测试、Go 全量测试、race、vet 和根 `init.ps1` + 通过。隔离 Gin/SQLite + 本地伪 Connector 验证成功/失败水位;Playwright 在 + 1440×900、390×844、360×800 无横向溢出、页头遮挡、过小命令控件或 PII 文本。 + 未访问真实 ERP。 diff --git a/erp-connector/README.md b/erp-connector/README.md index b1430c1..466eddd 100644 --- a/erp-connector/README.md +++ b/erp-connector/README.md @@ -1,7 +1,8 @@ # 顺运宝 ERP Connector -内部只读适配器:保持顺运宝验证码/Cookie 会话,按完整单号查询货运列表和详情, -再输出不含收件人、电话、地址、Cookie、JWT 或原始响应的规范化 JSON。 +内部只读适配器:保持顺运宝验证码/Cookie 会话,按完整单号或最多 7 个自然日的 +创建日期闭区间查询货运列表和详情,再输出不含收件人、电话、地址、Cookie、JWT +或原始响应的规范化 JSON。 协议客户端基于同一开发机上已经用本地 HAR 验证的 `shunyunbaoerp 0.1.0` 代码收敛; 本目录不包含 HAR、真实响应、账号、密码或订单号。 @@ -40,7 +41,17 @@ $env:SHUNYUNBAO_SERVICE_API_KEY = "at-least-32-utf8-bytes" 3. 调用 `POST /v1/freight/query`: ```json - {"order_number":"完整单号"} + {"mode":"ORDER_NUMBER","order_number":"完整单号"} + ``` + + 日期查询严格使用 Asia/Shanghai 的 `YYYY-MM-DD`: + + ```json + { + "mode":"CREATED_RANGE", + "created_from":"2026-07-22", + "created_to":"2026-07-28" + } ``` Connector 不自动 OCR 绕过验证码。服务重启或 ERP 会话过期后需重新执行前两步。 diff --git a/erp-connector/src/shunyunbaoerp/api.py b/erp-connector/src/shunyunbaoerp/api.py index 65627fc..2b23c4c 100644 --- a/erp-connector/src/shunyunbaoerp/api.py +++ b/erp-connector/src/shunyunbaoerp/api.py @@ -35,7 +35,10 @@ class LoginRequest(BaseModel): class FreightQueryRequest(BaseModel): - order_number: str = Field(min_length=1, max_length=128) + mode: str = Field(default="ORDER_NUMBER", max_length=32) + order_number: str | None = Field(default=None, min_length=1, max_length=128) + created_from: str | None = Field(default=None, min_length=10, max_length=10) + created_to: str | None = Field(default=None, min_length=10, max_length=10) def get_client() -> ERPClient: @@ -125,7 +128,22 @@ def query_freight( _: APIKeyDependency, ) -> dict[str, object]: try: - result = get_client().get_freight_details(body.order_number) + if body.mode == "ORDER_NUMBER" and body.order_number: + if body.created_from is not None or body.created_to is not None: + raise ValueError("ORDER_NUMBER 不能包含日期范围") + result = get_client().get_freight_details(body.order_number) + elif ( + body.mode == "CREATED_RANGE" + and body.order_number is None + and body.created_from + and body.created_to + ): + result = get_client().get_freight_details_by_created_range( + body.created_from, + body.created_to, + ) + else: + raise ValueError("查询模式与参数不匹配") return normalize_freight_result(result) except ERPAuthenticationError as exc: raise HTTPException( diff --git a/erp-connector/src/shunyunbaoerp/client.py b/erp-connector/src/shunyunbaoerp/client.py index 38ca1a5..cdb806f 100644 --- a/erp-connector/src/shunyunbaoerp/client.py +++ b/erp-connector/src/shunyunbaoerp/client.py @@ -8,7 +8,7 @@ import random import threading import time from collections.abc import Iterable, Mapping -from datetime import datetime, timezone +from datetime import date, datetime, timedelta, timezone from typing import Any from urllib.parse import urljoin, urlparse @@ -20,6 +20,7 @@ from .constants import ( STOCK_DETAIL_PATH, STOCK_LIST_PATH, STOCK_LIST_TOTAL_PATH, + created_range_query, order_number_query, stock_columns, ) @@ -202,8 +203,21 @@ class ERPClient: """ normalized = self._validate_order_number(order_number) + return self._query_stock(order_number_query(normalized)) + + def query_stock_by_created_range( + self, + created_from: date | str, + created_to: date | str, + ) -> list[dict[str, Any]]: + """按 Asia/Shanghai 自然日创建时间闭区间查询货运列表。""" + + start, end = self._validate_created_range(created_from, created_to) + return self._query_stock(created_range_query(start.isoformat(), end.isoformat())) + + def _query_stock(self, query: Mapping[str, Any]) -> list[dict[str, Any]]: with self._lock: - first_payload = self._stock_payload(normalized, start=0, page_index=1) + first_payload = self._stock_payload(query, start=0, page_index=1) total_raw = self._request_json( "POST", STOCK_LIST_TOTAL_PATH, @@ -228,7 +242,7 @@ class ERPClient: for start in range(0, total, self.page_size): page_index = start // self.page_size + 1 payload = self._stock_payload( - normalized, + query, start=start, page_index=page_index, ) @@ -297,36 +311,65 @@ class ERPClient: stocks = self.query_stock_by_order_number(normalized) if not stocks: raise ERPNotFoundError("未找到对应货运记录") - - ids = [row.get("id") for row in stocks if row.get("id") is not None] - if len(ids) != len(stocks): - raise ERPProtocolError("货运列表存在缺少 id 的记录") - details = self.get_stock_details(ids) - details_by_id: dict[object, dict[str, Any]] = {} - for item in details: - item_id = item.get("id") - if item_id is None: - raise ERPProtocolError("货运详情存在缺少 id 的记录") - existing = details_by_id.get(item_id) - if existing is not None and existing != item: - raise ERPProtocolError("货运详情存在冲突的重复 id") - details_by_id[item_id] = item - - records = [ + return self._join_freight_details( + stocks, { - "stock": stock, - "detail": details_by_id.get(stock["id"]), - } - for stock in stocks - ] - return { - "query": { + "mode": "ORDER_NUMBER", "orderNumber": normalized, "matchField": "allcode", }, - "count": len(records), - "records": records, + ) + + def get_freight_details_by_created_range( + self, + created_from: date | str, + created_to: date | str, + ) -> dict[str, Any]: + """查询一个最多七天的创建日期窗口并合并详情。""" + + start, end = self._validate_created_range(created_from, created_to) + with self._lock: + stocks = self.query_stock_by_created_range(start, end) + return self._join_freight_details( + stocks, + { + "mode": "CREATED_RANGE", + "createdFrom": start.isoformat(), + "createdTo": end.isoformat(), + }, + ) + + def _join_freight_details( + self, + stocks: list[dict[str, Any]], + query: Mapping[str, Any], + ) -> dict[str, Any]: + ids = [row.get("id") for row in stocks if row.get("id") is not None] + if len(ids) != len(stocks): + raise ERPProtocolError("货运列表存在缺少 id 的记录") + details = self.get_stock_details(ids) + details_by_id: dict[object, dict[str, Any]] = {} + for item in details: + item_id = item.get("id") + if item_id is None: + raise ERPProtocolError("货运详情存在缺少 id 的记录") + existing = details_by_id.get(item_id) + if existing is not None and existing != item: + raise ERPProtocolError("货运详情存在冲突的重复 id") + details_by_id[item_id] = item + + records = [ + { + "stock": stock, + "detail": details_by_id.get(stock["id"]), } + for stock in stocks + ] + return { + "query": dict(query), + "count": len(records), + "records": records, + } def close(self) -> None: self.session.close() @@ -339,7 +382,7 @@ class ERPClient: def _stock_payload( self, - order_number: str, + query: Mapping[str, Any], *, start: int, page_index: int, @@ -352,7 +395,7 @@ class ERPClient: "pageIndex": page_index, "store": False, "columns": stock_columns(), - "queries": [order_number_query(order_number)], + "queries": [dict(query)], } def _request_json( @@ -455,6 +498,34 @@ class ERPClient: raise ValueError("order_number 不能包含控制字符") return normalized + @staticmethod + def _validate_created_range( + created_from: date | str, + created_to: date | str, + ) -> tuple[date, date]: + def parse(value: date | str, field: str) -> date: + if isinstance(value, datetime): + raise ValueError(f"{field} 必须是 YYYY-MM-DD 日期") + if isinstance(value, date): + return value + if not isinstance(value, str): + raise ValueError(f"{field} 必须是 YYYY-MM-DD 日期") + try: + parsed = date.fromisoformat(value) + except ValueError as exc: + raise ValueError(f"{field} 必须是 YYYY-MM-DD 日期") from exc + if parsed.isoformat() != value: + raise ValueError(f"{field} 必须是 YYYY-MM-DD 日期") + return parsed + + start = parse(created_from, "created_from") + end = parse(created_to, "created_to") + if end < start: + raise ValueError("created_to 不能早于 created_from") + if end - start > timedelta(days=6): + raise ValueError("创建日期闭区间不能超过 7 天") + return start, end + @staticmethod def _decode_jwt_claims(token: str | None) -> dict[str, Any]: """仅解码 JWT payload 供过期时间展示,不验证其真实性。""" diff --git a/erp-connector/src/shunyunbaoerp/constants.py b/erp-connector/src/shunyunbaoerp/constants.py index 008f741..9a51ec5 100644 --- a/erp-connector/src/shunyunbaoerp/constants.py +++ b/erp-connector/src/shunyunbaoerp/constants.py @@ -115,3 +115,17 @@ def order_number_query(order_number: str) -> dict[str, Any]: "tableAlias": "t", "optType": 1, } + + +def created_range_query(created_from: str, created_to: str) -> dict[str, Any]: + """生成 HAR 中“创建时间”闭区间查询条件。""" + + return { + "dvalue": f"{created_from},{created_to}", + "tableName": "t_stock", + "colName": "created", + "op": 0, + "type": 3, + "tableAlias": "t", + "optType": 0, + } diff --git a/erp-connector/src/shunyunbaoerp/normalizer.py b/erp-connector/src/shunyunbaoerp/normalizer.py index 513e235..4f00161 100644 --- a/erp-connector/src/shunyunbaoerp/normalizer.py +++ b/erp-connector/src/shunyunbaoerp/normalizer.py @@ -9,6 +9,7 @@ from .errors import ERPProtocolError def normalize_freight_result(result: Mapping[str, Any]) -> dict[str, Any]: + normalized_query = _normalize_query(result.get("query")) records = result.get("records") if not isinstance(records, list): raise ERPProtocolError("货运查询结果缺少 records 数组") @@ -84,11 +85,30 @@ def normalize_freight_result(result: Mapping[str, Any]) -> dict[str, Any]: ) return { "schema_version": 1, - "query": {"mode": "ORDER_NUMBER"}, + "query": normalized_query, "orders": orders, } +def _normalize_query(value: Any) -> dict[str, Any]: + if not isinstance(value, Mapping): + raise ERPProtocolError("货运查询结果缺少 query 对象") + mode = value.get("mode", "ORDER_NUMBER") + if mode == "ORDER_NUMBER": + return {"mode": "ORDER_NUMBER"} + if mode != "CREATED_RANGE": + raise ERPProtocolError("货运查询模式无效") + created_from = _text(value.get("createdFrom")) + created_to = _text(value.get("createdTo")) + if len(created_from) != 10 or len(created_to) != 10: + raise ERPProtocolError("货运日期范围无效") + return { + "mode": "CREATED_RANGE", + "created_from": created_from, + "created_to": created_to, + } + + def _normalize_item(item: Mapping[str, Any]) -> dict[str, Any]: title = _text(item.get("productTitle")) if not title: diff --git a/erp-connector/tests/test_api.py b/erp-connector/tests/test_api.py index 32a314b..1f3db25 100644 --- a/erp-connector/tests/test_api.py +++ b/erp-connector/tests/test_api.py @@ -21,6 +21,20 @@ class FakeClient: assert order_number == "SOURCE-12" return raw_result() + def get_freight_details_by_created_range( + self, + created_from: str, + created_to: str, + ) -> dict: + assert (created_from, created_to) == ("2026-07-22", "2026-07-28") + result = raw_result() + result["query"] = { + "mode": "CREATED_RANGE", + "createdFrom": created_from, + "createdTo": created_to, + } + return result + def test_query_requires_configured_service_key(monkeypatch) -> None: monkeypatch.delenv("SHUNYUNBAO_SERVICE_API_KEY", raising=False) @@ -56,3 +70,42 @@ def test_query_returns_only_normalized_fields(monkeypatch) -> None: assert "receiverTel" not in encoded assert "receiverAddr" not in encoded assert "PRIVATE-QUERY" not in encoded + + +def test_created_range_query_returns_allowlisted_range(monkeypatch) -> None: + monkeypatch.setenv("SHUNYUNBAO_SERVICE_API_KEY", VALID_KEY) + monkeypatch.setattr(api, "get_client", lambda: FakeClient()) + + response = api.query_freight( + api.FreightQueryRequest( + mode="CREATED_RANGE", + created_from="2026-07-22", + created_to="2026-07-28", + ), + None, + ) + + assert response["query"] == { + "mode": "CREATED_RANGE", + "created_from": "2026-07-22", + "created_to": "2026-07-28", + } + + +def test_created_range_rejects_mixed_mode_parameters(monkeypatch) -> None: + monkeypatch.setenv("SHUNYUNBAO_SERVICE_API_KEY", VALID_KEY) + monkeypatch.setattr(api, "get_client", lambda: FakeClient()) + + with pytest.raises(HTTPException) as raised: + api.query_freight( + api.FreightQueryRequest( + mode="CREATED_RANGE", + order_number="SOURCE-12", + created_from="2026-07-22", + created_to="2026-07-28", + ), + None, + ) + + assert raised.value.status_code == 422 + assert raised.value.detail == "ERP_QUERY_INVALID" diff --git a/erp-connector/tests/test_client.py b/erp-connector/tests/test_client.py index 175b95c..fa52758 100644 --- a/erp-connector/tests/test_client.py +++ b/erp-connector/tests/test_client.py @@ -2,6 +2,7 @@ from __future__ import annotations import base64 import json +from datetime import date from urllib.parse import urlparse import pytest @@ -153,6 +154,62 @@ def test_query_and_detail_reproduce_har_contract() -> None: assert session.calls[2]["json"] == {"ids": [99001122]} +def test_created_range_query_reproduces_har_contract_and_deduplicates() -> None: + stock = {"id": 99001122, "code": "FREIGHT-001"} + detail = {"id": 99001122, "details": []} + session = FakeSession( + [ + FakeResponse(envelope(2)), + FakeResponse(envelope({"total": 2, "list": [stock, stock]})), + FakeResponse(envelope({"total": 1, "list": [detail]})), + ] + ) + client = ERPClient(session=session) + + result = client.get_freight_details_by_created_range( + date(2026, 7, 22), + date(2026, 7, 28), + ) + + assert result["count"] == 1 + assert result["query"] == { + "mode": "CREATED_RANGE", + "createdFrom": "2026-07-22", + "createdTo": "2026-07-28", + } + assert session.calls[0]["json"]["queries"] == [ + { + "dvalue": "2026-07-22,2026-07-28", + "tableName": "t_stock", + "colName": "created", + "op": 0, + "type": 3, + "tableAlias": "t", + "optType": 0, + } + ] + + +def test_created_range_rejects_more_than_seven_inclusive_days() -> None: + client = ERPClient(session=FakeSession([])) + + with pytest.raises(ValueError, match="7 天"): + client.query_stock_by_created_range("2026-07-21", "2026-07-28") + + +def test_empty_created_range_does_not_request_details() -> None: + session = FakeSession([FakeResponse(envelope(0))]) + client = ERPClient(session=session) + + result = client.get_freight_details_by_created_range( + "2026-07-28", + "2026-07-28", + ) + + assert result["records"] == [] + assert len(session.calls) == 1 + + def test_not_found_stops_before_detail_request() -> None: session = FakeSession([FakeResponse(envelope(0))]) client = ERPClient(session=session) diff --git a/erp-connector/tests/test_normalizer.py b/erp-connector/tests/test_normalizer.py index 8ed3dff..5ac0124 100644 --- a/erp-connector/tests/test_normalizer.py +++ b/erp-connector/tests/test_normalizer.py @@ -110,6 +110,25 @@ def test_invalid_procurement_fields_remain_reviewable() -> None: assert result["product_thumb_ref"] is None +def test_normalizes_created_range_without_exposing_raw_query() -> None: + source = raw_result() + source["query"] = { + "mode": "CREATED_RANGE", + "createdFrom": "2026-07-22", + "createdTo": "2026-07-28", + "private": "must-not-leak", + } + + normalized = normalize_freight_result(source) + + assert normalized["query"] == { + "mode": "CREATED_RANGE", + "created_from": "2026-07-22", + "created_to": "2026-07-28", + } + assert "must-not-leak" not in json.dumps(normalized) + + def test_conflicting_duplicate_item_is_rejected() -> None: source = raw_result() duplicate = dict(source["records"][0]["detail"]["details"][0])