diff --git a/main.go b/main.go index bd587c4..0968dd5 100644 --- a/main.go +++ b/main.go @@ -7,6 +7,7 @@ import ( "nearle/facade" "nearle/messaging" "nearle/models" + "nearle/repositories" "nearle/routes" "os" "os/signal" @@ -62,6 +63,21 @@ func main() { log.Fatal("staff shift schema migration failed:", err) } + // The rest of a product's photos. + // + // An explicit ALTER rather than AutoMigrate on `models.Products`: that model + // has drifted from the live table over time, and letting GORM reconcile the + // whole thing to add one column would rewrite far more than anyone asked + // for. `IF NOT EXISTS` makes it a no-op on every boot after the first. + // + // Not fatal on failure. A missing column costs the extra images and nothing + // else — `productimage` still carries the first — and refusing to start the + // API over a gallery would be the worse trade. + if err := db.DB.Exec( + `ALTER TABLE products ADD COLUMN IF NOT EXISTS productimages jsonb`).Error; err != nil { + log.Println("⚠️ could not add products.productimages, extra photos will not be stored:", err) + } + f := facade.NewFacade(db.DB, db.CatalogueDB) routes.RegisterRoutes(app, f) @@ -72,10 +88,19 @@ func main() { // // A broker that is configured but unreachable is fatal on purpose: coming // up healthy while every till quietly queues is the worse failure. + // The consumer is also the publisher for `nearle/pos/{loc}/catalogue`. + // Registering it here inverts the dependency: the repositories that move + // stock cannot import `messaging` — wiring runs this way and the reverse + // would be a cycle — so they hold an interface and this supplies it. Left + // unset the notification is simply skipped, which is what happens when the + // broker is unavailable and must not stop the API booting. posMqtt, err := messaging.StartPosMqttConsumer(f.PosService()) if err != nil { log.Fatal("POS MQTT consumer failed to start:", err) } + if posMqtt != nil { + repositories.SetCatalogueNotifier(posMqtt) + } // Start server go func() { diff --git a/models/product.go b/models/product.go index 7cf4f8a..9021cd4 100644 --- a/models/product.go +++ b/models/product.go @@ -62,6 +62,23 @@ type Products struct { Pricingid int `json:"pricingid,omitempty"` Productname string `json:"productname,omitempty"` Productimage string `json:"productimage,omitempty"` + + // Every photo the product has, as a JSON array of URLs. + // + // `Productimage` above stays the first of these and is left untouched: every + // existing reader — the store catalogue, the customer app, the POS catalogue + // pull, each order line — reads that column, and repointing them all at an + // array is a far larger change than giving the extra photos somewhere to live. + // + // Before this the import kept `Images[0]` and discarded the rest, so a product + // with ten photos in the global catalogue arrived in a shop with one. 90 of + // nestle's 123 products have more than one. + // + // Held as a string rather than a []string because GORM's raw scan-into-struct + // silently drops slice-kind destination fields — the same reason + // `catalogueProductColumns` casts its text[] columns to text. + Productimages string `json:"productimages,omitempty" gorm:"column:productimages;type:jsonb"` + Productdesc string `json:"productdesc,omitempty"` Productsku string `json:"productsku,omitempty"` Brandid int `json:"brandid,omitempty"` diff --git a/repositories/catalogueNotify.go b/repositories/catalogueNotify.go new file mode 100644 index 0000000..2ba0832 --- /dev/null +++ b/repositories/catalogueNotify.go @@ -0,0 +1,85 @@ +package repositories + +import ( + "log" + "strconv" + "sync" + "time" +) + +// Telling tills that a shop's shelf has moved. +// +// The broker topic `nearle/pos/{loc}/catalogue` and its publisher have existed +// since the POS integration landed, retained flag and all — and **nothing ever +// called it**. Every terminal's stock figure was therefore whatever it last +// pulled, with no signal that an app order had just sold the last of something. +// +// It could not be called from here, which is why it never was: `repositories` +// does not import `messaging`, and must not — wiring runs the other way, from +// main.go, and reversing it would be an import cycle. So the dependency is +// inverted through this interface. `messaging.PosMqttConsumer` already +// satisfies it; main.go registers it once at startup. +// +// A package-level registration rather than a constructor parameter because the +// two write paths that need it — POS bill ingest and createOrderTx — are +// reached through several layers that would each have to grow a field for a +// notification neither of them cares about the result of. + +// CatalogueNotifier tells every till at a store to pull the catalogue again. +type CatalogueNotifier interface { + PublishCatalogueChanged(storeID, revision string) error +} + +var ( + catalogueNotifierMu sync.RWMutex + catalogueNotifier CatalogueNotifier +) + +// SetCatalogueNotifier registers the publisher. Safe to leave unset — the API +// must run with the broker down, and a missed notification only means a till is +// stale until its next periodic pull. +func SetCatalogueNotifier(n CatalogueNotifier) { + catalogueNotifierMu.Lock() + defer catalogueNotifierMu.Unlock() + catalogueNotifier = n +} + +// notifyCatalogueChanged asks the tills at one outlet to refresh. +// +// Three rules, all of them learned the hard way elsewhere in this codebase: +// +// 1. **Call it after the transaction commits, never inside.** A rolled-back +// sale must not announce a change that did not happen. +// 2. **It cannot fail a sale.** A publish error is logged and swallowed — +// exactly as syncProductLocationStatus already does for the availability +// flag. Stock is committed; a terminal being told about it late is not +// worth losing a bill over. +// 3. **Only the revision travels, never quantities.** The till pulls. Pushed +// numbers race with concurrent sales, and the catalogue endpoint already +// answers a delta from a revision. +// +// Fire-and-forget on its own goroutine: the publisher waits for the broker's +// acknowledgement, and a bill's response should not sit behind that. +func notifyCatalogueChanged(locationID int) { + if locationID <= 0 { + return + } + + catalogueNotifierMu.RLock() + n := catalogueNotifier + catalogueNotifierMu.RUnlock() + if n == nil { + return + } + + // Stamped a second in the past for the same reason Catalogue() does it: a + // revision must not claim to include a write that is still landing. + revision := posRevisionFor(locationID, time.Now().Add(-time.Second)) + storeID := strconv.Itoa(locationID) + + go func() { + if err := n.PublishCatalogueChanged(storeID, revision); err != nil { + log.Printf("catalogue notify: store %s: %v", storeID, err) + } + }() +} diff --git a/repositories/catalogueRepository.go b/repositories/catalogueRepository.go index d576676..2983d83 100644 --- a/repositories/catalogueRepository.go +++ b/repositories/catalogueRepository.go @@ -8,6 +8,7 @@ import ( "nearle/models" "sort" "strings" + "sync" "time" "gorm.io/gorm" @@ -105,8 +106,100 @@ func NewCatalogueRepository(db *gorm.DB) CatalogueRepository { return &catalogueRepository{db: db} } -func tableForBrand(brand string) (string, error) { - table, ok := catalogueBrandTables[strings.ToLower(strings.TrimSpace(brand))] +// brandTables is the allowlist actually used, discovered from the catalogue +// database rather than hardcoded. +// +// The literal map above only ever exposed six brands, so a `brand_*` table +// added to the catalogue DB was invisible to the whole product until somebody +// shipped a release. Discovery removes that, and keeps the property the +// allowlist existed for: **table names still never come from a request.** They +// come from the database's own catalog, and the `brand` parameter is only ever +// looked up in the result. +// +// Two things it checks that a name alone would not: +// +// - the table has every column `catalogueProductColumns` selects. A +// `brand_*` table of a different shape would error on read, which is +// exactly the `brand_sakthi` failure — a table that exists in the map and +// cannot be queried. +// - discovery returning nothing falls back to the literal map, so a +// permissions problem on information_schema cannot take the catalogue dark. +// +// Cached with a TTL so a brand added today appears without a restart. +var ( + brandTablesMu sync.RWMutex + brandTablesData map[string]string + brandTablesAt time.Time +) + +const brandTablesTTL = 5 * time.Minute + +func (r *catalogueRepository) brandTables() map[string]string { + brandTablesMu.RLock() + if brandTablesData != nil && time.Since(brandTablesAt) < brandTablesTTL { + defer brandTablesMu.RUnlock() + return brandTablesData + } + brandTablesMu.RUnlock() + + discovered := catalogueBrandTables + if r.db != nil { + if found, err := r.discoverBrandTables(); err != nil { + log.Printf("catalogue: brand discovery failed, using the built-in list: %v", err) + } else if len(found) > 0 { + discovered = found + } else { + log.Println("catalogue: brand discovery found no usable tables, using the built-in list") + } + } + + brandTablesMu.Lock() + brandTablesData, brandTablesAt = discovered, time.Now() + brandTablesMu.Unlock() + return discovered +} + +// discoverBrandTables lists `brand_*` tables that carry every column the reader +// needs, in one query. +func (r *catalogueRepository) discoverBrandTables() (map[string]string, error) { + required := []string{ + "id", "product_name", "title", "description", "category", "image_id", "size", + "variant_key", "product_sku", "sku_source", "price_range", "providers", + "fssai_license", "highlights", "nutrients", "search_query", "created_at", "updated_at", + } + + // `IN (?)` with a slice rather than `= ANY(?)` with a driver array type: + // GORM expands the former itself, and the latter would pull in lib/pq for + // one call in a codebase that has no other use for it. + var names []string + err := r.db.Raw(` + SELECT c.table_name + FROM information_schema.columns c + WHERE c.table_schema = 'public' + AND c.table_name LIKE 'brand\_%' + AND c.column_name IN (?) + GROUP BY c.table_name + HAVING COUNT(DISTINCT c.column_name) = ? + ORDER BY c.table_name`, + required, len(required), + ).Scan(&names).Error + if err != nil { + return nil, err + } + + out := make(map[string]string, len(names)) + for _, table := range names { + brand := strings.TrimPrefix(table, "brand_") + if brand == "" || brand == table { + continue + } + out[strings.ToLower(brand)] = table + } + return out, nil +} + +func (r *catalogueRepository) tableForBrand(brand string) (string, error) { + table, ok := r.brandTables()[strings.ToLower(strings.TrimSpace(brand))] if !ok { return "", ErrUnknownBrand } @@ -121,7 +214,7 @@ func (r *catalogueRepository) GetBrands() ([]models.CatalogueBrand, error) { var brands []models.CatalogueBrand var failures int - for brand, table := range catalogueBrandTables { + for brand, table := range r.brandTables() { var count int64 if err := r.db.Table(table).Count(&count).Error; err != nil { // One brand's table being absent or unreadable must not hide the @@ -138,7 +231,7 @@ func (r *catalogueRepository) GetBrands() ([]models.CatalogueBrand, error) { // Every table failing is a different thing from every table being empty: // the first is an outage and must be reported, the second is a fact. - if failures == len(catalogueBrandTables) { + if failures == len(r.brandTables()) { return nil, fmt.Errorf("no catalogue brand table could be read (%d configured)", failures) } @@ -150,7 +243,7 @@ func (r *catalogueRepository) GetCategories(brand string) ([]string, error) { return nil, ErrCatalogueDBUnavailable } - table, err := tableForBrand(brand) + table, err := r.tableForBrand(brand) if err != nil { return nil, err } @@ -187,7 +280,7 @@ func (r *catalogueRepository) GetProducts(brand, category, keyword string, pagen } func (r *catalogueRepository) getProductsForBrand(brand, category, keyword string, pageno, pagesize int) ([]models.CatalogueProduct, int64, error) { - table, err := tableForBrand(brand) + table, err := r.tableForBrand(brand) if err != nil { return nil, 0, err } @@ -229,7 +322,8 @@ func (r *catalogueRepository) getProductsAllBrands(category, keyword string, pag whereClause, args := catalogueWhereClause(category, keyword) brands := make([]string, 0, len(catalogueBrandTables)) - for brand := range catalogueBrandTables { + tables := r.brandTables() + for brand := range tables { brands = append(brands, brand) } sort.Strings(brands) @@ -238,7 +332,7 @@ func (r *catalogueRepository) getProductsAllBrands(category, keyword string, pag var failures int for _, brand := range brands { - table := catalogueBrandTables[brand] + table := tables[brand] var rows []catalogueProductRow dataQuery := fmt.Sprintf(`SELECT %s FROM %s WHERE %s ORDER BY id`, catalogueProductColumns, table, whereClause) @@ -302,7 +396,7 @@ func (r *catalogueRepository) GetProductBySKU(brand, sku string) (*models.Catalo return nil, ErrCatalogueDBUnavailable } - table, err := tableForBrand(brand) + table, err := r.tableForBrand(brand) if err != nil { return nil, err } @@ -329,7 +423,7 @@ func (r *catalogueRepository) GetProductByID(brand string, id int64) (*models.Ca return nil, ErrCatalogueDBUnavailable } - table, err := tableForBrand(brand) + table, err := r.tableForBrand(brand) if err != nil { return nil, err } diff --git a/repositories/orderRepository.go b/repositories/orderRepository.go index 052e4f8..89221e6 100644 --- a/repositories/orderRepository.go +++ b/repositories/orderRepository.go @@ -1354,6 +1354,13 @@ func (r *orderRepository) CreateOrder(data models.Orders) (models.Orders, error) return models.Orders{}, err } + // An app order depletes the same shelf the counter sells from, so the tills + // at that outlet need to know as much as they do after a counter sale. This + // half was the one that would have been forgotten: it is easy to think of + // the catalogue as a POS concern, and a till overselling something the app + // just took is exactly the failure the push exists to prevent. + notifyCatalogueChanged(created.Locationid) + return r.reloadOrder(created.Orderheaderid) } diff --git a/repositories/posRepository.go b/repositories/posRepository.go index a758c25..2b007b3 100644 --- a/repositories/posRepository.go +++ b/repositories/posRepository.go @@ -386,6 +386,11 @@ func (r *posRepository) importPosOrder( return fmt.Sprintf("could not commit bill %s: %v", order.Invoicenumber, err) } + // Stock has moved at this outlet, so every other till standing at the same + // counter is now holding a figure that is one sale out of date. After the + // commit, never before — a rolled-back bill must not announce itself. + notifyCatalogueChanged(ctx.Locationid) + return "" } diff --git a/repositories/productRepository.go b/repositories/productRepository.go index d1de908..fcff05d 100644 --- a/repositories/productRepository.go +++ b/repositories/productRepository.go @@ -255,7 +255,7 @@ func (r *productRepository) GetProductStocks(tenantID, locationID string) ([]mod MAX(a.stocktype) AS stocktype, MAX(a.maxquantity) AS maxquantity, MAX(a.minquantity) AS minquantity, MAX(a.status) AS status, b.applocationid, b.categoryid, b.subcategoryid, b.catalogueid, b.addonid, b.discountid, b.pricingid, - b.productname, b.productimage, b.productdesc, b.productsku, b.brandid, b.productbrand, b.productunit, + b.productname, b.productimage, COALESCE(b.productimages::text, '') AS productimages, b.productdesc, b.productsku, b.brandid, b.productbrand, b.productunit, b.unitvalue, b.toppicks, b.productcost, b.taxamount, b.taxpercent, b.producttax, b.productstock, b.productcombo, b.variants, b.retailprice, b.diffprice, b.diffpercent, b.othercost, b.approve, b.productstatus, b.created, b.updated, c.subcatname AS subcategoryname, diff --git a/scratch/productimages/main.go b/scratch/productimages/main.go new file mode 100644 index 0000000..c0cb51f --- /dev/null +++ b/scratch/productimages/main.go @@ -0,0 +1,151 @@ +// Backfills products.productimages for everything imported before the column +// existed. +// +// The import used to keep `Images[0]` and discard the rest, so roughly six +// thousand products carry one photo where the global catalogue has up to ten. +// Re-importing would work but rewrites price, category and stock links; this +// only fills the gallery. +// +// The images are re-listed rather than reconstructed. Filenames are +// `image_000.jpeg`, `image_001.jpg` … — the extension varies within a single +// product, so the sibling URLs cannot be derived from the first one. The +// `image_id` is taken out of the stored URL and the same S3 listing the API +// uses is asked for the rest, which is why this boots the app's own +// connections instead of talking to Postgres alone. +// +// Dry by default; pass --apply to write. +// +// go run ./scratch/productimages # report only +// go run ./scratch/productimages --apply # write +// go run ./scratch/productimages --apply --tenant 1087 +package main + +import ( + "encoding/json" + "flag" + "fmt" + "log" + "strings" + + "nearle/db" + + "github.com/joho/godotenv" +) + +type row struct { + Productid int + Tenantid int + Productbrand string + Productimage string +} + +func main() { + apply := flag.Bool("apply", false, "write the changes; otherwise report only") + tenant := flag.Int("tenant", 0, "restrict to one tenant") + flag.Parse() + + _ = godotenv.Load() + db.Connect() + + if db.ImageStore == nil { + log.Fatal("the image store is not configured (USE_S3 / S3_*), so there is nothing to list images from") + } + + q := ` + SELECT productid, tenantid, + COALESCE(productbrand,'') AS productbrand, + COALESCE(productimage,'') AS productimage + FROM products + WHERE COALESCE(productimage,'') <> '' + AND COALESCE(productimages::text,'') IN ('', 'null', '[]')` + args := []interface{}{} + if *tenant > 0 { + q += ` AND tenantid = ?` + args = append(args, *tenant) + } + q += ` ORDER BY productid` + + var rows []row + if err := db.DB.Raw(q, args...).Scan(&rows).Error; err != nil { + log.Fatal(err) + } + + fmt.Printf("candidates: %d\n\n", len(rows)) + + var filled, single, unparsed, nolist int + for _, r := range rows { + brand, imageID := parseImageRef(r.Productimage) + if brand == "" || imageID == "" { + unparsed++ + continue + } + // The stored brand is authoritative where it exists; the URL is only a + // fallback, because a product renamed in the back office keeps its + // original image path. + if b := strings.ToLower(strings.TrimSpace(r.Productbrand)); b != "" { + brand = b + } + + images := db.GetImages(brand, imageID) + if len(images) == 0 { + nolist++ + continue + } + if len(images) == 1 { + // Nothing gained — leave the column null rather than writing a + // one-element array that says the same as productimage. + single++ + continue + } + + encoded, err := json.Marshal(images) + if err != nil { + log.Printf("product %d: encoding failed: %v", r.Productid, err) + continue + } + + if *apply { + if err := db.DB.Exec( + `UPDATE products SET productimages = ?::jsonb, updated = NOW() WHERE productid = ?`, + string(encoded), r.Productid).Error; err != nil { + log.Printf("product %d: update failed: %v", r.Productid, err) + continue + } + } + filled++ + if filled <= 10 { + fmt.Printf(" product %-7d %-10s %-34s %d images\n", r.Productid, brand, imageID, len(images)) + } + } + + fmt.Printf(` +would fill (>1 image) : %d +only one image : %d +image_id unparseable : %d +nothing listed in S3 : %d +`, filled, single, unparsed, nolist) + + if !*apply { + fmt.Println("\ndry run — nothing written. Re-run with --apply.") + } +} + +// parseImageRef pulls the brand and image_id out of a stored image URL. +// +// https://…/daily/brands/nestle/nestle_maggi_615a9fa9/image_000.jpeg +// ^brand ^image_id +// +// Returns empty strings for anything that does not match, so a hand-uploaded +// image URL is skipped rather than guessed at. +func parseImageRef(url string) (brand, imageID string) { + const marker = "/daily/brands/" + i := strings.Index(url, marker) + if i < 0 { + return "", "" + } + parts := strings.Split(strings.Trim(url[i+len(marker):], "/"), "/") + if len(parts) < 3 { + return "", "" + } + return strings.ToLower(parts[0]), parts[1] +} diff --git a/services/productService.go b/services/productService.go index ce3d6d8..311e0e1 100644 --- a/services/productService.go +++ b/services/productService.go @@ -1,7 +1,9 @@ package services import ( + "encoding/json" "fmt" + "log" "nearle/models" "nearle/repositories" "time" @@ -284,7 +286,22 @@ func (s *productService) ImportCatalogueProduct(reqs []models.ImportCataloguePro Approve: 1, } if len(catalogueProduct.Images) > 0 { + // The first stays where every reader already looks for it. snapshot.Productimage = catalogueProduct.Images[0] + + // The rest used to be dropped on the floor here — a product with + // ten photos in the global catalogue arrived in a shop with one, + // and there was no way to get them back short of re-importing. + // + // Encoding failure is swallowed: the product itself is fine + // without a gallery, and refusing an import over a photo list + // would be the wrong trade. + if encoded, err := json.Marshal(catalogueProduct.Images); err == nil { + snapshot.Productimages = string(encoded) + } else { + log.Printf("import: could not encode images for catalogue %s/%d: %v", + req.Brand, req.Catalogueid, err) + } } productID, err = s.repo.CreateProductReturningID(snapshot)