Files

866 lines
24 KiB
Go
Raw Permalink Normal View History

package feeds
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"io"
"strings"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
)
const upsertChunkSize = 100
type syncStats struct {
Total int
Synced int
Skipped int
Unchanged int
Progress int
ContentHash string
UnchangedFeed bool
Deltas syncDeltaCounts
}
type pendingProduct struct {
GTIN string
RawData map[string]string
MappedData map[string]any
ContentHash string
}
// Sync downloads the feed URL, parses CSV/XML, applies mappings, and upserts raw_products
// in chunks with progress updates on feed_sync_jobs. Replaces SyncStub.
func (s *Service) Sync(ctx context.Context, companyID, feedID uuid.UUID) (map[string]any, error) {
feed, err := s.Get(ctx, companyID, feedID)
if err != nil {
if IsNotFound(err) {
return nil, ErrNotFound
}
return nil, err
}
if err := s.ensureMappingsReadyForSync(ctx, companyID, feedID); err != nil {
return nil, err
}
jobID, err := s.createSyncJob(ctx, companyID, feedID)
if err != nil {
return nil, err
}
if err := s.markJobRunning(ctx, jobID); err != nil {
return nil, err
}
stats, runErr := s.runSync(ctx, companyID, feedID, jobID, feed)
if runErr != nil {
_ = s.failJob(ctx, jobID, runErr.Error(), stats)
return nil, runErr
}
if err := s.completeJob(ctx, jobID, stats); err != nil {
return nil, err
}
_ = s.persistFeedSyncMeta(ctx, companyID, feedID, jobID, stats)
job, err := s.GetSyncJob(ctx, companyID, feedID, jobID)
if err != nil {
return nil, err
}
return enrichJobWithDeltas(job, stats.Deltas), nil
}
// EnqueueSync creates a pending feed_sync_jobs row and wakes the worker via NOTIFY.
// Same-feed pending jobs are reused (no duplicate pending stack). The API process
// does not run sync work (no unbound goroutines). Returns the job id immediately
// for 202/poll (dashboard) or legacy 200 + jobId (v1).
func (s *Service) EnqueueSync(ctx context.Context, companyID, feedID uuid.UUID) (uuid.UUID, error) {
if _, err := s.Get(ctx, companyID, feedID); err != nil {
if IsNotFound(err) {
return uuid.Nil, ErrNotFound
}
return uuid.Nil, err
}
if err := s.ensureMappingsReadyForSync(ctx, companyID, feedID); err != nil {
return uuid.Nil, err
}
jobID, err := s.findOrCreatePendingSyncJob(ctx, companyID, feedID)
if err != nil {
return uuid.Nil, err
}
_, _ = s.Pool.Exec(ctx, `SELECT pg_notify('feed_sync_jobs', $1)`, jobID.String())
return jobID, nil
}
// ClaimNextPendingSyncJob claims one pending feed_sync_jobs row (FOR UPDATE SKIP LOCKED)
// and marks it running for the worker.
func (s *Service) ClaimNextPendingSyncJob(ctx context.Context) (jobID, companyID, feedID uuid.UUID, err error) {
err = s.Pool.QueryRow(ctx, `
WITH candidate AS (
SELECT id FROM feed_sync_jobs
WHERE status = 'pending'
ORDER BY created_at ASC, id ASC
LIMIT 1
FOR UPDATE SKIP LOCKED
)
UPDATE feed_sync_jobs j
SET status = 'running', started_at = now(), updated_at = now()
FROM candidate
WHERE j.id = candidate.id
RETURNING j.id, j.company_id, j.feed_id`).Scan(&jobID, &companyID, &feedID)
if errors.Is(err, pgx.ErrNoRows) {
return uuid.Nil, uuid.Nil, uuid.Nil, pgx.ErrNoRows
}
return jobID, companyID, feedID, err
}
// ProcessSyncJob runs sync for a job already claimed (status=running) by ClaimNextPendingSyncJob.
func (s *Service) ProcessSyncJob(ctx context.Context, companyID, feedID, jobID uuid.UUID) error {
feed, err := s.Get(ctx, companyID, feedID)
if err != nil {
msg := err.Error()
if IsNotFound(err) {
msg = "feed not found"
}
_ = s.failJob(ctx, jobID, msg, syncStats{})
return err
}
stats, runErr := s.runSync(ctx, companyID, feedID, jobID, feed)
if runErr != nil {
_ = s.failJob(ctx, jobID, runErr.Error(), stats)
return runErr
}
if err := s.completeJob(ctx, jobID, stats); err != nil {
_ = s.failJob(ctx, jobID, err.Error(), stats)
return err
}
_ = s.persistFeedSyncMeta(ctx, companyID, feedID, jobID, stats)
return nil
}
// SyncStub is retained as a compatibility alias for Sync.
func (s *Service) SyncStub(ctx context.Context, companyID, feedID uuid.UUID) (map[string]any, error) {
return s.Sync(ctx, companyID, feedID)
}
func (s *Service) GetSyncJob(ctx context.Context, companyID, feedID, jobID uuid.UUID) (map[string]any, error) {
row := s.Pool.QueryRow(ctx, `
SELECT id, feed_id, company_id, status, started_at, completed_at,
products_synced, products_total, products_skipped, products_unchanged,
progress, content_hash, error, created_at, updated_at
FROM feed_sync_jobs
WHERE id = $1 AND company_id = $2 AND feed_id = $3`, jobID, companyID, feedID)
job, err := scanMap(row, []string{
"id", "feed_id", "company_id", "status", "started_at", "completed_at",
"products_synced", "products_total", "products_skipped", "products_unchanged",
"progress", "content_hash", "error", "created_at", "updated_at",
})
if err != nil {
return nil, err
}
if deltas, ok := s.loadFeedLastSyncDeltas(ctx, companyID, feedID); ok {
if matchJobID(deltas["job_id"], jobID) {
return enrichJobWithDeltasMap(job, deltas), nil
}
}
return job, nil
}
func (s *Service) persistFeedSyncMeta(ctx context.Context, companyID, feedID, jobID uuid.UUID, stats syncStats) error {
stats.Deltas.JobID = jobID.String()
stats.Deltas.Unchanged = stats.Unchanged
stats.Deltas.Skipped = stats.Skipped
deltaJSON, err := json.Marshal(stats.Deltas.asMap())
if err != nil {
return err
}
_, err = s.Pool.Exec(ctx, `
UPDATE input_feeds SET last_synced_at = now(), updated_at = now(),
options = COALESCE(options, '{}'::jsonb) || jsonb_build_object(
'last_content_hash', to_jsonb($2::text),
'last_sync_deltas', $3::jsonb
)
WHERE id = $1 AND company_id = $4`, feedID, stats.ContentHash, deltaJSON, companyID)
return err
}
func (s *Service) loadFeedLastSyncDeltas(ctx context.Context, companyID, feedID uuid.UUID) (map[string]any, bool) {
var raw []byte
err := s.Pool.QueryRow(ctx, `
SELECT options->'last_sync_deltas' FROM input_feeds
WHERE id = $1 AND company_id = $2`, feedID, companyID).Scan(&raw)
if err != nil || len(raw) == 0 || string(raw) == "null" {
return nil, false
}
var m map[string]any
if err := json.Unmarshal(raw, &m); err != nil || m == nil {
return nil, false
}
return m, true
}
func enrichJobWithDeltas(job map[string]any, d syncDeltaCounts) map[string]any {
return enrichJobWithDeltasMap(job, d.asMap())
}
func enrichJobWithDeltasMap(job map[string]any, deltas map[string]any) map[string]any {
if job == nil {
return nil
}
out := make(map[string]any, len(job)+1)
for k, v := range job {
out[k] = v
}
out["deltas"] = deltas
out["price_changed"] = deltas["price_changed"]
out["stock_changed"] = deltas["stock_changed"]
out["availability_changed"] = deltas["availability_changed"]
out["title_changed"] = deltas["title_changed"]
out["other_changed"] = deltas["other_changed"]
out["products_new"] = deltas["new"]
return out
}
func matchJobID(v any, id uuid.UUID) bool {
switch t := v.(type) {
case string:
return strings.EqualFold(strings.TrimSpace(t), id.String())
case uuid.UUID:
return t == id
default:
return false
}
}
func (s *Service) ListSyncJobs(ctx context.Context, companyID, feedID uuid.UUID, limit int) ([]map[string]any, error) {
if limit <= 0 {
limit = 50
}
if limit > 200 {
limit = 200
}
rows, err := s.Pool.Query(ctx, `
SELECT id, feed_id, status, started_at, completed_at,
products_synced, products_total, products_skipped, products_unchanged,
progress, content_hash, error, created_at
FROM feed_sync_jobs
WHERE company_id = $1 AND feed_id = $2
ORDER BY created_at DESC LIMIT $3`, companyID, feedID, limit)
if err != nil {
return nil, err
}
defer rows.Close()
return scanMaps(rows, []string{
"id", "feed_id", "status", "started_at", "completed_at",
"products_synced", "products_total", "products_skipped", "products_unchanged",
"progress", "content_hash", "error", "created_at",
})
}
func (s *Service) runSync(ctx context.Context, companyID, feedID, jobID uuid.UUID, feed map[string]any) (syncStats, error) {
var stats syncStats
urlStr, _ := feed["url"].(string)
feedType, _ := feed["feed_type"].(string)
src, err := s.loadFeedSource(ctx, companyID, feed)
if err != nil {
return stats, err
}
defer src.Close()
hash, err := sha256HexFile(src)
if err != nil {
return stats, err
}
stats.ContentHash = hash
prev, _ := s.lastContentHash(ctx, feedID)
if prev != "" && prev == hash {
stats.UnchangedFeed = true
stats.Progress = 100
_ = s.updateJobProgress(ctx, jobID, stats)
return stats, nil
}
mappingsRaw, mapErr := s.loadMappingsRaw(ctx, companyID, feedID)
if mapErr != nil && !errors.Is(mapErr, pgx.ErrNoRows) {
return stats, mapErr
}
mappings := activeMappings(parseMappings(mappingsRaw))
if len(mappings) == 0 {
return stats, ClientMsg("no field mappings defined for feed")
}
itemPath := "item"
if path := itemPathFromMappings(mappingsRaw); path != "" {
itemPath = itemLocalFromPath(path)
} else if opts, ok := feed["options"].(map[string]any); ok {
if v, ok := opts["item_path"].(string); ok && strings.TrimSpace(v) != "" {
itemPath = itemLocalFromPath(v)
}
}
if itemPath == "" {
itemPath = "item"
}
sample, sniffErr := src.Sniff(4096)
if sniffErr != nil {
return stats, sniffErr
}
format := detectFeedFormat(feedType, contentTypeFromBlob(src), urlStr, sample)
chunk := make([]pendingProduct, 0, upsertChunkSize)
flush := func(force bool) error {
if len(chunk) == 0 {
return nil
}
if !force && len(chunk) < upsertChunkSize {
return nil
}
synced, unchanged, skipped, deltas, err := s.upsertChunk(ctx, companyID, feedID, jobID, chunk)
if err != nil {
return err
}
stats.Synced += synced
stats.Unchanged += unchanged
stats.Skipped += skipped
stats.Deltas.New += deltas.New
stats.Deltas.PriceChanged += deltas.PriceChanged
stats.Deltas.StockChanged += deltas.StockChanged
stats.Deltas.AvailabilityChanged += deltas.AvailabilityChanged
stats.Deltas.TitleChanged += deltas.TitleChanged
stats.Deltas.OtherChanged += deltas.OtherChanged
chunk = chunk[:0]
done := stats.Synced + stats.Unchanged + stats.Skipped
if stats.Total > 0 {
stats.Progress = done * 100 / stats.Total
if stats.Progress > 99 {
stats.Progress = 99
}
}
_ = s.updateJobProgress(ctx, jobID, stats)
return nil
}
onRow := func(row feedRow) error {
stats.Total++
mapped, gtin := applyMappings(row, mappings)
if gtin == "" {
stats.Skipped++
return nil
}
rawCopy := make(map[string]string, len(row))
for k, v := range row {
rawCopy[k] = v
}
mh := sha256Hex([]byte(mustJSON(mapped)))
chunk = append(chunk, pendingProduct{
GTIN: gtin, RawData: rawCopy, MappedData: mapped, ContentHash: mh,
})
return flush(false)
}
body, err := src.Open()
if err != nil {
return stats, err
}
defer body.Close()
var parseCount int
switch format {
case "xml":
parseCount, err = parseXMLItems(body, itemPath, onRow)
default:
parseCount, err = parseCSV(body, onRow)
}
if err != nil {
return stats, err
}
if parseCount == 0 {
return stats, ClientMsg("feed contained no rows")
}
if err := flush(true); err != nil {
return stats, err
}
stats.Progress = 100
_ = s.updateJobProgress(ctx, jobID, stats)
return stats, nil
}
type existingProduct struct {
ID uuid.UUID
MappedData []byte
}
const (
// Set-based upserts (UNNEST / ANY) — one statement per op kind, queued in a single
// pgx.Batch round-trip. Mirrors catalog/import_csv.go patterns.
upsertSQLInsertSet = `
INSERT INTO raw_products (
company_id, gtin, feed_id, raw_data, mapped_data, sync_job_id,
is_processed, processing_status, updated_at
)
SELECT $1, v.gtin, $2, v.raw_data::jsonb, v.mapped_data::jsonb, $3, false, 'unprocessed', now()
FROM unnest($4::text[], $5::text[], $6::text[]) AS v(gtin, raw_data, mapped_data)
ON CONFLICT (company_id, gtin) DO UPDATE SET
feed_id = EXCLUDED.feed_id,
raw_data = EXCLUDED.raw_data,
mapped_data = EXCLUDED.mapped_data,
sync_job_id = EXCLUDED.sync_job_id,
is_processed = false,
processing_status = 'unprocessed',
updated_at = now()`
upsertSQLTouchSet = `
UPDATE raw_products SET sync_job_id = $3, feed_id = $4, updated_at = now()
WHERE company_id = $1 AND id = ANY($2::uuid[])`
upsertSQLUpdateSet = `
UPDATE raw_products AS r SET
feed_id = $2,
raw_data = v.raw_data::jsonb,
mapped_data = v.mapped_data::jsonb,
sync_job_id = $3,
is_processed = false,
processing_status = 'unprocessed',
updated_at = now()
FROM unnest($4::uuid[], $5::text[], $6::text[]) AS v(id, raw_data, mapped_data)
WHERE r.id = v.id AND r.company_id = $1`
)
type upsertOpKind int
const (
upsertOpInsert upsertOpKind = iota
upsertOpTouch
upsertOpUpdate
)
type upsertOp struct {
kind upsertOpKind
gtin string
id uuid.UUID
rawJSON []byte
mappedJSON []byte
changes []string
}
// classifyUpsertOps decides insert / touch / update per row without DB I/O so a chunk
// can be applied with set-based UNNEST statements (≤3) in one pgx.Batch round-trip.
// Touch keeps is_processed/processing_status when mapped_data is equal (canonical JSON).
// Updates/inserts stamp mapped_data._sync_changes for seller filters (price/stock/…).
func classifyUpsertOps(chunk []pendingProduct, existing map[string]existingProduct) (ops []upsertOp, skipped int) {
ops = make([]upsertOp, 0, len(chunk))
for _, p := range chunk {
rawJSON, err := json.Marshal(p.RawData)
if err != nil {
skipped++
continue
}
ex, found := existing[p.GTIN]
cleanMapped := stripSyncChanges(p.MappedData)
if !found {
withFlags := withSyncChanges(cleanMapped, []string{"new"})
mappedJSON, err := json.Marshal(withFlags)
if err != nil {
skipped++
continue
}
ops = append(ops, upsertOp{
kind: upsertOpInsert, gtin: p.GTIN, rawJSON: rawJSON, mappedJSON: mappedJSON,
changes: []string{"new"},
})
continue
}
// Category (and similar extras) live outside feed mappings — keep them.
cleanMapped = preserveSyncedExtras(ex.MappedData, cleanMapped)
plainJSON, err := json.Marshal(cleanMapped)
if err != nil {
skipped++
continue
}
if bytesEqualJSON(stripSyncChangesBytes(ex.MappedData), plainJSON) {
ops = append(ops, upsertOp{kind: upsertOpTouch, id: ex.ID})
continue
}
changes := detectMappedChanges(ex.MappedData, cleanMapped)
withFlags := withSyncChanges(cleanMapped, changes)
mappedJSON, err := json.Marshal(withFlags)
if err != nil {
skipped++
continue
}
ops = append(ops, upsertOp{
kind: upsertOpUpdate, id: ex.ID, rawJSON: rawJSON, mappedJSON: mappedJSON,
changes: changes,
})
}
return ops, skipped
}
func stripSyncChangesBytes(raw []byte) []byte {
if len(raw) == 0 {
return raw
}
var m map[string]any
if err := json.Unmarshal(raw, &m); err != nil {
return raw
}
out, err := json.Marshal(stripSyncChanges(m))
if err != nil {
return raw
}
return out
}
// dedupePendingByGTIN keeps first-seen order but last-seen payload per GTIN so a single
// UNNEST INSERT cannot hit "cannot affect row a second time".
func dedupePendingByGTIN(chunk []pendingProduct) []pendingProduct {
if len(chunk) < 2 {
return chunk
}
by := make(map[string]pendingProduct, len(chunk))
order := make([]string, 0, len(chunk))
for _, p := range chunk {
if _, ok := by[p.GTIN]; !ok {
order = append(order, p.GTIN)
}
by[p.GTIN] = p
}
if len(order) == len(chunk) {
return chunk
}
out := make([]pendingProduct, 0, len(order))
for _, gtin := range order {
out = append(out, by[gtin])
}
return out
}
func partitionUpsertOps(ops []upsertOp) (inserts, touches, updates []upsertOp) {
for _, op := range ops {
switch op.kind {
case upsertOpInsert:
inserts = append(inserts, op)
case upsertOpTouch:
touches = append(touches, op)
default:
updates = append(updates, op)
}
}
return inserts, touches, updates
}
func (s *Service) upsertChunk(ctx context.Context, companyID, feedID, jobID uuid.UUID, chunk []pendingProduct) (synced, unchanged, skipped int, deltas syncDeltaCounts, err error) {
chunk = dedupePendingByGTIN(chunk)
if len(chunk) == 0 {
return 0, 0, 0, deltas, nil
}
gtins := make([]string, 0, len(chunk))
for _, p := range chunk {
gtins = append(gtins, p.GTIN)
}
existing, err := s.loadExistingByGTIN(ctx, companyID, gtins)
if err != nil {
return 0, 0, 0, deltas, err
}
ops, skipped := classifyUpsertOps(chunk, existing)
if len(ops) == 0 {
deltas.Skipped = skipped
return 0, 0, skipped, deltas, nil
}
inserts, touches, updates := partitionUpsertOps(ops)
batch := &pgx.Batch{}
queued := 0
if len(inserts) > 0 {
gtinCol := make([]string, len(inserts))
rawCol := make([]string, len(inserts))
mappedCol := make([]string, len(inserts))
for i, op := range inserts {
gtinCol[i] = op.gtin
rawCol[i] = string(op.rawJSON)
mappedCol[i] = string(op.mappedJSON)
deltas.addChanges(op.changes)
}
batch.Queue(upsertSQLInsertSet, companyID, feedID, jobID, gtinCol, rawCol, mappedCol)
queued++
}
if len(touches) > 0 {
ids := make([]uuid.UUID, len(touches))
for i, op := range touches {
ids[i] = op.id
}
batch.Queue(upsertSQLTouchSet, companyID, ids, jobID, feedID)
queued++
}
if len(updates) > 0 {
ids := make([]uuid.UUID, len(updates))
rawCol := make([]string, len(updates))
mappedCol := make([]string, len(updates))
for i, op := range updates {
ids[i] = op.id
rawCol[i] = string(op.rawJSON)
mappedCol[i] = string(op.mappedJSON)
deltas.addChanges(op.changes)
}
batch.Queue(upsertSQLUpdateSet, companyID, feedID, jobID, ids, rawCol, mappedCol)
queued++
}
br := s.Pool.SendBatch(ctx, batch)
defer br.Close()
for i := 0; i < queued; i++ {
if _, err := br.Exec(); err != nil {
return synced, unchanged, skipped, deltas, err
}
}
// Content inserts/updates stamp raw as unprocessed — drop catalog rows so
// processed_products cannot outlive that reset (matches resetRawProducts).
// Touches keep processing_status and must not invalidate catalog.
if len(inserts) > 0 || len(updates) > 0 {
invalidateIDs := make([]uuid.UUID, 0, len(updates))
for _, op := range updates {
invalidateIDs = append(invalidateIDs, op.id)
}
invalidateGTINs := make([]string, 0, len(inserts))
for _, op := range inserts {
invalidateGTINs = append(invalidateGTINs, op.gtin)
}
if _, err := s.Pool.Exec(ctx, `
DELETE FROM processed_products
WHERE company_id = $1
AND (
raw_product_id = ANY($2::uuid[])
OR raw_product_id IN (
SELECT id FROM raw_products
WHERE company_id = $1 AND gtin = ANY($3::text[])
)
)`, companyID, invalidateIDs, invalidateGTINs); err != nil {
return synced, unchanged, skipped, deltas, err
}
}
synced = len(inserts) + len(updates)
unchanged = len(touches)
deltas.Unchanged = unchanged
deltas.Skipped = skipped
return synced, unchanged, skipped, deltas, nil
}
func (s *Service) loadExistingByGTIN(ctx context.Context, companyID uuid.UUID, gtins []string) (map[string]existingProduct, error) {
out := make(map[string]existingProduct, len(gtins))
if len(gtins) == 0 {
return out, nil
}
rows, err := s.Pool.Query(ctx, `
SELECT id, gtin, mapped_data FROM raw_products
WHERE company_id = $1 AND gtin = ANY($2::text[])`, companyID, gtins)
if err != nil {
return nil, err
}
defer rows.Close()
for rows.Next() {
var ex existingProduct
var gtin string
if err := rows.Scan(&ex.ID, &gtin, &ex.MappedData); err != nil {
return nil, err
}
out[gtin] = ex
}
return out, rows.Err()
}
func (s *Service) createSyncJob(ctx context.Context, companyID, feedID uuid.UUID) (uuid.UUID, error) {
var id uuid.UUID
err := s.Pool.QueryRow(ctx, `
INSERT INTO feed_sync_jobs (feed_id, company_id, status)
VALUES ($1, $2, 'pending') RETURNING id`, feedID, companyID).Scan(&id)
return id, err
}
// findOrCreatePendingSyncJob returns an existing pending job for the feed, or inserts one.
// Uses a transaction advisory lock so concurrent enqueues do not stack duplicates.
func (s *Service) findOrCreatePendingSyncJob(ctx context.Context, companyID, feedID uuid.UUID) (uuid.UUID, error) {
tx, err := s.Pool.Begin(ctx)
if err != nil {
return uuid.Nil, err
}
defer tx.Rollback(ctx)
if _, err := tx.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtextextended($1::text, 0))`, feedID.String()); err != nil {
return uuid.Nil, err
}
var id uuid.UUID
err = tx.QueryRow(ctx, `
SELECT id FROM feed_sync_jobs
WHERE feed_id = $1 AND company_id = $2 AND status = 'pending'
ORDER BY created_at
LIMIT 1`, feedID, companyID).Scan(&id)
if err == nil {
if err := tx.Commit(ctx); err != nil {
return uuid.Nil, err
}
return id, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return uuid.Nil, err
}
err = tx.QueryRow(ctx, `
INSERT INTO feed_sync_jobs (feed_id, company_id, status)
VALUES ($1, $2, 'pending') RETURNING id`, feedID, companyID).Scan(&id)
if err != nil {
return uuid.Nil, err
}
if err := tx.Commit(ctx); err != nil {
return uuid.Nil, err
}
return id, nil
}
func (s *Service) markJobRunning(ctx context.Context, jobID uuid.UUID) error {
_, err := s.Pool.Exec(ctx, `
UPDATE feed_sync_jobs SET status = 'running', started_at = now(), updated_at = now()
WHERE id = $1`, jobID)
return err
}
func (s *Service) updateJobProgress(ctx context.Context, jobID uuid.UUID, stats syncStats) error {
_, err := s.Pool.Exec(ctx, `
UPDATE feed_sync_jobs SET
products_synced = $2,
products_total = $3,
products_skipped = $4,
products_unchanged = $5,
progress = $6,
content_hash = NULLIF($7, ''),
updated_at = now()
WHERE id = $1`,
jobID, stats.Synced, stats.Total, stats.Skipped, stats.Unchanged, stats.Progress, stats.ContentHash)
return err
}
func (s *Service) completeJob(ctx context.Context, jobID uuid.UUID, stats syncStats) error {
_, err := s.Pool.Exec(ctx, `
UPDATE feed_sync_jobs SET
status = 'completed',
completed_at = now(),
products_synced = $2,
products_total = $3,
products_skipped = $4,
products_unchanged = $5,
progress = 100,
content_hash = NULLIF($6, ''),
updated_at = now()
WHERE id = $1`,
jobID, stats.Synced, stats.Total, stats.Skipped, stats.Unchanged, stats.ContentHash)
return err
}
func (s *Service) failJob(ctx context.Context, jobID uuid.UUID, msg string, stats syncStats) error {
_, err := s.Pool.Exec(ctx, `
UPDATE feed_sync_jobs SET
status = 'failed',
error = $2,
completed_at = now(),
products_synced = $3,
products_total = $4,
products_skipped = $5,
products_unchanged = $6,
progress = $7,
content_hash = NULLIF($8, ''),
updated_at = now()
WHERE id = $1`,
jobID, truncateErr(msg), stats.Synced, stats.Total, stats.Skipped, stats.Unchanged, stats.Progress, stats.ContentHash)
return err
}
func (s *Service) lastContentHash(ctx context.Context, feedID uuid.UUID) (string, error) {
var hash *string
err := s.Pool.QueryRow(ctx, `
SELECT content_hash FROM feed_sync_jobs
WHERE feed_id = $1 AND status = 'completed' AND content_hash IS NOT NULL AND content_hash <> ''
ORDER BY completed_at DESC NULLS LAST LIMIT 1`, feedID).Scan(&hash)
if errors.Is(err, pgx.ErrNoRows) {
return "", nil
}
if err != nil {
return "", err
}
if hash == nil {
return "", nil
}
return *hash, nil
}
func (s *Service) loadMappingsRaw(ctx context.Context, companyID, feedID uuid.UUID) (any, error) {
var raw []byte
err := s.Pool.QueryRow(ctx, `
SELECT mappings FROM feed_mappings
WHERE feed_id = $1 AND company_id = $2 AND is_active = true
ORDER BY version DESC LIMIT 1`, feedID, companyID).Scan(&raw)
if err != nil {
return nil, err
}
var m any
if err := json.Unmarshal(raw, &m); err != nil {
return nil, err
}
return m, nil
}
func sha256Hex(b []byte) string {
sum := sha256.Sum256(b)
return hex.EncodeToString(sum[:])
}
func sha256HexFile(src *feedBlob) (string, error) {
f, err := src.Open()
if err != nil {
return "", err
}
defer f.Close()
h := sha256.New()
if _, err := io.Copy(h, f); err != nil {
return "", err
}
return hex.EncodeToString(h.Sum(nil)), nil
}
func contentTypeFromBlob(src *feedBlob) string {
if src == nil {
return ""
}
return src.contentType
}
func mustJSON(v any) string {
b, err := json.Marshal(v)
if err != nil {
return ""
}
return string(b)
}
func bytesEqualJSON(a, b []byte) bool {
if len(a) == 0 && len(b) == 0 {
return true
}
var xa, xb any
if json.Unmarshal(a, &xa) != nil || json.Unmarshal(b, &xb) != nil {
return string(a) == string(b)
}
ba, _ := json.Marshal(xa)
bb, _ := json.Marshal(xb)
return string(ba) == string(bb)
}
func truncateErr(msg string) string {
if len(msg) > 2000 {
return msg[:2000]
}
return msg
}