636 lines
20 KiB
Go
636 lines
20 KiB
Go
package feeds
|
|||
|
|
|
||
|
|
import (
|
||
|
|
"context"
|
||
|
|
"encoding/json"
|
||
|
|
"errors"
|
||
|
|
"fmt"
|
||
|
|
"strings"
|
||
|
|
|
||
|
|
"github.com/google/uuid"
|
||
|
|
"github.com/jackc/pgx/v5"
|
||
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
||
|
|
)
|
||
|
|
|
||
|
|
type Service struct {
|
||
|
|
Pool *pgxpool.Pool
|
||
|
|
UploadDir string
|
||
|
|
}
|
||
|
|
|
||
|
|
// CreateInput is the payload for creating an input feed (URL and/or uploaded CSV).
|
||
|
|
// Legacy clients send name + item_path (url optional). V2 clients send name + url/file
|
||
|
|
// plus optional feed_type / sync_interval_minutes (or legacy sync_frequency in hours).
|
||
|
|
type CreateInput struct {
|
||
|
|
Name string
|
||
|
|
URL string
|
||
|
|
ItemPath string
|
||
|
|
FeedType string
|
||
|
|
SyncIntervalMinutes int
|
||
|
|
SyncFrequencyHours int // legacy alias; converted to minutes when SyncIntervalMinutes unset
|
||
|
|
Options map[string]any
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) List(ctx context.Context, companyID uuid.UUID, limit, offset int, q string) ([]map[string]any, int64, int64, int64, error) {
|
||
|
|
q = strings.TrimSpace(q)
|
||
|
|
where := `company_id = $1`
|
||
|
|
args := []any{companyID}
|
||
|
|
if q != "" {
|
||
|
|
where += ` AND (
|
||
|
|
COALESCE(name, '') ILIKE '%' || $2 || '%' OR
|
||
|
|
COALESCE(url, '') ILIKE '%' || $2 || '%' OR
|
||
|
|
COALESCE(feed_type, '') ILIKE '%' || $2 || '%' OR
|
||
|
|
COALESCE(status, '') ILIKE '%' || $2 || '%' OR
|
||
|
|
COALESCE(options->>'source_filename', '') ILIKE '%' || $2 || '%'
|
||
|
|
)`
|
||
|
|
args = append(args, q)
|
||
|
|
}
|
||
|
|
var total int64
|
||
|
|
if err := s.Pool.QueryRow(ctx, `SELECT count(*) FROM input_feeds WHERE `+where, args...).Scan(&total); err != nil {
|
||
|
|
return nil, 0, 0, 0, err
|
||
|
|
}
|
||
|
|
// active_total = truly syncing; mapped_total = fields saved but not activated.
|
||
|
|
var activeTotal, mappedTotal int64
|
||
|
|
if err := s.Pool.QueryRow(ctx,
|
||
|
|
`SELECT
|
||
|
|
count(*) FILTER (WHERE lower(COALESCE(status, '')) = 'active'),
|
||
|
|
count(*) FILTER (WHERE lower(COALESCE(status, '')) = 'mapped')
|
||
|
|
FROM input_feeds WHERE `+where,
|
||
|
|
args...,
|
||
|
|
).Scan(&activeTotal, &mappedTotal); err != nil {
|
||
|
|
return nil, 0, 0, 0, err
|
||
|
|
}
|
||
|
|
limitArg := len(args) + 1
|
||
|
|
offsetArg := len(args) + 2
|
||
|
|
query := fmt.Sprintf(`
|
||
|
|
SELECT id, name, url, feed_type, status, sync_interval_minutes, last_synced_at, options, created_at, updated_at,
|
||
|
|
(SELECT count(*)::int FROM raw_products rp WHERE rp.feed_id = input_feeds.id AND rp.company_id = input_feeds.company_id) AS product_count,
|
||
|
|
(SELECT max(rp.updated_at) FROM raw_products rp WHERE rp.feed_id = input_feeds.id AND rp.company_id = input_feeds.company_id) AS products_updated_at
|
||
|
|
FROM input_feeds WHERE %s
|
||
|
|
ORDER BY created_at DESC LIMIT $%d OFFSET $%d`, where, limitArg, offsetArg)
|
||
|
|
queryArgs := append(append([]any{}, args...), limit, offset)
|
||
|
|
rows, err := s.Pool.Query(ctx, query, queryArgs...)
|
||
|
|
if err != nil {
|
||
|
|
return nil, 0, 0, 0, err
|
||
|
|
}
|
||
|
|
defer rows.Close()
|
||
|
|
items, err := scanMaps(rows, []string{"id", "name", "url", "feed_type", "status", "sync_interval_minutes", "last_synced_at", "options", "created_at", "updated_at", "product_count", "products_updated_at"})
|
||
|
|
if err != nil {
|
||
|
|
return nil, 0, 0, 0, err
|
||
|
|
}
|
||
|
|
if err := s.attachMappingFieldCounts(ctx, companyID, items); err != nil {
|
||
|
|
return nil, 0, 0, 0, err
|
||
|
|
}
|
||
|
|
return items, total, activeTotal, mappedTotal, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// ProductTotals is company-scoped catalog counts across all feeds.
|
||
|
|
type ProductTotals struct {
|
||
|
|
Total int64
|
||
|
|
Processed int64
|
||
|
|
Unprocessed int64
|
||
|
|
}
|
||
|
|
|
||
|
|
// CompanyProductTotals returns company-scoped catalog counts for feeds/dashboard cards.
|
||
|
|
// ASSUMPTION: Total = count(raw_products); Processed = count(processed_products);
|
||
|
|
// Unprocessed = count(raw where processing_status='unprocessed'). These are not a
|
||
|
|
// partition of Total (P+U≠Total by design). Do not redefine without documenting a new ASSUMPTION.
|
||
|
|
func (s *Service) CompanyProductTotals(ctx context.Context, companyID uuid.UUID) (ProductTotals, error) {
|
||
|
|
var t ProductTotals
|
||
|
|
err := s.Pool.QueryRow(ctx, `
|
||
|
|
SELECT
|
||
|
|
(SELECT count(*)::bigint FROM raw_products WHERE company_id = $1),
|
||
|
|
(SELECT count(*)::bigint FROM processed_products WHERE company_id = $1),
|
||
|
|
(SELECT count(*)::bigint FROM raw_products
|
||
|
|
WHERE company_id = $1 AND lower(COALESCE(processing_status, '')) = 'unprocessed')`,
|
||
|
|
companyID,
|
||
|
|
).Scan(&t.Total, &t.Processed, &t.Unprocessed)
|
||
|
|
return t, err
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) Create(ctx context.Context, companyID uuid.UUID, in CreateInput) (map[string]any, error) {
|
||
|
|
name := strings.TrimSpace(in.Name)
|
||
|
|
if name == "" {
|
||
|
|
return nil, ClientMsg("name required")
|
||
|
|
}
|
||
|
|
url := strings.TrimSpace(in.URL)
|
||
|
|
if err := ValidateFeedURL(ctx, url); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
opts := in.Options
|
||
|
|
if opts == nil {
|
||
|
|
opts = map[string]any{}
|
||
|
|
}
|
||
|
|
itemPath := strings.TrimSpace(in.ItemPath)
|
||
|
|
if itemPath == "" {
|
||
|
|
if p, ok := opts["item_path"].(string); ok {
|
||
|
|
itemPath = strings.TrimSpace(p)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if itemPath != "" {
|
||
|
|
opts["item_path"] = itemPath
|
||
|
|
}
|
||
|
|
hasLocal := sourcePathFromOptions(opts) != ""
|
||
|
|
// Dual-support: legacy create allows name + item_path without url/file.
|
||
|
|
if url == "" && !hasLocal && itemPath == "" {
|
||
|
|
return nil, errSourceRequired
|
||
|
|
}
|
||
|
|
feedType := strings.ToLower(strings.TrimSpace(in.FeedType))
|
||
|
|
if feedType == "" {
|
||
|
|
if hasLocal {
|
||
|
|
feedType = "csv"
|
||
|
|
} else {
|
||
|
|
feedType = "xml"
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if feedType != "xml" && feedType != "csv" {
|
||
|
|
return nil, ClientMsg("feed_type must be xml or csv")
|
||
|
|
}
|
||
|
|
interval := in.SyncIntervalMinutes
|
||
|
|
if interval <= 0 && in.SyncFrequencyHours > 0 {
|
||
|
|
interval = in.SyncFrequencyHours * 60
|
||
|
|
}
|
||
|
|
if interval <= 0 {
|
||
|
|
interval = 60
|
||
|
|
}
|
||
|
|
optsBytes, err := json.Marshal(opts)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
var id uuid.UUID
|
||
|
|
err = s.Pool.QueryRow(ctx, `
|
||
|
|
INSERT INTO input_feeds (company_id, name, url, feed_type, status, sync_interval_minutes, options)
|
||
|
|
VALUES ($1, $2, $3, $4, 'unmapped', $5, $6::jsonb) RETURNING id`,
|
||
|
|
companyID, name, nullStr(url), feedType, interval, optsBytes).Scan(&id)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
return s.Get(ctx, companyID, id)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) Get(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
|
||
|
|
row := s.Pool.QueryRow(ctx, `
|
||
|
|
SELECT id, name, url, feed_type, status, sync_interval_minutes, last_synced_at, options, created_at, updated_at,
|
||
|
|
(SELECT count(*)::int FROM raw_products rp WHERE rp.feed_id = input_feeds.id AND rp.company_id = input_feeds.company_id) AS product_count,
|
||
|
|
(SELECT max(rp.updated_at) FROM raw_products rp WHERE rp.feed_id = input_feeds.id AND rp.company_id = input_feeds.company_id) AS products_updated_at
|
||
|
|
FROM input_feeds WHERE id = $1 AND company_id = $2`, id, companyID)
|
||
|
|
item, err := scanMap(row, []string{"id", "name", "url", "feed_type", "status", "sync_interval_minutes", "last_synced_at", "options", "created_at", "updated_at", "product_count", "products_updated_at"})
|
||
|
|
if err != nil {
|
||
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
||
|
|
return nil, ErrNotFound
|
||
|
|
}
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
if err := s.attachMappingFieldCounts(ctx, companyID, []map[string]any{item}); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
return item, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// attachMappingFieldCounts sets mapping_field_count / has_mappings / mapping_incomplete
|
||
|
|
// on each feed row from the active feed_mappings document (one query for the page).
|
||
|
|
// mapping_incomplete mirrors list-chip blocking preflight (empty/required/item_path).
|
||
|
|
func (s *Service) attachMappingFieldCounts(ctx context.Context, companyID uuid.UUID, items []map[string]any) error {
|
||
|
|
if len(items) == 0 {
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
ids := make([]uuid.UUID, 0, len(items))
|
||
|
|
index := make(map[uuid.UUID]map[string]any, len(items))
|
||
|
|
for _, item := range items {
|
||
|
|
id, ok := asUUID(item["id"])
|
||
|
|
if !ok {
|
||
|
|
continue
|
||
|
|
}
|
||
|
|
ids = append(ids, id)
|
||
|
|
index[id] = item
|
||
|
|
item["mapping_field_count"] = 0
|
||
|
|
item["has_mappings"] = false
|
||
|
|
item["mapping_incomplete"] = true
|
||
|
|
}
|
||
|
|
if len(ids) == 0 {
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
required, err := s.loadRequiredStandardFields(ctx, companyID)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
rows, err := s.Pool.Query(ctx, `
|
||
|
|
SELECT DISTINCT ON (feed_id) feed_id, mappings
|
||
|
|
FROM feed_mappings
|
||
|
|
WHERE company_id = $1 AND is_active = true AND feed_id = ANY($2::uuid[])
|
||
|
|
ORDER BY feed_id, version DESC`, companyID, ids)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
defer rows.Close()
|
||
|
|
for rows.Next() {
|
||
|
|
var feedID uuid.UUID
|
||
|
|
var raw []byte
|
||
|
|
if err := rows.Scan(&feedID, &raw); err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
item := index[feedID]
|
||
|
|
if item == nil {
|
||
|
|
continue
|
||
|
|
}
|
||
|
|
var parsed any
|
||
|
|
if err := json.Unmarshal(raw, &parsed); err != nil {
|
||
|
|
continue
|
||
|
|
}
|
||
|
|
mappings := parseMappings(parsed)
|
||
|
|
n := len(mappings)
|
||
|
|
item["mapping_field_count"] = n
|
||
|
|
item["has_mappings"] = n > 0
|
||
|
|
item["mapping_incomplete"] = mappingDocIncomplete(item, parsed, mappings, required)
|
||
|
|
}
|
||
|
|
return rows.Err()
|
||
|
|
}
|
||
|
|
|
||
|
|
func isCSVFeedType(feedType string) bool {
|
||
|
|
t := strings.ToLower(strings.TrimSpace(feedType))
|
||
|
|
return t == "csv" || t == "excel"
|
||
|
|
}
|
||
|
|
|
||
|
|
func feedItemPathHint(item map[string]any, mappingsRaw any) string {
|
||
|
|
if p := itemPathFromMappings(mappingsRaw); p != "" {
|
||
|
|
return p
|
||
|
|
}
|
||
|
|
if opts, ok := item["options"].(map[string]any); ok {
|
||
|
|
if v, ok := opts["item_path"].(string); ok {
|
||
|
|
if p := strings.TrimSpace(v); p != "" {
|
||
|
|
return p
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if v, ok := item["item_path"].(string); ok {
|
||
|
|
return strings.TrimSpace(v)
|
||
|
|
}
|
||
|
|
return ""
|
||
|
|
}
|
||
|
|
|
||
|
|
// mappingDocIncomplete reports list-chip blocking gaps (empty mappings, required targets, XML item_path).
|
||
|
|
func mappingDocIncomplete(item map[string]any, mappingsRaw any, mappings []FieldMapping, required []requiredStandardField) bool {
|
||
|
|
if err := validateMappingsForSync(mappings, required); err != nil {
|
||
|
|
return true
|
||
|
|
}
|
||
|
|
feedType, _ := item["feed_type"].(string)
|
||
|
|
if !isCSVFeedType(feedType) && feedItemPathHint(item, mappingsRaw) == "" {
|
||
|
|
return true
|
||
|
|
}
|
||
|
|
return false
|
||
|
|
}
|
||
|
|
|
||
|
|
func asUUID(v any) (uuid.UUID, bool) {
|
||
|
|
switch t := v.(type) {
|
||
|
|
case uuid.UUID:
|
||
|
|
return t, true
|
||
|
|
case [16]byte:
|
||
|
|
return uuid.UUID(t), true
|
||
|
|
case string:
|
||
|
|
id, err := uuid.Parse(t)
|
||
|
|
return id, err == nil
|
||
|
|
default:
|
||
|
|
return uuid.Nil, false
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) Update(ctx context.Context, companyID, id uuid.UUID, body map[string]any) (map[string]any, error) {
|
||
|
|
name, _ := body["name"].(string)
|
||
|
|
url, _ := body["url"].(string)
|
||
|
|
status, _ := body["status"].(string)
|
||
|
|
feedType, _ := body["feed_type"].(string)
|
||
|
|
itemPath, _ := body["item_path"].(string)
|
||
|
|
if err := ValidateFeedURL(ctx, url); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
feedType = strings.ToLower(strings.TrimSpace(feedType))
|
||
|
|
if feedType != "" && feedType != "xml" && feedType != "csv" {
|
||
|
|
return nil, ClientMsg("feed_type must be xml or csv")
|
||
|
|
}
|
||
|
|
interval := 0
|
||
|
|
switch v := body["sync_interval_minutes"].(type) {
|
||
|
|
case float64:
|
||
|
|
interval = int(v)
|
||
|
|
case int:
|
||
|
|
interval = v
|
||
|
|
case json.Number:
|
||
|
|
n, _ := v.Int64()
|
||
|
|
interval = int(n)
|
||
|
|
}
|
||
|
|
if interval <= 0 {
|
||
|
|
switch v := body["sync_frequency"].(type) {
|
||
|
|
case float64:
|
||
|
|
interval = int(v) * 60
|
||
|
|
case int:
|
||
|
|
interval = v * 60
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
ct, err := s.Pool.Exec(ctx, `
|
||
|
|
UPDATE input_feeds SET
|
||
|
|
name = CASE WHEN $3 <> '' THEN $3 ELSE name END,
|
||
|
|
url = CASE WHEN $4 <> '' THEN $4 ELSE url END,
|
||
|
|
status = CASE WHEN $5 <> '' THEN $5 ELSE status END,
|
||
|
|
feed_type = CASE WHEN $6 <> '' THEN $6 ELSE feed_type END,
|
||
|
|
sync_interval_minutes = CASE WHEN $7 > 0 THEN $7 ELSE sync_interval_minutes END,
|
||
|
|
options = CASE
|
||
|
|
WHEN $8 <> '' THEN COALESCE(options, '{}'::jsonb) || jsonb_build_object('item_path', to_jsonb($8::text))
|
||
|
|
ELSE options
|
||
|
|
END,
|
||
|
|
updated_at = now()
|
||
|
|
WHERE id = $1 AND company_id = $2`,
|
||
|
|
id, companyID, name, url, status, feedType, interval, strings.TrimSpace(itemPath))
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
if ct.RowsAffected() == 0 {
|
||
|
|
return nil, ErrNotFound
|
||
|
|
}
|
||
|
|
return s.Get(ctx, companyID, id)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) Delete(ctx context.Context, companyID, id uuid.UUID) error {
|
||
|
|
ct, err := s.Pool.Exec(ctx, `DELETE FROM input_feeds WHERE id = $1 AND company_id = $2`, id, companyID)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
if ct.RowsAffected() == 0 {
|
||
|
|
return ErrNotFound
|
||
|
|
}
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) GetMappings(ctx context.Context, companyID, feedID uuid.UUID) (map[string]any, error) {
|
||
|
|
var id uuid.UUID
|
||
|
|
var version int
|
||
|
|
var mappings []byte
|
||
|
|
err := s.Pool.QueryRow(ctx, `
|
||
|
|
SELECT id, version, 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(&id, &version, &mappings)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
var m any
|
||
|
|
_ = json.Unmarshal(mappings, &m)
|
||
|
|
return map[string]any{"id": id, "version": version, "mappings": m}, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) PutMappings(ctx context.Context, companyID, feedID uuid.UUID, mappings any) (map[string]any, error) {
|
||
|
|
b, err := json.Marshal(mappings)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
var version int
|
||
|
|
_ = s.Pool.QueryRow(ctx, `
|
||
|
|
SELECT COALESCE(MAX(version), 0) FROM feed_mappings WHERE feed_id = $1`, feedID).Scan(&version)
|
||
|
|
version++
|
||
|
|
_, _ = s.Pool.Exec(ctx, `UPDATE feed_mappings SET is_active = false WHERE feed_id = $1`, feedID)
|
||
|
|
var id uuid.UUID
|
||
|
|
err = s.Pool.QueryRow(ctx, `
|
||
|
|
INSERT INTO feed_mappings (feed_id, company_id, version, mappings, is_active)
|
||
|
|
VALUES ($1, $2, $3, $4, true) RETURNING id`, feedID, companyID, version, b).Scan(&id)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
// Keep feed.options.item_path in sync for Sync() XML item selection, and
|
||
|
|
// flip unmapped -> mapped whenever at least one field mapping is saved.
|
||
|
|
path := itemPathFromMappings(mappings)
|
||
|
|
hasFields := len(parseMappings(mappings)) > 0
|
||
|
|
switch {
|
||
|
|
case path != "":
|
||
|
|
_, _ = s.Pool.Exec(ctx, `
|
||
|
|
UPDATE input_feeds SET
|
||
|
|
options = COALESCE(options, '{}'::jsonb) || jsonb_build_object('item_path', to_jsonb($3::text)),
|
||
|
|
status = CASE WHEN status = 'unmapped' THEN 'mapped' ELSE status END,
|
||
|
|
updated_at = now()
|
||
|
|
WHERE id = $1 AND company_id = $2`, feedID, companyID, path)
|
||
|
|
case hasFields:
|
||
|
|
_, _ = s.Pool.Exec(ctx, `
|
||
|
|
UPDATE input_feeds SET
|
||
|
|
status = CASE WHEN status = 'unmapped' THEN 'mapped' ELSE status END,
|
||
|
|
updated_at = now()
|
||
|
|
WHERE id = $1 AND company_id = $2`, feedID, companyID)
|
||
|
|
}
|
||
|
|
return map[string]any{"id": id, "version": version, "mappings": mappings}, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) ListExportFeeds(ctx context.Context, companyID uuid.UUID, limit, offset int) ([]map[string]any, int64, error) {
|
||
|
|
var total int64
|
||
|
|
if err := s.Pool.QueryRow(ctx, `SELECT count(*) FROM export_feeds WHERE company_id = $1`, companyID).Scan(&total); err != nil {
|
||
|
|
return nil, 0, err
|
||
|
|
}
|
||
|
|
rows, err := s.Pool.Query(ctx, `
|
||
|
|
SELECT id, name, source_feed_id, format, public_token, is_active, last_generated_at, created_at, updated_at
|
||
|
|
FROM export_feeds WHERE company_id = $1
|
||
|
|
ORDER BY created_at DESC LIMIT $2 OFFSET $3`, companyID, limit, offset)
|
||
|
|
if err != nil {
|
||
|
|
return nil, 0, err
|
||
|
|
}
|
||
|
|
defer rows.Close()
|
||
|
|
items, err := scanMaps(rows, []string{"id", "name", "source_feed_id", "format", "public_token", "is_active", "last_generated_at", "created_at", "updated_at"})
|
||
|
|
if err != nil {
|
||
|
|
return nil, 0, err
|
||
|
|
}
|
||
|
|
return items, total, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) GetExportFeed(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
|
||
|
|
rows, err := s.Pool.Query(ctx, `
|
||
|
|
SELECT id, name, source_feed_id, format, public_token, template, filters, is_active, last_generated_at, created_at, updated_at
|
||
|
|
FROM export_feeds WHERE id = $1 AND company_id = $2`, id, companyID)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
defer rows.Close()
|
||
|
|
items, err := scanMaps(rows, []string{"id", "name", "source_feed_id", "format", "public_token", "template", "filters", "is_active", "last_generated_at", "created_at", "updated_at"})
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
if len(items) == 0 {
|
||
|
|
return nil, errors.New("not found")
|
||
|
|
}
|
||
|
|
return items[0], nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) DeleteExportFeed(ctx context.Context, companyID, id uuid.UUID) error {
|
||
|
|
ct, err := s.Pool.Exec(ctx, `DELETE FROM export_feeds WHERE id = $1 AND company_id = $2`, id, companyID)
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
if ct.RowsAffected() == 0 {
|
||
|
|
return errors.New("not found")
|
||
|
|
}
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) UpdateExportFeed(ctx context.Context, companyID, id uuid.UUID, name *string, isActive *bool, template, filters any) (map[string]any, error) {
|
||
|
|
current, err := s.GetExportFeed(ctx, companyID, id)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
if name != nil {
|
||
|
|
n := strings.TrimSpace(*name)
|
||
|
|
if n == "" {
|
||
|
|
return nil, ClientMsg("name required")
|
||
|
|
}
|
||
|
|
if _, err := s.Pool.Exec(ctx, `
|
||
|
|
UPDATE export_feeds SET name = $3, updated_at = now() WHERE id = $1 AND company_id = $2`, id, companyID, n); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if isActive != nil {
|
||
|
|
if _, err := s.Pool.Exec(ctx, `
|
||
|
|
UPDATE export_feeds SET is_active = $3, updated_at = now() WHERE id = $1 AND company_id = $2`, id, companyID, *isActive); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
}
|
||
|
|
if template != nil || filters != nil {
|
||
|
|
tpl := template
|
||
|
|
flt := filters
|
||
|
|
if tpl == nil {
|
||
|
|
tpl = current["template"]
|
||
|
|
}
|
||
|
|
if flt == nil {
|
||
|
|
flt = current["filters"]
|
||
|
|
}
|
||
|
|
if _, err := s.UpdateExportFeedTemplate(ctx, companyID, id, tpl, flt); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return s.GetExportFeed(ctx, companyID, id)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (s *Service) CreateExportFeed(ctx context.Context, companyID uuid.UUID, in CreateExportInput) (map[string]any, error) {
|
||
|
|
if strings.TrimSpace(in.Name) == "" {
|
||
|
|
return nil, ClientMsg("name required")
|
||
|
|
}
|
||
|
|
format := strings.ToLower(strings.TrimSpace(in.Format))
|
||
|
|
if format == "" {
|
||
|
|
format = "xml"
|
||
|
|
}
|
||
|
|
if format != "xml" && format != "csv" {
|
||
|
|
return nil, ClientMsg("format must be xml or csv")
|
||
|
|
}
|
||
|
|
var src *uuid.UUID
|
||
|
|
if in.SourceFeedID != nil && *in.SourceFeedID != "" {
|
||
|
|
id, err := uuid.Parse(*in.SourceFeedID)
|
||
|
|
if err != nil {
|
||
|
|
return nil, ClientMsg("invalid source_feed_id")
|
||
|
|
}
|
||
|
|
src = &id
|
||
|
|
}
|
||
|
|
tplBytes := []byte("{}")
|
||
|
|
if in.Template != nil {
|
||
|
|
b, err := json.Marshal(in.Template)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
tplBytes = b
|
||
|
|
}
|
||
|
|
filterBytes := []byte("{}")
|
||
|
|
if in.Filters != nil {
|
||
|
|
b, err := json.Marshal(in.Filters)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
filterBytes = b
|
||
|
|
}
|
||
|
|
token, err := newPublicExportToken()
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
var id uuid.UUID
|
||
|
|
err = s.Pool.QueryRow(ctx, `
|
||
|
|
INSERT INTO export_feeds (company_id, name, source_feed_id, format, template, filters, public_token)
|
||
|
|
VALUES ($1, $2, $3, $4, $5::jsonb, $6::jsonb, $7) RETURNING id, public_token`,
|
||
|
|
companyID, in.Name, src, format, tplBytes, filterBytes, token).Scan(&id, &token)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
return map[string]any{
|
||
|
|
"id": id, "name": in.Name, "format": format, "public_token": token,
|
||
|
|
"template": in.Template, "filters": in.Filters,
|
||
|
|
}, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// RotateExportFeedPublicToken replaces the public URL token (revokes the previous URL).
|
||
|
|
func (s *Service) RotateExportFeedPublicToken(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
|
||
|
|
token, err := newPublicExportToken()
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
ct, err := s.Pool.Exec(ctx, `
|
||
|
|
UPDATE export_feeds SET public_token = $3, updated_at = now()
|
||
|
|
WHERE id = $1 AND company_id = $2`, id, companyID, token)
|
||
|
|
if err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
if ct.RowsAffected() == 0 {
|
||
|
|
return nil, errors.New("not found")
|
||
|
|
}
|
||
|
|
return s.GetExportFeed(ctx, companyID, id)
|
||
|
|
}
|
||
|
|
|
||
|
|
func nullStr(s string) *string {
|
||
|
|
if s == "" {
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
return &s
|
||
|
|
}
|
||
|
|
|
||
|
|
func xmlEscape(s string) string {
|
||
|
|
r := strings.NewReplacer("&", "&", "<", "<", ">", ">", `"`, """)
|
||
|
|
return r.Replace(s)
|
||
|
|
}
|
||
|
|
|
||
|
|
func scanMaps(rows pgx.Rows, cols []string) ([]map[string]any, error) {
|
||
|
|
out := make([]map[string]any, 0)
|
||
|
|
for rows.Next() {
|
||
|
|
vals := make([]any, len(cols))
|
||
|
|
ptrs := make([]any, len(cols))
|
||
|
|
for i := range vals {
|
||
|
|
ptrs[i] = &vals[i]
|
||
|
|
}
|
||
|
|
if err := rows.Scan(ptrs...); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
m := make(map[string]any, len(cols))
|
||
|
|
for i, c := range cols {
|
||
|
|
m[c] = normalize(vals[i])
|
||
|
|
}
|
||
|
|
out = append(out, m)
|
||
|
|
}
|
||
|
|
return out, rows.Err()
|
||
|
|
}
|
||
|
|
|
||
|
|
func scanMap(row pgx.Row, cols []string) (map[string]any, error) {
|
||
|
|
vals := make([]any, len(cols))
|
||
|
|
ptrs := make([]any, len(cols))
|
||
|
|
for i := range vals {
|
||
|
|
ptrs[i] = &vals[i]
|
||
|
|
}
|
||
|
|
if err := row.Scan(ptrs...); err != nil {
|
||
|
|
return nil, err
|
||
|
|
}
|
||
|
|
m := make(map[string]any, len(cols))
|
||
|
|
for i, c := range cols {
|
||
|
|
m[c] = normalize(vals[i])
|
||
|
|
}
|
||
|
|
return m, nil
|
||
|
|
}
|
||
|
|
|
||
|
|
func normalize(v any) any {
|
||
|
|
switch t := v.(type) {
|
||
|
|
case []byte:
|
||
|
|
var j any
|
||
|
|
if json.Unmarshal(t, &j) == nil {
|
||
|
|
return j
|
||
|
|
}
|
||
|
|
return string(t)
|
||
|
|
case [16]byte:
|
||
|
|
return uuid.UUID(t).String()
|
||
|
|
default:
|
||
|
|
return v
|
||
|
|
}
|
||
|
|
}
|