ui: add B2B olmayan stok (orphans) page

This commit is contained in:
M_Kececi
2026-07-02 15:33:49 +03:00
parent 810c0e5eff
commit cbb853d433
19 changed files with 3365 additions and 395 deletions
+963
View File
@@ -0,0 +1,963 @@
package queries
import (
"bssapp-backend/db"
"bssapp-backend/models"
"context"
"database/sql"
"fmt"
"log"
"strings"
"time"
)
type ProductPerformanceFilters struct {
Search string
ProductCode string
MarketKey string
Kategori string
Seri string
Bucket string
Limit int
Page int
}
type ProductPerformanceRefreshRequest struct {
Mode string
Stage string
ResumeAfter int
SkipDelete bool
StartDate time.Time
EndDate time.Time
}
type ProductPerformanceRefreshResult struct {
Mode string `json:"mode"`
Stage string `json:"stage"`
StartDate string `json:"start_date"`
EndDate string `json:"end_date"`
SalesRows int `json:"sales_rows"`
StockRows int `json:"stock_rows"`
KpiRows int `json:"kpi_rows"`
DurationMS int64 `json:"duration_ms"`
}
func EnsureProductPerformanceTables(pg *sql.DB) error {
stmts := []string{
`
CREATE TABLE IF NOT EXISTS mk_product_performance_sales_daily (
sales_date DATE NOT NULL,
product_code TEXT NOT NULL,
color_code TEXT NOT NULL DEFAULT '',
yaka_kodu TEXT NOT NULL DEFAULT '',
item_description TEXT NOT NULL DEFAULT '',
kategori TEXT NOT NULL DEFAULT '',
seri TEXT NOT NULL DEFAULT '',
yas_grubu TEXT NOT NULL DEFAULT '',
askili_yan TEXT NOT NULL DEFAULT '',
urun_ilk_grubu TEXT NOT NULL DEFAULT '',
urun_ana_grubu TEXT NOT NULL DEFAULT '',
urun_alt_grubu TEXT NOT NULL DEFAULT '',
market_key TEXT NOT NULL DEFAULT '',
channel_code TEXT NOT NULL DEFAULT '',
customer_country TEXT NOT NULL DEFAULT '',
customer_segment TEXT NOT NULL DEFAULT '',
customer_code TEXT NOT NULL DEFAULT '',
customer_name TEXT NOT NULL DEFAULT '',
sales_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_tl NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_usd NUMERIC(18,4) NOT NULL DEFAULT 0,
avg_price_usd NUMERIC(18,6) NOT NULL DEFAULT 0,
invoice_line_count INTEGER NOT NULL DEFAULT 0,
invoice_count INTEGER NOT NULL DEFAULT 0,
customer_count INTEGER NOT NULL DEFAULT 0,
last_ref_number TEXT NOT NULL DEFAULT '',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
CONSTRAINT pk_mk_product_performance_sales_daily PRIMARY KEY
(sales_date, product_code, color_code, yaka_kodu, market_key, customer_country, customer_segment, customer_code)
)`,
`ALTER TABLE mk_product_performance_sales_daily ADD COLUMN IF NOT EXISTS customer_code TEXT NOT NULL DEFAULT ''`,
`ALTER TABLE mk_product_performance_sales_daily ADD COLUMN IF NOT EXISTS customer_name TEXT NOT NULL DEFAULT ''`,
`ALTER TABLE mk_product_performance_sales_daily ADD COLUMN IF NOT EXISTS urun_ilk_grubu TEXT NOT NULL DEFAULT ''`,
`ALTER TABLE mk_product_performance_sales_daily ADD COLUMN IF NOT EXISTS urun_ana_grubu TEXT NOT NULL DEFAULT ''`,
`ALTER TABLE mk_product_performance_sales_daily ADD COLUMN IF NOT EXISTS urun_alt_grubu TEXT NOT NULL DEFAULT ''`,
`
DO $$
BEGIN
IF EXISTS (
SELECT 1
FROM pg_constraint
WHERE conname = 'pk_mk_product_performance_sales_daily'
AND conrelid = 'mk_product_performance_sales_daily'::regclass
) THEN
ALTER TABLE mk_product_performance_sales_daily DROP CONSTRAINT pk_mk_product_performance_sales_daily;
END IF;
ALTER TABLE mk_product_performance_sales_daily
ADD CONSTRAINT pk_mk_product_performance_sales_daily PRIMARY KEY
(sales_date, product_code, color_code, yaka_kodu, market_key, customer_country, customer_segment, customer_code);
END $$`,
`CREATE INDEX IF NOT EXISTS ix_mk_product_perf_sales_product_date ON mk_product_performance_sales_daily (product_code, sales_date DESC)`,
`CREATE INDEX IF NOT EXISTS ix_mk_product_perf_sales_market ON mk_product_performance_sales_daily (market_key, sales_date DESC)`,
`CREATE INDEX IF NOT EXISTS ix_mk_product_perf_sales_customer ON mk_product_performance_sales_daily (customer_code, sales_date DESC)`,
`
CREATE TABLE IF NOT EXISTS mk_product_performance_stock_daily (
stock_date DATE NOT NULL,
product_code TEXT NOT NULL,
color_code TEXT NOT NULL DEFAULT '',
yaka_kodu TEXT NOT NULL DEFAULT '',
stock_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
in_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
out_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
kpi_in_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
kpi_out_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_movement_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
production_in_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
purchase_in_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
consumption_out_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
count_diff_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
CONSTRAINT pk_mk_product_performance_stock_daily PRIMARY KEY
(stock_date, product_code, color_code, yaka_kodu)
)`,
`CREATE INDEX IF NOT EXISTS ix_mk_product_perf_stock_product_date ON mk_product_performance_stock_daily (product_code, stock_date DESC)`,
`
CREATE TABLE IF NOT EXISTS mk_product_performance_price_dim (
product_code TEXT PRIMARY KEY,
cost_price_usd NUMERIC(18,6) NOT NULL DEFAULT 0,
base_price_usd NUMERIC(18,6) NOT NULL DEFAULT 0,
base_price_try NUMERIC(18,6) NOT NULL DEFAULT 0,
last_pricing_date DATE,
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)`,
`
CREATE TABLE IF NOT EXISTS mk_product_performance_kpi_daily (
kpi_date DATE NOT NULL,
product_code TEXT NOT NULL,
color_code TEXT NOT NULL DEFAULT '',
yaka_kodu TEXT NOT NULL DEFAULT '',
item_description TEXT NOT NULL DEFAULT '',
kategori TEXT NOT NULL DEFAULT '',
seri TEXT NOT NULL DEFAULT '',
yas_grubu TEXT NOT NULL DEFAULT '',
askili_yan TEXT NOT NULL DEFAULT '',
urun_ilk_grubu TEXT NOT NULL DEFAULT '',
urun_ana_grubu TEXT NOT NULL DEFAULT '',
urun_alt_grubu TEXT NOT NULL DEFAULT '',
market_key TEXT NOT NULL DEFAULT '',
stock_qty NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_qty_30d NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_qty_90d NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_qty_365d NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_qty_730d NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_usd_30d NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_usd_90d NUMERIC(18,4) NOT NULL DEFAULT 0,
sales_usd_365d NUMERIC(18,4) NOT NULL DEFAULT 0,
avg_daily_sales_90d NUMERIC(18,6) NOT NULL DEFAULT 0,
stock_days_90d NUMERIC(18,6) NOT NULL DEFAULT 0,
avg_price_usd_90d NUMERIC(18,6) NOT NULL DEFAULT 0,
cost_price_usd NUMERIC(18,6) NOT NULL DEFAULT 0,
base_price_usd NUMERIC(18,6) NOT NULL DEFAULT 0,
base_price_try NUMERIC(18,6) NOT NULL DEFAULT 0,
gross_profit_usd_90d NUMERIC(18,4) NOT NULL DEFAULT 0,
gross_margin_90d NUMERIC(18,6) NOT NULL DEFAULT 0,
customer_count_90d INTEGER NOT NULL DEFAULT 0,
sales_index_90d NUMERIC(18,6) NOT NULL DEFAULT 0,
price_index_90d NUMERIC(18,6) NOT NULL DEFAULT 0,
margin_index_90d NUMERIC(18,6) NOT NULL DEFAULT 0,
performance_score NUMERIC(18,6) NOT NULL DEFAULT 0,
performance_bucket TEXT NOT NULL DEFAULT '',
recommendation TEXT NOT NULL DEFAULT '',
last_sale_date DATE,
last_ref_number TEXT NOT NULL DEFAULT '',
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
CONSTRAINT pk_mk_product_performance_kpi_daily PRIMARY KEY
(kpi_date, product_code, color_code, yaka_kodu, market_key)
)`,
`CREATE INDEX IF NOT EXISTS ix_mk_product_perf_kpi_bucket ON mk_product_performance_kpi_daily (kpi_date DESC, performance_bucket)`,
`CREATE INDEX IF NOT EXISTS ix_mk_product_perf_kpi_market ON mk_product_performance_kpi_daily (kpi_date DESC, market_key, sales_index_90d DESC)`,
`ALTER TABLE mk_product_performance_kpi_daily ADD COLUMN IF NOT EXISTS urun_ilk_grubu TEXT NOT NULL DEFAULT ''`,
`ALTER TABLE mk_product_performance_kpi_daily ADD COLUMN IF NOT EXISTS urun_ana_grubu TEXT NOT NULL DEFAULT ''`,
`ALTER TABLE mk_product_performance_kpi_daily ADD COLUMN IF NOT EXISTS urun_alt_grubu TEXT NOT NULL DEFAULT ''`,
}
for _, stmt := range stmts {
if _, err := pg.Exec(stmt); err != nil {
return err
}
}
return nil
}
func RefreshProductPerformance(ctx context.Context, pg *sql.DB, req ProductPerformanceRefreshRequest) (ProductPerformanceRefreshResult, error) {
started := time.Now()
stage := productPerformanceRefreshStage(req.Stage)
log.Printf("[ProductPerformanceRefresh] start mode=%s stage=%s start=%s end=%s", req.Mode, stage, req.StartDate.Format("2006-01-02"), req.EndDate.Format("2006-01-02"))
if pg == nil {
return ProductPerformanceRefreshResult{}, fmt.Errorf("postgres db nil")
}
if db.MssqlDB == nil {
return ProductPerformanceRefreshResult{}, fmt.Errorf("mssql db nil")
}
log.Printf("[ProductPerformanceRefresh] ensure tables start")
if err := EnsureProductPerformanceTables(pg); err != nil {
return ProductPerformanceRefreshResult{}, err
}
log.Printf("[ProductPerformanceRefresh] ensure tables done")
mode := strings.ToLower(strings.TrimSpace(req.Mode))
if mode == "" {
mode = "delta"
}
if req.EndDate.IsZero() {
req.EndDate = time.Now()
}
if req.StartDate.IsZero() {
if mode == "full" {
req.StartDate = time.Date(2022, 1, 1, 0, 0, 0, 0, req.EndDate.Location())
} else {
req.StartDate = req.EndDate.AddDate(0, 0, -45)
}
}
req.StartDate = dateOnly(req.StartDate)
req.EndDate = dateOnly(req.EndDate)
shouldRun := func(name string) bool {
if stage == "all" {
return true
}
order := map[string]int{
"sales": 1,
"stock": 2,
"price": 3,
"kpi": 4,
}
return order[name] >= order[stage]
}
runStage := func(stage string, fn func(*sql.Tx) error) error {
log.Printf("[ProductPerformanceRefresh] %s tx begin", stage)
tx, err := pg.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
if err := fn(tx); err != nil {
return err
}
log.Printf("[ProductPerformanceRefresh] %s tx commit start", stage)
if err := tx.Commit(); err != nil {
return err
}
log.Printf("[ProductPerformanceRefresh] %s tx commit done elapsed=%s", stage, time.Since(started).Round(time.Second))
return nil
}
var salesRows int
if shouldRun("sales") {
if err := runStage("sales", func(tx *sql.Tx) error {
log.Printf("[ProductPerformanceRefresh] delete sales cache start")
if _, err := tx.ExecContext(ctx, `DELETE FROM mk_product_performance_sales_daily WHERE sales_date BETWEEN $1 AND $2`, req.StartDate, req.EndDate); err != nil {
return err
}
log.Printf("[ProductPerformanceRefresh] delete sales cache done")
log.Printf("[ProductPerformanceRefresh] sales refresh start")
rows, err := refreshProductPerformanceSales(ctx, tx, req.StartDate, req.EndDate)
if err != nil {
return err
}
salesRows = rows
log.Printf("[ProductPerformanceRefresh] sales refresh done rows=%d elapsed=%s", salesRows, time.Since(started).Round(time.Second))
return nil
}); err != nil {
return ProductPerformanceRefreshResult{}, err
}
} else {
log.Printf("[ProductPerformanceRefresh] sales skipped stage=%s", stage)
}
var stockRows int
if shouldRun("stock") {
rows, err := refreshProductPerformanceStockChunked(ctx, pg, req.StartDate, req.EndDate, started, req.SkipDelete || req.ResumeAfter > 0, req.ResumeAfter)
if err != nil {
return ProductPerformanceRefreshResult{}, err
}
stockRows = rows
} else {
log.Printf("[ProductPerformanceRefresh] stock skipped stage=%s", stage)
}
if shouldRun("price") {
if err := runStage("price", func(tx *sql.Tx) error {
log.Printf("[ProductPerformanceRefresh] price refresh start")
if err := refreshProductPerformancePrices(ctx, tx); err != nil {
return err
}
log.Printf("[ProductPerformanceRefresh] price refresh done elapsed=%s", time.Since(started).Round(time.Second))
return nil
}); err != nil {
return ProductPerformanceRefreshResult{}, err
}
} else {
log.Printf("[ProductPerformanceRefresh] price skipped stage=%s", stage)
}
var kpiRows int
if shouldRun("kpi") {
if err := runStage("kpi", func(tx *sql.Tx) error {
log.Printf("[ProductPerformanceRefresh] kpi rebuild start")
rows, err := RebuildProductPerformanceKPI(ctx, tx, req.EndDate)
if err != nil {
return err
}
kpiRows = rows
log.Printf("[ProductPerformanceRefresh] kpi rebuild done rows=%d elapsed=%s", kpiRows, time.Since(started).Round(time.Second))
return nil
}); err != nil {
return ProductPerformanceRefreshResult{}, err
}
} else {
log.Printf("[ProductPerformanceRefresh] kpi skipped stage=%s", stage)
}
log.Printf("[ProductPerformanceRefresh] refresh done total_elapsed=%s", time.Since(started).Round(time.Second))
return ProductPerformanceRefreshResult{
Mode: mode,
Stage: stage,
StartDate: req.StartDate.Format("2006-01-02"),
EndDate: req.EndDate.Format("2006-01-02"),
SalesRows: salesRows,
StockRows: stockRows,
KpiRows: kpiRows,
DurationMS: time.Since(started).Milliseconds(),
}, nil
}
func productPerformanceRefreshStage(raw string) string {
switch strings.ToLower(strings.TrimSpace(raw)) {
case "", "all", "full":
return "all"
case "sales", "stock", "price", "kpi":
return strings.ToLower(strings.TrimSpace(raw))
default:
return "all"
}
}
func refreshProductPerformanceSales(ctx context.Context, tx *sql.Tx, startDate, endDate time.Time) (int, error) {
log.Printf("[ProductPerformanceRefresh] sales mssql query start start=%s end=%s", startDate.Format("2006-01-02"), endDate.Format("2006-01-02"))
rows, err := db.MssqlDB.QueryContext(ctx, productPerformanceSalesSQL(), startDate, endDate)
if err != nil {
return 0, err
}
defer rows.Close()
log.Printf("[ProductPerformanceRefresh] sales mssql query returned, postgres insert start")
count := 0
for rows.Next() {
var r productPerformanceSalesDaily
if err := rows.Scan(
&r.SalesDate, &r.ProductCode, &r.ColorCode, &r.YakaKodu, &r.ItemDescription,
&r.Kategori, &r.Seri, &r.YasGrubu, &r.AskiliYan, &r.UrunIlkGrubu, &r.UrunAnaGrubu, &r.UrunAltGrubu, &r.MarketKey, &r.ChannelCode,
&r.CustomerCountry, &r.CustomerSegment, &r.CustomerCode, &r.CustomerName, &r.SalesQty, &r.SalesTL, &r.SalesUSD,
&r.AvgPriceUSD, &r.InvoiceLineCount, &r.InvoiceCount, &r.CustomerCount, &r.LastRefNumber,
); err != nil {
return count, err
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO mk_product_performance_sales_daily (
sales_date, product_code, color_code, yaka_kodu, item_description,
kategori, seri, yas_grubu, askili_yan, urun_ilk_grubu, urun_ana_grubu, urun_alt_grubu, market_key, channel_code,
customer_country, customer_segment, customer_code, customer_name, sales_qty, sales_tl, sales_usd,
avg_price_usd, invoice_line_count, invoice_count, customer_count, last_ref_number, updated_at
) VALUES (
$1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23,$24,$25,$26,now()
)
ON CONFLICT (sales_date, product_code, color_code, yaka_kodu, market_key, customer_country, customer_segment, customer_code)
DO UPDATE SET
item_description=EXCLUDED.item_description,
kategori=EXCLUDED.kategori,
seri=EXCLUDED.seri,
yas_grubu=EXCLUDED.yas_grubu,
askili_yan=EXCLUDED.askili_yan,
urun_ilk_grubu=EXCLUDED.urun_ilk_grubu,
urun_ana_grubu=EXCLUDED.urun_ana_grubu,
urun_alt_grubu=EXCLUDED.urun_alt_grubu,
channel_code=EXCLUDED.channel_code,
customer_name=EXCLUDED.customer_name,
sales_qty=EXCLUDED.sales_qty,
sales_tl=EXCLUDED.sales_tl,
sales_usd=EXCLUDED.sales_usd,
avg_price_usd=EXCLUDED.avg_price_usd,
invoice_line_count=EXCLUDED.invoice_line_count,
invoice_count=EXCLUDED.invoice_count,
customer_count=EXCLUDED.customer_count,
last_ref_number=EXCLUDED.last_ref_number,
updated_at=now()
`, r.SalesDate, r.ProductCode, r.ColorCode, r.YakaKodu, r.ItemDescription,
r.Kategori, r.Seri, r.YasGrubu, r.AskiliYan, r.UrunIlkGrubu, r.UrunAnaGrubu, r.UrunAltGrubu, r.MarketKey, r.ChannelCode,
r.CustomerCountry, r.CustomerSegment, r.CustomerCode, r.CustomerName, r.SalesQty, r.SalesTL, r.SalesUSD,
r.AvgPriceUSD, r.InvoiceLineCount, r.InvoiceCount, r.CustomerCount, r.LastRefNumber); err != nil {
return count, err
}
count++
if count%5000 == 0 {
log.Printf("[ProductPerformanceRefresh] sales inserted rows=%d", count)
}
}
return count, rows.Err()
}
func refreshProductPerformanceStock(ctx context.Context, tx *sql.Tx, startDate, endDate time.Time) (int, error) {
log.Printf("[ProductPerformanceRefresh] stock mssql query start start=%s end=%s", startDate.Format("2006-01-02"), endDate.Format("2006-01-02"))
rows, err := db.MssqlDB.QueryContext(ctx, productPerformanceStockSQL(), startDate, endDate)
if err != nil {
return 0, err
}
defer rows.Close()
log.Printf("[ProductPerformanceRefresh] stock mssql query returned, postgres insert start")
count := 0
for rows.Next() {
var r productPerformanceStockDaily
if err := rows.Scan(
&r.StockDate, &r.ProductCode, &r.ColorCode, &r.YakaKodu, &r.StockQty,
&r.InQty, &r.OutQty, &r.KpiInQty, &r.KpiOutQty, &r.SalesMovementQty,
&r.ProductionInQty, &r.PurchaseInQty, &r.ConsumptionOutQty, &r.CountDiffQty,
); err != nil {
return count, err
}
if err := insertProductPerformanceStock(ctx, tx, r); err != nil {
return count, err
}
count++
if count%5000 == 0 {
log.Printf("[ProductPerformanceRefresh] stock inserted rows=%d", count)
}
}
return count, rows.Err()
}
func refreshProductPerformanceStockChunked(ctx context.Context, pg *sql.DB, startDate, endDate, started time.Time, skipDelete bool, resumeAfter int) (int, error) {
if skipDelete {
log.Printf("[ProductPerformanceRefresh] stock delete skipped resume_after=%d", resumeAfter)
} else {
log.Printf("[ProductPerformanceRefresh] stock delete tx begin")
deleteTx, err := pg.BeginTx(ctx, nil)
if err != nil {
return 0, err
}
defer deleteTx.Rollback()
log.Printf("[ProductPerformanceRefresh] delete stock cache start")
if _, err := deleteTx.ExecContext(ctx, `DELETE FROM mk_product_performance_stock_daily WHERE stock_date BETWEEN $1 AND $2`, startDate, endDate); err != nil {
return 0, err
}
log.Printf("[ProductPerformanceRefresh] delete stock cache done")
log.Printf("[ProductPerformanceRefresh] stock delete tx commit start")
if err := deleteTx.Commit(); err != nil {
return 0, err
}
log.Printf("[ProductPerformanceRefresh] stock delete tx commit done elapsed=%s", time.Since(started).Round(time.Second))
}
log.Printf("[ProductPerformanceRefresh] stock refresh start")
log.Printf("[ProductPerformanceRefresh] stock mssql query start start=%s end=%s", startDate.Format("2006-01-02"), endDate.Format("2006-01-02"))
rows, err := db.MssqlDB.QueryContext(ctx, productPerformanceStockSQL(), startDate, endDate)
if err != nil {
return 0, err
}
defer rows.Close()
log.Printf("[ProductPerformanceRefresh] stock mssql query returned, postgres chunk insert start")
const chunkSize = 5000
count := 0
seen := 0
chunkRows := 0
tx, err := pg.BeginTx(ctx, nil)
if err != nil {
return 0, err
}
defer tx.Rollback()
for rows.Next() {
var r productPerformanceStockDaily
if err := rows.Scan(
&r.StockDate, &r.ProductCode, &r.ColorCode, &r.YakaKodu, &r.StockQty,
&r.InQty, &r.OutQty, &r.KpiInQty, &r.KpiOutQty, &r.SalesMovementQty,
&r.ProductionInQty, &r.PurchaseInQty, &r.ConsumptionOutQty, &r.CountDiffQty,
); err != nil {
return count, err
}
seen++
if resumeAfter > 0 && seen <= resumeAfter {
if seen%50000 == 0 || seen == resumeAfter {
log.Printf("[ProductPerformanceRefresh] stock resume skipped rows=%d", seen)
}
continue
}
if err := insertProductPerformanceStock(ctx, tx, r); err != nil {
return count, err
}
count++
chunkRows++
if chunkRows >= chunkSize {
log.Printf("[ProductPerformanceRefresh] stock chunk commit start inserted_rows=%d scanned_rows=%d", count, seen)
if err := tx.Commit(); err != nil {
return count, err
}
log.Printf("[ProductPerformanceRefresh] stock chunk commit done inserted_rows=%d scanned_rows=%d elapsed=%s", count, seen, time.Since(started).Round(time.Second))
tx, err = pg.BeginTx(ctx, nil)
if err != nil {
return count, err
}
chunkRows = 0
}
}
if err := rows.Err(); err != nil {
return count, err
}
if chunkRows > 0 {
log.Printf("[ProductPerformanceRefresh] stock final chunk commit start inserted_rows=%d scanned_rows=%d", count, seen)
if err := tx.Commit(); err != nil {
return count, err
}
log.Printf("[ProductPerformanceRefresh] stock final chunk commit done inserted_rows=%d scanned_rows=%d elapsed=%s", count, seen, time.Since(started).Round(time.Second))
} else {
_ = tx.Rollback()
}
log.Printf("[ProductPerformanceRefresh] stock refresh done inserted_rows=%d scanned_rows=%d elapsed=%s", count, seen, time.Since(started).Round(time.Second))
return count, nil
}
type productPerformanceStockExec interface {
ExecContext(context.Context, string, ...any) (sql.Result, error)
}
func insertProductPerformanceStock(ctx context.Context, exec productPerformanceStockExec, r productPerformanceStockDaily) error {
_, err := exec.ExecContext(ctx, `
INSERT INTO mk_product_performance_stock_daily (
stock_date, product_code, color_code, yaka_kodu, stock_qty, in_qty, out_qty,
kpi_in_qty, kpi_out_qty, sales_movement_qty, production_in_qty, purchase_in_qty,
consumption_out_qty, count_diff_qty, updated_at
) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,now())
ON CONFLICT (stock_date, product_code, color_code, yaka_kodu)
DO UPDATE SET
stock_qty=EXCLUDED.stock_qty,
in_qty=EXCLUDED.in_qty,
out_qty=EXCLUDED.out_qty,
kpi_in_qty=EXCLUDED.kpi_in_qty,
kpi_out_qty=EXCLUDED.kpi_out_qty,
sales_movement_qty=EXCLUDED.sales_movement_qty,
production_in_qty=EXCLUDED.production_in_qty,
purchase_in_qty=EXCLUDED.purchase_in_qty,
consumption_out_qty=EXCLUDED.consumption_out_qty,
count_diff_qty=EXCLUDED.count_diff_qty,
updated_at=now()
`, r.StockDate, r.ProductCode, r.ColorCode, r.YakaKodu, r.StockQty,
r.InQty, r.OutQty, r.KpiInQty, r.KpiOutQty, r.SalesMovementQty,
r.ProductionInQty, r.PurchaseInQty, r.ConsumptionOutQty, r.CountDiffQty)
return err
}
func refreshProductPerformancePrices(ctx context.Context, tx *sql.Tx) error {
log.Printf("[ProductPerformanceRefresh] price mssql query start")
rows, err := db.MssqlDB.QueryContext(ctx, productPerformancePriceSQL())
if err != nil {
return err
}
defer rows.Close()
log.Printf("[ProductPerformanceRefresh] price mssql query returned, postgres upsert start")
count := 0
for rows.Next() {
var code string
var cost, usd, tryPrice float64
var lastPricing sql.NullTime
if err := rows.Scan(&code, &cost, &usd, &tryPrice, &lastPricing); err != nil {
return err
}
var lp any
if lastPricing.Valid {
lp = lastPricing.Time
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO mk_product_performance_price_dim (product_code, cost_price_usd, base_price_usd, base_price_try, last_pricing_date, updated_at)
VALUES ($1,$2,$3,$4,$5,now())
ON CONFLICT (product_code) DO UPDATE SET
cost_price_usd=EXCLUDED.cost_price_usd,
base_price_usd=EXCLUDED.base_price_usd,
base_price_try=EXCLUDED.base_price_try,
last_pricing_date=EXCLUDED.last_pricing_date,
updated_at=now()
`, strings.TrimSpace(code), cost, usd, tryPrice, lp); err != nil {
return err
}
count++
if count%5000 == 0 {
log.Printf("[ProductPerformanceRefresh] price upserted rows=%d", count)
}
}
log.Printf("[ProductPerformanceRefresh] price upsert done rows=%d", count)
return rows.Err()
}
func RebuildProductPerformanceKPI(ctx context.Context, exec interface {
ExecContext(context.Context, string, ...any) (sql.Result, error)
}, kpiDate time.Time) (int, error) {
kpiDate = dateOnly(kpiDate)
if _, err := exec.ExecContext(ctx, `DELETE FROM mk_product_performance_kpi_daily WHERE kpi_date=$1`, kpiDate); err != nil {
return 0, err
}
res, err := exec.ExecContext(ctx, productPerformanceKPISQL(), kpiDate)
if err != nil {
return 0, err
}
n, _ := res.RowsAffected()
return int(n), nil
}
func ListProductPerformance(ctx context.Context, pg *sql.DB, f ProductPerformanceFilters) ([]models.ProductPerformanceRow, int, error) {
if err := EnsureProductPerformanceTables(pg); err != nil {
return nil, 0, err
}
limit := f.Limit
if limit <= 0 || limit > 500 {
limit = 100
}
page := f.Page
if page <= 0 {
page = 1
}
where, args := productPerformanceWhere(f)
countQuery := `SELECT COUNT(*) FROM mk_product_performance_kpi_daily WHERE kpi_date = (SELECT MAX(kpi_date) FROM mk_product_performance_kpi_daily)` + where
var total int
if err := pg.QueryRowContext(ctx, countQuery, args...).Scan(&total); err != nil {
return nil, 0, err
}
args = append(args, limit, (page-1)*limit)
rows, err := pg.QueryContext(ctx, `
SELECT
to_char(kpi_date,'YYYY-MM-DD'), product_code, color_code, yaka_kodu, item_description,
kategori, seri, yas_grubu, askili_yan, urun_ilk_grubu, urun_ana_grubu, urun_alt_grubu, market_key, stock_qty,
sales_qty_30d, sales_qty_90d, sales_qty_365d, sales_qty_730d,
sales_usd_30d, sales_usd_90d, sales_usd_365d, avg_daily_sales_90d,
stock_days_90d, avg_price_usd_90d, cost_price_usd, base_price_usd, base_price_try,
gross_profit_usd_90d, gross_margin_90d, customer_count_90d, sales_index_90d,
price_index_90d, margin_index_90d, performance_score, performance_bucket,
recommendation, COALESCE(to_char(last_sale_date,'YYYY-MM-DD'),''), last_ref_number,
to_char(updated_at,'YYYY-MM-DD HH24:MI:SS')
FROM mk_product_performance_kpi_daily
WHERE kpi_date = (SELECT MAX(kpi_date) FROM mk_product_performance_kpi_daily)`+where+`
ORDER BY performance_score DESC, sales_qty_90d DESC, product_code
LIMIT $`+fmt.Sprint(len(args)-1)+` OFFSET $`+fmt.Sprint(len(args)), args...)
if err != nil {
return nil, 0, err
}
defer rows.Close()
out := make([]models.ProductPerformanceRow, 0, limit)
for rows.Next() {
var r models.ProductPerformanceRow
if err := rows.Scan(
&r.KpiDate, &r.ProductCode, &r.ColorCode, &r.YakaKodu, &r.ItemDescription,
&r.Kategori, &r.Seri, &r.YasGrubu, &r.AskiliYan, &r.UrunIlkGrubu, &r.UrunAnaGrubu, &r.UrunAltGrubu, &r.MarketKey, &r.StockQty,
&r.SalesQty30, &r.SalesQty90, &r.SalesQty365, &r.SalesQty730,
&r.SalesUSD30, &r.SalesUSD90, &r.SalesUSD365, &r.AvgDailySales90,
&r.StockDays90, &r.AvgPriceUSD90, &r.CostPriceUSD, &r.BasePriceUSD, &r.BasePriceTRY,
&r.GrossProfitUSD90, &r.GrossMargin90, &r.CustomerCount90, &r.SalesIndex90,
&r.PriceIndex90, &r.MarginIndex90, &r.PerformanceScore, &r.PerformanceBucket,
&r.Recommendation, &r.LastSaleDate, &r.LastRefNumber, &r.UpdatedAt,
); err != nil {
return nil, 0, err
}
out = append(out, r)
}
return out, total, rows.Err()
}
func GetProductPerformanceSummary(ctx context.Context, pg *sql.DB) (models.ProductPerformanceSummary, error) {
if err := EnsureProductPerformanceTables(pg); err != nil {
return models.ProductPerformanceSummary{}, err
}
var s models.ProductPerformanceSummary
err := pg.QueryRowContext(ctx, `
SELECT
COALESCE(to_char(MAX(kpi_date),'YYYY-MM-DD'),''),
COUNT(*),
COUNT(*) FILTER (WHERE performance_bucket='YILDIZ_URUN'),
COUNT(*) FILTER (WHERE performance_bucket='STOK_RISKI'),
COUNT(*) FILTER (WHERE performance_bucket='STOKSUZ_TALEP'),
COALESCE(SUM(stock_qty),0),
COALESCE(SUM(stock_qty * cost_price_usd),0),
COALESCE(SUM(CASE WHEN performance_bucket IN ('STOK_RISKI','TAKIP') AND COALESCE(sales_qty_90d,0)=0 THEN stock_qty * cost_price_usd ELSE 0 END),0),
COALESCE(SUM(sales_qty_90d),0),
COALESCE(SUM(sales_usd_90d),0),
COALESCE(SUM(gross_profit_usd_90d),0),
COALESCE(to_char(MAX(updated_at),'YYYY-MM-DD HH24:MI:SS'),'')
FROM mk_product_performance_kpi_daily
WHERE kpi_date = (SELECT MAX(kpi_date) FROM mk_product_performance_kpi_daily)
`).Scan(&s.KpiDate, &s.TotalRows, &s.StarCount, &s.StockRisk, &s.NoStockDemand, &s.TotalStock, &s.StockCostUSD, &s.RiskCostUSD, &s.SalesQty90, &s.SalesUSD90, &s.GrossProfit90, &s.UpdatedAt)
return s, err
}
func ListProductPerformanceMarkets(ctx context.Context, pg *sql.DB, limit int) ([]models.ProductPerformanceMarketRow, error) {
if err := EnsureProductPerformanceTables(pg); err != nil {
return nil, err
}
if limit <= 0 || limit > 500 {
limit = 100
}
rows, err := pg.QueryContext(ctx, `
SELECT
market_key,
MAX(kategori) AS kategori,
MAX(seri) AS seri,
MAX(yas_grubu) AS yas_grubu,
MAX(askili_yan) AS askili_yan,
MAX(urun_ilk_grubu) AS urun_ilk_grubu,
MAX(urun_ana_grubu) AS urun_ana_grubu,
MAX(urun_alt_grubu) AS urun_alt_grubu,
COUNT(DISTINCT product_code) AS product_count,
COUNT(*) FILTER (WHERE performance_bucket='YILDIZ_URUN') AS star_count,
COUNT(*) FILTER (WHERE performance_bucket='STOK_RISKI') AS stock_risk_count,
COALESCE(SUM(stock_qty),0) AS stock_qty,
COALESCE(SUM(stock_qty * cost_price_usd),0) AS stock_cost_value_usd,
COALESCE(SUM(CASE WHEN performance_bucket IN ('STOK_RISKI','TAKIP') AND COALESCE(sales_qty_90d,0)=0 THEN stock_qty * cost_price_usd ELSE 0 END),0) AS risk_stock_cost_value_usd,
COALESCE(SUM(sales_qty_90d),0) AS sales_qty_90d,
COALESCE(SUM(sales_usd_90d),0) AS sales_usd_90d,
COALESCE(SUM(gross_profit_usd_90d),0) AS gross_profit_usd_90d,
COALESCE(AVG(NULLIF(gross_margin_90d,0)),0) AS avg_gross_margin_90d,
COALESCE(AVG(NULLIF(stock_days_90d,0)),0) AS avg_stock_days_90d
FROM mk_product_performance_kpi_daily
WHERE kpi_date = (SELECT MAX(kpi_date) FROM mk_product_performance_kpi_daily)
GROUP BY market_key
ORDER BY risk_stock_cost_value_usd DESC, stock_cost_value_usd DESC, sales_usd_90d DESC
LIMIT $1
`, limit)
if err != nil {
return nil, err
}
defer rows.Close()
out := make([]models.ProductPerformanceMarketRow, 0, limit)
for rows.Next() {
var r models.ProductPerformanceMarketRow
if err := rows.Scan(
&r.MarketKey, &r.Kategori, &r.Seri, &r.YasGrubu, &r.AskiliYan, &r.UrunIlkGrubu, &r.UrunAnaGrubu, &r.UrunAltGrubu,
&r.ProductCount, &r.StarCount, &r.StockRiskCount, &r.StockQty,
&r.StockCostValueUSD, &r.RiskStockCostValueUSD, &r.SalesQty90,
&r.SalesUSD90, &r.GrossProfitUSD90, &r.AvgGrossMargin90, &r.AvgStockDays90,
); err != nil {
return nil, err
}
out = append(out, r)
}
return out, rows.Err()
}
func ListProductPerformanceCountries(ctx context.Context, pg *sql.DB, limit int) ([]models.ProductPerformanceCountryRow, error) {
if err := EnsureProductPerformanceTables(pg); err != nil {
return nil, err
}
if limit <= 0 || limit > 500 {
limit = 100
}
rows, err := pg.QueryContext(ctx, `
SELECT
customer_country,
customer_segment,
market_key,
MAX(kategori) AS kategori,
MAX(seri) AS seri,
COUNT(DISTINCT product_code) AS product_count,
COALESCE(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0) AS sales_qty_90d,
COALESCE(SUM(sales_usd) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0) AS sales_usd_90d,
CASE
WHEN COALESCE(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)=0 THEN 0
ELSE COALESCE(SUM(sales_usd) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)
/ NULLIF(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)
END AS avg_price_usd_90d,
COALESCE(SUM(customer_count) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)::integer AS customer_count_90d,
COALESCE(SUM(invoice_count) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)::integer AS invoice_count_90d,
COALESCE(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '364 days'),0) AS sales_qty_365d,
COALESCE(SUM(sales_usd) FILTER (WHERE sales_date >= current_date - INTERVAL '364 days'),0) AS sales_usd_365d
FROM mk_product_performance_sales_daily
WHERE sales_date >= current_date - INTERVAL '364 days'
GROUP BY customer_country, customer_segment, market_key
ORDER BY sales_usd_90d DESC, sales_qty_90d DESC
LIMIT $1
`, limit)
if err != nil {
return nil, err
}
defer rows.Close()
out := make([]models.ProductPerformanceCountryRow, 0, limit)
for rows.Next() {
var r models.ProductPerformanceCountryRow
if err := rows.Scan(
&r.Country, &r.CustomerSegment, &r.MarketKey, &r.Kategori, &r.Seri,
&r.ProductCount, &r.SalesQty90, &r.SalesUSD90, &r.AvgPriceUSD90,
&r.CustomerCount90, &r.InvoiceCount90, &r.SalesQty365, &r.SalesUSD365,
); err != nil {
return nil, err
}
out = append(out, r)
}
return out, rows.Err()
}
func ListProductPerformanceCustomers(ctx context.Context, pg *sql.DB, breakdown string, limit int) ([]models.ProductPerformanceCustomerRow, error) {
if err := EnsureProductPerformanceTables(pg); err != nil {
return nil, err
}
if limit <= 0 || limit > 500 {
limit = 100
}
mode := strings.ToLower(strings.TrimSpace(breakdown))
if mode == "" {
mode = "market_customer"
}
selectMarket := "''"
selectCountry := "''"
selectSegment := "''"
groupCols := []string{"customer_code", "customer_name"}
switch mode {
case "country_customer":
selectCountry = "customer_country"
selectSegment = "customer_segment"
groupCols = append(groupCols, "customer_country", "customer_segment")
case "market_country_customer":
selectMarket = "market_key"
selectCountry = "customer_country"
selectSegment = "customer_segment"
groupCols = append(groupCols, "market_key", "customer_country", "customer_segment")
default:
mode = "market_customer"
selectMarket = "market_key"
groupCols = append(groupCols, "market_key")
}
query := fmt.Sprintf(`
SELECT
$2::text AS breakdown,
%s AS market_key,
%s AS country,
%s AS customer_segment,
customer_code,
customer_name,
COUNT(DISTINCT product_code) AS product_count,
COALESCE(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0) AS sales_qty_90d,
COALESCE(SUM(sales_usd) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0) AS sales_usd_90d,
CASE
WHEN COALESCE(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)=0 THEN 0
ELSE COALESCE(SUM(sales_usd) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)
/ NULLIF(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)
END AS avg_price_usd_90d,
COALESCE(SUM(invoice_count) FILTER (WHERE sales_date >= current_date - INTERVAL '89 days'),0)::integer AS invoice_count_90d,
COALESCE(SUM(sales_qty) FILTER (WHERE sales_date >= current_date - INTERVAL '364 days'),0) AS sales_qty_365d,
COALESCE(SUM(sales_usd) FILTER (WHERE sales_date >= current_date - INTERVAL '364 days'),0) AS sales_usd_365d,
COALESCE(to_char(MAX(sales_date),'YYYY-MM-DD'),'') AS last_sale_date
FROM mk_product_performance_sales_daily
WHERE sales_date >= current_date - INTERVAL '364 days'
AND COALESCE(customer_code,'') <> ''
GROUP BY %s
ORDER BY sales_usd_90d DESC, sales_qty_90d DESC
LIMIT $1
`, selectMarket, selectCountry, selectSegment, strings.Join(groupCols, ", "))
rows, err := pg.QueryContext(ctx, query, limit, mode)
if err != nil {
return nil, err
}
defer rows.Close()
out := make([]models.ProductPerformanceCustomerRow, 0, limit)
for rows.Next() {
var r models.ProductPerformanceCustomerRow
if err := rows.Scan(
&r.Breakdown, &r.MarketKey, &r.Country, &r.CustomerSegment, &r.CustomerCode, &r.CustomerName,
&r.ProductCount, &r.SalesQty90, &r.SalesUSD90, &r.AvgPriceUSD90, &r.InvoiceCount90,
&r.SalesQty365, &r.SalesUSD365, &r.LastSaleDate,
); err != nil {
return nil, err
}
out = append(out, r)
}
return out, rows.Err()
}
func productPerformanceWhere(f ProductPerformanceFilters) (string, []any) {
parts := make([]string, 0, 6)
args := make([]any, 0, 6)
add := func(cond string, value any) {
args = append(args, value)
parts = append(parts, fmt.Sprintf(cond, len(args)))
}
if q := strings.TrimSpace(f.Search); q != "" {
args = append(args, q)
idx := len(args)
parts = append(parts, fmt.Sprintf(" AND (product_code ILIKE '%%' || $%d || '%%' OR item_description ILIKE '%%' || $%d || '%%')", idx, idx))
}
if v := strings.TrimSpace(f.ProductCode); v != "" {
add(" AND product_code = $%d", v)
}
if v := strings.TrimSpace(f.MarketKey); v != "" {
add(" AND market_key = $%d", v)
}
if v := strings.TrimSpace(f.Kategori); v != "" {
add(" AND kategori = $%d", v)
}
if v := strings.TrimSpace(f.Seri); v != "" {
args = append(args, v)
idx := len(args)
parts = append(parts, fmt.Sprintf(" AND (seri = $%d OR urun_ana_grubu = $%d)", idx, idx))
}
if v := strings.TrimSpace(f.Bucket); v != "" {
add(" AND performance_bucket = $%d", v)
}
return strings.Join(parts, ""), args
}
func dateOnly(t time.Time) time.Time {
y, m, d := t.Date()
return time.Date(y, m, d, 0, 0, 0, 0, t.Location())
}
type productPerformanceSalesDaily struct {
SalesDate time.Time
ProductCode string
ColorCode string
YakaKodu string
ItemDescription string
Kategori string
Seri string
YasGrubu string
AskiliYan string
UrunIlkGrubu string
UrunAnaGrubu string
UrunAltGrubu string
MarketKey string
ChannelCode string
CustomerCountry string
CustomerSegment string
CustomerCode string
CustomerName string
SalesQty float64
SalesTL float64
SalesUSD float64
AvgPriceUSD float64
InvoiceLineCount int
InvoiceCount int
CustomerCount int
LastRefNumber string
}
type productPerformanceStockDaily struct {
StockDate time.Time
ProductCode string
ColorCode string
YakaKodu string
StockQty float64
InQty float64
OutQty float64
KpiInQty float64
KpiOutQty float64
SalesMovementQty float64
ProductionInQty float64
PurchaseInQty float64
ConsumptionOutQty float64
CountDiffQty float64
}