feat(t240): cache freight item images

This commit is contained in:
QiuSW
2026-07-29 15:12:05 +08:00
parent 5f17fc11aa
commit c7e152aaa1
28 changed files with 1728 additions and 30 deletions
@@ -0,0 +1,42 @@
package usecase
import (
"context"
"io"
"cmroubao/backend-api/internal/domain"
)
type FreightSourceImage struct {
Content io.ReadCloser
MediaType string
}
type FreightImageSource interface {
FetchProductImage(
context.Context,
string,
) (FreightSourceImage, error)
}
type FreightImageRepository interface {
ListFreightImageJobs(
context.Context,
string,
string,
int,
) ([]domain.FreightItemImageJob, error)
SaveFreightItemImage(
context.Context,
domain.FreightItemImage,
) (*string, error)
GetReadyFreightItemImage(
context.Context,
string,
string,
) (domain.FreightItemImage, error)
}
type FreightImageCache interface {
CacheRun(context.Context, string, string) error
}
@@ -0,0 +1,270 @@
package usecase
import (
"context"
"errors"
"io"
"sync"
"time"
"cmroubao/backend-api/internal/domain"
)
const (
maxFreightImageJobsPerRun = 1000
freightImageWorkers = 4
freightImageCleanup = 5 * time.Second
)
type FreightImageService struct {
repository FreightImageRepository
source FreightImageSource
store ReferenceImageStore
clock Clock
}
type FreightItemImageContent struct {
Image domain.FreightItemImage
Content io.ReadCloser
}
func NewFreightImageService(
repository FreightImageRepository,
source FreightImageSource,
store ReferenceImageStore,
clock Clock,
) (*FreightImageService, error) {
if repository == nil || source == nil || store == nil || clock == nil {
return nil, errors.New("freight image service dependencies are required")
}
return &FreightImageService{
repository: repository,
source: source,
store: store,
clock: clock,
}, nil
}
func (service *FreightImageService) CacheRun(
ctx context.Context,
creatorSubject, runID string,
) error {
jobs, err := service.repository.ListFreightImageJobs(
ctx,
creatorSubject,
runID,
maxFreightImageJobsPerRun,
)
if err != nil {
return wrapRepositoryError(err)
}
if len(jobs) == 0 {
return nil
}
workerCount := freightImageWorkers
if len(jobs) < workerCount {
workerCount = len(jobs)
}
queue := make(chan domain.FreightItemImageJob)
var workers sync.WaitGroup
workers.Add(workerCount)
for range workerCount {
go func() {
defer workers.Done()
for job := range queue {
if ctx.Err() != nil {
return
}
service.cacheJob(ctx, job)
}
}()
}
for _, job := range jobs {
select {
case queue <- job:
case <-ctx.Done():
close(queue)
workers.Wait()
return nil
}
}
close(queue)
workers.Wait()
return nil
}
func (service *FreightImageService) cacheJob(
ctx context.Context,
job domain.FreightItemImageJob,
) {
sourceImage, err := service.source.FetchProductImage(
ctx,
job.ProductThumbRef,
)
if err != nil {
if ctx.Err() == nil {
service.saveFailure(ctx, job, freightImageSourceFailure(err))
}
return
}
defer sourceImage.Content.Close()
normalized, err := service.store.Put(
ctx,
job.FreightOrderItemID,
sourceImage.MediaType,
sourceImage.Content,
)
if err != nil {
if ctx.Err() == nil {
service.saveFailure(ctx, job, freightImageStoreFailure(err))
}
return
}
candidate := domain.FreightItemImage{
CreatorSubject: job.CreatorSubject,
FreightOrderItemID: job.FreightOrderItemID,
ProductThumbRef: job.ProductThumbRef,
Status: domain.FreightItemImageReady,
MediaType: normalized.MediaType,
SizeBytes: normalized.SizeBytes,
SHA256: normalized.SHA256,
StorageKey: normalized.StorageKey,
UpdatedAt: service.clock.Now().UTC(),
}
replaced, err := service.repository.SaveFreightItemImage(ctx, candidate)
if err != nil {
service.deleteStoredImage(normalized.StorageKey)
return
}
if replaced != nil {
service.deleteStoredImage(*replaced)
}
}
func (service *FreightImageService) saveFailure(
ctx context.Context,
job domain.FreightItemImageJob,
failure freightImageFailure,
) {
errorCode := failure.code
replaced, err := service.repository.SaveFreightItemImage(
ctx,
domain.FreightItemImage{
CreatorSubject: job.CreatorSubject,
FreightOrderItemID: job.FreightOrderItemID,
ProductThumbRef: job.ProductThumbRef,
Status: failure.status,
ErrorCode: &errorCode,
UpdatedAt: service.clock.Now().UTC(),
},
)
if err == nil && replaced != nil {
service.deleteStoredImage(*replaced)
}
}
func (service *FreightImageService) deleteStoredImage(storageKey string) {
ctx, cancel := context.WithTimeout(context.Background(), freightImageCleanup)
defer cancel()
_ = service.store.Delete(ctx, storageKey)
}
func (service *FreightImageService) OpenItemImage(
ctx context.Context,
creatorSubject, itemID string,
) (FreightItemImageContent, error) {
if creatorSubject == "" || !isUUID(itemID) {
return FreightItemImageContent{}, freightImageNotFoundError(nil)
}
image, err := service.repository.GetReadyFreightItemImage(
ctx,
creatorSubject,
itemID,
)
if err != nil {
wrapped := wrapRepositoryError(err)
var typed *Error
if errors.As(wrapped, &typed) &&
typed.Kind == ErrorKindNotFound {
return FreightItemImageContent{},
freightImageNotFoundError(err)
}
return FreightItemImageContent{}, wrapped
}
content, err := service.store.Open(ctx, image.StorageKey)
if err != nil {
var storeError *ImageStoreError
if errors.As(err, &storeError) &&
storeError.Kind == ImageStoreErrorNotFound {
return FreightItemImageContent{},
freightImageNotFoundError(err)
}
return FreightItemImageContent{}, mapImageStoreError(err)
}
return FreightItemImageContent{Image: image, Content: content}, nil
}
type freightImageFailure struct {
status domain.FreightItemImageStatus
code string
}
func freightImageSourceFailure(err error) freightImageFailure {
switch {
case errors.Is(err, domain.ErrFreightImageNotFound):
return freightImageFailure{
status: domain.FreightItemImageMissing,
code: "ERP_IMAGE_NOT_FOUND",
}
case errors.Is(err, domain.ErrFreightImageInvalid):
return freightImageFailure{
status: domain.FreightItemImageFailed,
code: "ERP_IMAGE_INVALID",
}
case errors.Is(err, domain.ErrFreightSourceSessionNeeded):
return freightImageFailure{
status: domain.FreightItemImageFailed,
code: "ERP_SESSION_REQUIRED",
}
case errors.Is(err, domain.ErrFreightSourceNotConfigured):
return freightImageFailure{
status: domain.FreightItemImageFailed,
code: "ERP_NOT_CONFIGURED",
}
default:
return freightImageFailure{
status: domain.FreightItemImageFailed,
code: "ERP_IMAGE_UNAVAILABLE",
}
}
}
func freightImageStoreFailure(err error) freightImageFailure {
var storeError *ImageStoreError
if errors.As(err, &storeError) {
switch storeError.Kind {
case ImageStoreErrorTooLarge,
ImageStoreErrorUnsupported,
ImageStoreErrorInvalid:
return freightImageFailure{
status: domain.FreightItemImageFailed,
code: "ERP_IMAGE_INVALID",
}
}
}
return freightImageFailure{
status: domain.FreightItemImageFailed,
code: "IMAGE_STORE_UNAVAILABLE",
}
}
func freightImageNotFoundError(cause error) error {
return newError(
ErrorKindNotFound,
"FREIGHT_IMAGE_NOT_FOUND",
"freight item image not found",
cause,
)
}
var _ FreightImageCache = (*FreightImageService)(nil)
@@ -0,0 +1,219 @@
package usecase
import (
"bytes"
"context"
"errors"
"io"
"strings"
"sync"
"testing"
"cmroubao/backend-api/internal/domain"
)
func TestFreightImageCacheIsBestEffortAndCleansReplacedFiles(t *testing.T) {
repository := &freightImageRepositoryFake{
jobs: []domain.FreightItemImageJob{
{
CreatorSubject: "local-admin",
FreightOrderItemID: "00000000-0000-4000-8000-000000000001",
ProductThumbRef: "190",
},
{
CreatorSubject: "local-admin",
FreightOrderItemID: "00000000-0000-4000-8000-000000000002",
ProductThumbRef: "191",
},
},
replaced: "old/image.jpg",
}
source := &freightImageSourceFake{}
store := &freightImageStoreFake{}
service, err := NewFreightImageService(
repository,
source,
store,
fakeClock{},
)
if err != nil {
t.Fatalf("NewFreightImageService() error = %v", err)
}
if err := service.CacheRun(
context.Background(),
"local-admin",
"00000000-0000-4000-8000-000000000099",
); err != nil {
t.Fatalf("CacheRun() error = %v", err)
}
repository.mu.Lock()
defer repository.mu.Unlock()
if len(repository.saved) != 2 {
t.Fatalf("saved images = %+v", repository.saved)
}
statusByRef := map[string]domain.FreightItemImage{}
for _, image := range repository.saved {
statusByRef[image.ProductThumbRef] = image
}
if statusByRef["190"].Status != domain.FreightItemImageReady ||
statusByRef["190"].StorageKey != "new/image.jpg" ||
statusByRef["191"].Status != domain.FreightItemImageMissing ||
statusByRef["191"].ErrorCode == nil ||
*statusByRef["191"].ErrorCode != "ERP_IMAGE_NOT_FOUND" {
t.Fatalf("saved image statuses = %+v", statusByRef)
}
store.mu.Lock()
defer store.mu.Unlock()
if len(store.deleted) != 1 || store.deleted[0] != "old/image.jpg" {
t.Fatalf("deleted storage keys = %#v", store.deleted)
}
}
func TestFreightImageOpenDoesNotRevealUnavailableItems(t *testing.T) {
repository := &freightImageRepositoryFake{
ready: domain.FreightItemImage{
CreatorSubject: "local-admin",
FreightOrderItemID: "00000000-0000-4000-8000-000000000001",
Status: domain.FreightItemImageReady,
StorageKey: "ready/image.jpg",
},
}
store := &freightImageStoreFake{
openContent: []byte("normalized-jpeg"),
}
service, err := NewFreightImageService(
repository,
&freightImageSourceFake{},
store,
fakeClock{},
)
if err != nil {
t.Fatalf("NewFreightImageService() error = %v", err)
}
result, err := service.OpenItemImage(
context.Background(),
"local-admin",
repository.ready.FreightOrderItemID,
)
if err != nil {
t.Fatalf("OpenItemImage() error = %v", err)
}
defer result.Content.Close()
content, _ := io.ReadAll(result.Content)
if string(content) != "normalized-jpeg" {
t.Fatalf("image content = %q", content)
}
if _, err := service.OpenItemImage(
context.Background(),
"another-subject",
repository.ready.FreightOrderItemID,
); err == nil {
t.Fatal("cross-subject OpenItemImage() error = nil")
} else {
var typed *Error
if !errors.As(err, &typed) ||
typed.Code != "FREIGHT_IMAGE_NOT_FOUND" {
t.Fatalf("cross-subject error = %v", err)
}
}
}
type freightImageRepositoryFake struct {
mu sync.Mutex
jobs []domain.FreightItemImageJob
saved []domain.FreightItemImage
replaced string
ready domain.FreightItemImage
}
func (repository *freightImageRepositoryFake) ListFreightImageJobs(
context.Context,
string,
string,
int,
) ([]domain.FreightItemImageJob, error) {
return append([]domain.FreightItemImageJob(nil), repository.jobs...), nil
}
func (repository *freightImageRepositoryFake) SaveFreightItemImage(
_ context.Context,
image domain.FreightItemImage,
) (*string, error) {
repository.mu.Lock()
defer repository.mu.Unlock()
repository.saved = append(repository.saved, image)
if image.Status == domain.FreightItemImageReady &&
repository.replaced != "" {
value := repository.replaced
repository.replaced = ""
return &value, nil
}
return nil, nil
}
func (repository *freightImageRepositoryFake) GetReadyFreightItemImage(
_ context.Context,
creatorSubject, itemID string,
) (domain.FreightItemImage, error) {
if repository.ready.CreatorSubject != creatorSubject ||
repository.ready.FreightOrderItemID != itemID {
return domain.FreightItemImage{}, ErrRepositoryNotFound
}
return repository.ready, nil
}
type freightImageSourceFake struct{}
func (*freightImageSourceFake) FetchProductImage(
_ context.Context,
productThumbRef string,
) (FreightSourceImage, error) {
if productThumbRef == "191" {
return FreightSourceImage{}, domain.ErrFreightImageNotFound
}
return FreightSourceImage{
Content: io.NopCloser(bytes.NewReader([]byte("source-image"))),
MediaType: "image/png",
}, nil
}
type freightImageStoreFake struct {
mu sync.Mutex
deleted []string
openContent []byte
}
func (*freightImageStoreFake) Put(
context.Context,
string,
string,
io.Reader,
) (NormalizedReferenceImage, error) {
return NormalizedReferenceImage{
StorageKey: "new/image.jpg",
MediaType: "image/jpeg",
SizeBytes: 100,
SHA256: repeatSHA("a"),
}, nil
}
func (store *freightImageStoreFake) Open(
context.Context,
string,
) (io.ReadCloser, error) {
return io.NopCloser(bytes.NewReader(store.openContent)), nil
}
func (store *freightImageStoreFake) Delete(
_ context.Context,
storageKey string,
) error {
store.mu.Lock()
defer store.mu.Unlock()
store.deleted = append(store.deleted, storageKey)
return nil
}
func repeatSHA(value string) string {
return strings.Repeat(value, 64)
}
@@ -22,6 +22,7 @@ const (
freightWatermarkOverlap = 10 * time.Minute
orderFreightSyncTimeout = 55 * time.Second
freightCleanupTimeout = 5 * time.Second
freightImageCacheBudget = 15 * time.Second
)
type FreightService struct {
@@ -33,6 +34,19 @@ type FreightService struct {
orderTimeout time.Duration
cleanupTimeout time.Duration
orderSyncGate chan struct{}
imageCache FreightImageCache
}
type FreightServiceOption func(*FreightService) error
func WithFreightImageCache(cache FreightImageCache) FreightServiceOption {
return func(service *FreightService) error {
if cache == nil {
return errors.New("freight image cache is required")
}
service.imageCache = cache
return nil
}
}
type FreightSourcePreflight interface {
@@ -66,12 +80,13 @@ func NewFreightService(
clock Clock,
ids IDGenerator,
timeout time.Duration,
options ...FreightServiceOption,
) (*FreightService, error) {
if repository == nil || source == nil || clock == nil || ids == nil ||
timeout <= 0 {
return nil, errors.New("freight service dependencies are required")
}
return &FreightService{
service := &FreightService{
repository: repository,
source: source,
clock: clock,
@@ -80,7 +95,16 @@ func NewFreightService(
orderTimeout: orderFreightSyncTimeout,
cleanupTimeout: freightCleanupTimeout,
orderSyncGate: make(chan struct{}, 1),
}, nil
}
for _, option := range options {
if option == nil {
return nil, errors.New("freight service option is required")
}
if err := option(service); err != nil {
return nil, err
}
}
return service, nil
}
func (service *FreightService) CreateOrderSync(
@@ -424,6 +448,18 @@ func (service *FreightService) executeRun(
run.OrderCount = len(batch.Orders)
run.ItemCount = itemCount
run.FinishedAt = &finishedAt
if service.imageCache != nil {
cacheCtx, cancel := context.WithTimeout(
ctx,
freightImageCacheBudget,
)
_ = service.imageCache.CacheRun(
cacheCtx,
run.CreatorSubject,
run.ID,
)
cancel()
}
return run, nil
}
@@ -184,12 +184,16 @@ func TestCreateFreightOrderSyncStopsBeforePersistingWhenOCRIsInvalid(t *testing.
func TestCreateFreightOrderSyncReturnsCommittedResult(t *testing.T) {
repository := &syncTrackingRepository{}
imageCache := &recordingFreightImageCache{
err: errors.New("best effort image failure"),
}
service, err := NewFreightService(
repository,
&recordingDateSource{},
fakeClock{},
&sequenceIDs{},
time.Minute,
WithFreightImageCache(imageCache),
)
if err != nil {
t.Fatalf("NewFreightService() error = %v", err)
@@ -209,6 +213,11 @@ func TestCreateFreightOrderSyncReturnsCommittedResult(t *testing.T) {
result.Run.StartedAt == nil || result.Run.FinishedAt == nil {
t.Fatalf("CreateOrderSync() result = %+v", result)
}
if imageCache.calls != 1 ||
imageCache.creatorSubject != "local-admin" ||
imageCache.runID != result.Run.ID {
t.Fatalf("image cache calls = %+v", imageCache)
}
repository.mu.Lock()
defer repository.mu.Unlock()
if repository.startCalls != 1 || repository.completeCalls != 1 ||
@@ -222,6 +231,23 @@ func TestCreateFreightOrderSyncReturnsCommittedResult(t *testing.T) {
}
}
type recordingFreightImageCache struct {
calls int
creatorSubject string
runID string
err error
}
func (cache *recordingFreightImageCache) CacheRun(
_ context.Context,
creatorSubject, runID string,
) error {
cache.calls++
cache.creatorSubject = creatorSubject
cache.runID = runID
return cache.err
}
func TestCreateFreightOrderSyncTimeoutUsesLiveCleanupContext(t *testing.T) {
repository := &syncTrackingRepository{}
source := &blockingOrderSource{}