Files
descrybe/apps/api/internal/catalog/service.go
T

1466 lines
52 KiB
Go
Raw Normal View History

package catalog
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"github.com/descrybe/descrybe-v2/apps/api/internal/company"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
// exactTotalFromPage returns an exact total when the page itself proves the
// result set size: a short page (fewer rows than limit) means there are no
// further rows, except an empty page past offset 0 which may be beyond EOF.
func exactTotalFromPage(offset, limit, pageLen int) (int64, bool) {
if limit <= 0 {
return 0, false
}
if pageLen < limit && (offset == 0 || pageLen > 0) {
return int64(offset + pageLen), true
}
return 0, false
}
// parallelCountAndList runs count and page queries concurrently. When the
// page is short enough to prove the exact total, the count query is cancelled
// so huge-table COUNT(*) can abort early. Safe when both share the same WHERE
// args and do not mutate shared slices.
func parallelCountAndList(
ctx context.Context,
limit, offset int,
countFn func(context.Context) (int64, error),
listFn func(context.Context) ([]map[string]any, error),
) ([]map[string]any, int64, error) {
type countRes struct {
n int64
err error
}
type listRes struct {
items []map[string]any
err error
}
countCtx, cancelCount := context.WithCancel(ctx)
defer cancelCount()
countCh := make(chan countRes, 1)
listCh := make(chan listRes, 1)
go func() {
n, err := countFn(countCtx)
countCh <- countRes{n: n, err: err}
}()
go func() {
items, err := listFn(ctx)
listCh <- listRes{items: items, err: err}
}()
lr := <-listCh
if lr.err != nil {
cancelCount()
<-countCh
return nil, 0, lr.err
}
if total, ok := exactTotalFromPage(offset, limit, len(lr.items)); ok {
cancelCount()
<-countCh
return lr.items, total, nil
}
cr := <-countCh
if cr.err != nil {
return nil, 0, cr.err
}
return lr.items, cr.n, nil
}
type Service struct {
Pool *pgxpool.Pool
}
func (s *Service) ListCategories(ctx context.Context, companyID uuid.UUID, f ListFilter) ([]map[string]any, int64, error) {
f = NormalizeListFilter(f)
args := []any{companyID}
where := []string{"company_id = $1"}
if f.Query != "" {
args = append(args, "%"+f.Query+"%")
n := len(args)
where = append(where, fmt.Sprintf("(name ILIKE $%d OR unique_id ILIKE $%d OR COALESCE(path, '') ILIKE $%d)", n, n, n))
}
wSQL := strings.Join(where, " AND ")
args = append(args, f.Limit, f.Offset)
lim := len(args) - 1
off := len(args)
rows, err := s.Pool.Query(ctx, fmt.Sprintf(`
SELECT id, name, unique_id, parent_unique_id, path, level, position, is_active, description,
title_template, description_template,
(title_template IS NOT NULL AND jsonb_typeof(title_template) = 'object'
AND COALESCE(jsonb_array_length(title_template->'elements'), 0) > 0) AS has_title_formula,
(description_template IS NOT NULL AND jsonb_typeof(description_template) = 'object'
AND COALESCE(jsonb_array_length(description_template->'sections'), 0) > 0) AS has_description_formula,
(EXISTS (
SELECT 1 FROM jsonb_each_text(COALESCE(prompt, '{}'::jsonb)) kv
WHERE length(trim(kv.value)) > 0
)) AS has_prompt,
created_at, updated_at
FROM categories WHERE %s
ORDER BY path NULLS LAST, position, name
LIMIT $%d OFFSET $%d`, wSQL, lim, off), args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
items, err := scanMaps(rows, []string{
"id", "name", "unique_id", "parent_unique_id", "path", "level", "position", "is_active", "description",
"title_template", "description_template", "has_title_formula", "has_description_formula", "has_prompt",
"created_at", "updated_at",
})
if err != nil {
return nil, 0, err
}
if total, ok := exactTotalFromPage(f.Offset, f.Limit, len(items)); ok {
return items, total, nil
}
countArgs := args[:len(args)-2]
var total int64
if err := s.Pool.QueryRow(ctx, `SELECT count(*) FROM categories WHERE `+wSQL, countArgs...).Scan(&total); err != nil {
return nil, 0, err
}
return items, total, nil
}
func (s *Service) CreateCategory(ctx context.Context, companyID uuid.UUID, name, uniqueID string, parent *string, desc *string) (map[string]any, error) {
name = strings.TrimSpace(name)
uniqueID = strings.TrimSpace(uniqueID)
if name == "" || uniqueID == "" {
return nil, ClientMsg("name and unique_id required")
}
path := uniqueID
level := 0
if parent != nil && strings.TrimSpace(*parent) != "" {
p := strings.TrimSpace(*parent)
var parentPath *string
var parentLevel int
err := s.Pool.QueryRow(ctx, `
SELECT path, level FROM categories
WHERE company_id = $1 AND unique_id = $2`, companyID, p).
Scan(&parentPath, &parentLevel)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, ClientMsg("parent category not found")
}
return nil, err
}
if parentPath != nil && *parentPath != "" {
path = *parentPath + "/" + uniqueID
} else {
path = p + "/" + uniqueID
}
level = parentLevel + 1
parent = &p
}
var id uuid.UUID
err := s.Pool.QueryRow(ctx, `
INSERT INTO categories (company_id, name, unique_id, parent_unique_id, description, path, level)
VALUES ($1, $2, $3, $4, $5, $6, $7)
RETURNING id`, companyID, name, uniqueID, parent, desc, path, level).Scan(&id)
if err != nil {
return nil, err
}
return s.GetCategory(ctx, companyID, id)
}
func (s *Service) GetCategory(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
row := s.Pool.QueryRow(ctx, `
SELECT id, name, unique_id, parent_unique_id, path, level, position, is_active, description,
COALESCE(prompt, '{}'::jsonb) AS prompts,
(EXISTS (
SELECT 1 FROM jsonb_each_text(COALESCE(prompt, '{}'::jsonb)) kv
WHERE length(trim(kv.value)) > 0
)) AS has_prompt,
title_template, description_template, created_at, updated_at
FROM categories WHERE id = $1 AND company_id = $2`, id, companyID)
item, err := scanMap(row, []string{
"id", "name", "unique_id", "parent_unique_id", "path", "level", "position", "is_active", "description",
"prompts", "has_prompt", "title_template", "description_template", "created_at", "updated_at",
})
if err != nil {
return nil, err
}
return enrichCategoryPrompts(ctx, s.Pool, companyID, item)
}
func (s *Service) UpdateTitleFormula(ctx context.Context, companyID, id uuid.UUID, template any) (map[string]any, error) {
b, err := json.Marshal(template)
if err != nil {
return nil, ClientMsg("invalid title_template")
}
ct, err := s.Pool.Exec(ctx, `
UPDATE categories SET title_template = $3::jsonb, updated_at = now()
WHERE id = $1 AND company_id = $2`, id, companyID, string(b))
if err != nil {
return nil, err
}
if ct.RowsAffected() == 0 {
return nil, ErrNotFound
}
return s.GetCategory(ctx, companyID, id)
}
func (s *Service) UpdateDescriptionFormula(ctx context.Context, companyID, id uuid.UUID, template any) (map[string]any, error) {
b, err := json.Marshal(template)
if err != nil {
return nil, ClientMsg("invalid description_template")
}
ct, err := s.Pool.Exec(ctx, `
UPDATE categories SET description_template = $3::jsonb, updated_at = now()
WHERE id = $1 AND company_id = $2`, id, companyID, string(b))
if err != nil {
return nil, err
}
if ct.RowsAffected() == 0 {
return nil, ErrNotFound
}
return s.GetCategory(ctx, companyID, id)
}
// MaxCategoryPromptRunes bounds per-category AI generation prompts.
const MaxCategoryPromptRunes = 8000
// UpdateCategoryPrompt sets per-language AI enhance user prompts for a category.
// Empty map (or all-empty values) clears overrides (company/built-in template applies).
func (s *Service) UpdateCategoryPrompt(ctx context.Context, companyID, id uuid.UUID, prompts map[string]string) (map[string]any, error) {
cleaned, err := company.SanitizeLangPromptMap(prompts, MaxCategoryPromptRunes)
if err != nil {
return nil, ClientMsg(err.Error())
}
raw, err := company.EncodeLangPromptMap(cleaned)
if err != nil {
return nil, err
}
ct, err := s.Pool.Exec(ctx, `
UPDATE categories SET prompt = $3::jsonb, updated_at = now()
WHERE id = $1 AND company_id = $2`, id, companyID, string(raw))
if err != nil {
return nil, err
}
if ct.RowsAffected() == 0 {
return nil, ErrNotFound
}
return s.GetCategory(ctx, companyID, id)
}
func enrichCategoryPrompts(ctx context.Context, pool *pgxpool.Pool, companyID uuid.UUID, item map[string]any) (map[string]any, error) {
if item == nil {
return nil, ErrNotFound
}
prompts, err := company.DecodeLangPromptMap(item["prompts"])
if err != nil {
prompts = company.LangPromptMap{}
}
primary := company.LoadLanguage(ctx, pool, companyID)
item["prompts"] = prompts
item["prompt"] = company.PromptForLanguage(prompts, primary)
item["has_prompt"] = company.HasAnyPrompt(prompts)
return item, nil
}
func (s *Service) ListVariables(ctx context.Context, companyID uuid.UUID, f ListFilter) ([]map[string]any, int64, error) {
f = NormalizeListFilter(f)
rows, err := s.Pool.Query(ctx, `
SELECT id, name, value, description, created_at, updated_at
FROM custom_variables WHERE company_id = $1
ORDER BY name LIMIT $2 OFFSET $3`, companyID, f.Limit, f.Offset)
if err != nil {
return nil, 0, err
}
defer rows.Close()
items, err := scanMaps(rows, []string{"id", "name", "value", "description", "created_at", "updated_at"})
if err != nil {
return nil, 0, err
}
if total, ok := exactTotalFromPage(f.Offset, f.Limit, len(items)); ok {
return items, total, nil
}
var total int64
if err := s.Pool.QueryRow(ctx, `SELECT count(*) FROM custom_variables WHERE company_id = $1`, companyID).Scan(&total); err != nil {
return nil, 0, err
}
return items, total, nil
}
func (s *Service) CreateVariable(ctx context.Context, companyID uuid.UUID, name, value string, description *string) (map[string]any, error) {
name = strings.TrimSpace(name)
if name == "" {
return nil, ClientMsg("name required")
}
var id uuid.UUID
err := s.Pool.QueryRow(ctx, `
INSERT INTO custom_variables (company_id, name, value, description)
VALUES ($1, $2, $3, $4) RETURNING id`, companyID, name, value, description).Scan(&id)
if err != nil {
return nil, err
}
return s.getVariable(ctx, companyID, id)
}
func (s *Service) DeleteVariable(ctx context.Context, companyID, id uuid.UUID) error {
ct, err := s.Pool.Exec(ctx, `DELETE FROM custom_variables 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) getVariable(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
row := s.Pool.QueryRow(ctx, `
SELECT id, name, value, description, created_at, updated_at
FROM custom_variables WHERE id = $1 AND company_id = $2`, id, companyID)
return scanMap(row, []string{"id", "name", "value", "description", "created_at", "updated_at"})
}
func (s *Service) UpdateCategory(ctx context.Context, companyID, id uuid.UUID, body map[string]any) (map[string]any, error) {
name, _ := body["name"].(string)
desc, _ := body["description"].(string)
var isActive *bool
if v, ok := body["is_active"].(bool); ok {
isActive = &v
}
_, err := s.Pool.Exec(ctx, `
UPDATE categories SET
name = CASE WHEN $3 <> '' THEN $3 ELSE name END,
description = CASE WHEN $4 <> '' THEN $4 ELSE description END,
is_active = COALESCE($5, is_active),
updated_at = now()
WHERE id = $1 AND company_id = $2`, id, companyID, name, desc, isActive)
if err != nil {
return nil, err
}
return s.GetCategory(ctx, companyID, id)
}
func (s *Service) DeleteCategory(ctx context.Context, companyID, id uuid.UUID) error {
_, err := s.Pool.Exec(ctx, `DELETE FROM categories WHERE id = $1 AND company_id = $2`, id, companyID)
return err
}
// DeleteCategoryByUniqueID deletes a category by its unique_id (legacy public DELETE path).
func (s *Service) DeleteCategoryByUniqueID(ctx context.Context, companyID uuid.UUID, uniqueID string) error {
uniqueID = strings.TrimSpace(uniqueID)
if uniqueID == "" {
return ClientMsg("invalid category id")
}
ct, err := s.Pool.Exec(ctx, `
DELETE FROM categories WHERE company_id = $1 AND unique_id = $2`, companyID, uniqueID)
if err != nil {
return err
}
if ct.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
func (s *Service) ListAttributes(ctx context.Context, companyID uuid.UUID, f ListFilter) ([]map[string]any, int64, error) {
f = NormalizeListFilter(f)
args := []any{companyID}
where := []string{"a.company_id = $1"}
from := "attributes a"
selectCols := "a.id, a.attribute_key, a.name, a.value_type, a.unit, a.example, a.parent_key, a.created_at, a.updated_at"
scanCols := []string{"id", "attribute_key", "name", "value_type", "unit", "example", "parent_key", "created_at", "updated_at"}
if f.Category != "" {
from = `attributes a
INNER JOIN category_attributes ca
ON ca.attribute_id = a.id AND ca.company_id = a.company_id`
args = append(args, f.Category)
where = append(where, fmt.Sprintf("ca.category_unique_id = $%d", len(args)))
selectCols += ", ca.required, ca.category_unique_id"
scanCols = append(scanCols, "required", "category_unique_id")
}
if f.RootsOnly {
where = append(where, "a.parent_key IS NULL")
}
if f.ParentKey != "" {
args = append(args, f.ParentKey)
where = append(where, fmt.Sprintf("a.parent_key = $%d", len(args)))
}
if f.Query != "" {
args = append(args, "%"+f.Query+"%")
n := len(args)
where = append(where, fmt.Sprintf("(a.name ILIKE $%d OR a.attribute_key ILIKE $%d OR COALESCE(a.parent_key, '') ILIKE $%d)", n, n, n))
}
wSQL := strings.Join(where, " AND ")
args = append(args, f.Limit, f.Offset)
lim := len(args) - 1
off := len(args)
rows, err := s.Pool.Query(ctx, fmt.Sprintf(`
SELECT %s
FROM %s WHERE %s
ORDER BY a.name
LIMIT $%d OFFSET $%d`, selectCols, from, wSQL, lim, off), args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
items, err := scanMaps(rows, scanCols)
if err != nil {
return nil, 0, err
}
if total, ok := exactTotalFromPage(f.Offset, f.Limit, len(items)); ok {
return items, total, nil
}
countArgs := args[:len(args)-2]
var total int64
if err := s.Pool.QueryRow(ctx, fmt.Sprintf(`SELECT count(*) FROM %s WHERE %s`, from, wSQL), countArgs...).Scan(&total); err != nil {
return nil, 0, err
}
return items, total, nil
}
var attributeValueTypes = map[string]struct{}{
"string": {}, "number": {}, "boolean": {}, "date": {}, "list": {}, "multiselect": {},
}
func (s *Service) CreateAttribute(ctx context.Context, companyID uuid.UUID, key, name, valueType string, unit, example, parent *string) (map[string]any, error) {
if key == "" || name == "" {
return nil, ClientMsg("attribute_key and name required")
}
if valueType == "" {
valueType = "string"
}
if _, ok := attributeValueTypes[valueType]; !ok {
return nil, ClientMsg("value_type must be one of: string, number, boolean, date, list, multiselect")
}
var id uuid.UUID
err := s.Pool.QueryRow(ctx, `
INSERT INTO attributes (company_id, attribute_key, name, value_type, unit, example, parent_key)
VALUES ($1, $2, $3, $4, $5, $6, $7) RETURNING id`,
companyID, key, name, valueType, unit, example, parent).Scan(&id)
if err != nil {
return nil, err
}
return s.getAttribute(ctx, companyID, id)
}
func (s *Service) getAttribute(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
row := s.Pool.QueryRow(ctx, `
SELECT id, attribute_key, name, value_type, unit, example, parent_key, created_at, updated_at
FROM attributes WHERE id = $1 AND company_id = $2`, id, companyID)
return scanMap(row, []string{"id", "attribute_key", "name", "value_type", "unit", "example", "parent_key", "created_at", "updated_at"})
}
func (s *Service) UpdateAttribute(ctx context.Context, companyID, id uuid.UUID, body map[string]any) (map[string]any, error) {
name, _ := body["name"].(string)
valueType, _ := body["value_type"].(string)
valueType = strings.TrimSpace(valueType)
if valueType != "" {
if _, ok := attributeValueTypes[valueType]; !ok {
return nil, ClientMsg("value_type must be one of: string, number, boolean, date, list, multiselect")
}
}
unit := optionalStringPtr(body, "unit")
example := optionalStringPtr(body, "example")
_, err := s.Pool.Exec(ctx, `
UPDATE attributes SET
name = CASE WHEN $3 <> '' THEN $3 ELSE name END,
value_type = CASE WHEN $4 <> '' THEN $4 ELSE value_type END,
unit = CASE WHEN $5::boolean THEN $6 ELSE unit END,
example = CASE WHEN $7::boolean THEN $8 ELSE example END,
updated_at = now()
WHERE id = $1 AND company_id = $2`,
id, companyID, name, valueType,
unit != nil, nullableString(unit),
example != nil, nullableString(example))
if err != nil {
return nil, err
}
return s.getAttribute(ctx, companyID, id)
}
func optionalStringPtr(body map[string]any, key string) *string {
v, ok := body[key]
if !ok {
return nil
}
if v == nil {
empty := ""
return &empty
}
s, ok := v.(string)
if !ok {
return nil
}
return &s
}
func nullableString(p *string) any {
if p == nil {
return nil
}
if *p == "" {
return nil
}
return *p
}
func (s *Service) DeleteAttribute(ctx context.Context, companyID, id uuid.UUID) error {
ct, err := s.Pool.Exec(ctx, `DELETE FROM attributes WHERE id = $1 AND company_id = $2`, id, companyID)
if err != nil {
return err
}
if ct.RowsAffected() == 0 {
return ErrNotFound
}
return nil
}
// ListFilter is the shared SQL pagination/search filter for catalog list APIs
// (categories, attributes, products). Product-only fields may be left empty.
type ListFilter struct {
Query string
Status string
Category string
FeedID string // optional UUID; filters raw/processed products by feed
// Coverage filters processed products by enrichment completeness:
// complete | incomplete | missing_name | missing_description | missing_attributes | missing_category.
Coverage string
// Eprel filters processed products by mapped EPREL id presence: has_eprel | no_eprel.
Eprel string
// SyncChange filters by raw mapped_data._sync_changes from the latest modifying feed sync:
// price | stock | availability | title | other | new | any.
SyncChange string
SortBy string // updatedAt | name (products list)
SortOrder string // asc | desc
Limit int
Offset int
// Cursor is an opaque keyset bookmark (preferred over Offset for deep pages).
Cursor string
// AfterID is a product UUID keyset bookmark; resolved to sort keys server-side.
// When Cursor is also set, Cursor wins. Offset is ignored when either is set.
AfterID string
// RootsOnly limits attributes to top-level definitions (parent_key IS NULL).
RootsOnly bool
// ParentKey limits attributes to children of a list/multiselect attribute.
ParentKey string
}
// ProductFilter is an alias kept for existing call sites.
type ProductFilter = ListFilter
// Product list pagination bounds (categories/attrs still use NormalizeListFilter's higher cap).
const (
MaxProductPageLimit = 200
MaxOffsetWithoutCursor = 5000
)
func NormalizeListFilter(f ListFilter) ListFilter {
if f.Limit <= 0 {
f.Limit = 50
}
if f.Limit > 2000 {
f.Limit = 2000
}
if f.Offset < 0 {
f.Offset = 0
}
f.Query = strings.TrimSpace(f.Query)
f.Status = strings.TrimSpace(f.Status)
f.Category = strings.TrimSpace(f.Category)
f.FeedID = strings.TrimSpace(f.FeedID)
f.Coverage = normalizeCoverageFilter(f.Coverage)
f.Eprel = normalizeEprelFilter(f.Eprel)
f.SyncChange = normalizeSyncChangeFilter(f.SyncChange)
f.ParentKey = strings.TrimSpace(f.ParentKey)
f.Cursor = strings.TrimSpace(f.Cursor)
f.AfterID = strings.TrimSpace(f.AfterID)
f.SortBy = strings.TrimSpace(f.SortBy)
f.SortOrder = strings.ToLower(strings.TrimSpace(f.SortOrder))
switch f.SortBy {
case "name", "updatedAt", "createdAt":
// keep
default:
f.SortBy = "updatedAt"
}
if f.SortOrder != "asc" {
f.SortOrder = "desc"
}
if HasProductCursor(f) {
f.Offset = 0
}
return f
}
// normalizeProductListFilter caps product pages and rejects deep OFFSET without a keyset cursor.
func normalizeProductListFilter(f ListFilter) (ListFilter, error) {
f = NormalizeListFilter(f)
if f.Limit > MaxProductPageLimit {
f.Limit = MaxProductPageLimit
}
if !HasProductCursor(f) && f.Offset > MaxOffsetWithoutCursor {
return f, ClientMsg("offset too large; use cursor or after_id for deep pages")
}
return f, nil
}
func productSortDir(order string) string {
if strings.EqualFold(order, "asc") {
return "ASC"
}
return "DESC"
}
func processedProductsOrderBy(f ListFilter) string {
dir := productSortDir(f.SortOrder)
if f.SortBy == "name" {
return fmt.Sprintf(
`(CASE WHEN COALESCE(NULLIF(p.processed_name, ''), NULLIF(p.name, '')) IS NULL THEN 1 ELSE 0 END),
LOWER(COALESCE(NULLIF(p.processed_name, ''), NULLIF(p.name, ''), p.product_id)) %s, p.id %s`,
dir, dir,
)
}
if f.SortBy == "createdAt" {
return fmt.Sprintf("p.created_at %s, p.id %s", dir, dir)
}
return fmt.Sprintf("p.updated_at %s, p.id %s", dir, dir)
}
func rawProductsOrderBy(f ListFilter) string {
dir := productSortDir(f.SortOrder)
if f.SortBy == "name" {
return fmt.Sprintf(
`(CASE WHEN COALESCE(NULLIF(rp.mapped_data->>'name', ''), NULLIF(rp.mapped_data->>'title', '')) IS NULL THEN 1 ELSE 0 END),
LOWER(COALESCE(NULLIF(rp.mapped_data->>'name', ''), NULLIF(rp.mapped_data->>'title', ''), rp.gtin)) %s, rp.id %s`,
dir, dir,
)
}
if f.SortBy == "updatedAt" {
return fmt.Sprintf("rp.updated_at %s, rp.id %s", dir, dir)
}
return fmt.Sprintf("rp.created_at %s, rp.id %s", dir, dir)
}
func normalizeProductFilter(f ProductFilter) (ProductFilter, error) {
return normalizeProductListFilter(f)
}
func rawProductsCountFromSQL(needsFeedJoin bool) string {
if needsFeedJoin {
return `
FROM raw_products rp
LEFT JOIN input_feeds f ON f.id = rp.feed_id`
}
return `
FROM raw_products rp`
}
func processedProductsCountFromSQL(needsRawJoin bool) string {
if needsRawJoin {
return `
FROM processed_products p
LEFT JOIN raw_products r ON r.id = p.raw_product_id`
}
return `
FROM processed_products p`
}
// Enrichment coverage SQL predicates (alias p = processed_products, r = raw_products).
const (
processedHasNameSQL = `(COALESCE(NULLIF(p.processed_name, ''), NULLIF(p.name, ''), NULLIF(r.mapped_data->>'name', ''), NULLIF(r.mapped_data->>'title', ''), '') <> '')`
processedHasDescriptionSQL = `(COALESCE(NULLIF(p.processed_description, ''), NULLIF(p.description, ''), NULLIF(r.mapped_data->>'description', ''), '') <> '')`
processedHasCategorySQL = `(COALESCE(NULLIF(p.category, ''), '') <> '' AND lower(p.category) <> 'none')`
// Resolve display name / canonical unique_id when products store unique_id, UUID id, or name.
processedCategoryResolveJoin = `
LEFT JOIN LATERAL (
SELECT c.name, c.unique_id
FROM categories c
WHERE c.company_id = p.company_id
AND NULLIF(BTRIM(p.category), '') IS NOT NULL
AND lower(BTRIM(p.category)) <> 'none'
AND (
c.unique_id = BTRIM(p.category)
OR c.id::text = BTRIM(p.category)
OR lower(c.name) = lower(BTRIM(p.category))
)
ORDER BY
CASE
WHEN c.unique_id = BTRIM(p.category) THEN 0
WHEN c.id::text = BTRIM(p.category) THEN 1
ELSE 2
END
LIMIT 1
) cat ON true`
// Raw inventory (A1 clean / unprocessed tab): category lives on mapped_data only.
rawCategoryResolveJoin = `
LEFT JOIN LATERAL (
SELECT c.name, c.unique_id
FROM categories c
WHERE c.company_id = rp.company_id
AND NULLIF(BTRIM(rp.mapped_data->>'category'), '') IS NOT NULL
AND lower(BTRIM(rp.mapped_data->>'category')) <> 'none'
AND (
c.unique_id = BTRIM(rp.mapped_data->>'category')
OR c.id::text = BTRIM(rp.mapped_data->>'category')
OR lower(c.name) = lower(BTRIM(rp.mapped_data->>'category'))
)
ORDER BY
CASE
WHEN c.unique_id = BTRIM(rp.mapped_data->>'category') THEN 0
WHEN c.id::text = BTRIM(rp.mapped_data->>'category') THEN 1
ELSE 2
END
LIMIT 1
) cat ON true`
)
// Attribute presence: AI bag, original bag, or mapped feed specs/dimensions/eprel/warranty.
// (Assembled as vars so we can OR the pieces without repeating the EXISTS body.)
var (
processedHasProcessedAttributesSQL = `EXISTS (
SELECT 1
FROM jsonb_each(COALESCE(p.processed_attributes, '{}'::jsonb)) AS kv(key, value)
WHERE (jsonb_typeof(kv.value) = 'string' AND length(trim(both '"' from kv.value::text)) > 0)
OR jsonb_typeof(kv.value) IN ('number', 'boolean')
OR (jsonb_typeof(kv.value) = 'object' AND COALESCE(NULLIF(kv.value->>'name', ''), NULLIF(kv.value->>'value', ''), '') <> '')
OR (jsonb_typeof(kv.value) = 'array' AND jsonb_array_length(kv.value) > 0)
)`
processedHasOriginalAttributesSQL = `EXISTS (
SELECT 1
FROM jsonb_each(COALESCE(p.attributes, '{}'::jsonb)) AS kv(key, value)
WHERE (jsonb_typeof(kv.value) = 'string' AND length(trim(both '"' from kv.value::text)) > 0)
OR jsonb_typeof(kv.value) IN ('number', 'boolean')
OR (jsonb_typeof(kv.value) = 'object' AND COALESCE(NULLIF(kv.value->>'name', ''), NULLIF(kv.value->>'value', ''), '') <> '')
OR (jsonb_typeof(kv.value) = 'array' AND jsonb_array_length(kv.value) > 0)
)`
processedHasFeedAttributesSQL = `(
CASE jsonb_typeof(r.mapped_data->'specifications')
WHEN 'string' THEN length(trim(r.mapped_data->>'specifications')) > 0
WHEN 'object' THEN r.mapped_data->'specifications' <> '{}'::jsonb
WHEN 'array' THEN jsonb_array_length(r.mapped_data->'specifications') > 0
ELSE false
END
OR CASE jsonb_typeof(r.mapped_data->'specs')
WHEN 'string' THEN length(trim(r.mapped_data->>'specs')) > 0
WHEN 'object' THEN r.mapped_data->'specs' <> '{}'::jsonb
WHEN 'array' THEN jsonb_array_length(r.mapped_data->'specs') > 0
ELSE false
END
OR COALESCE(NULLIF(trim(r.mapped_data->>'eprel_id'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'eprel'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'netwidth'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'net_width'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'netheight'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'net_height'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'netdepth'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'net_depth'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'netmass'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'net_mass'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'warranty'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'productmodel'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'product_model'), ''), '') <> ''
)`
processedHasAttributesSQL = `(` + processedHasProcessedAttributesSQL + ` OR ` + processedHasOriginalAttributesSQL + ` OR ` + processedHasFeedAttributesSQL + `)`
processedHasEprelSQL = `(
COALESCE(NULLIF(trim(r.mapped_data->>'eprel_id'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'eprel'), ''), '') <> ''
OR COALESCE(NULLIF(trim(r.mapped_data->>'EPRELID'), ''), '') <> ''
)`
)
func normalizeCoverageFilter(raw string) string {
c := strings.ToLower(strings.TrimSpace(raw))
c = strings.ReplaceAll(c, "-", "_")
switch c {
case "", "all", "any":
return ""
case "complete", "full", "ok":
return "complete"
case "incomplete", "partial":
return "incomplete"
case "missing_name", "name":
return "missing_name"
case "missing_description", "description":
return "missing_description"
case "missing_attributes", "attributes", "attrs":
return "missing_attributes"
case "missing_category", "category":
return "missing_category"
default:
return ""
}
}
func normalizeEprelFilter(raw string) string {
c := strings.ToLower(strings.TrimSpace(raw))
c = strings.ReplaceAll(c, "-", "_")
switch c {
case "", "all", "any":
return ""
case "has_eprel", "eprel", "with_eprel", "yes", "true", "1":
return "has_eprel"
case "no_eprel", "without_eprel", "missing_eprel", "none", "no", "false", "0":
return "no_eprel"
default:
return ""
}
}
func normalizeSyncChangeFilter(raw string) string {
c := strings.ToLower(strings.TrimSpace(raw))
c = strings.ReplaceAll(c, "-", "_")
switch c {
case "", "all":
return ""
case "any", "changed", "has_change", "has_changes":
return "any"
case "price", "price_changed":
return "price"
case "stock", "stock_changed", "qty", "quantity":
return "stock"
case "availability", "availability_changed", "stock_status":
return "availability"
case "title", "name", "title_changed":
return "title"
case "other", "other_changed":
return "other"
case "new", "inserted":
return "new"
default:
return ""
}
}
func appendSyncChangeFilter(alias, syncChange string, where []string) []string {
col := alias + `.mapped_data->'_sync_changes'`
switch syncChange {
case "any":
return append(where, `jsonb_typeof(`+col+`) = 'array' AND jsonb_array_length(`+col+`) > 0`)
case "price", "stock", "availability", "title", "other", "new":
return append(where, col+` ? '`+syncChange+`'`)
default:
return where
}
}
func appendProcessedCoverageFilter(coverage string, where []string) []string {
switch coverage {
case "complete":
return append(where, processedHasNameSQL+" AND "+processedHasDescriptionSQL+" AND "+processedHasCategorySQL+" AND "+processedHasAttributesSQL)
case "incomplete":
return append(where, "NOT ("+processedHasNameSQL+" AND "+processedHasDescriptionSQL+" AND "+processedHasCategorySQL+" AND "+processedHasAttributesSQL+")")
case "missing_name":
return append(where, "NOT "+processedHasNameSQL)
case "missing_description":
return append(where, "NOT "+processedHasDescriptionSQL)
case "missing_attributes":
return append(where, "NOT "+processedHasAttributesSQL)
case "missing_category":
return append(where, "NOT "+processedHasCategorySQL)
default:
return where
}
}
func appendProcessedEprelFilter(eprel string, where []string) []string {
switch eprel {
case "has_eprel":
return append(where, processedHasEprelSQL)
case "no_eprel":
return append(where, "NOT "+processedHasEprelSQL)
default:
return where
}
}
func processedListNeedsRawJoin(f ListFilter) bool {
return f.Query != "" || f.Coverage != "" || f.Eprel != "" || f.SyncChange != ""
}
func appendRawProductFilters(f ListFilter, args []any, where []string) ([]any, []string) {
// Search keyed JSON paths + gtin/feed only. Avoid CAST(jsonb AS text) ILIKE:
// it forces full-document scans and cannot use btree/trigram expression indexes usefully.
// Defer pg_trgm/GIN until leading-wildcard ILIKE is measured hot after this shape.
if f.Query != "" {
args = append(args, "%"+f.Query+"%")
n := len(args)
where = append(where, fmt.Sprintf(`(
rp.gtin ILIKE $%d
OR COALESCE(rp.mapped_data->>'name', '') ILIKE $%d
OR COALESCE(rp.mapped_data->>'title', '') ILIKE $%d
OR COALESCE(f.name, '') ILIKE $%d
)`, n, n, n, n))
}
if f.Status != "" {
args = append(args, f.Status)
where = append(where, fmt.Sprintf("rp.processing_status = $%d", len(args)))
}
if feedID, err := uuid.Parse(strings.TrimSpace(f.FeedID)); err == nil {
args = append(args, feedID)
where = append(where, fmt.Sprintf("rp.feed_id = $%d", len(args)))
}
where = appendSyncChangeFilter("rp", f.SyncChange, where)
return args, where
}
func (s *Service) ListRawProducts(ctx context.Context, companyID uuid.UUID, f ListFilter) ([]map[string]any, int64, error) {
f, err := normalizeProductListFilter(f)
if err != nil {
return nil, 0, err
}
args := []any{companyID}
where := []string{"rp.company_id = $1"}
args, where = appendRawProductFilters(f, args, where)
countArgs := append([]any{}, args...)
countSQL := strings.Join(where, " AND ")
cur, useCursor, missing, err := s.resolveRawCursor(ctx, companyID, f)
if err != nil {
return nil, 0, err
}
if missing {
where = append(where, "FALSE")
} else if useCursor {
args, where, err = appendRawKeyset(f, cur, args, where)
if err != nil {
return nil, 0, ClientMsg("invalid cursor")
}
}
wSQL := strings.Join(where, " AND ")
listFromSQL := `
FROM raw_products rp
LEFT JOIN input_feeds f ON f.id = rp.feed_id`
countFromSQL := rawProductsCountFromSQL(f.Query != "")
orderBy := rawProductsOrderBy(f)
listArgs := append(append([]any{}, args...), f.Limit, f.Offset)
lim := len(listArgs) - 1
off := len(listArgs)
return parallelCountAndList(ctx,
f.Limit, f.Offset,
func(ctx context.Context) (int64, error) {
var total int64
err := s.Pool.QueryRow(ctx, `SELECT count(*) `+countFromSQL+` WHERE `+countSQL, countArgs...).Scan(&total)
return total, err
},
func(ctx context.Context) ([]map[string]any, error) {
rows, err := s.Pool.Query(ctx, fmt.Sprintf(`
SELECT rp.id, rp.gtin, rp.feed_id, rp.is_processed, rp.processing_status,
COALESCE(NULLIF(rp.mapped_data->>'name', ''), NULLIF(rp.mapped_data->>'title', ''), '') AS name,
NULLIF(BTRIM(rp.mapped_data->>'category'), '') AS category,
COALESCE(cat.name, NULLIF(BTRIM(rp.mapped_data->>'category'), '')) AS category_name,
COALESCE(cat.unique_id, NULLIF(BTRIM(rp.mapped_data->>'category'), '')) AS category_unique_id,
f.name AS feed_name,
rp.mapped_data->'_sync_changes' AS sync_changes,
(COALESCE(NULLIF(trim(rp.mapped_data->>'name'), ''), NULLIF(trim(rp.mapped_data->>'title'), ''), '') <> '') AS has_name,
(COALESCE(NULLIF(trim(rp.mapped_data->>'description'), ''), '') <> '') AS has_description,
(COALESCE(NULLIF(trim(rp.mapped_data->>'category'), ''), '') <> '' AND lower(trim(rp.mapped_data->>'category')) <> 'none') AS has_category,
`+strings.ReplaceAll(processedHasFeedAttributesSQL, "r.mapped_data", "rp.mapped_data")+` AS has_attributes,
rp.created_at, rp.updated_at
%s%s WHERE %s
ORDER BY %s
LIMIT $%d OFFSET $%d`, listFromSQL, rawCategoryResolveJoin, wSQL, orderBy, lim, off), listArgs...)
if err != nil {
return nil, err
}
defer rows.Close()
return scanMaps(rows, []string{
"id", "gtin", "feed_id", "is_processed", "processing_status",
"name", "category", "category_name", "category_unique_id",
"feed_name", "sync_changes",
"has_name", "has_description", "has_category", "has_attributes",
"created_at", "updated_at",
})
},
)
}
// ListRawProductIDsByFeed returns up to limit raw product UUIDs for a company-scoped feed.
func (s *Service) ListRawProductIDsByFeed(ctx context.Context, companyID, feedID uuid.UUID, limit int) ([]uuid.UUID, error) {
if limit <= 0 {
limit = 10
}
if limit > 100 {
limit = 100
}
rows, err := s.Pool.Query(ctx, `
SELECT id FROM raw_products
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()
ids := make([]uuid.UUID, 0, limit)
for rows.Next() {
var id uuid.UUID
if err := rows.Scan(&id); err != nil {
return nil, err
}
ids = append(ids, id)
}
return ids, rows.Err()
}
func appendProcessedProductFilters(f ProductFilter, args []any, where []string) ([]any, []string) {
if f.Query != "" {
args = append(args, "%"+f.Query+"%")
n := len(args)
where = append(where, fmt.Sprintf(
`(p.name ILIKE $%d OR COALESCE(p.processed_name, '') ILIKE $%d OR p.product_id ILIKE $%d OR p.category ILIKE $%d OR COALESCE(r.gtin, '') ILIKE $%d)`,
n, n, n, n, n))
}
if f.Status != "" {
// needs_review includes legacy pipeline status "processed" (pre-P0-8).
if f.Status == "needs_review" {
where = append(where, "p.status IN ('needs_review', 'processed')")
} else {
args = append(args, f.Status)
where = append(where, fmt.Sprintf("p.status = $%d", len(args)))
}
}
if f.Category != "" {
args = append(args, f.Category)
n := len(args)
// Match stored unique_id, UUID id, or display name for the selected category.
where = append(where, fmt.Sprintf(`(
p.category = $%d
OR EXISTS (
SELECT 1 FROM categories c
WHERE c.company_id = p.company_id
AND (
c.unique_id = $%d
OR c.id::text = $%d
OR lower(c.name) = lower($%d)
)
AND (
p.category = c.unique_id
OR p.category = c.id::text
OR lower(p.category) = lower(c.name)
)
)
)`, n, n, n, n))
}
if feedID, err := uuid.Parse(f.FeedID); err == nil {
args = append(args, feedID)
where = append(where, fmt.Sprintf("p.feed_id = $%d", len(args)))
}
where = appendProcessedCoverageFilter(f.Coverage, where)
where = appendProcessedEprelFilter(f.Eprel, where)
where = appendSyncChangeFilter("r", f.SyncChange, where)
return args, where
}
// ListProcessedProducts returns a lean page without heavy JSONB columns
// (attributes, descriptions, mapped_data). Prefer this for UI tables;
// use ListProcessedProductsDetailed when quality scoring or full attrs are needed.
func (s *Service) ListProcessedProducts(ctx context.Context, companyID uuid.UUID, f ProductFilter) ([]map[string]any, int64, error) {
f, err := normalizeProductFilter(f)
if err != nil {
return nil, 0, err
}
args := []any{companyID}
where := []string{"p.company_id = $1"}
args, where = appendProcessedProductFilters(f, args, where)
countArgs := append([]any{}, args...)
countSQL := strings.Join(where, " AND ")
cur, useCursor, missing, err := s.resolveProcessedCursor(ctx, companyID, f)
if err != nil {
return nil, 0, err
}
if missing {
where = append(where, "FALSE")
} else if useCursor {
args, where, err = appendProcessedKeyset(f, cur, args, where)
if err != nil {
return nil, 0, ClientMsg("invalid cursor")
}
}
wSQL := strings.Join(where, " AND ")
orderBy := processedProductsOrderBy(f)
listArgs := append(append([]any{}, args...), f.Limit, f.Offset)
lim := len(listArgs) - 1
off := len(listArgs)
return parallelCountAndList(ctx,
f.Limit, f.Offset,
func(ctx context.Context) (int64, error) {
var total int64
err := s.Pool.QueryRow(ctx, `
SELECT count(*) `+processedProductsCountFromSQL(processedListNeedsRawJoin(f))+`
WHERE `+countSQL, countArgs...).Scan(&total)
return total, err
},
func(ctx context.Context) ([]map[string]any, error) {
rows, err := s.Pool.Query(ctx, fmt.Sprintf(`
SELECT p.id, p.product_id,
COALESCE(NULLIF(p.name, ''), NULLIF(r.mapped_data->>'name', ''), NULLIF(r.mapped_data->>'title', ''), '') AS name,
COALESCE(NULLIF(p.processed_name, ''), '') AS processed_name,
p.category,
COALESCE(cat.name, NULLIF(BTRIM(p.category), '')) AS category_name,
COALESCE(cat.unique_id, NULLIF(BTRIM(p.category), '')) AS category_unique_id,
p.status, p.raw_product_id, COALESCE(p.feed_id, r.feed_id) AS feed_id, r.gtin,
f.name AS feed_name,
f.last_synced_at AS feed_last_synced_at,
r.updated_at AS raw_updated_at,
(COALESCE(NULLIF(p.processed_name, ''), NULLIF(p.name, ''), NULLIF(r.mapped_data->>'name', ''), NULLIF(r.mapped_data->>'title', ''), '') <> '') AS has_name,
(COALESCE(NULLIF(p.processed_name, ''), '') <> '') AS has_processed_name,
(COALESCE(NULLIF(p.processed_description, ''), NULLIF(p.description, ''), NULLIF(r.mapped_data->>'description', ''), '') <> '') AS has_description,
(COALESCE(NULLIF(p.processed_description, ''), '') <> '') AS has_processed_description,
(COALESCE(NULLIF(p.category, ''), '') <> '' AND lower(p.category) <> 'none') AS has_category,
`+processedHasAttributesSQL+` AS has_attributes,
`+processedHasProcessedAttributesSQL+` AS has_processed_attributes,
`+processedHasEprelSQL+` AS has_eprel,
p.created_at, p.updated_at
FROM processed_products p
LEFT JOIN raw_products r ON r.id = p.raw_product_id
LEFT JOIN input_feeds f ON f.id = COALESCE(p.feed_id, r.feed_id)`+processedCategoryResolveJoin+`
WHERE %s
ORDER BY %s LIMIT $%d OFFSET $%d`, wSQL, orderBy, lim, off), listArgs...)
if err != nil {
return nil, err
}
defer rows.Close()
return scanMaps(rows, []string{
"id", "product_id", "name", "processed_name", "category", "category_name", "category_unique_id",
"status", "raw_product_id", "feed_id", "gtin",
"feed_name", "feed_last_synced_at", "raw_updated_at",
"has_name", "has_processed_name", "has_description", "has_processed_description", "has_category", "has_attributes", "has_processed_attributes",
"has_eprel",
"created_at", "updated_at",
})
},
)
}
// ListProcessedProductsDetailed includes fields needed for quality scoring.
func (s *Service) ListProcessedProductsDetailed(ctx context.Context, companyID uuid.UUID, f ProductFilter) ([]map[string]any, int64, error) {
f, err := normalizeProductFilter(f)
if err != nil {
return nil, 0, err
}
args := []any{companyID}
where := []string{"p.company_id = $1"}
args, where = appendProcessedProductFilters(f, args, where)
countArgs := append([]any{}, args...)
countSQL := strings.Join(where, " AND ")
cur, useCursor, missing, err := s.resolveProcessedCursor(ctx, companyID, f)
if err != nil {
return nil, 0, err
}
if missing {
where = append(where, "FALSE")
} else if useCursor {
args, where, err = appendProcessedKeyset(f, cur, args, where)
if err != nil {
return nil, 0, ClientMsg("invalid cursor")
}
}
wSQL := strings.Join(where, " AND ")
orderBy := processedProductsOrderBy(f)
listArgs := append(append([]any{}, args...), f.Limit, f.Offset)
lim := len(listArgs) - 1
off := len(listArgs)
return parallelCountAndList(ctx,
f.Limit, f.Offset,
func(ctx context.Context) (int64, error) {
var total int64
err := s.Pool.QueryRow(ctx, `
SELECT count(*) `+processedProductsCountFromSQL(processedListNeedsRawJoin(f))+`
WHERE `+countSQL, countArgs...).Scan(&total)
return total, err
},
func(ctx context.Context) ([]map[string]any, error) {
rows, err := s.Pool.Query(ctx, fmt.Sprintf(`
SELECT p.id, p.product_id,
COALESCE(NULLIF(p.name, ''), NULLIF(r.mapped_data->>'name', ''), NULLIF(r.mapped_data->>'title', ''), '') AS name,
p.category,
COALESCE(cat.name, NULLIF(BTRIM(p.category), '')) AS category_name,
COALESCE(cat.unique_id, NULLIF(BTRIM(p.category), '')) AS category_unique_id,
p.status, p.raw_product_id, p.feed_id, r.gtin,
COALESCE(NULLIF(p.description, ''), NULLIF(r.mapped_data->>'description', ''), '') AS description,
p.processed_name, p.processed_description,
COALESCE(p.meta_title, ''), COALESCE(p.meta_description, ''),
p.attributes, p.processed_attributes, r.mapped_data,
p.created_at, p.updated_at
FROM processed_products p
LEFT JOIN raw_products r ON r.id = p.raw_product_id`+processedCategoryResolveJoin+`
WHERE %s
ORDER BY %s LIMIT $%d OFFSET $%d`, wSQL, orderBy, lim, off), listArgs...)
if err != nil {
return nil, err
}
defer rows.Close()
return scanMaps(rows, []string{
"id", "product_id", "name", "category", "category_name", "category_unique_id",
"status", "raw_product_id", "feed_id", "gtin",
"description", "processed_name", "processed_description",
"meta_title", "meta_description",
"attributes", "processed_attributes", "mapped_data",
"created_at", "updated_at",
})
},
)
}
func (s *Service) GetProcessedProduct(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
row := s.Pool.QueryRow(ctx, `
SELECT p.id, p.product_id,
COALESCE(NULLIF(p.name, ''), NULLIF(r.mapped_data->>'name', ''), NULLIF(r.mapped_data->>'title', ''), '') AS name,
p.category,
COALESCE(cat.name, NULLIF(BTRIM(p.category), '')) AS category_name,
COALESCE(cat.unique_id, NULLIF(BTRIM(p.category), '')) AS category_unique_id,
COALESCE(NULLIF(p.description, ''), NULLIF(r.mapped_data->>'description', ''), '') AS description,
p.processed_name, p.processed_description,
p.status, p.attributes, p.processed_attributes, r.gtin, r.mapped_data,
COALESCE(p.feed_id, r.feed_id) AS feed_id,
f.name AS feed_name,
f.last_synced_at AS feed_last_synced_at,
r.updated_at AS raw_updated_at,
(COALESCE(NULLIF(p.processed_name, ''), '') <> '') AS has_processed_name,
(COALESCE(NULLIF(p.processed_description, ''), '') <> '') AS has_processed_description,
`+processedHasAttributesSQL+` AS has_attributes,
`+processedHasProcessedAttributesSQL+` AS has_processed_attributes,
`+processedHasEprelSQL+` AS has_eprel,
COALESCE(p.localized_content, '{}'::jsonb) AS localized_content,
p.created_at, p.updated_at
FROM processed_products p
LEFT JOIN raw_products r ON r.id = p.raw_product_id
LEFT JOIN input_feeds f ON f.id = COALESCE(p.feed_id, r.feed_id)`+processedCategoryResolveJoin+`
WHERE p.id = $1 AND p.company_id = $2`, id, companyID)
item, err := scanMap(row, []string{
"id", "product_id", "name", "category", "category_name", "category_unique_id",
"description", "processed_name", "processed_description",
"status", "attributes", "processed_attributes", "gtin", "mapped_data",
"feed_id", "feed_name", "feed_last_synced_at", "raw_updated_at",
"has_processed_name", "has_processed_description", "has_attributes", "has_processed_attributes",
"has_eprel",
"localized_content",
"created_at", "updated_at",
})
if err != nil {
return nil, err
}
linkFeedSpecificationsIntoProduct(item)
primary := company.LoadLanguage(ctx, s.Pool, companyID)
item["content_language"] = primary
item["content_languages"] = company.LoadContentLanguages(ctx, s.Pool, companyID)
if loc, err := company.DecodeLocalizedContent(item["localized_content"]); err == nil {
item["localized_content"] = loc
}
return item, nil
}
// GetRawProduct returns a raw inventory row with feed-origin mapped_data so the
// product panel can show original description, category, and attributes when
// processed_products is empty (A1 demo seed keeps processed=0).
func (s *Service) GetRawProduct(ctx context.Context, companyID, id uuid.UUID) (map[string]any, error) {
row := s.Pool.QueryRow(ctx, `
SELECT rp.id,
COALESCE(NULLIF(rp.gtin, ''), '') AS product_id,
COALESCE(NULLIF(rp.mapped_data->>'name', ''), NULLIF(rp.mapped_data->>'title', ''), '') AS name,
NULLIF(BTRIM(rp.mapped_data->>'category'), '') AS category,
COALESCE(cat.name, NULLIF(BTRIM(rp.mapped_data->>'category'), '')) AS category_name,
COALESCE(cat.unique_id, NULLIF(BTRIM(rp.mapped_data->>'category'), '')) AS category_unique_id,
COALESCE(NULLIF(rp.mapped_data->>'description', ''), '') AS description,
''::text AS processed_name,
''::text AS processed_description,
COALESCE(NULLIF(rp.processing_status, ''), 'unprocessed') AS status,
'{}'::jsonb AS attributes,
'{}'::jsonb AS processed_attributes,
rp.gtin,
rp.mapped_data,
rp.feed_id,
f.name AS feed_name,
f.last_synced_at AS feed_last_synced_at,
rp.updated_at AS raw_updated_at,
false AS has_processed_name,
false AS has_processed_description,
`+strings.ReplaceAll(processedHasFeedAttributesSQL, "r.mapped_data", "rp.mapped_data")+` AS has_attributes,
false AS has_processed_attributes,
`+strings.ReplaceAll(processedHasEprelSQL, "r.mapped_data", "rp.mapped_data")+` AS has_eprel,
'{}'::jsonb AS localized_content,
rp.created_at, rp.updated_at,
(COALESCE(NULLIF(trim(rp.mapped_data->>'name'), ''), NULLIF(trim(rp.mapped_data->>'title'), ''), '') <> '') AS has_name,
(COALESCE(NULLIF(trim(rp.mapped_data->>'description'), ''), '') <> '') AS has_description,
(COALESCE(NULLIF(trim(rp.mapped_data->>'category'), ''), '') <> '' AND lower(trim(rp.mapped_data->>'category')) <> 'none') AS has_category,
rp.is_processed, rp.processing_status
FROM raw_products rp
LEFT JOIN input_feeds f ON f.id = rp.feed_id`+rawCategoryResolveJoin+`
WHERE rp.id = $1 AND rp.company_id = $2`, id, companyID)
item, err := scanMap(row, []string{
"id", "product_id", "name", "category", "category_name", "category_unique_id",
"description", "processed_name", "processed_description",
"status", "attributes", "processed_attributes", "gtin", "mapped_data",
"feed_id", "feed_name", "feed_last_synced_at", "raw_updated_at",
"has_processed_name", "has_processed_description", "has_attributes", "has_processed_attributes",
"has_eprel",
"localized_content",
"created_at", "updated_at",
"has_name", "has_description", "has_category",
"is_processed", "processing_status",
})
if err != nil {
return nil, err
}
linkFeedSpecificationsIntoProduct(item)
primary := company.LoadLanguage(ctx, s.Pool, companyID)
item["content_language"] = primary
item["content_languages"] = company.LoadContentLanguages(ctx, s.Pool, companyID)
if loc, err := company.DecodeLocalizedContent(item["localized_content"]); err == nil {
item["localized_content"] = loc
}
return item, nil
}
func (s *Service) UpdateProcessedProduct(ctx context.Context, companyID, id uuid.UUID, body map[string]any) (map[string]any, error) {
name, _ := body["name"].(string)
desc, _ := body["description"].(string)
status, _ := body["status"].(string)
category, _ := body["category"].(string)
productID, _ := body["product_id"].(string)
processedName, _ := body["processed_name"].(string)
processedDesc, _ := body["processed_description"].(string)
langRaw, _ := body["language"].(string)
var attrsJSON *string
if v, ok := body["attributes"]; ok && v != nil {
b, err := json.Marshal(v)
if err != nil {
return nil, ClientMsg("invalid attributes")
}
s := string(b)
attrsJSON = &s
}
primary := company.LoadLanguage(ctx, s.Pool, companyID)
lang := primary
if strings.TrimSpace(langRaw) != "" {
parsed, err := company.ParseLanguage(langRaw, false)
if err != nil {
return nil, ClientMsg("unsupported language")
}
lang = parsed
}
// Load existing localized_content and merge this language's fields.
var existingRaw []byte
_ = s.Pool.QueryRow(ctx, `
SELECT COALESCE(localized_content, '{}'::jsonb)
FROM processed_products WHERE id = $1 AND company_id = $2`, id, companyID).Scan(&existingRaw)
localized, _ := company.DecodeLocalizedContent(existingRaw)
fields := company.FieldsForLanguage(localized, lang)
if _, ok := body["processed_name"]; ok {
fields.ProcessedName = processedName
}
if _, ok := body["processed_description"]; ok {
fields.ProcessedDescription = processedDesc
}
if mt, ok := body["meta_title"].(string); ok {
fields.MetaTitle = mt
}
if md, ok := body["meta_description"].(string); ok {
fields.MetaDescription = md
}
localized = company.SetFieldsForLanguage(localized, lang, fields)
locJSON, err := company.EncodeLocalizedContent(localized)
if err != nil {
return nil, err
}
// Denormalized columns always reflect primary language.
primaryFields := company.FieldsForLanguage(localized, primary)
denormName := primaryFields.ProcessedName
denormDesc := primaryFields.ProcessedDescription
if lang == primary {
if _, ok := body["processed_name"]; ok {
denormName = processedName
}
if _, ok := body["processed_description"]; ok {
denormDesc = processedDesc
}
}
ct, err := s.Pool.Exec(ctx, `
UPDATE processed_products SET
name = CASE WHEN $3 <> '' THEN $3 ELSE name END,
description = CASE WHEN $4 <> '' THEN $4 ELSE description END,
status = CASE WHEN $5 <> '' THEN $5 ELSE status END,
category = CASE WHEN $14::boolean THEN $6 ELSE category END,
product_id = CASE WHEN $7 <> '' THEN $7 ELSE product_id END,
processed_name = CASE WHEN $11::boolean THEN $8 ELSE processed_name END,
processed_description = CASE WHEN $12::boolean THEN $9 ELSE processed_description END,
attributes = CASE WHEN $10::jsonb IS NOT NULL THEN $10::jsonb ELSE attributes END,
localized_content = $13::jsonb,
updated_at = now()
WHERE id = $1 AND company_id = $2`,
id, companyID, name, desc, status, category, productID, denormName, denormDesc, attrsJSON,
lang == primary && hasKey(body, "processed_name"),
lang == primary && hasKey(body, "processed_description"),
string(locJSON),
hasKey(body, "category"))
if err != nil {
return nil, err
}
if ct.RowsAffected() == 0 {
return nil, ErrNotFound
}
return s.GetProcessedProduct(ctx, companyID, id)
}
func hasKey(m map[string]any, key string) bool {
_, ok := m[key]
return ok
}
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
}
}
var _ = fmt.Sprintf