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

276 lines
7.1 KiB
Go

package sqlite
import (
"context"
"database/sql"
"errors"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/usecase"
)
func (store *Store) ListFreightImageJobs(
ctx context.Context,
creatorSubject, runID string,
limit int,
) ([]domain.FreightItemImageJob, error) {
rows, err := store.db.QueryContext(
ctx,
`SELECT item.id, item.product_thumb_ref
FROM freight_order_items AS item
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
LEFT JOIN freight_item_images AS image
ON image.freight_order_item_id = item.id
WHERE freight.creator_subject = ?
AND item.last_sync_run_id = ?
AND item.is_present = 1
AND item.product_thumb_ref IS NOT NULL
AND (
image.freight_order_item_id IS NULL
OR image.product_thumb_ref != item.product_thumb_ref
OR image.status != 'READY'
)
ORDER BY item.id
LIMIT ?`,
creatorSubject,
runID,
limit,
)
if err != nil {
return nil, repositoryFailure(err)
}
defer rows.Close()
jobs := make([]domain.FreightItemImageJob, 0)
for rows.Next() {
var job domain.FreightItemImageJob
job.CreatorSubject = creatorSubject
if err := rows.Scan(
&job.FreightOrderItemID,
&job.ProductThumbRef,
); err != nil {
return nil, repositoryFailure(err)
}
jobs = append(jobs, job)
}
if err := rows.Err(); err != nil {
return nil, repositoryFailure(err)
}
return jobs, nil
}
func (store *Store) SaveFreightItemImage(
ctx context.Context,
candidate domain.FreightItemImage,
) (*string, error) {
tx, err := store.db.BeginTx(ctx, nil)
if err != nil {
return nil, repositoryFailure(err)
}
defer tx.Rollback()
var currentThumb sql.NullString
err = tx.QueryRowContext(
ctx,
`SELECT item.product_thumb_ref
FROM freight_order_items AS item
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
WHERE freight.creator_subject = ?
AND item.id = ?
AND item.is_present = 1`,
candidate.CreatorSubject,
candidate.FreightOrderItemID,
).Scan(&currentThumb)
if errors.Is(err, sql.ErrNoRows) {
return nil, usecase.ErrRepositoryNotFound
}
if err != nil {
return nil, repositoryFailure(err)
}
if !currentThumb.Valid ||
currentThumb.String != candidate.ProductThumbRef {
return nil, usecase.ErrTaskStateConflict
}
var existingThumb string
var existingStorage sql.NullString
var existingAttempts int
err = tx.QueryRowContext(
ctx,
`SELECT product_thumb_ref, storage_key, attempt_count
FROM freight_item_images
WHERE freight_order_item_id = ?`,
candidate.FreightOrderItemID,
).Scan(&existingThumb, &existingStorage, &existingAttempts)
switch {
case errors.Is(err, sql.ErrNoRows):
existingAttempts = 0
case err != nil:
return nil, repositoryFailure(err)
case existingThumb != candidate.ProductThumbRef:
existingAttempts = 0
}
attemptCount := existingAttempts + 1
_, err = tx.ExecContext(
ctx,
`INSERT INTO freight_item_images (
freight_order_item_id, product_thumb_ref, status, media_type,
size_bytes, sha256, storage_key, error_code, attempt_count,
updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT (freight_order_item_id)
DO UPDATE SET
product_thumb_ref = excluded.product_thumb_ref,
status = excluded.status,
media_type = excluded.media_type,
size_bytes = excluded.size_bytes,
sha256 = excluded.sha256,
storage_key = excluded.storage_key,
error_code = excluded.error_code,
attempt_count = excluded.attempt_count,
updated_at = excluded.updated_at`,
candidate.FreightOrderItemID,
candidate.ProductThumbRef,
candidate.Status,
nullableFreightImageText(candidate.MediaType),
nullableFreightImageSize(candidate),
nullableFreightImageText(candidate.SHA256),
nullableFreightImageText(candidate.StorageKey),
nullableString(candidate.ErrorCode),
attemptCount,
formatTimestamp(candidate.UpdatedAt),
)
if err != nil {
return nil, repositoryFailure(err)
}
if err := tx.Commit(); err != nil {
return nil, repositoryFailure(err)
}
if existingStorage.Valid &&
existingStorage.String != candidate.StorageKey {
return &existingStorage.String, nil
}
return nil, nil
}
func (store *Store) GetReadyFreightItemImage(
ctx context.Context,
creatorSubject, itemID string,
) (domain.FreightItemImage, error) {
var image domain.FreightItemImage
var updatedAt string
err := store.db.QueryRowContext(
ctx,
`SELECT freight.creator_subject, image.freight_order_item_id,
image.product_thumb_ref, image.status, image.media_type,
image.size_bytes, image.sha256, image.storage_key,
image.attempt_count, image.updated_at
FROM freight_item_images AS image
JOIN freight_order_items AS item
ON item.id = image.freight_order_item_id
JOIN freight_orders AS freight
ON freight.id = item.freight_order_id
WHERE freight.creator_subject = ?
AND item.id = ?
AND item.is_present = 1
AND item.product_thumb_ref = image.product_thumb_ref
AND image.status = 'READY'`,
creatorSubject,
itemID,
).Scan(
&image.CreatorSubject,
&image.FreightOrderItemID,
&image.ProductThumbRef,
&image.Status,
&image.MediaType,
&image.SizeBytes,
&image.SHA256,
&image.StorageKey,
&image.AttemptCount,
&updatedAt,
)
if errors.Is(err, sql.ErrNoRows) {
return domain.FreightItemImage{}, usecase.ErrRepositoryNotFound
}
if err != nil {
return domain.FreightItemImage{}, repositoryFailure(err)
}
image.UpdatedAt, err = parseTimestamp(updatedAt)
if err != nil {
return domain.FreightItemImage{}, repositoryFailure(err)
}
return image, nil
}
func (store *Store) applyFreightItemImageStatuses(
ctx context.Context,
orderID string,
items []domain.FreightOrderItem,
) error {
byID := make(map[string]int, len(items))
for index := range items {
byID[items[index].ID] = index
if items[index].ProductThumbRef == nil {
items[index].ImageStatus = domain.FreightItemImageNone
} else {
items[index].ImageStatus = domain.FreightItemImagePending
}
}
rows, err := store.db.QueryContext(
ctx,
`SELECT item.id, image.product_thumb_ref, image.status,
image.error_code
FROM freight_order_items AS item
JOIN freight_item_images AS image
ON image.freight_order_item_id = item.id
WHERE item.freight_order_id = ?
AND item.is_present = 1`,
orderID,
)
if err != nil {
return repositoryFailure(err)
}
defer rows.Close()
for rows.Next() {
var itemID, imageThumb string
var status domain.FreightItemImageStatus
var errorCode sql.NullString
if err := rows.Scan(
&itemID,
&imageThumb,
&status,
&errorCode,
); err != nil {
return repositoryFailure(err)
}
index, exists := byID[itemID]
if !exists || items[index].ProductThumbRef == nil ||
*items[index].ProductThumbRef != imageThumb {
continue
}
items[index].ImageStatus = status
items[index].ImageErrorCode = optionalString(errorCode)
}
if err := rows.Err(); err != nil {
return repositoryFailure(err)
}
return nil
}
func nullableFreightImageText(value string) any {
if value == "" {
return nil
}
return value
}
func nullableFreightImageSize(image domain.FreightItemImage) any {
if image.Status != domain.FreightItemImageReady {
return nil
}
return image.SizeBytes
}
var _ usecase.FreightImageRepository = (*Store)(nil)