From 72907dae7461575596c9a9b50e57528c7a896dff Mon Sep 17 00:00:00 2001 From: Suriyakumarvijayanayagam Date: Tue, 15 Sep 2026 17:04:34 +0530 Subject: [PATCH] Scan-to-order: label from the customer's camera to "buy it here" MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit POST /v1/mob/scan/lookup label + customer → catalogue match, sizes, and every registered store that sells it with live stock, in-stock first / nearest first, one recommended POST /v1/mob/scan/confirm chosen store + size + qty → re-read the ledger; ok, or the next-nearest store with enough of the same product GET /v1/mob/scan/stores registered stores nearest first Recognition is pgvector cosine search over every brand_* table (each with its own index, merged) plus a word match that settles near-ties and works alone when no model is configured. The embedder is chosen by EMBEDDING_PROVIDER (OpenAI-compatible or Gemini) and must be the model that indexed the catalogue: verified 2026-09-15 as all-MiniLM-L6-v2 over search_query, served by the cluster's Ollama as `all-minilm`; the first search refuses a width mismatch by name. Customer, stores and catalogue are read concurrently under a 5 s cap; a slow model degrades to a text answer. Vectors and ranked hits are cached in Redis and in-process; live stock never is. Availability uses the same rules as the customer catalogue (approve, publishedat, ledger balance, outlet price else retail). No stock reservation: confirm re-reads. scratch/cataloguedims reports the catalogue's embedding width and fill. Co-Authored-By: Claude Opus 5 (1M context) --- controllers/scanController.go | 101 +++++ docs/SCAN_TO_ORDER.md | 188 +++++++++ facade/container.go | 13 +- models/scan.go | 141 +++++++ repositories/scanRepository.go | 633 +++++++++++++++++++++++++++++ routes/routes.go | 1 + routes/scanroutes.go | 17 + scratch/cataloguedims/main.go | 89 ++++ services/scanService.go | 721 +++++++++++++++++++++++++++++++++ services/scan_test.go | 438 ++++++++++++++++++++ utils/embedding.go | 225 ++++++++++ utils/embedding_test.go | 106 +++++ utils/geo.go | 104 +++++ utils/geo_test.go | 65 +++ 14 files changed, 2841 insertions(+), 1 deletion(-) create mode 100644 controllers/scanController.go create mode 100644 docs/SCAN_TO_ORDER.md create mode 100644 models/scan.go create mode 100644 repositories/scanRepository.go create mode 100644 routes/scanroutes.go create mode 100644 scratch/cataloguedims/main.go create mode 100644 services/scanService.go create mode 100644 services/scan_test.go create mode 100644 utils/embedding.go create mode 100644 utils/embedding_test.go create mode 100644 utils/geo.go create mode 100644 utils/geo_test.go diff --git a/controllers/scanController.go b/controllers/scanController.go new file mode 100644 index 0000000..302d899 --- /dev/null +++ b/controllers/scanController.go @@ -0,0 +1,101 @@ +package controllers + +import ( + "errors" + "nearle/models" + "nearle/services" + "net/http" + + "github.com/gofiber/fiber/v2" +) + +// ScanController is the scan-to-order surface for the customer app: +// +// POST /v1/mob/scan/lookup a label → the product, and which of my stores has it +// POST /v1/mob/scan/confirm I picked a store and a size → still there? else where? +// GET /v1/mob/scan/stores my stores, nearest first +// +// Business outcomes ("out of stock", "not registered with that store") are +// 200s with a reason in the body: the app renders them, it does not retry +// them. HTTP errors are reserved for a request that cannot be served at all. +type ScanController struct { + scanService services.ScanService +} + +func NewScanController(scanService services.ScanService) *ScanController { + return &ScanController{scanService: scanService} +} + +func (ctl *ScanController) Lookup(c *fiber.Ctx) error { + var req models.ScanLookupRequest + if err := c.BodyParser(&req); err != nil { + return scanBadRequest(c, "Invalid request body") + } + resp, err := ctl.scanService.Lookup(c.Context(), req) + if err != nil { + return scanError(c, err, "Could not look up that product") + } + return c.Status(http.StatusOK).JSON(fiber.Map{ + "code": http.StatusOK, + "status": true, + "message": resp.Message, + "details": resp, + }) +} + +func (ctl *ScanController) Confirm(c *fiber.Ctx) error { + var req models.ScanConfirmRequest + if err := c.BodyParser(&req); err != nil { + return scanBadRequest(c, "Invalid request body") + } + resp, err := ctl.scanService.Confirm(c.Context(), req) + if err != nil { + return scanError(c, err, "Could not check that store") + } + return c.Status(http.StatusOK).JSON(fiber.Map{ + "code": http.StatusOK, + "status": true, + "message": resp.Message, + "details": resp, + }) +} + +func (ctl *ScanController) Stores(c *fiber.Ctx) error { + customerid, _ := c.QueryInt("customerid"), 0 + stores, err := ctl.scanService.Stores(c.Context(), customerid, + models.FlexibleString(c.Query("latitude")), models.FlexibleString(c.Query("longitude"))) + if err != nil { + return scanError(c, err, "Could not list your stores") + } + return c.Status(http.StatusOK).JSON(fiber.Map{ + "code": http.StatusOK, + "status": true, + "message": "Success", + "details": stores, + }) +} + +func scanBadRequest(c *fiber.Ctx, msg string) error { + return c.Status(http.StatusBadRequest).JSON(fiber.Map{ + "code": http.StatusBadRequest, + "status": false, + "message": msg, + }) +} + +func scanError(c *fiber.Ctx, err error, fallback string) error { + code, msg := http.StatusInternalServerError, fallback + switch { + case errors.Is(err, services.ErrScanBadRequest): + code, msg = http.StatusBadRequest, err.Error() + case errors.Is(err, services.ErrScanCustomerNotFound): + code, msg = http.StatusNotFound, "Customer not found" + case errors.Is(err, services.ErrScanCatalogueDown): + code, msg = http.StatusServiceUnavailable, "Product search is temporarily unavailable" + } + return c.Status(code).JSON(fiber.Map{ + "code": code, + "status": false, + "message": msg, + }) +} diff --git a/docs/SCAN_TO_ORDER.md b/docs/SCAN_TO_ORDER.md new file mode 100644 index 0000000..6508f89 --- /dev/null +++ b/docs/SCAN_TO_ORDER.md @@ -0,0 +1,188 @@ +# Scan-to-order — mobile integration + +A customer photographs a product. Google Lens (on the phone) turns the photo +into a label — `"Milk Bikis"`, `"Dabur Honey 500g"`. The app sends that label +here and gets back: what the product is, which of the customer's stores sell +it, in which sizes, with live stock, nearest first, and which store we +recommend. When the customer taps a store and a size, a second call confirms +the shelf still has it — and if it does not, names the next-nearest store +that does. + +Base path: `/live/api/v1/mob/scan`. Every response uses the usual envelope +`{ code, status, message, details }`; the shapes below are `details`. + +## The flow + +``` +photo ──Lens──▶ label + │ + ▼ + POST /lookup ───▶ match + stores[] (recommended first) + │ + customer taps a store + a size + │ + ▼ + POST /confirm ───▶ ok:true → add to basket with existing order APIs + ok:false + alternative → offer the other store +``` + +`GET /stores` is for the "choose another shop" sheet: the customer's +registered stores, nearest first, independent of any product. + +## `POST /lookup` + +```json +{ + "customerid": 5123, + "label": "Milk Bikis", + "latitude": 11.0290, // phone fix; optional — saved address is used without it + "longitude": 77.0290, + "tenantids": [1135, 1140], // optional: what the app THINKS the customer joined + "limit": 0 // optional: max stores, 0 = all +} +``` + +`tenantids` is verified, never trusted: the server intersects it with the +`tenantcustomers` table. Ids the customer is not actually registered with +come back in `unregistered_tenantids` — treat that as "refresh the local +list". A list that matches nothing at all is treated as stale and all +registered stores are used. + +Response: + +```json +{ + "label": "Milk Bikis", + "match": { + "brand": "britannia", "catalogueid": 7, "imageid": "britannia_milk_bikis_100g", + "product_name": "Milk Bikis", "size": "100 g", "variant_key": "milk_bikis", + "image": "https://…", "score": 0.94, "method": "vector+text" + }, + "catalogue_variants": [ { "…same shape…": "100 g" }, { "…": "200 g" } ], + "confidence": 0.94, + "available": true, + "recommended_locationid": 20, + "stores": [ + { + "tenantid": 2, "tenantname": "R Mart", "locationid": 20, "locationname": "Hopes", + "latitude": 11.01, "longitude": 77.0, "distance_km": 3.8, "open": true, + "deliveryradius": 5, "deliverymins": 30, + "recommended": true, "available": true, + "options": [ + { "productid": 200, "productname": "Milk Bikis 100g", "size": "100 g", "price": 12, "stock": 6, + "available": true, "is_variant": false, "matched_by": "imageid", "image": "…" }, + { "productid": 201, "productname": "Milk Bikis 200g", "size": "200 g", "price": 22, "stock": 3, + "available": true, "is_variant": true, "variantname": "200 g", "matched_by": "variant-of:200" } + ] + }, + { "locationid": 10, "locationname": "Peelamedu", "distance_km": 0.9, "available": false, "recommended": false, + "options": [ { "productid": 100, "stock": 0, "available": false, "…": "…" } ] } + ], + "unregistered_tenantids": [], + "message": "Available at 1 of your stores." +} +``` + +How to read it: + +- `match == null` → nothing recognised; show `message` and let them retry. + `confidence` below ~0.5 → recognised but unsure; confirm the name with the + customer before showing prices. `method: "text"` means no embedding model + was involved (not configured, or it timed out) — be a little more cautious. +- `stores` is ordered **in-stock first, then nearest**. Exactly one store has + `recommended: true` — the nearest with stock — and only when `available` + is true. Stores that sell it but have nothing on the shelf are still listed + (so the customer understands why they are not recommended); stores that do + not sell it are not. +- `options` are the things that can actually go in a basket at that store — + the matched product and each of its sizes — each a real product with its + own `productid`, price and live `stock`. Use `productid` in the existing + cart/order calls exactly as you would from the catalogue screen. +- `distance_km: -1` means the distance is unknown (no fix from the phone and + no saved address, or the store has no coordinates). Do not render it as 0. + +## `POST /confirm` + +Sent when the customer taps a store and an option. Re-reads live stock — +nothing is cached on this path. + +```json +{ "customerid": 5123, "tenantid": 1, "locationid": 10, "productid": 100, "quantity": 2, + "latitude": 11.029, "longitude": 77.029 } +``` + +```json +{ + "ok": false, + "reason": "out_of_stock", // in_stock | insufficient_stock | out_of_stock | not_sold_here | store_not_registered + "store": { "…the store they tapped…" }, + "option": { "productid": 100, "stock": 0, "…": "…" }, + "requested": 2, + "alternative": { // absent when nobody has enough + "locationid": 20, "locationname": "Hopes", "distance_km": 3.8, "recommended": true, "available": true, + "options": [ { "productid": 200, "stock": 6, "price": 12, "…": "…" } ] + }, + "message": "Out of stock at Peelamedu. Hopes has it (3.8 km away)." +} +``` + +`ok: true` → proceed to the basket. `ok: false` → show `message`; if +`alternative` is present offer it as a one-tap switch (it is the **same +product**, not another size — the customer chose a size and we do not +substitute). These are HTTP 200s: they are answers, not errors. + +## `GET /stores?customerid=5123&latitude=11.029&longitude=77.029` + +The customer's registered stores, nearest first, `distance_km: -1` last. +Same `ScanStore` shape as inside `stores[]` above, without options. + +## Errors (HTTP status ≠ 200) + +| Status | When | +|---|---| +| 400 | Missing `customerid`/`label`/ids, or a body that is not JSON. `message` says which. | +| 404 | `customerid` does not exist. | +| 503 | The catalogue database is not reachable. Retry later; the rest of the app is unaffected. | +| 500 | Anything else. Logged server-side. | + +## Behind the curtain (for whoever operates it) + +- **Recognition** = pgvector cosine search over every `brand_*` table in the + catalogue (each with its own index, merged), plus a word match on + `product_name`/`title`/`search_query` that settles near-ties and works on + its own when no embedding model is configured. The model is set by + `EMBEDDING_PROVIDER/MODEL/API_KEY` and **must** be the one that indexed + the catalogue — the first search checks the vector width and refuses a + mismatch by name. +- **The catalogue's model** (verified 2026-09-15 by cosine against a stored + row: 1.0000): `all-MiniLM-L6-v2`, 384-d, unit-normalised, embedding the + `search_query` column (brand + name + category + blurb + price range). + Ollama ships it as `all-minilm`; the cluster's `ollama.krow` service serves + it, so production is: + ``` + EMBEDDING_PROVIDER=openai + EMBEDDING_BASE_URL=http://ollama.krow.svc.cluster.local:11434/v1 + EMBEDDING_MODEL=all-minilm + EMBEDDING_API_KEY=ollama # any non-empty value; Ollama ignores it + EMBEDDING_DIMENSIONS=384 + ``` + A bare label ("Milk Bikis") scores ~0.92 against its product's stored + vector and ~0.23 against an unrelated one, which is what the 0.30 floor in + `scanService.go` is set against. If the catalogue team ever re-embeds + with another model, change `EMBEDDING_MODEL`/`DIMENSIONS` here and + nothing else. +- **Speed**: the label's vector (7 days) and the ranked catalogue hits + (30 min) are cached in Redis and in-process, so a popular product costs + one model call platform-wide. Customer, stores and catalogue are read + concurrently; the whole lookup is capped at 5 s and a slow model degrades + to a text answer instead of a spinner. Live stock is one indexed query and + is never cached. +- **Availability** is the same rule the app's catalogue screen uses: + `products.approve = 1`, `productlocations.publishedat IS NOT NULL`, stock = + live `SUM(in) − SUM(out)` of `productstocks` at that outlet, price = the + outlet's own price else the tenant's retail price. +- **No reservation.** Confirm re-reads the ledger; a hold would give the + same answer with a timer to babysit. If contention becomes real, a + Redis-backed short hold slots in at `Confirm` without changing the API. +- **Identity** is the `customerid` in the body, like every other mobile + endpoint here — there is no auth layer yet (see `SECURITY_HANDOFF.md`). diff --git a/facade/container.go b/facade/container.go index fba786e..058eace 100644 --- a/facade/container.go +++ b/facade/container.go @@ -4,6 +4,7 @@ import ( "nearle/controllers" "nearle/repositories" "nearle/services" + "nearle/utils" "gorm.io/gorm" ) @@ -22,6 +23,7 @@ type Facade struct { PosController *controllers.PosController LiveController *controllers.LiveController CatalogueUploadController *controllers.CatalogueUploadController + ScanController *controllers.ScanController // Held so the NATS consumer can reach the ingest without going through // HTTP. Unexported: everything else should use the controller. @@ -32,7 +34,8 @@ type Facade struct { // catalogueDB is a separate connection to the pgvector catalogue database; // it may be nil if catalogue env vars are not configured, in which case // catalogue endpoints will error at query time rather than at startup. -func NewFacade(db *gorm.DB, catalogueDB *gorm.DB) *Facade { +// embedder may be nil too: scan-to-order then matches on words alone. +func NewFacade(db *gorm.DB, catalogueDB *gorm.DB, embedder utils.Embedder) *Facade { // User Module userRepo := repositories.NewUserRepository(db) @@ -109,6 +112,13 @@ func NewFacade(db *gorm.DB, catalogueDB *gorm.DB) *Facade { catalogueUploadService := services.NewCatalogueUploadService(catalogueUploadRepo) catalogueUploadController := controllers.NewCatalogueUploadController(catalogueUploadService) + // Scan Module — a label from the customer's camera to "buy it here". + // Reads both databases: the catalogue to recognise the product, nearledb + // for who the customer is and what their outlets have on the shelf. + scanRepo := repositories.NewScanRepository(db, catalogueDB) + scanService := services.NewScanService(scanRepo, embedder) + scanController := controllers.NewScanController(scanService) + return &Facade{ UserController: userController, ProductController: productController, @@ -123,6 +133,7 @@ func NewFacade(db *gorm.DB, catalogueDB *gorm.DB) *Facade { PosController: posController, LiveController: liveController, CatalogueUploadController: catalogueUploadController, + ScanController: scanController, posService: posService, } } diff --git a/models/scan.go b/models/scan.go new file mode 100644 index 0000000..fc61af2 --- /dev/null +++ b/models/scan.go @@ -0,0 +1,141 @@ +package models + +// Scan-to-order. +// +// A customer points the app at a packet, Google Lens (on the phone) reads a +// label off it, and the app asks: "which of MY shops has this, in what sizes, +// and which one should I buy from?" These are the shapes on both sides of +// that conversation. + +// ScanLookupRequest is what the app sends once Lens has produced a label. +type ScanLookupRequest struct { + Customerid int `json:"customerid"` + // What Lens read: "Milk Bikis", "Dabur Honey 500g". Free text, trimmed + // and capped by the service. + Label string `json:"label"` + // Where the customer is right now. Optional: without it the customer's + // saved primary address is used, and without that stores are listed in + // registration order with no distance. + Latitude FlexibleString `json:"latitude"` + Longitude FlexibleString `json:"longitude"` + // The tenants the app believes the customer has scanned into. Optional and + // never trusted on its own: the server intersects it with the + // tenantcustomers table and reports anything it dropped. + Tenantids []int `json:"tenantids"` + // How many stores to return. 0 = all registered stores that stock it. + Limit int `json:"limit"` +} + +// ScanStore is one of the customer's registered outlets. +type ScanStore struct { + Tenantid int `json:"tenantid"` + Tenantname string `json:"tenantname"` + Locationid int `json:"locationid"` + Locationname string `json:"locationname"` + Address string `json:"address,omitempty"` + Latitude float64 `json:"latitude"` + Longitude float64 `json:"longitude"` + // Kilometres from the customer, or -1 when either side has no usable + // coordinates. Never omitted: a missing number is easy to misread as 0. + DistanceKm float64 `json:"distance_km"` + // Delivery reach in the outlet's own units, straight from tenantlocations. + Deliveryradius int `json:"deliveryradius"` + Deliverymins int `json:"deliverymins"` + Open bool `json:"open"` +} + +// ScanOption is one thing the customer can actually put in the basket at one +// store: the matched product itself, or one of its sizes. Each is a real +// product row with its own price and stock, which is why they are flat. +type ScanOption struct { + Productid int `json:"productid"` + Productname string `json:"productname"` + Size string `json:"size"` // "500 g", "1 kg" — unitvalue + productunit + Price float64 `json:"price"` + Stock int `json:"stock"` + Available bool `json:"available"` + Image string `json:"image,omitempty"` + // Is this the product that matched, or a size hanging under it? + IsVariant bool `json:"is_variant"` + Variantname string `json:"variantname,omitempty"` + // How the row was tied back to the catalogue: "imageid", + // "brand+catalogueid", "name" or "variant-of:". + MatchedBy string `json:"matched_by"` +} + +// ScanStoreOffer is one store and what it can sell. +type ScanStoreOffer struct { + ScanStore + // Nearest store with at least one option in stock. Exactly one offer + // carries this, and only when something is in stock somewhere. + Recommended bool `json:"recommended"` + // Any option in stock here. + Available bool `json:"available"` + Options []ScanOption `json:"options"` +} + +// ScanCatalogueMatch is what the catalogue search settled on. +type ScanCatalogueMatch struct { + Brand string `json:"brand"` + Catalogueid int64 `json:"catalogueid"` + Imageid string `json:"imageid,omitempty"` + ProductName string `json:"product_name"` + Title string `json:"title,omitempty"` + Category string `json:"category,omitempty"` + Size string `json:"size,omitempty"` + VariantKey string `json:"variant_key,omitempty"` + Image string `json:"image,omitempty"` + Score float64 `json:"score"` + // "vector", "vector+text" or "text" — how the score was produced. The app + // can be more cautious with a text-only match. + Method string `json:"method"` +} + +// ScanLookupResponse is the answer to a scan. +type ScanLookupResponse struct { + Label string `json:"label"` + // The best catalogue product for the label, and the sizes of it the + // catalogue knows about (each a separate catalogue row). + Match *ScanCatalogueMatch `json:"match"` + Variants []ScanCatalogueMatch `json:"catalogue_variants"` + // 0..1. Below ~0.5 the app should confirm with the customer before + // showing prices. + Confidence float64 `json:"confidence"` + // Registered stores that stock the product, nearest first, in-stock + // first. Empty with Available=false when none does. + Stores []ScanStoreOffer `json:"stores"` + Available bool `json:"available"` + // The locationid of the store marked Recommended, or 0. + RecommendedLocationid int `json:"recommended_locationid"` + // Tenant ids the app sent that the customer is not actually registered + // with. Empty normally; non-empty means the app's local list is stale. + UnregisteredTenantids []int `json:"unregistered_tenantids,omitempty"` + Message string `json:"message"` +} + +// ScanConfirmRequest is sent when the customer taps a store and a size. +type ScanConfirmRequest struct { + Customerid int `json:"customerid"` + Tenantid int `json:"tenantid"` + Locationid int `json:"locationid"` + Productid int `json:"productid"` + Quantity int `json:"quantity"` + Latitude FlexibleString `json:"latitude"` + Longitude FlexibleString `json:"longitude"` +} + +// ScanConfirmResponse says whether the pick still holds, and where to go if +// it does not. +type ScanConfirmResponse struct { + Ok bool `json:"ok"` + // "in_stock", "insufficient_stock", "out_of_stock", "not_sold_here", + // "store_not_registered". + Reason string `json:"reason"` + Store *ScanStore `json:"store,omitempty"` + Option *ScanOption `json:"option,omitempty"` + Requested int `json:"requested"` + // The next-nearest registered store with enough of the same product, when + // the chosen one has run out. Nil when there is none. + Alternative *ScanStoreOffer `json:"alternative,omitempty"` + Message string `json:"message"` +} diff --git a/repositories/scanRepository.go b/repositories/scanRepository.go new file mode 100644 index 0000000..7e52cba --- /dev/null +++ b/repositories/scanRepository.go @@ -0,0 +1,633 @@ +package repositories + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "log" + "nearle/db" + "nearle/models" + "nearle/utils" + "sort" + "strings" + "sync" + "time" + + "gorm.io/gorm" +) + +/* +Scan-to-order reads from three places and this file is the only one that +knows which is which: + + - the catalogue database (pgvector, one table per brand) — to turn a label + into a catalogue product; + - nearledb — who the customer is, which outlets they scanned into, and what + those outlets have on the shelf right now; + - Redis — a cache for the expensive and stable half (the label's vector and + its catalogue hits). Live stock is never cached. + +Two connections are held rather than one because the catalogue must never be +reachable through the nearledb handle: the comment on db.CatalogueDB is +explicit about that and every catalogue reader in this package honours it. +*/ + +// CatalogueKey is how a tenant's product row points back at the catalogue. +// Imageid is the stable one; brand+catalogueid is kept for rows imported +// before imageid existed (see models.Products.Imageid). +type CatalogueKey struct { + Brand string + Catalogueid int64 + Imageid string +} + +// CatalogueHit is one catalogue row the search considered. +type CatalogueHit struct { + Brand string + ID int64 + ProductName string + Title string + Category string + Size string + VariantKey string + ImageID string + ImageURL string + // Cosine distance from pgvector (0 = identical); -1 for a text-only hit. + Distance float64 +} + +// StoreOptionRow is one sellable product at one outlet, with its live stock. +type StoreOptionRow struct { + Tenantid int + Locationid int + Productid int + Productname string + Productbrand string + Catalogueid int64 + Imageid string + Productimage string + Productunit string + Unitvalue string + Price float64 + Stock int + // For a size row: the product it hangs under and the label given to it. + Parentid int + Variantname string +} + +type ScanRepository interface { + // nearledb + CustomerExists(ctx context.Context, customerid int) (bool, error) + CustomerHome(ctx context.Context, customerid int) (lat, lng float64, ok bool, err error) + RegisteredStores(ctx context.Context, customerid int) ([]models.ScanStore, error) + // StoreOptions finds, at the given outlets, every published product tied + // to one of the catalogue keys (or, for hand-made products, one of the + // names) — and every size hanging under those products. + StoreOptions(ctx context.Context, locationids []int, keys []CatalogueKey, names []string) ([]StoreOptionRow, error) + // ProductAt is one product at one outlet with its live stock, or nil. + ProductAt(ctx context.Context, tenantid, locationid, productid int) (*StoreOptionRow, error) + + // catalogue + VectorSearch(ctx context.Context, vector []float32, limit int) ([]CatalogueHit, error) + TextSearch(ctx context.Context, label string, limit int) ([]CatalogueHit, error) + VectorSearchAvailable() bool + + // cache + CachedVector(ctx context.Context, model, label string) ([]float32, bool) + CacheVector(ctx context.Context, model, label string, v []float32) + CachedHits(ctx context.Context, method, label string) ([]CatalogueHit, bool) + CacheHits(ctx context.Context, method, label string, hits []CatalogueHit) +} + +type scanRepository struct { + db *gorm.DB + catalogue *gorm.DB + + // Catalogue tables and their columns, discovered once and refreshed on a + // timer — the catalogue pipeline adds brands without telling anyone. + tablesMu sync.Mutex + tables map[string]map[string]bool // table -> column set + tablesAt time.Time + embeddingDim int + + // Process-local cache in front of Redis, bounded, so a hot label costs + // nothing even when Redis is not configured. + memMu sync.Mutex + memVecs map[string][]float32 + memHits map[string][]CatalogueHit +} + +const ( + scanTablesTTL = 10 * time.Minute + scanVectorTTL = 7 * 24 * time.Hour // a label's vector never changes for a given model + scanHitsTTL = 30 * time.Minute // the catalogue is rebuilt by scrape; not for long + scanMemCacheMax = 2000 +) + +func NewScanRepository(nearle, catalogue *gorm.DB) ScanRepository { + return &scanRepository{ + db: nearle, + catalogue: catalogue, + memVecs: make(map[string][]float32), + memHits: make(map[string][]CatalogueHit), + } +} + +// ── nearledb ──────────────────────────────────────────────────────────────── + +func (r *scanRepository) CustomerExists(ctx context.Context, customerid int) (bool, error) { + var n int64 + err := r.db.WithContext(ctx).Raw( + `SELECT COUNT(1) FROM customers WHERE customerid = ?`, customerid).Scan(&n).Error + return n > 0, err +} + +// CustomerHome is the saved primary address, falling back to the customers +// row itself. Either may be blank or unparsable — a customer created from a +// phone number alone has neither — and that is reported as ok=false rather +// than as (0, 0), which is a real place in the Gulf of Guinea. +func (r *scanRepository) CustomerHome(ctx context.Context, customerid int) (float64, float64, bool, error) { + var row struct { + Lat string + Lng string + } + err := r.db.WithContext(ctx).Raw(` + SELECT COALESCE(NULLIF(l.latitude, ''), c.latitude, '') AS lat, + COALESCE(NULLIF(l.longitude, ''), c.longitude, '') AS lng + FROM customers c + LEFT JOIN customerlocations l ON l.customerid = c.customerid AND l.primaryaddress = 1 + WHERE c.customerid = ? + LIMIT 1`, customerid).Scan(&row).Error + if err != nil { + return 0, 0, false, err + } + lat, lng, ok := utils.ParseLatLng(row.Lat, row.Lng) + return lat, lng, ok, nil +} + +// RegisteredStores is every active outlet of every tenant the customer has +// scanned into. A tenantcustomers row with locationid 0 means "the tenant", +// i.e. all of its outlets; a non-zero one pins a single outlet. +func (r *scanRepository) RegisteredStores(ctx context.Context, customerid int) ([]models.ScanStore, error) { + var rows []struct { + Tenantid int + Tenantname string + Locationid int + Locationname string + Address string + Latitude string + Longitude string + Deliveryradius int + Deliverymins int + Opentime string + Closetime string + } + err := r.db.WithContext(ctx).Raw(` + SELECT DISTINCT + tl.tenantid, COALESCE(t.tenantname, '') AS tenantname, + tl.locationid, COALESCE(tl.locationname, '') AS locationname, + COALESCE(tl.address, '') AS address, + COALESCE(tl.latitude, '') AS latitude, COALESCE(tl.longitude, '') AS longitude, + COALESCE(tl.deliveryradius, 0) AS deliveryradius, COALESCE(tl.deliverymins, 0) AS deliverymins, + COALESCE(tl.opentime, '') AS opentime, COALESCE(tl.closetime, '') AS closetime + FROM tenantcustomers tc + INNER JOIN tenantlocations tl + ON tl.tenantid = tc.tenantid + AND (COALESCE(tc.locationid, 0) = 0 OR tc.locationid = tl.locationid) + LEFT JOIN tenants t ON t.tenantid = tl.tenantid + WHERE tc.customerid = ? + AND LOWER(COALESCE(tl.status, 'active')) <> 'inactive' + ORDER BY tl.tenantid, tl.locationid`, customerid).Scan(&rows).Error + if err != nil { + return nil, err + } + + now := time.Now() + stores := make([]models.ScanStore, 0, len(rows)) + for _, row := range rows { + lat, lng, _ := utils.ParseLatLng(row.Latitude, row.Longitude) + stores = append(stores, models.ScanStore{ + Tenantid: row.Tenantid, + Tenantname: row.Tenantname, + Locationid: row.Locationid, + Locationname: row.Locationname, + Address: row.Address, + Latitude: lat, + Longitude: lng, + DistanceKm: -1, + Deliveryradius: row.Deliveryradius, + Deliverymins: row.Deliverymins, + Open: utils.OpenNow(row.Opentime, row.Closetime, now), + }) + } + return stores, nil +} + +// storeOptionSelect is the projection every outlet read shares, so the +// price and stock rules cannot differ between the lookup and the confirm. +// +// Price: the outlet's own price when it set one, else the tenant's retail +// price — the same rule GetProducts applies. Stock: the live IN−OUT balance +// of the ledger at that outlet, the same expression the app displays, so a +// product can never be offered here and show 0 on the next screen. +const storeOptionSelect = ` + SELECT a.tenantid, b.locationid, a.productid, + COALESCE(a.productname, '') AS productname, + LOWER(COALESCE(a.productbrand, '')) AS productbrand, + COALESCE(a.catalogueid, 0) AS catalogueid, + COALESCE(a.imageid, '') AS imageid, + COALESCE(a.productimage, '') AS productimage, + COALESCE(a.productunit, '') AS productunit, + COALESCE(a.unitvalue, '') AS unitvalue, + CASE WHEN COALESCE(b.price, 0) > 0 THEN b.price ELSE COALESCE(a.retailprice, 0) END AS price, + COALESCE(( + SELECT SUM(CASE WHEN LOWER(c.stocktype) = 'in' THEN c.quantity + WHEN LOWER(c.stocktype) = 'out' THEN -c.quantity + ELSE 0 END) + FROM productstocks c + WHERE c.productid = a.productid AND c.locationid = b.locationid AND c.tenantid = a.tenantid + ), 0) AS stock, + COALESCE(v.productid, 0) AS parentid, + COALESCE(v.variantname, '') AS variantname + FROM products a + INNER JOIN productlocations b ON b.productid = a.productid AND b.tenantid = a.tenantid + LEFT JOIN productvariants v ON v.variantproductid = a.productid AND v.tenantid = a.tenantid + AND LOWER(COALESCE(v.status, 'active')) <> 'inactive'` + +func (r *scanRepository) StoreOptions(ctx context.Context, locationids []int, keys []CatalogueKey, names []string) ([]StoreOptionRow, error) { + if len(locationids) == 0 || (len(keys) == 0 && len(names) == 0) { + return nil, nil + } + + // The products that ARE the catalogue match, at these outlets. + var matchConds []string + var args []interface{} + args = append(args, locationids) + for _, k := range keys { + if k.Imageid != "" { + matchConds = append(matchConds, "a.imageid = ?") + args = append(args, k.Imageid) + } + if k.Brand != "" && k.Catalogueid > 0 { + matchConds = append(matchConds, "(LOWER(a.productbrand) = ? AND a.catalogueid = ?)") + args = append(args, strings.ToLower(k.Brand), k.Catalogueid) + } + } + for _, n := range names { + if n = strings.ToLower(strings.TrimSpace(n)); n != "" { + matchConds = append(matchConds, "LOWER(a.productname) = ?") + args = append(args, n) + } + } + if len(matchConds) == 0 { + return nil, nil + } + + // Two reads rather than one recursive query: the second is keyed on the + // first's product ids, and a variant of a variant is not a thing here. + query := storeOptionSelect + ` + WHERE a.approve = 1 AND b.publishedat IS NOT NULL + AND b.locationid IN (?) + AND (` + strings.Join(matchConds, " OR ") + `)` + + var parents []StoreOptionRow + if err := r.db.WithContext(ctx).Raw(query, args...).Scan(&parents).Error; err != nil { + return nil, err + } + if len(parents) == 0 { + return nil, nil + } + + parentIDs := make([]int, 0, len(parents)) + for _, p := range parents { + parentIDs = append(parentIDs, p.Productid) + } + + // The sizes hanging under those products, at the same outlets. Only the + // rows whose parent is one of ours — the LEFT JOIN in the select can + // attach any parent, so it is pinned here. + var sizes []StoreOptionRow + err := r.db.WithContext(ctx).Raw(storeOptionSelect+` + WHERE a.approve = 1 AND b.publishedat IS NOT NULL + AND b.locationid IN (?) + AND v.productid IN (?)`, locationids, parentIDs).Scan(&sizes).Error + if err != nil { + return nil, err + } + + return append(parents, sizes...), nil +} + +func (r *scanRepository) ProductAt(ctx context.Context, tenantid, locationid, productid int) (*StoreOptionRow, error) { + var rows []StoreOptionRow + err := r.db.WithContext(ctx).Raw(storeOptionSelect+` + WHERE a.approve = 1 AND b.publishedat IS NOT NULL + AND a.tenantid = ? AND b.locationid = ? AND a.productid = ? + LIMIT 1`, tenantid, locationid, productid).Scan(&rows).Error + if err != nil || len(rows) == 0 { + return nil, err + } + return &rows[0], nil +} + +// ── catalogue ─────────────────────────────────────────────────────────────── + +// brandTables is every `brand_*` table and its columns, cached briefly. +func (r *scanRepository) brandTables(ctx context.Context) (map[string]map[string]bool, error) { + if r.catalogue == nil { + return nil, ErrCatalogueDBUnavailable + } + r.tablesMu.Lock() + defer r.tablesMu.Unlock() + if r.tables != nil && time.Since(r.tablesAt) < scanTablesTTL { + return r.tables, nil + } + + var rows []struct { + TableName string + ColumnName string + } + err := r.catalogue.WithContext(ctx).Raw(` + SELECT c.table_name, c.column_name + FROM information_schema.columns c + WHERE c.table_schema = 'public' AND c.table_name LIKE 'brand\_%'`).Scan(&rows).Error + if err != nil { + return nil, err + } + tables := make(map[string]map[string]bool) + for _, row := range rows { + if tables[row.TableName] == nil { + tables[row.TableName] = make(map[string]bool) + } + tables[row.TableName][row.ColumnName] = true + } + for name, cols := range tables { + if !cols["id"] || !cols["product_name"] { + delete(tables, name) + } + } + + // The vector width, read from the first embedding column found. pgvector + // stores it as the type modifier, so a mismatch with the model can be + // named in the error instead of surfacing as a bare "different vector + // dimensions" from the driver. + if r.embeddingDim == 0 { + for name, cols := range tables { + if !cols["embedding"] { + continue + } + var dim int + r.catalogue.WithContext(ctx).Raw(` + SELECT a.atttypmod FROM pg_attribute a + JOIN pg_class c ON c.oid = a.attrelid + WHERE c.relname = ? AND a.attname = 'embedding'`, name).Scan(&dim) + if dim > 0 { + r.embeddingDim = dim + } + break + } + } + + r.tables, r.tablesAt = tables, time.Now() + return tables, nil +} + +// VectorSearchAvailable is whether any catalogue table carries a vector. +func (r *scanRepository) VectorSearchAvailable() bool { + tables, err := r.brandTables(context.Background()) + if err != nil { + return false + } + for _, cols := range tables { + if cols["embedding"] { + return true + } + } + return false +} + +// hitColumns is the projection each search returns, with NULL stand-ins for +// columns a particular brand table lacks — the same tolerance +// catalogueRepository applies, for the same reason: a newer table missing +// one enrichment column is still a perfectly good catalogue of products. +func hitColumns(brand string, cols map[string]bool) string { + opt := func(name string) string { + if cols[name] { + return "COALESCE(" + name + ", '') AS " + name + } + return "'' AS " + name + } + return fmt.Sprintf(`'%s' AS brand, id, COALESCE(product_name, '') AS product_name, %s, %s, %s, %s, %s, %s`, + brand, opt("title"), opt("category"), opt("size"), opt("variant_key"), opt("image_id"), opt("image_url")) +} + +// VectorSearch ranks every brand table by cosine distance to the label's +// vector and merges the top of each. +// +// One branch per table, each with its own ORDER BY and LIMIT inside +// parentheses, so Postgres can use the per-table vector index instead of +// scanning the union. The literal is bound as a parameter and cast — never +// concatenated — and table names come from information_schema, never from +// the request. +func (r *scanRepository) VectorSearch(ctx context.Context, vector []float32, limit int) ([]CatalogueHit, error) { + tables, err := r.brandTables(ctx) + if err != nil { + return nil, err + } + if r.embeddingDim > 0 && len(vector) != r.embeddingDim { + return nil, fmt.Errorf("embedding is %d wide but the catalogue's embedding column is %d: EMBEDDING_MODEL/EMBEDDING_DIMENSIONS do not match the model that indexed the catalogue", len(vector), r.embeddingDim) + } + + literal := utils.VectorLiteral(vector) + var branches []string + var args []interface{} + for _, table := range sortedKeys(tables) { + cols := tables[table] + if !cols["embedding"] { + continue + } + brand := strings.TrimPrefix(table, "brand_") + branches = append(branches, fmt.Sprintf( + `(SELECT %s, (embedding <=> ?::vector) AS distance FROM %s WHERE embedding IS NOT NULL ORDER BY embedding <=> ?::vector LIMIT %d)`, + hitColumns(brand, cols), table, limit)) + args = append(args, literal, literal) + } + if len(branches) == 0 { + return nil, errors.New("no catalogue table has an embedding column") + } + + query := strings.Join(branches, " UNION ALL ") + fmt.Sprintf(" ORDER BY distance LIMIT %d", limit) + var hits []CatalogueHit + if err := r.catalogue.WithContext(ctx).Raw(query, args...).Scan(&hits).Error; err != nil { + return nil, err + } + return hits, nil +} + +// TextSearch is the fallback when there is no embedder, and the tie-breaker +// beside it when there is: rows whose name or title contains the label, or +// contains every word of it. +func (r *scanRepository) TextSearch(ctx context.Context, label string, limit int) ([]CatalogueHit, error) { + tables, err := r.brandTables(ctx) + if err != nil { + return nil, err + } + label = strings.ToLower(strings.TrimSpace(label)) + tokens := utils.SearchTokens(label) + if label == "" || len(tokens) == 0 { + return nil, nil + } + + var branches []string + var args []interface{} + for _, table := range sortedKeys(tables) { + cols := tables[table] + brand := strings.TrimPrefix(table, "brand_") + + hay := "LOWER(COALESCE(product_name, ''))" + if cols["title"] { + hay = "LOWER(COALESCE(product_name, '') || ' ' || COALESCE(title, ''))" + } + if cols["search_query"] { + hay = "LOWER(COALESCE(product_name, '') || ' ' || COALESCE(title, '') || ' ' || COALESCE(search_query, ''))" + } + + conds := []string{hay + " LIKE ?"} + args = append(args, "%"+label+"%") + all := make([]string, 0, len(tokens)) + for _, tok := range tokens { + all = append(all, hay+" LIKE ?") + args = append(args, "%"+tok+"%") + } + conds = append(conds, "("+strings.Join(all, " AND ")+")") + + branches = append(branches, fmt.Sprintf( + `(SELECT %s, -1::float8 AS distance FROM %s WHERE %s LIMIT %d)`, + hitColumns(brand, cols), table, strings.Join(conds, " OR "), limit)) + } + if len(branches) == 0 { + return nil, nil + } + + var hits []CatalogueHit + if err := r.catalogue.WithContext(ctx).Raw(strings.Join(branches, " UNION ALL "), args...).Scan(&hits).Error; err != nil { + return nil, err + } + return hits, nil +} + +func sortedKeys(m map[string]map[string]bool) []string { + keys := make([]string, 0, len(m)) + for k := range m { + keys = append(keys, k) + } + sort.Strings(keys) + return keys +} + +// ── cache ─────────────────────────────────────────────────────────────────── +// +// Two tiers. Redis is shared across replicas and survives a restart; the +// in-process map is there so the request after a cache hit costs no network +// round trip at all, and so a deployment without Redis still gets the +// benefit within one process. Neither tier ever holds stock. + +func scanCacheKey(kind, scope, label string) string { + sum := sha256.Sum256([]byte(strings.ToLower(strings.TrimSpace(label)))) + return "scan:" + kind + ":v1:" + scope + ":" + hex.EncodeToString(sum[:16]) +} + +func (r *scanRepository) CachedVector(ctx context.Context, model, label string) ([]float32, bool) { + key := scanCacheKey("emb", model, label) + + r.memMu.Lock() + v, ok := r.memVecs[key] + r.memMu.Unlock() + if ok { + return v, true + } + + if db.Rdb == nil { + return nil, false + } + raw, err := db.Rdb.Get(ctx, key).Bytes() + if err != nil { + return nil, false + } + if json.Unmarshal(raw, &v) != nil || len(v) == 0 { + return nil, false + } + r.remember(key, v, nil) + return v, true +} + +func (r *scanRepository) CacheVector(ctx context.Context, model, label string, v []float32) { + key := scanCacheKey("emb", model, label) + r.remember(key, v, nil) + if db.Rdb == nil { + return + } + if raw, err := json.Marshal(v); err == nil { + if err := db.Rdb.Set(ctx, key, raw, scanVectorTTL).Err(); err != nil { + log.Printf("scan: could not cache vector: %v", err) + } + } +} + +func (r *scanRepository) CachedHits(ctx context.Context, method, label string) ([]CatalogueHit, bool) { + key := scanCacheKey("hits", method, label) + + r.memMu.Lock() + h, ok := r.memHits[key] + r.memMu.Unlock() + if ok { + return h, true + } + + if db.Rdb == nil { + return nil, false + } + raw, err := db.Rdb.Get(ctx, key).Bytes() + if err != nil { + return nil, false + } + if json.Unmarshal(raw, &h) != nil { + return nil, false + } + r.remember(key, nil, h) + return h, true +} + +func (r *scanRepository) CacheHits(ctx context.Context, method, label string, hits []CatalogueHit) { + key := scanCacheKey("hits", method, label) + r.remember(key, nil, hits) + if db.Rdb == nil { + return + } + if raw, err := json.Marshal(hits); err == nil { + if err := db.Rdb.Set(ctx, key, raw, scanHitsTTL).Err(); err != nil { + log.Printf("scan: could not cache hits: %v", err) + } + } +} + +// remember writes one entry into the process-local tier. Eviction is the +// simplest thing that bounds memory: when full, drop everything. Labels are +// short-lived popularity, not a working set worth an LRU. +func (r *scanRepository) remember(key string, v []float32, h []CatalogueHit) { + r.memMu.Lock() + defer r.memMu.Unlock() + if len(r.memVecs)+len(r.memHits) >= scanMemCacheMax { + r.memVecs = make(map[string][]float32) + r.memHits = make(map[string][]CatalogueHit) + } + if v != nil { + r.memVecs[key] = v + } + if h != nil { + r.memHits[key] = h + } +} diff --git a/routes/routes.go b/routes/routes.go index 688fe5c..e7c8e57 100644 --- a/routes/routes.go +++ b/routes/routes.go @@ -21,4 +21,5 @@ func RegisterRoutes(app *fiber.App, f *facade.Facade) { RegisterCatalogueRoutes(api, f) RegisterPosRoutes(api, f) RegisterUploadRoutes(api, f) + RegisterScanRoutes(api, f) } diff --git a/routes/scanroutes.go b/routes/scanroutes.go new file mode 100644 index 0000000..8277ba2 --- /dev/null +++ b/routes/scanroutes.go @@ -0,0 +1,17 @@ +package routes + +import ( + "nearle/facade" + + "github.com/gofiber/fiber/v2" +) + +// Scan-to-order, customer app only. See controllers/scanController.go for +// the three calls and services/scanService.go for the pipeline behind them. +func RegisterScanRoutes(api fiber.Router, f *facade.Facade) { + scan := api.Group("/v1/mob/scan") + + scan.Post("/lookup", f.ScanController.Lookup) + scan.Post("/confirm", f.ScanController.Confirm) + scan.Get("/stores", f.ScanController.Stores) +} diff --git a/scratch/cataloguedims/main.go b/scratch/cataloguedims/main.go new file mode 100644 index 0000000..1b39c1f --- /dev/null +++ b/scratch/cataloguedims/main.go @@ -0,0 +1,89 @@ +// Reports how the catalogue's `embedding` columns are shaped — width, how +// many rows are filled, and a sample norm — so the embedding model Fiesta +// calls can be matched to the one that indexed the catalogue. Metadata and +// counts only, on a read-only transaction; it never writes. +// +// go run ./scratch/cataloguedims # reads CATALOGUE_DB_* from .env.production +package main + +import ( + "flag" + "fmt" + "log" + "net/url" + "os" + + "github.com/joho/godotenv" + "gorm.io/driver/postgres" + "gorm.io/gorm" +) + +func main() { + sample := flag.String("sample", "", "print one row's texts and stored vector from this table, to check which model produced it") + flag.Parse() + _ = godotenv.Load(".env.production") + dsn := url.URL{ + Scheme: "postgres", + User: url.UserPassword(os.Getenv("CATALOGUE_DB_USER"), os.Getenv("CATALOGUE_DB_PASSWORD")), + Host: os.Getenv("CATALOGUE_DB_HOST") + ":" + os.Getenv("CATALOGUE_DB_PORT"), + Path: "/" + os.Getenv("CATALOGUE_DB_NAME"), + } + q := dsn.Query() + q.Set("sslmode", "disable") + q.Set("default_transaction_read_only", "on") + dsn.RawQuery = q.Encode() + + db, err := gorm.Open(postgres.Open(dsn.String()), &gorm.Config{}) + if err != nil { + log.Fatal(err) + } + + if *sample != "" { + var row struct { + ProductName string + Title string + SearchQuery string + Embedding string + } + db.Raw(fmt.Sprintf(`SELECT product_name, COALESCE(title, '') AS title, COALESCE(search_query, '') AS search_query, + embedding::text AS embedding FROM %s WHERE embedding IS NOT NULL ORDER BY id LIMIT 1`, *sample)).Scan(&row) + fmt.Printf("product_name: %s\ntitle: %s\nsearch_query: %s\nembedding: %s\n", row.ProductName, row.Title, row.SearchQuery, row.Embedding) + return + } + + var cols []struct { + Relname string + Typname string + Atttypmod int + } + if err := db.Raw(` + SELECT c.relname, t.typname, a.atttypmod + FROM pg_attribute a + JOIN pg_class c ON c.oid = a.attrelid + JOIN pg_type t ON t.oid = a.atttypid + WHERE a.attname = 'embedding' AND c.relname LIKE 'brand\_%' + ORDER BY c.relname`).Scan(&cols).Error; err != nil { + log.Fatal(err) + } + if len(cols) == 0 { + fmt.Println("no brand_* table has an embedding column") + return + } + fmt.Printf("%-28s %-8s %5s %5s %5s\n", "table", "type", "dims", "rows", "embd") + for _, c := range cols { + var total, filled int64 + db.Raw(fmt.Sprintf(`SELECT COUNT(1) FROM %s`, c.Relname)).Scan(&total) + db.Raw(fmt.Sprintf(`SELECT COUNT(1) FROM %s WHERE embedding IS NOT NULL`, c.Relname)).Scan(&filled) + fmt.Printf("%-28s %-8s %5d %5d %5d\n", c.Relname, c.Typname, c.Atttypmod, total, filled) + } + + // nomic/bge emit unit vectors; a norm far from 1 means another pipeline. + for _, c := range cols { + var norm float64 + db.Raw(fmt.Sprintf(`SELECT vector_norm(embedding) FROM %s WHERE embedding IS NOT NULL LIMIT 1`, c.Relname)).Scan(&norm) + if norm > 0 { + fmt.Printf("sample vector norm (%s): %.4f\n", c.Relname, norm) + break + } + } +} diff --git a/services/scanService.go b/services/scanService.go new file mode 100644 index 0000000..4695212 --- /dev/null +++ b/services/scanService.go @@ -0,0 +1,721 @@ +package services + +import ( + "context" + "errors" + "fmt" + "log" + "math" + "nearle/models" + "nearle/repositories" + "nearle/utils" + "sort" + "strings" + "sync" + "time" +) + +/* +Scan-to-order: from a label Lens read off a packet to "buy it here". + +The pipeline, in the order it runs: + + 1. Who is asking, and where can they buy? The customer's registered outlets + (tenantcustomers → tenantlocations) and their position — the phone's + fix if it sent one, else the saved address. + 2. What did they scan? The label goes to the catalogue: embedded and ranked + by pgvector when a model is configured, matched on words when it is not, + and both when it is (the text match settles near-ties). The vector and + the hits are cached; a second person scanning the same packet today does + not pay for the model call. + 3. Which of THEIR outlets sell it, in which sizes, with how many on the + shelf right now? One read of nearledb keyed on the catalogue links every + imported product carries. Live stock is never cached. + 4. Rank: in-stock outlets first, nearest first, and the first of those is + the recommendation. The customer may still tap any other. + +Steps 1 and 2 touch different databases and run concurrently; the whole +lookup is bounded by scanLookupTimeout so a slow model degrades to a text +answer rather than a spinner. + +What this deliberately does not do: reserve stock. The confirm call +re-reads the ledger and, if the shelf emptied in between, points at the next +outlet — the same answer a hold would give, without a hold to expire. +*/ + +const ( + scanLookupTimeout = 5 * time.Second + scanMaxLabelLen = 200 + scanCatalogueTopK = 15 + // Below this the best hit is not shown as a match at all. + scanMinScore = 0.30 +) + +// ScanErrors the controller maps to statuses. Everything else is a 500. +var ( + ErrScanBadRequest = errors.New("scan: bad request") + ErrScanCustomerNotFound = errors.New("scan: customer not found") + ErrScanCatalogueDown = errors.New("scan: catalogue unavailable") +) + +type ScanService interface { + Lookup(ctx context.Context, req models.ScanLookupRequest) (*models.ScanLookupResponse, error) + Confirm(ctx context.Context, req models.ScanConfirmRequest) (*models.ScanConfirmResponse, error) + Stores(ctx context.Context, customerid int, lat, lng models.FlexibleString) ([]models.ScanStore, error) +} + +type scanService struct { + repo repositories.ScanRepository + embedder utils.Embedder // nil → text matching only +} + +func NewScanService(repo repositories.ScanRepository, embedder utils.Embedder) ScanService { + return &scanService{repo: repo, embedder: embedder} +} + +// ── Lookup ────────────────────────────────────────────────────────────────── + +func (s *scanService) Lookup(ctx context.Context, req models.ScanLookupRequest) (*models.ScanLookupResponse, error) { + label := strings.TrimSpace(req.Label) + if req.Customerid <= 0 { + return nil, fmt.Errorf("%w: customerid is required", ErrScanBadRequest) + } + if label == "" { + return nil, fmt.Errorf("%w: label is required", ErrScanBadRequest) + } + if len(label) > scanMaxLabelLen { + label = label[:scanMaxLabelLen] + } + + ctx, cancel := context.WithTimeout(ctx, scanLookupTimeout) + defer cancel() + + // Steps 1 and 2 in parallel — different databases, no dependency. + var ( + wg sync.WaitGroup + stores []models.ScanStore + exists bool + homeLat float64 + homeLng float64 + homeOK bool + custErr error + hits []scoredHit + method string + matchErr error + ) + wg.Add(2) + go func() { + defer wg.Done() + exists, custErr = s.repo.CustomerExists(ctx, req.Customerid) + if custErr != nil || !exists { + return + } + stores, custErr = s.repo.RegisteredStores(ctx, req.Customerid) + if custErr != nil { + return + } + // Only read the saved address when the phone sent nothing usable. + if _, _, ok := utils.ParseLatLng(string(req.Latitude), string(req.Longitude)); !ok { + homeLat, homeLng, homeOK, custErr = s.repo.CustomerHome(ctx, req.Customerid) + } + }() + go func() { + defer wg.Done() + hits, method, matchErr = s.searchCatalogue(ctx, label) + }() + wg.Wait() + + if custErr != nil { + return nil, custErr + } + if !exists { + return nil, ErrScanCustomerNotFound + } + if matchErr != nil { + return nil, matchErr + } + + lat, lng, hasPos := utils.ParseLatLng(string(req.Latitude), string(req.Longitude)) + if !hasPos && homeOK { + lat, lng, hasPos = homeLat, homeLng, true + } + + resp := &models.ScanLookupResponse{ + Label: label, + Stores: []models.ScanStoreOffer{}, + Variants: []models.ScanCatalogueMatch{}, + } + + // Verify the app's idea of the customer's tenants against the truth. + stores, resp.UnregisteredTenantids = restrictToTenants(stores, req.Tenantids) + + if len(hits) == 0 || hits[0].score < scanMinScore { + resp.Message = "We couldn't recognise that product. Try a clearer photo of the front of the pack." + return resp, nil + } + + best := hits[0] + family := catalogueFamily(hits) + resp.Match = ptr(best.toMatch(method)) + resp.Confidence = round3(best.score) + for _, h := range family { + resp.Variants = append(resp.Variants, h.toMatch(method)) + } + + if len(stores) == 0 { + resp.Message = "You haven't joined a store yet. Scan a store's QR code in the app to shop from it." + return resp, nil + } + + // Step 3: what those outlets have. + keys := make([]repositories.CatalogueKey, 0, len(family)) + names := make([]string, 0, len(family)) + for _, h := range family { + keys = append(keys, repositories.CatalogueKey{Brand: h.Brand, Catalogueid: h.ID, Imageid: h.ImageID}) + names = append(names, h.ProductName) + } + locationids := make([]int, 0, len(stores)) + for _, st := range stores { + locationids = append(locationids, st.Locationid) + } + rows, err := s.repo.StoreOptions(ctx, locationids, keys, names) + if err != nil { + return nil, err + } + + // Step 4: rank. + offers := buildOffers(stores, rows, keys, lat, lng, hasPos) + if req.Limit > 0 && len(offers) > req.Limit { + offers = offers[:req.Limit] + } + resp.Stores = offers + for i := range offers { + if offers[i].Recommended { + resp.Available = true + resp.RecommendedLocationid = offers[i].Locationid + break + } + } + + switch { + case resp.Available: + resp.Message = fmt.Sprintf("Available at %d of your stores.", countAvailable(offers)) + case len(offers) > 0: + resp.Message = "Your stores sell this but it's out of stock right now." + default: + resp.Message = "None of your stores sell this product yet." + } + return resp, nil +} + +// ── Confirm ───────────────────────────────────────────────────────────────── + +func (s *scanService) Confirm(ctx context.Context, req models.ScanConfirmRequest) (*models.ScanConfirmResponse, error) { + if req.Customerid <= 0 || req.Tenantid <= 0 || req.Locationid <= 0 || req.Productid <= 0 { + return nil, fmt.Errorf("%w: customerid, tenantid, locationid and productid are required", ErrScanBadRequest) + } + if req.Quantity <= 0 { + req.Quantity = 1 + } + + ctx, cancel := context.WithTimeout(ctx, scanLookupTimeout) + defer cancel() + + stores, err := s.repo.RegisteredStores(ctx, req.Customerid) + if err != nil { + return nil, err + } + resp := &models.ScanConfirmResponse{Requested: req.Quantity} + + var chosen *models.ScanStore + for i := range stores { + if stores[i].Tenantid == req.Tenantid && stores[i].Locationid == req.Locationid { + chosen = &stores[i] + break + } + } + if chosen == nil { + resp.Reason = "store_not_registered" + resp.Message = "You're not registered with that store. Scan its QR code first." + return resp, nil + } + resp.Store = chosen + + row, err := s.repo.ProductAt(ctx, req.Tenantid, req.Locationid, req.Productid) + if err != nil { + return nil, err + } + if row == nil { + resp.Reason = "not_sold_here" + resp.Message = "That store doesn't sell this product." + return resp, nil + } + opt := optionFromRow(*row, nil) + resp.Option = &opt + + if row.Stock >= req.Quantity { + resp.Ok = true + resp.Reason = "in_stock" + resp.Message = "In stock." + return resp, nil + } + if row.Stock > 0 { + resp.Reason = "insufficient_stock" + resp.Message = fmt.Sprintf("Only %d left at %s.", row.Stock, chosen.Locationname) + } else { + resp.Reason = "out_of_stock" + resp.Message = fmt.Sprintf("Out of stock at %s.", chosen.Locationname) + } + + // The same product elsewhere, nearest first, with enough of it. + lat, lng, hasPos := utils.ParseLatLng(string(req.Latitude), string(req.Longitude)) + if !hasPos { + if hl, hg, ok, err := s.repo.CustomerHome(ctx, req.Customerid); err == nil && ok { + lat, lng, hasPos = hl, hg, true + } + } + others := make([]models.ScanStore, 0, len(stores)) + locationids := make([]int, 0, len(stores)) + for _, st := range stores { + if st.Locationid == req.Locationid { + continue + } + others = append(others, st) + locationids = append(locationids, st.Locationid) + } + if len(others) == 0 { + return resp, nil + } + + key := repositories.CatalogueKey{Brand: row.Productbrand, Catalogueid: row.Catalogueid, Imageid: row.Imageid} + rows, err := s.repo.StoreOptions(ctx, locationids, []repositories.CatalogueKey{key}, []string{row.Productname}) + if err != nil { + return nil, err + } + // Only the product itself — not its other sizes: the customer chose a + // size and a different one is not a substitute they asked for. + var same []repositories.StoreOptionRow + for _, r := range rows { + if r.Stock >= req.Quantity && sameStoreProduct(*row, r) { + same = append(same, r) + } + } + offers := buildOffers(others, same, []repositories.CatalogueKey{key}, lat, lng, hasPos) + for i := range offers { + if offers[i].Available { + offers[i].Recommended = true + resp.Alternative = &offers[i] + resp.Message += fmt.Sprintf(" %s has it", offers[i].Locationname) + if offers[i].DistanceKm >= 0 { + resp.Message += fmt.Sprintf(" (%.1f km away)", offers[i].DistanceKm) + } + resp.Message += "." + break + } + } + return resp, nil +} + +// ── Stores ────────────────────────────────────────────────────────────────── + +func (s *scanService) Stores(ctx context.Context, customerid int, latStr, lngStr models.FlexibleString) ([]models.ScanStore, error) { + if customerid <= 0 { + return nil, fmt.Errorf("%w: customerid is required", ErrScanBadRequest) + } + ctx, cancel := context.WithTimeout(ctx, scanLookupTimeout) + defer cancel() + + exists, err := s.repo.CustomerExists(ctx, customerid) + if err != nil { + return nil, err + } + if !exists { + return nil, ErrScanCustomerNotFound + } + stores, err := s.repo.RegisteredStores(ctx, customerid) + if err != nil { + return nil, err + } + lat, lng, hasPos := utils.ParseLatLng(string(latStr), string(lngStr)) + if !hasPos { + if hl, hg, ok, err := s.repo.CustomerHome(ctx, customerid); err == nil && ok { + lat, lng, hasPos = hl, hg, true + } + } + for i := range stores { + stores[i].DistanceKm = distanceKm(stores[i], lat, lng, hasPos) + } + sort.SliceStable(stores, func(i, j int) bool { return nearer(stores[i].DistanceKm, stores[j].DistanceKm) }) + if stores == nil { + stores = []models.ScanStore{} + } + return stores, nil +} + +// ── Catalogue search ──────────────────────────────────────────────────────── + +type scoredHit struct { + repositories.CatalogueHit + score float64 +} + +func (h scoredHit) toMatch(method string) models.ScanCatalogueMatch { + return models.ScanCatalogueMatch{ + Brand: h.Brand, + Catalogueid: h.ID, + Imageid: h.ImageID, + ProductName: h.ProductName, + Title: h.Title, + Category: h.Category, + Size: h.Size, + VariantKey: h.VariantKey, + Image: h.ImageURL, + Score: round3(h.score), + Method: method, + } +} + +// searchCatalogue returns hits best-first and the method that produced them. +// +// Vector and text are combined as max(vector, text) with a small bonus when +// both agree. Max rather than a weighted sum so that a text-only hit — the +// exact product name typed on the pack — is never dragged below a vaguely +// similar vector neighbour, and a vector hit is never punished for a label +// Lens spelled slightly differently from the catalogue. +func (s *scanService) searchCatalogue(ctx context.Context, label string) ([]scoredHit, string, error) { + method := "text" + useVector := s.embedder != nil && s.repo.VectorSearchAvailable() + if useVector { + method = "vector+text" + } + + if cached, ok := s.repo.CachedHits(ctx, method+":"+s.modelName(), label); ok { + return scoreCachedHits(cached), method, nil + } + + tokens := utils.SearchTokens(label) + byKey := make(map[string]*scoredHit) + keyOf := func(h repositories.CatalogueHit) string { return h.Brand + "#" + fmt.Sprint(h.ID) } + + if useVector { + vec, err := s.embed(ctx, label) + if err != nil { + // Degrade, loudly in the log and quietly to the customer: a text + // answer now beats a vector answer never. + log.Printf("scan: embedding %q failed, falling back to text: %v", label, err) + method = "text" + } else { + vhits, err := s.repo.VectorSearch(ctx, vec, scanCatalogueTopK) + if err != nil { + log.Printf("scan: vector search failed, falling back to text: %v", err) + method = "text" + } + for _, h := range vhits { + sim := 1 - h.Distance // cosine distance → similarity + if sim < 0 { + sim = 0 + } + byKey[keyOf(h)] = &scoredHit{CatalogueHit: h, score: sim} + } + } + } + + thits, err := s.repo.TextSearch(ctx, label, scanCatalogueTopK) + if err != nil { + if errors.Is(err, repositories.ErrCatalogueDBUnavailable) { + return nil, method, ErrScanCatalogueDown + } + if len(byKey) == 0 { + return nil, method, err + } + log.Printf("scan: text search failed, vector only: %v", err) + } + for _, h := range thits { + ts := textScore(h, label, tokens) + if existing, ok := byKey[keyOf(h)]; ok { + existing.score = math.Min(1, math.Max(existing.score, ts)+0.10) + continue + } + byKey[keyOf(h)] = &scoredHit{CatalogueHit: h, score: ts} + } + + hits := make([]scoredHit, 0, len(byKey)) + for _, h := range byKey { + hits = append(hits, *h) + } + sortHits(hits) + + // Remember the ranked rows with their score folded into Distance, so the + // cache does not need a second shape. + toCache := make([]repositories.CatalogueHit, 0, len(hits)) + for _, h := range hits { + c := h.CatalogueHit + c.Distance = 1 - h.score + toCache = append(toCache, c) + } + s.repo.CacheHits(ctx, method+":"+s.modelName(), label, toCache) + + return hits, method, nil +} + +func scoreCachedHits(cached []repositories.CatalogueHit) []scoredHit { + hits := make([]scoredHit, 0, len(cached)) + for _, c := range cached { + hits = append(hits, scoredHit{CatalogueHit: c, score: 1 - c.Distance}) + } + sortHits(hits) + return hits +} + +func sortHits(hits []scoredHit) { + sort.SliceStable(hits, func(i, j int) bool { + if hits[i].score != hits[j].score { + return hits[i].score > hits[j].score + } + return hits[i].ProductName < hits[j].ProductName + }) +} + +func (s *scanService) modelName() string { + if s.embedder == nil { + return "none" + } + return s.embedder.Model() +} + +func (s *scanService) embed(ctx context.Context, label string) ([]float32, error) { + if v, ok := s.repo.CachedVector(ctx, s.embedder.Model(), label); ok { + return v, nil + } + v, err := s.embedder.Embed(ctx, label) + if err != nil { + return nil, err + } + s.repo.CacheVector(ctx, s.embedder.Model(), label, v) + return v, nil +} + +// textScore is how well a catalogue row's name matches the words Lens read. +// The whole label as a substring of the name is near-certain; otherwise the +// share of label words found in name+title, scaled so that "all of them" +// stops short of the substring case. +func textScore(h repositories.CatalogueHit, label string, tokens []string) float64 { + name := strings.ToLower(h.ProductName) + hay := name + " " + strings.ToLower(h.Title) + label = strings.ToLower(strings.TrimSpace(label)) + if label != "" && strings.Contains(name, label) { + return 0.95 + } + if len(tokens) == 0 { + return 0 + } + found := 0 + for _, t := range tokens { + if strings.Contains(hay, t) { + found++ + } + } + return 0.8 * float64(found) / float64(len(tokens)) +} + +// catalogueFamily is the best hit and its other pack sizes: same brand, and +// the same variant_key when the catalogue assigned one, else the same name. +// Every member is a separate catalogue row a shop may have imported. +func catalogueFamily(hits []scoredHit) []scoredHit { + if len(hits) == 0 { + return nil + } + best := hits[0] + family := []scoredHit{best} + for _, h := range hits[1:] { + if h.Brand != best.Brand { + continue + } + switch { + case best.VariantKey != "" && h.VariantKey != "": + if h.VariantKey == best.VariantKey { + family = append(family, h) + } + case strings.EqualFold(strings.TrimSpace(h.ProductName), strings.TrimSpace(best.ProductName)): + family = append(family, h) + } + } + return family +} + +// ── Offers ────────────────────────────────────────────────────────────────── + +// buildOffers turns outlet rows into ranked offers: one per outlet that had +// any row, in-stock outlets first, nearest first, the first in-stock one +// recommended. +func buildOffers(stores []models.ScanStore, rows []repositories.StoreOptionRow, keys []repositories.CatalogueKey, lat, lng float64, hasPos bool) []models.ScanStoreOffer { + byLocation := make(map[int][]repositories.StoreOptionRow) + for _, r := range rows { + byLocation[r.Locationid] = append(byLocation[r.Locationid], r) + } + + offers := make([]models.ScanStoreOffer, 0, len(byLocation)) + for _, st := range stores { + rs := byLocation[st.Locationid] + if len(rs) == 0 { + continue + } + st.DistanceKm = distanceKm(st, lat, lng, hasPos) + offer := models.ScanStoreOffer{ScanStore: st, Options: []models.ScanOption{}} + + // A product can arrive twice — as a direct match and as a size of + // another match. Keep the direct one; it carries the better label. + seen := make(map[int]int) + for _, r := range rs { + opt := optionFromRow(r, keys) + if idx, dup := seen[r.Productid]; dup { + if offer.Options[idx].IsVariant && !opt.IsVariant { + offer.Options[idx] = opt + } + continue + } + seen[r.Productid] = len(offer.Options) + offer.Options = append(offer.Options, opt) + if opt.Available { + offer.Available = true + } + } + // Direct matches first, then sizes; in stock before out. + sort.SliceStable(offer.Options, func(i, j int) bool { + a, b := offer.Options[i], offer.Options[j] + if a.Available != b.Available { + return a.Available + } + if a.IsVariant != b.IsVariant { + return !a.IsVariant + } + return a.Productname < b.Productname + }) + offers = append(offers, offer) + } + + sort.SliceStable(offers, func(i, j int) bool { + a, b := offers[i], offers[j] + if a.Available != b.Available { + return a.Available + } + if a.DistanceKm != b.DistanceKm { + return nearer(a.DistanceKm, b.DistanceKm) + } + return a.Locationname < b.Locationname + }) + for i := range offers { + if offers[i].Available { + offers[i].Recommended = true + break + } + } + return offers +} + +func optionFromRow(r repositories.StoreOptionRow, keys []repositories.CatalogueKey) models.ScanOption { + opt := models.ScanOption{ + Productid: r.Productid, + Productname: r.Productname, + Size: strings.TrimSpace(r.Unitvalue + " " + r.Productunit), + Price: r.Price, + Stock: r.Stock, + Available: r.Stock > 0, + Image: r.Productimage, + IsVariant: r.Parentid > 0, + Variantname: r.Variantname, + MatchedBy: "name", + } + if r.Parentid > 0 { + opt.MatchedBy = fmt.Sprintf("variant-of:%d", r.Parentid) + return opt + } + for _, k := range keys { + if k.Imageid != "" && k.Imageid == r.Imageid { + opt.MatchedBy = "imageid" + return opt + } + if k.Brand != "" && strings.EqualFold(k.Brand, r.Productbrand) && k.Catalogueid == r.Catalogueid && k.Catalogueid > 0 { + opt.MatchedBy = "brand+catalogueid" + return opt + } + } + return opt +} + +// sameStoreProduct is the cross-tenant identity of a product: the catalogue key +// when both rows carry one, the name otherwise. +func sameStoreProduct(a, b repositories.StoreOptionRow) bool { + if a.Imageid != "" && b.Imageid != "" { + return a.Imageid == b.Imageid + } + if a.Catalogueid > 0 && b.Catalogueid > 0 && a.Productbrand != "" { + return a.Catalogueid == b.Catalogueid && strings.EqualFold(a.Productbrand, b.Productbrand) + } + return strings.EqualFold(strings.TrimSpace(a.Productname), strings.TrimSpace(b.Productname)) +} + +// restrictToTenants keeps the outlets whose tenant the app named, and +// reports the names it got wrong. An app list that matches nothing is +// treated as stale rather than as "no stores": all registered outlets are +// used and every id it sent is reported. +func restrictToTenants(stores []models.ScanStore, tenantids []int) ([]models.ScanStore, []int) { + if len(tenantids) == 0 { + return stores, nil + } + registered := make(map[int]bool) + for _, st := range stores { + registered[st.Tenantid] = true + } + wanted := make(map[int]bool) + var unregistered []int + for _, id := range tenantids { + if registered[id] { + wanted[id] = true + } else { + unregistered = append(unregistered, id) + } + } + if len(wanted) == 0 { + return stores, unregistered + } + kept := make([]models.ScanStore, 0, len(stores)) + for _, st := range stores { + if wanted[st.Tenantid] { + kept = append(kept, st) + } + } + return kept, unregistered +} + +func distanceKm(st models.ScanStore, lat, lng float64, hasPos bool) float64 { + if !hasPos || (st.Latitude == 0 && st.Longitude == 0) { + return -1 + } + return math.Round(utils.HaversineKm(lat, lng, st.Latitude, st.Longitude)*100) / 100 +} + +// nearer orders distances with -1 (unknown) last. +func nearer(a, b float64) bool { + if a < 0 { + return false + } + if b < 0 { + return true + } + return a < b +} + +func countAvailable(offers []models.ScanStoreOffer) int { + n := 0 + for _, o := range offers { + if o.Available { + n++ + } + } + return n +} + +func round3(f float64) float64 { return math.Round(f*1000) / 1000 } + +func ptr[T any](v T) *T { return &v } diff --git a/services/scan_test.go b/services/scan_test.go new file mode 100644 index 0000000..6c5fe98 --- /dev/null +++ b/services/scan_test.go @@ -0,0 +1,438 @@ +package services + +import ( + "context" + "errors" + "testing" + + "nearle/models" + "nearle/repositories" +) + +/* +The scan pipeline has three decisions worth defending: which catalogue rows +count as "the product", which of the customer's outlets get shown and in what +order, and what happens when the outlet they tapped has run out. Everything +below drives those through a fake repository; the SQL itself is exercised +against a real database in scratch/ when there is one. +*/ + +// fakeScanRepo answers from fixtures and records what it was asked. +type fakeScanRepo struct { + exists bool + homeLat float64 + homeLng float64 + homeOK bool + stores []models.ScanStore + text []repositories.CatalogueHit + vector []repositories.CatalogueHit + hasVec bool + options []repositories.StoreOptionRow + at map[int]*repositories.StoreOptionRow // productid → row + + askedKeys []repositories.CatalogueKey + askedNames []string + askedLocs []int + cachedHits map[string][]repositories.CatalogueHit +} + +func (f *fakeScanRepo) CustomerExists(context.Context, int) (bool, error) { return f.exists, nil } +func (f *fakeScanRepo) CustomerHome(context.Context, int) (float64, float64, bool, error) { + return f.homeLat, f.homeLng, f.homeOK, nil +} +func (f *fakeScanRepo) RegisteredStores(context.Context, int) ([]models.ScanStore, error) { + out := make([]models.ScanStore, len(f.stores)) + copy(out, f.stores) + return out, nil +} +func (f *fakeScanRepo) StoreOptions(_ context.Context, locs []int, keys []repositories.CatalogueKey, names []string) ([]repositories.StoreOptionRow, error) { + f.askedLocs, f.askedKeys, f.askedNames = locs, keys, names + allowed := make(map[int]bool) + for _, l := range locs { + allowed[l] = true + } + var out []repositories.StoreOptionRow + for _, o := range f.options { + if allowed[o.Locationid] { + out = append(out, o) + } + } + return out, nil +} +func (f *fakeScanRepo) ProductAt(_ context.Context, _, _, productid int) (*repositories.StoreOptionRow, error) { + return f.at[productid], nil +} +func (f *fakeScanRepo) VectorSearch(context.Context, []float32, int) ([]repositories.CatalogueHit, error) { + return f.vector, nil +} +func (f *fakeScanRepo) TextSearch(context.Context, string, int) ([]repositories.CatalogueHit, error) { + return f.text, nil +} +func (f *fakeScanRepo) VectorSearchAvailable() bool { return f.hasVec } +func (f *fakeScanRepo) CachedVector(context.Context, string, string) ([]float32, bool) { + return nil, false +} +func (f *fakeScanRepo) CacheVector(context.Context, string, string, []float32) {} +func (f *fakeScanRepo) CachedHits(_ context.Context, method, label string) ([]repositories.CatalogueHit, bool) { + h, ok := f.cachedHits[method+"|"+label] + return h, ok +} +func (f *fakeScanRepo) CacheHits(_ context.Context, method, label string, hits []repositories.CatalogueHit) { + if f.cachedHits == nil { + f.cachedHits = make(map[string][]repositories.CatalogueHit) + } + f.cachedHits[method+"|"+label] = hits +} + +type fakeEmbedder struct { + vec []float32 + err error +} + +func (e fakeEmbedder) Embed(context.Context, string) ([]float32, error) { return e.vec, e.err } +func (e fakeEmbedder) Model() string { return "fake-model" } + +// A customer in Peelamedu with three outlets: one 1 km away, one 4 km away, +// one across town with no coordinates on file. +func fixtureStores() []models.ScanStore { + return []models.ScanStore{ + {Tenantid: 1, Tenantname: "Suriya Store", Locationid: 10, Locationname: "Peelamedu", Latitude: 11.030, Longitude: 77.030}, + {Tenantid: 2, Tenantname: "R Mart", Locationid: 20, Locationname: "Hopes", Latitude: 11.010, Longitude: 77.000}, + {Tenantid: 3, Tenantname: "Daily Needs", Locationid: 30, Locationname: "Gandhipuram"}, + } +} + +var milkBikis = repositories.CatalogueHit{Brand: "britannia", ID: 7, ProductName: "Milk Bikis", Size: "100 g", VariantKey: "milk_bikis", ImageID: "britannia_milk_bikis_100g", Distance: 0.05} +var milkBikis200 = repositories.CatalogueHit{Brand: "britannia", ID: 8, ProductName: "Milk Bikis", Size: "200 g", VariantKey: "milk_bikis", ImageID: "britannia_milk_bikis_200g", Distance: 0.12} +var goodDay = repositories.CatalogueHit{Brand: "britannia", ID: 9, ProductName: "Good Day Butter", VariantKey: "good_day", ImageID: "britannia_good_day", Distance: 0.40} + +func newLookupFixture() *fakeScanRepo { + return &fakeScanRepo{ + exists: true, + stores: fixtureStores(), + hasVec: true, + vector: []repositories.CatalogueHit{milkBikis, milkBikis200, goodDay}, + text: []repositories.CatalogueHit{milkBikis}, + options: []repositories.StoreOptionRow{ + // Nearest outlet: sells it, but the shelf is empty. + {Tenantid: 1, Locationid: 10, Productid: 100, Productname: "Milk Bikis 100g", Productbrand: "britannia", Catalogueid: 7, Imageid: "britannia_milk_bikis_100g", Unitvalue: "100", Productunit: "g", Price: 10, Stock: 0}, + // 4 km away: has both sizes. + {Tenantid: 2, Locationid: 20, Productid: 200, Productname: "Milk Bikis 100g", Productbrand: "britannia", Catalogueid: 7, Imageid: "britannia_milk_bikis_100g", Unitvalue: "100", Productunit: "g", Price: 12, Stock: 6}, + {Tenantid: 2, Locationid: 20, Productid: 201, Productname: "Milk Bikis 200g", Productbrand: "britannia", Catalogueid: 8, Imageid: "britannia_milk_bikis_200g", Unitvalue: "200", Productunit: "g", Price: 22, Stock: 3}, + // Same product listed again as a size under the 100g row — must not + // appear twice. + {Tenantid: 2, Locationid: 20, Productid: 201, Productname: "Milk Bikis 200g", Productbrand: "britannia", Catalogueid: 8, Imageid: "britannia_milk_bikis_200g", Unitvalue: "200", Productunit: "g", Price: 22, Stock: 3, Parentid: 200, Variantname: "200 g"}, + // Unknown distance, in stock. + {Tenantid: 3, Locationid: 30, Productid: 300, Productname: "Milk Bikis", Productbrand: "britannia", Catalogueid: 7, Imageid: "britannia_milk_bikis_100g", Price: 11, Stock: 2}, + }, + } +} + +func TestLookupRecommendsTheNearestOutletWithStock(t *testing.T) { + repo := newLookupFixture() + svc := NewScanService(repo, fakeEmbedder{vec: []float32{0.1, 0.2}}) + + resp, err := svc.Lookup(context.Background(), models.ScanLookupRequest{ + Customerid: 5, Label: "Milk Bikis", Latitude: "11.035", Longitude: "77.035", + }) + if err != nil { + t.Fatal(err) + } + if resp.Match == nil || resp.Match.ProductName != "Milk Bikis" || resp.Match.Catalogueid != 7 { + t.Fatalf("expected Milk Bikis 100 g as the match, got %+v", resp.Match) + } + if resp.Match.Method != "vector+text" { + t.Errorf("method = %q, want vector+text", resp.Match.Method) + } + if len(resp.Variants) != 2 { + t.Errorf("the catalogue family should be the two Milk Bikis sizes, got %d: %+v", len(resp.Variants), resp.Variants) + } + if !resp.Available || resp.RecommendedLocationid != 20 { + t.Fatalf("Hopes (4 km, in stock) should be recommended over Peelamedu (1 km, empty); got available=%v recommended=%d", resp.Available, resp.RecommendedLocationid) + } + + // Order: in stock first (Hopes, then Gandhipuram with no distance), then + // the empty nearest outlet. + var order []int + for _, o := range resp.Stores { + order = append(order, o.Locationid) + } + if len(order) != 3 || order[0] != 20 || order[1] != 30 || order[2] != 10 { + t.Fatalf("store order = %v, want [20 30 10]", order) + } + if !resp.Stores[0].Recommended || resp.Stores[1].Recommended || resp.Stores[2].Recommended { + t.Error("exactly the first in-stock offer should be recommended") + } + if resp.Stores[2].Available { + t.Error("Peelamedu has no stock and must not be available") + } + if resp.Stores[1].DistanceKm != -1 { + t.Errorf("an outlet with no coordinates reports distance -1, got %v", resp.Stores[1].DistanceKm) + } + if resp.Stores[0].DistanceKm <= 0 || resp.Stores[0].DistanceKm > 10 { + t.Errorf("Hopes should be a few km away, got %v", resp.Stores[0].DistanceKm) + } + + hopes := resp.Stores[0] + if len(hopes.Options) != 2 { + t.Fatalf("Hopes should offer two sizes once, got %d: %+v", len(hopes.Options), hopes.Options) + } + if hopes.Options[0].Productid != 200 || hopes.Options[0].MatchedBy != "imageid" || hopes.Options[0].Size != "100 g" { + t.Errorf("first option should be the direct 100 g match by imageid, got %+v", hopes.Options[0]) + } + if hopes.Options[1].Productid != 201 || hopes.Options[1].IsVariant { + t.Errorf("the 200 g row seen both directly and as a size keeps the direct form, got %+v", hopes.Options[1]) + } + + // Every catalogue size was asked for at every registered outlet. + if len(repo.askedKeys) != 2 || len(repo.askedLocs) != 3 { + t.Errorf("asked keys=%d locs=%d, want 2 and 3", len(repo.askedKeys), len(repo.askedLocs)) + } +} + +func TestLookupFallsBackToTextWhenTheModelFails(t *testing.T) { + repo := newLookupFixture() + svc := NewScanService(repo, fakeEmbedder{err: errors.New("429 rate limited")}) + + resp, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Customerid: 5, Label: "Milk Bikis"}) + if err != nil { + t.Fatal(err) + } + if resp.Match == nil || resp.Match.Method != "text" || resp.Match.Catalogueid != 7 { + t.Fatalf("expected a text-only match on Milk Bikis, got %+v", resp.Match) + } + if resp.Confidence < 0.9 { + t.Errorf("the label is the whole product name; confidence should be high, got %v", resp.Confidence) + } +} + +func TestLookupWithoutAnEmbedderIsTextOnly(t *testing.T) { + repo := newLookupFixture() + repo.hasVec = false + svc := NewScanService(repo, nil) + + resp, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Customerid: 5, Label: "milk bikis"}) + if err != nil { + t.Fatal(err) + } + if resp.Match == nil || resp.Match.Method != "text" { + t.Fatalf("expected text method, got %+v", resp.Match) + } +} + +func TestLookupUsesTheCachedRankingOnASecondScan(t *testing.T) { + repo := newLookupFixture() + svc := NewScanService(repo, fakeEmbedder{vec: []float32{0.1}}) + + if _, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Customerid: 5, Label: "Milk Bikis"}); err != nil { + t.Fatal(err) + } + // Take the catalogue away: the second scan must be served from cache. + repo.vector, repo.text = nil, nil + resp, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Customerid: 5, Label: "Milk Bikis"}) + if err != nil { + t.Fatal(err) + } + if resp.Match == nil || resp.Match.Catalogueid != 7 { + t.Fatalf("second scan should hit the cache, got %+v", resp.Match) + } +} + +func TestLookupRefusesAWeakMatch(t *testing.T) { + repo := newLookupFixture() + repo.vector = []repositories.CatalogueHit{{Brand: "x", ID: 1, ProductName: "Something Else", Distance: 0.9}} + repo.text = nil + svc := NewScanService(repo, fakeEmbedder{vec: []float32{0.1}}) + + resp, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Customerid: 5, Label: "zzz"}) + if err != nil { + t.Fatal(err) + } + if resp.Match != nil || resp.Available || len(resp.Stores) != 0 { + t.Fatalf("a 0.1 similarity is not a match; got %+v", resp) + } + if repo.askedLocs != nil { + t.Error("no outlet should be queried without a match") + } +} + +func TestLookupVerifiesTheAppsTenantList(t *testing.T) { + repo := newLookupFixture() + svc := NewScanService(repo, fakeEmbedder{vec: []float32{0.1}}) + + // The app says tenant 2 and tenant 99; 99 is not registered. + resp, err := svc.Lookup(context.Background(), models.ScanLookupRequest{ + Customerid: 5, Label: "Milk Bikis", Tenantids: []int{2, 99}, + }) + if err != nil { + t.Fatal(err) + } + if len(resp.UnregisteredTenantids) != 1 || resp.UnregisteredTenantids[0] != 99 { + t.Errorf("99 should be reported as unregistered, got %v", resp.UnregisteredTenantids) + } + if len(resp.Stores) != 1 || resp.Stores[0].Tenantid != 2 { + t.Errorf("only tenant 2's outlet should be offered, got %+v", resp.Stores) + } + + // A list that matches nothing is stale, not a request for nothing. + resp, _ = svc.Lookup(context.Background(), models.ScanLookupRequest{ + Customerid: 5, Label: "Milk Bikis", Tenantids: []int{98, 99}, + }) + if len(resp.Stores) != 3 || len(resp.UnregisteredTenantids) != 2 { + t.Errorf("a wholly stale list falls back to every registered outlet, got %d stores / %v", len(resp.Stores), resp.UnregisteredTenantids) + } +} + +func TestLookupWithNoStoresStillReturnsTheMatch(t *testing.T) { + repo := newLookupFixture() + repo.stores = nil + svc := NewScanService(repo, fakeEmbedder{vec: []float32{0.1}}) + + resp, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Customerid: 5, Label: "Milk Bikis"}) + if err != nil { + t.Fatal(err) + } + if resp.Match == nil || resp.Available || len(resp.Stores) != 0 { + t.Fatalf("match without stores, got %+v", resp) + } +} + +func TestLookupRejectsBadInput(t *testing.T) { + svc := NewScanService(newLookupFixture(), nil) + if _, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Label: "x"}); !errors.Is(err, ErrScanBadRequest) { + t.Errorf("missing customerid: %v", err) + } + if _, err := svc.Lookup(context.Background(), models.ScanLookupRequest{Customerid: 1, Label: " "}); !errors.Is(err, ErrScanBadRequest) { + t.Errorf("blank label: %v", err) + } + repo := newLookupFixture() + repo.exists = false + if _, err := NewScanService(repo, nil).Lookup(context.Background(), models.ScanLookupRequest{Customerid: 1, Label: "x"}); !errors.Is(err, ErrScanCustomerNotFound) { + t.Errorf("unknown customer: %v", err) + } +} + +func TestConfirmHoldsWhenStockIsThere(t *testing.T) { + repo := newLookupFixture() + repo.at = map[int]*repositories.StoreOptionRow{200: &repo.options[1]} + svc := NewScanService(repo, nil) + + resp, err := svc.Confirm(context.Background(), models.ScanConfirmRequest{ + Customerid: 5, Tenantid: 2, Locationid: 20, Productid: 200, Quantity: 4, + }) + if err != nil { + t.Fatal(err) + } + if !resp.Ok || resp.Reason != "in_stock" || resp.Option == nil || resp.Option.Stock != 6 { + t.Fatalf("4 of 6 should be fine, got %+v", resp) + } +} + +func TestConfirmPointsAtTheNextOutletWhenTheShelfIsEmpty(t *testing.T) { + repo := newLookupFixture() + repo.at = map[int]*repositories.StoreOptionRow{100: &repo.options[0]} // Peelamedu, stock 0 + svc := NewScanService(repo, nil) + + resp, err := svc.Confirm(context.Background(), models.ScanConfirmRequest{ + Customerid: 5, Tenantid: 1, Locationid: 10, Productid: 100, Quantity: 1, + Latitude: "11.035", Longitude: "77.035", + }) + if err != nil { + t.Fatal(err) + } + if resp.Ok || resp.Reason != "out_of_stock" { + t.Fatalf("expected out_of_stock, got %+v", resp) + } + if resp.Alternative == nil || resp.Alternative.Locationid != 20 { + t.Fatalf("Hopes is the nearest outlet with the same 100 g product, got %+v", resp.Alternative) + } + if len(resp.Alternative.Options) != 1 || resp.Alternative.Options[0].Productid != 200 { + t.Errorf("the alternative carries the same product, not its other sizes: %+v", resp.Alternative.Options) + } + if !resp.Alternative.Recommended { + t.Error("the alternative is the recommendation") + } + + // Asking for more than anyone has: no alternative, honest reason. + resp, _ = svc.Confirm(context.Background(), models.ScanConfirmRequest{ + Customerid: 5, Tenantid: 1, Locationid: 10, Productid: 100, Quantity: 50, + }) + if resp.Alternative != nil { + t.Errorf("nobody has 50; got alternative %+v", resp.Alternative) + } +} + +func TestConfirmReportsInsufficientRatherThanOut(t *testing.T) { + repo := newLookupFixture() + repo.at = map[int]*repositories.StoreOptionRow{300: &repo.options[4]} // Gandhipuram, stock 2 + svc := NewScanService(repo, nil) + + resp, err := svc.Confirm(context.Background(), models.ScanConfirmRequest{ + Customerid: 5, Tenantid: 3, Locationid: 30, Productid: 300, Quantity: 5, + }) + if err != nil { + t.Fatal(err) + } + if resp.Ok || resp.Reason != "insufficient_stock" { + t.Fatalf("2 in stock, 5 asked: want insufficient_stock, got %+v", resp) + } + if resp.Alternative == nil || resp.Alternative.Locationid != 20 { + t.Errorf("Hopes has 6 of the same product, got %+v", resp.Alternative) + } +} + +func TestConfirmRefusesAnUnregisteredStoreAndAnUnsoldProduct(t *testing.T) { + repo := newLookupFixture() + repo.at = map[int]*repositories.StoreOptionRow{} + svc := NewScanService(repo, nil) + + resp, _ := svc.Confirm(context.Background(), models.ScanConfirmRequest{Customerid: 5, Tenantid: 9, Locationid: 90, Productid: 1}) + if resp.Ok || resp.Reason != "store_not_registered" { + t.Errorf("got %+v", resp) + } + resp, _ = svc.Confirm(context.Background(), models.ScanConfirmRequest{Customerid: 5, Tenantid: 1, Locationid: 10, Productid: 424242}) + if resp.Ok || resp.Reason != "not_sold_here" { + t.Errorf("got %+v", resp) + } +} + +func TestStoresAreNearestFirstWithUnknownLast(t *testing.T) { + repo := newLookupFixture() + svc := NewScanService(repo, nil) + + stores, err := svc.Stores(context.Background(), 5, "11.035", "77.035") + if err != nil { + t.Fatal(err) + } + if len(stores) != 3 || stores[0].Locationid != 10 || stores[1].Locationid != 20 || stores[2].Locationid != 30 { + t.Fatalf("want [10 20 30], got %+v", stores) + } + + // No fix from the phone, saved address used instead. + repo.homeLat, repo.homeLng, repo.homeOK = 11.012, 77.001, true + stores, _ = svc.Stores(context.Background(), 5, "", "") + if stores[0].Locationid != 20 { + t.Errorf("from the saved address Hopes is nearest, got %+v", stores[0]) + } +} + +func TestCatalogueFamilyGroupsByVariantKeyThenName(t *testing.T) { + hits := []scoredHit{ + {CatalogueHit: milkBikis, score: 0.95}, + {CatalogueHit: goodDay, score: 0.6}, + {CatalogueHit: milkBikis200, score: 0.88}, + {CatalogueHit: repositories.CatalogueHit{Brand: "parle", ProductName: "Milk Bikis", VariantKey: "milk_bikis"}, score: 0.5}, + } + family := catalogueFamily(hits) + if len(family) != 2 || family[1].ID != 8 { + t.Fatalf("family should be the two britannia sizes, got %+v", family) + } + + // No variant keys: fall back to the name. + a := scoredHit{CatalogueHit: repositories.CatalogueHit{Brand: "b", ID: 1, ProductName: "Honey"}} + b := scoredHit{CatalogueHit: repositories.CatalogueHit{Brand: "b", ID: 2, ProductName: "honey "}} + c := scoredHit{CatalogueHit: repositories.CatalogueHit{Brand: "b", ID: 3, ProductName: "Honey Lite"}} + if family := catalogueFamily([]scoredHit{a, b, c}); len(family) != 2 { + t.Errorf("name match should join 1 and 2 only, got %+v", family) + } +} diff --git a/utils/embedding.go b/utils/embedding.go new file mode 100644 index 0000000..7d5a5fc --- /dev/null +++ b/utils/embedding.go @@ -0,0 +1,225 @@ +package utils + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "strings" + "time" + + "nearle/config" +) + +// Embedder turns a short piece of text — what Google Lens read off a packet — +// into the vector the catalogue was indexed with. +// +// One method on purpose. The scan pipeline needs exactly one thing from the +// model and nothing about which model it is; the provider is an environment +// decision (config.EmbeddingConfig) and the tests supply a fake. +type Embedder interface { + // Embed returns the vector for text. It must be the same length as the + // catalogue's `embedding` column or pgvector refuses the comparison. + Embed(ctx context.Context, text string) ([]float32, error) + // Model names what produced the vector, so a cache key can include it: a + // vector cached under one model must never be served for another. + Model() string +} + +// ErrEmbedderNotConfigured is what the scan search sees when no provider is +// set. It falls back to text matching rather than failing the request. +var ErrEmbedderNotConfigured = errors.New("embedding provider is not configured") + +// embedTimeout bounds one call to the provider. The scan endpoint has a +// customer waiting with a phone in their hand; a slow model is worse than a +// text-only answer, and the caller falls back on error. +const embedTimeout = 4 * time.Second + +// NewEmbedder builds the provider named in the config, or returns nil when +// none is configured. A nil Embedder is a supported state everywhere it is +// used: the search degrades to text matching and says so in the response. +func NewEmbedder(cfg config.EmbeddingConfig) (Embedder, error) { + if !cfg.Enabled() { + return nil, nil + } + client := &http.Client{Timeout: embedTimeout} + switch cfg.Provider { + case "openai": + base := strings.TrimRight(cfg.BaseURL, "/") + if base == "" { + base = "https://api.openai.com/v1" + } + return &openAIEmbedder{cfg: cfg, base: base, client: client}, nil + case "gemini": + base := strings.TrimRight(cfg.BaseURL, "/") + if base == "" { + base = "https://generativelanguage.googleapis.com/v1beta" + } + return &geminiEmbedder{cfg: cfg, base: base, client: client}, nil + } + return nil, fmt.Errorf("embedding provider %q is not supported", cfg.Provider) +} + +// ── OpenAI-compatible ─────────────────────────────────────────────────────── +// +// POST {base}/embeddings — the shape OpenAI, Azure OpenAI (with a base URL), +// Ollama, vLLM, LM Studio and most hosted models all accept. + +type openAIEmbedder struct { + cfg config.EmbeddingConfig + base string + client *http.Client +} + +func (e *openAIEmbedder) Model() string { return e.cfg.Model } + +func (e *openAIEmbedder) Embed(ctx context.Context, text string) ([]float32, error) { + body := map[string]interface{}{ + "model": e.cfg.Model, + "input": text, + } + if e.cfg.Dimensions > 0 { + body["dimensions"] = e.cfg.Dimensions + } + + var out struct { + Data []struct { + Embedding []float32 `json:"embedding"` + } `json:"data"` + Error *struct { + Message string `json:"message"` + } `json:"error"` + } + if err := postJSON(ctx, e.client, e.base+"/embeddings", "Bearer "+e.cfg.APIKey, body, &out); err != nil { + return nil, err + } + if out.Error != nil { + return nil, fmt.Errorf("embedding: %s", out.Error.Message) + } + if len(out.Data) == 0 || len(out.Data[0].Embedding) == 0 { + return nil, errors.New("embedding: provider returned no vector") + } + return out.Data[0].Embedding, nil +} + +// ── Gemini ────────────────────────────────────────────────────────────────── +// +// POST {base}/models/{model}:embedContent with the key as a header. + +type geminiEmbedder struct { + cfg config.EmbeddingConfig + base string + client *http.Client +} + +func (e *geminiEmbedder) Model() string { return e.cfg.Model } + +func (e *geminiEmbedder) Embed(ctx context.Context, text string) ([]float32, error) { + model := e.cfg.Model + if !strings.HasPrefix(model, "models/") { + model = "models/" + model + } + body := map[string]interface{}{ + "model": model, + "content": map[string]interface{}{"parts": []map[string]string{{"text": text}}}, + "taskType": "RETRIEVAL_QUERY", + } + if e.cfg.Dimensions > 0 { + body["outputDimensionality"] = e.cfg.Dimensions + } + + var out struct { + Embedding struct { + Values []float32 `json:"values"` + } `json:"embedding"` + Error *struct { + Message string `json:"message"` + } `json:"error"` + } + url := fmt.Sprintf("%s/%s:embedContent", e.base, model) + if err := postJSON(ctx, e.client, url, "", body, &out, "x-goog-api-key", e.cfg.APIKey); err != nil { + return nil, err + } + if out.Error != nil { + return nil, fmt.Errorf("embedding: %s", out.Error.Message) + } + if len(out.Embedding.Values) == 0 { + return nil, errors.New("embedding: provider returned no vector") + } + return out.Embedding.Values, nil +} + +// postJSON is the one HTTP call both providers make. Extra header pairs +// follow the body; `auth` is sent as Authorization when non-empty. +func postJSON(ctx context.Context, client *http.Client, url, auth string, body, out interface{}, headers ...string) error { + payload, err := json.Marshal(body) + if err != nil { + return err + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(payload)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/json") + if auth != "" { + req.Header.Set("Authorization", auth) + } + for i := 0; i+1 < len(headers); i += 2 { + req.Header.Set(headers[i], headers[i+1]) + } + + resp, err := client.Do(req) + if err != nil { + return fmt.Errorf("embedding: %w", err) + } + defer resp.Body.Close() + + // Bounded: an error page from a misconfigured proxy should not be read to + // the end of the internet. + raw, err := io.ReadAll(io.LimitReader(resp.Body, 4<<20)) + if err != nil { + return fmt.Errorf("embedding: %w", err) + } + if err := json.Unmarshal(raw, out); err != nil { + return fmt.Errorf("embedding: HTTP %d, unreadable body: %w", resp.StatusCode, err) + } + if resp.StatusCode/100 != 2 { + // The decoded body carries the provider's message where there is one; + // this is the fallback for a bare status. + if msg := extractMessage(raw); msg != "" { + return fmt.Errorf("embedding: HTTP %d: %s", resp.StatusCode, msg) + } + return fmt.Errorf("embedding: HTTP %d", resp.StatusCode) + } + return nil +} + +func extractMessage(raw []byte) string { + var e struct { + Error struct { + Message string `json:"message"` + } `json:"error"` + } + if json.Unmarshal(raw, &e) == nil { + return e.Error.Message + } + return "" +} + +// VectorLiteral renders a vector the way pgvector reads one: `[0.1,0.2,...]`. +func VectorLiteral(v []float32) string { + var b strings.Builder + b.Grow(len(v)*10 + 2) + b.WriteByte('[') + for i, f := range v { + if i > 0 { + b.WriteByte(',') + } + fmt.Fprintf(&b, "%g", f) + } + b.WriteByte(']') + return b.String() +} diff --git a/utils/embedding_test.go b/utils/embedding_test.go new file mode 100644 index 0000000..0a7653a --- /dev/null +++ b/utils/embedding_test.go @@ -0,0 +1,106 @@ +package utils + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "nearle/config" +) + +func TestOpenAIEmbedderSendsTheRequestTheAPIExpects(t *testing.T) { + var got map[string]interface{} + var auth, path string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + auth, path = r.Header.Get("Authorization"), r.URL.Path + json.NewDecoder(r.Body).Decode(&got) + w.Write([]byte(`{"data":[{"embedding":[0.1,0.2,0.3]}]}`)) + })) + defer srv.Close() + + e, err := NewEmbedder(config.EmbeddingConfig{Provider: "openai", Model: "text-embedding-3-small", APIKey: "sk-test", BaseURL: srv.URL + "/v1/", Dimensions: 3}) + if err != nil { + t.Fatal(err) + } + vec, err := e.Embed(context.Background(), "Milk Bikis") + if err != nil { + t.Fatal(err) + } + if len(vec) != 3 || vec[2] != 0.3 { + t.Errorf("vector = %v", vec) + } + if path != "/v1/embeddings" || auth != "Bearer sk-test" { + t.Errorf("path=%s auth=%s", path, auth) + } + if got["model"] != "text-embedding-3-small" || got["input"] != "Milk Bikis" || got["dimensions"] != float64(3) { + t.Errorf("body = %v", got) + } + if e.Model() != "text-embedding-3-small" { + t.Errorf("Model() = %q", e.Model()) + } +} + +func TestGeminiEmbedderSendsTheRequestTheAPIExpects(t *testing.T) { + var got map[string]interface{} + var key, path string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + key, path = r.Header.Get("x-goog-api-key"), r.URL.Path + json.NewDecoder(r.Body).Decode(&got) + w.Write([]byte(`{"embedding":{"values":[0.5,0.6]}}`)) + })) + defer srv.Close() + + e, err := NewEmbedder(config.EmbeddingConfig{Provider: "gemini", Model: "gemini-embedding-001", APIKey: "g-test", BaseURL: srv.URL}) + if err != nil { + t.Fatal(err) + } + vec, err := e.Embed(context.Background(), "Milk Bikis") + if err != nil { + t.Fatal(err) + } + if len(vec) != 2 || vec[0] != 0.5 { + t.Errorf("vector = %v", vec) + } + if path != "/models/gemini-embedding-001:embedContent" || key != "g-test" { + t.Errorf("path=%s key=%s", path, key) + } + if got["model"] != "models/gemini-embedding-001" || got["taskType"] != "RETRIEVAL_QUERY" { + t.Errorf("body = %v", got) + } +} + +func TestEmbedderSurfacesProviderErrors(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(429) + w.Write([]byte(`{"error":{"message":"Rate limit reached"}}`)) + })) + defer srv.Close() + + e, _ := NewEmbedder(config.EmbeddingConfig{Provider: "openai", Model: "m", APIKey: "k", BaseURL: srv.URL}) + _, err := e.Embed(context.Background(), "x") + if err == nil || !strings.Contains(err.Error(), "429") || !strings.Contains(err.Error(), "Rate limit") { + t.Fatalf("want a 429 with the provider's message, got %v", err) + } +} + +func TestNoProviderMeansNoEmbedder(t *testing.T) { + e, err := NewEmbedder(config.EmbeddingConfig{}) + if err != nil || e != nil { + t.Fatalf("got %v / %v", e, err) + } + if _, err := NewEmbedder(config.EmbeddingConfig{Provider: "cohere", Model: "m", APIKey: "k"}); err == nil { + t.Fatal("an unknown provider must be refused") + } +} + +func TestVectorLiteral(t *testing.T) { + if got := VectorLiteral([]float32{0.1, -2, 3.5}); got != "[0.1,-2,3.5]" { + t.Errorf("got %q", got) + } + if got := VectorLiteral(nil); got != "[]" { + t.Errorf("got %q", got) + } +} diff --git a/utils/geo.go b/utils/geo.go new file mode 100644 index 0000000..cc0a44c --- /dev/null +++ b/utils/geo.go @@ -0,0 +1,104 @@ +package utils + +import ( + "math" + "strconv" + "strings" + "time" +) + +// ParseLatLng reads the coordinate strings this schema stores — customers, +// customerlocations and tenantlocations all hold latitude/longitude as text +// — and says whether they name a real place. +// +// "0,0" is rejected along with blanks: nothing on this platform is in the +// Gulf of Guinea, and it is what an empty map picker saves. +func ParseLatLng(lat, lng string) (float64, float64, bool) { + la, err1 := strconv.ParseFloat(strings.TrimSpace(lat), 64) + lo, err2 := strconv.ParseFloat(strings.TrimSpace(lng), 64) + if err1 != nil || err2 != nil { + return 0, 0, false + } + if la < -90 || la > 90 || lo < -180 || lo > 180 || (la == 0 && lo == 0) { + return 0, 0, false + } + return la, lo, true +} + +// HaversineKm is the great-circle distance between two points. +func HaversineKm(lat1, lng1, lat2, lng2 float64) float64 { + const earthRadiusKm = 6371.0 + toRad := func(d float64) float64 { return d * math.Pi / 180 } + + dLat := toRad(lat2 - lat1) + dLng := toRad(lng2 - lng1) + a := math.Sin(dLat/2)*math.Sin(dLat/2) + + math.Cos(toRad(lat1))*math.Cos(toRad(lat2))*math.Sin(dLng/2)*math.Sin(dLng/2) + return 2 * earthRadiusKm * math.Asin(math.Sqrt(a)) +} + +// OpenNow reads tenantlocations.opentime/closetime ("09:00", "21:30", +// "9:00 AM") against the wall clock. Unparsable or blank hours are treated as +// open: a shop that never filled the field in should not vanish from the +// list, and the ordering flow re-checks at checkout anyway. +func OpenNow(open, closeAt string, now time.Time) bool { + o, ok1 := parseClock(open) + c, ok2 := parseClock(closeAt) + if !ok1 || !ok2 || o == c { + return true + } + cur := now.Hour()*60 + now.Minute() + if o < c { + return cur >= o && cur < c + } + // Past midnight: "20:00" – "02:00". + return cur >= o || cur < c +} + +func parseClock(s string) (int, bool) { + s = strings.TrimSpace(s) + if s == "" { + return 0, false + } + for _, layout := range []string{"15:04", "15:04:05", "3:04 PM", "3:04PM", "03:04 PM", "15.04"} { + if t, err := time.Parse(layout, strings.ToUpper(s)); err == nil { + return t.Hour()*60 + t.Minute(), true + } + } + return 0, false +} + +// SearchTokens splits a label into the words worth matching on: lowercased, +// punctuation stripped, single characters and pack-size noise dropped. "Milk +// Bikis 100g" → ["milk", "bikis"]; the size is matched separately, if at all. +func SearchTokens(label string) []string { + var tokens []string + seen := make(map[string]bool) + for _, raw := range strings.FieldsFunc(strings.ToLower(label), func(r rune) bool { + return !(r >= 'a' && r <= 'z' || r >= '0' && r <= '9') + }) { + if len(raw) < 2 || isPackSize(raw) || seen[raw] { + continue + } + seen[raw] = true + tokens = append(tokens, raw) + } + return tokens +} + +// isPackSize is "100g", "1kg", "500ml", "2l", "250gm" — a number with a unit +// glued on, or a bare number. +func isPackSize(tok string) bool { + digits := 0 + for digits < len(tok) && tok[digits] >= '0' && tok[digits] <= '9' { + digits++ + } + if digits == 0 { + return false + } + switch tok[digits:] { + case "", "g", "gm", "gms", "kg", "ml", "l", "ltr", "pcs", "pc", "x", "n": + return true + } + return false +} diff --git a/utils/geo_test.go b/utils/geo_test.go new file mode 100644 index 0000000..eb7c012 --- /dev/null +++ b/utils/geo_test.go @@ -0,0 +1,65 @@ +package utils + +import ( + "testing" + "time" +) + +func TestParseLatLng(t *testing.T) { + if _, _, ok := ParseLatLng("", ""); ok { + t.Error("blank must not parse") + } + if _, _, ok := ParseLatLng("0", "0"); ok { + t.Error("0,0 is an empty map picker, not a place") + } + if _, _, ok := ParseLatLng("91", "10"); ok { + t.Error("out of range") + } + lat, lng, ok := ParseLatLng(" 11.0168 ", "76.9558") + if !ok || lat != 11.0168 || lng != 76.9558 { + t.Errorf("got %v %v %v", lat, lng, ok) + } +} + +func TestHaversineKm(t *testing.T) { + // Coimbatore railway station to Peelamedu, roughly 8 km. + d := HaversineKm(11.0018, 76.9660, 11.0290, 77.0290) + if d < 7 || d > 9 { + t.Errorf("got %.2f km", d) + } + if HaversineKm(1, 1, 1, 1) != 0 { + t.Error("same point") + } +} + +func TestOpenNow(t *testing.T) { + at := func(h, m int) time.Time { return time.Date(2026, 9, 15, h, m, 0, 0, time.UTC) } + if !OpenNow("09:00", "21:00", at(12, 0)) || OpenNow("09:00", "21:00", at(22, 0)) { + t.Error("plain hours") + } + if !OpenNow("20:00", "02:00", at(1, 0)) || OpenNow("20:00", "02:00", at(12, 0)) { + t.Error("past midnight") + } + if !OpenNow("9:00 AM", "9:30 PM", at(21, 0)) || OpenNow("9:00 AM", "9:30 PM", at(21, 45)) { + t.Error("12-hour clock") + } + if !OpenNow("", "", at(3, 0)) || !OpenNow("always", "", at(3, 0)) { + t.Error("unknown hours mean open") + } +} + +func TestSearchTokens(t *testing.T) { + got := SearchTokens("Milk Bikis 100g, Britannia (2 x 50gm)") + want := []string{"milk", "bikis", "britannia"} + if len(got) != len(want) { + t.Fatalf("got %v want %v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("got %v want %v", got, want) + } + } + if len(SearchTokens("500ml 1kg 2")) != 0 { + t.Error("pack sizes alone are not searchable") + } +}