Files
cmroubao/backend-api/internal/repository/sqlite/freight_repository_test.go
T

403 lines
11 KiB
Go

package sqlite_test
import (
"context"
"database/sql"
"errors"
"testing"
"time"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/platform/migration"
repository "cmroubao/backend-api/internal/repository/sqlite"
"cmroubao/backend-api/internal/usecase"
)
func TestFreightImportIsAtomicIdempotentAndRevisioned(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, 28, 1, 2, 3, 0, time.UTC)
userID := uuid(900)
seedFreightUser(t, db, userID, now)
first := freightRun(901, userID, now)
stored, created, err := store.CreateFreightSync(
ctx,
first,
"freight-key-1",
repeatHex("1"),
)
if err != nil || !created || stored.ID != first.ID {
t.Fatalf("CreateFreightSync() = %+v, %t, %v", stored, created, err)
}
replayed, created, err := store.CreateFreightSync(
ctx,
freightRun(902, userID, now),
"freight-key-1",
repeatHex("1"),
)
if err != nil || created || replayed.ID != first.ID {
t.Fatalf("sync replay = %+v, %t, %v", replayed, created, err)
}
_, _, err = store.CreateFreightSync(
ctx,
freightRun(903, userID, now),
"freight-key-1",
repeatHex("2"),
)
if !errors.Is(err, usecase.ErrIdempotencyConflict) {
t.Fatalf("conflicting replay error = %v", err)
}
if err := store.StartFreightSync(ctx, first.ID, now.Add(time.Second)); err != nil {
t.Fatalf("StartFreightSync() error = %v", err)
}
batch := freightBatch(910, "a", "b")
price := int64(12950)
batch.Orders[0].Items[0].OriginalUnitPriceMinor = &price
if err := store.CompleteFreightSync(
ctx,
first,
batch,
now.Add(2*time.Second),
); err != nil {
t.Fatalf("CompleteFreightSync() error = %v", err)
}
orders, err := store.ListFreightOrders(ctx, "local-admin", 10)
if err != nil || len(orders) != 1 || orders[0].ItemCount != 2 ||
orders[0].Revision != 1 {
t.Fatalf("orders = %+v, error = %v", orders, err)
}
detail, err := store.GetFreightOrder(ctx, "local-admin", orders[0].ID)
if err != nil || len(detail.Items) != 2 {
t.Fatalf("detail = %+v, error = %v", detail, err)
}
if detail.Items[0].OriginalUnitPriceMinor == nil ||
*detail.Items[0].OriginalUnitPriceMinor != price ||
detail.Items[0].OriginalCurrency != domain.FreightCurrencyTWD {
t.Fatalf("item price metadata = %+v", detail.Items[0])
}
second := freightRun(904, userID, now.Add(time.Minute))
createAndStartFreightRun(t, store, second, "freight-key-2", "3")
unchanged := freightBatch(920, "a", "b")
unchanged.Orders[0].Items[0].OriginalUnitPriceMinor = &price
if err := store.CompleteFreightSync(
ctx,
second,
unchanged,
now.Add(time.Minute+time.Second),
); err != nil {
t.Fatalf("unchanged import error = %v", err)
}
detail, _ = store.GetFreightOrder(ctx, "local-admin", orders[0].ID)
if detail.Order.Revision != 1 ||
detail.Items[0].Revision != 1 ||
detail.Items[1].Revision != 1 {
t.Fatalf("unchanged revisions = %+v", detail)
}
third := freightRun(905, userID, now.Add(2*time.Minute))
createAndStartFreightRun(t, store, third, "freight-key-3", "4")
changed := freightBatch(930, "a", "c")
changed.Orders[0].Items[0].OriginalUnitPriceMinor = &price
changed.Orders[0].CanonicalSHA256 = repeatHex("d")
changed.Orders[0].Items[1].Title = "变化后的商品"
changed.Orders[0].Items[1].CanonicalSHA256 = repeatHex("e")
if err := store.CompleteFreightSync(
ctx,
third,
changed,
now.Add(2*time.Minute+time.Second),
); err != nil {
t.Fatalf("changed import error = %v", err)
}
detail, _ = store.GetFreightOrder(ctx, "local-admin", orders[0].ID)
if detail.Order.Revision != 2 ||
detail.Items[0].Revision != 1 ||
detail.Items[1].Revision != 2 {
t.Fatalf("changed revisions = %+v", detail)
}
if _, err := db.Exec(
`UPDATE freight_order_items SET original_unit_price_minor = NULL`,
); err != nil {
t.Fatalf("clear metadata before rollback guard test: %v", err)
}
runner, err := migration.New(db)
if err != nil {
t.Fatalf("migration.New() error = %v", err)
}
if err := runner.Down(ctx); err != nil {
t.Fatalf("metadata migration down: %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)
}
if err := runner.Down(ctx); err == nil {
t.Fatal("freight migration down succeeded with retained source data")
}
var retained int
if err := db.QueryRow(`SELECT count(*) FROM freight_orders`).Scan(
&retained,
); err != nil || retained != 1 {
t.Fatalf("retained freight orders = %d, error = %v", retained, err)
}
}
func TestFreightSyncRecoveryAndFailedBatchDoNotPersistOrders(t *testing.T) {
db := openDatabase(t)
store, _ := repository.New(db)
ctx := context.Background()
now := time.Date(2026, 7, 28, 2, 3, 4, 0, time.UTC)
userID := uuid(940)
seedFreightUser(t, db, userID, now)
run := freightRun(941, userID, now)
createAndStartFreightRun(t, store, run, "freight-recovery", "5")
count, err := store.RecoverFreightSyncs(ctx, now.Add(time.Minute))
if err != nil || count != 1 {
t.Fatalf("RecoverFreightSyncs() = %d, %v", count, err)
}
stored, err := store.GetFreightSync(ctx, "local-admin", run.ID)
if err != nil || stored.Status != domain.FreightSyncFailed ||
stored.ErrorCode == nil || *stored.ErrorCode != "PROCESS_INTERRUPTED" {
t.Fatalf("recovered run = %+v, error = %v", stored, err)
}
if err := store.CompleteFreightSync(
ctx,
run,
freightBatch(950, "a", "b"),
now.Add(2*time.Minute),
); !errors.Is(err, usecase.ErrTaskStateConflict) {
t.Fatalf("late completion error = %v", err)
}
orders, err := store.ListFreightOrders(ctx, "local-admin", 10)
if err != nil || len(orders) != 0 {
t.Fatalf("orders after failed batch = %+v, error = %v", orders, err)
}
}
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_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.Fatalf("metadata migration down: %v", err)
}
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,
userID string,
now time.Time,
) {
t.Helper()
_, err := db.Exec(
`INSERT INTO users (
id, username, password_hash, role, is_active, created_at, updated_at
) VALUES (?, 'freight-admin', 'hash', 'ADMIN', 1, ?, ?)`,
userID,
now.Format(time.RFC3339Nano),
now.Format(time.RFC3339Nano),
)
if err != nil {
t.Fatalf("seed freight user: %v", err)
}
}
func freightRun(
index int,
userID string,
now time.Time,
) domain.FreightSyncRun {
return domain.FreightSyncRun{
ID: uuid(index),
CreatorSubject: "local-admin",
CreatedByUserID: userID,
Mode: domain.FreightSyncOrderNumber,
OrderNumber: "SOURCE-12",
QuerySHA256: repeatHex("a"),
Status: domain.FreightSyncPending,
CreatedAt: now,
}
}
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
return domain.FreightImportBatch{Orders: []domain.FreightImportOrder{{
ID: uuid(index),
ExternalStockID: "12",
SourceCode: "SOURCE-12",
CanonicalSHA256: repeatHex("c"),
Items: []domain.FreightImportItem{
{
ID: uuid(index + 1),
ExternalItemID: "88",
Title: "商品一",
ProductSpec: "黑色,L",
SKU: "黑色,L",
Quantity: &quantityOne,
OriginalCurrency: domain.FreightCurrencyTWD,
CanonicalSHA256: repeatHex(firstHash),
},
{
ID: uuid(index + 2),
ExternalItemID: "89",
Title: "商品二",
ProductSpec: "白色,M",
SKU: "白色,M",
Quantity: &quantityTwo,
OriginalCurrency: domain.FreightCurrencyTWD,
CanonicalSHA256: repeatHex(secondHash),
},
},
}}}
}
func createAndStartFreightRun(
t *testing.T,
store *repository.Store,
run domain.FreightSyncRun,
key, hash string,
) {
t.Helper()
if _, _, err := store.CreateFreightSync(
context.Background(),
run,
key,
repeatHex(hash),
); err != nil {
t.Fatalf("CreateFreightSync() error = %v", err)
}
if err := store.StartFreightSync(
context.Background(),
run.ID,
run.CreatedAt.Add(time.Second),
); err != nil {
t.Fatalf("StartFreightSync() error = %v", err)
}
}