202 lines
5.2 KiB
Go
202 lines
5.2 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"net/http"
|
|
"path/filepath"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"git.k6n.net/mats/go-cart-actor/pkg/actor"
|
|
"git.k6n.net/mats/go-cart-actor/pkg/order"
|
|
)
|
|
|
|
type orderAlert struct {
|
|
Severity string `json:"severity"` // warn | error
|
|
Message string `json:"message"`
|
|
OrderId string `json:"orderId,omitempty"`
|
|
}
|
|
|
|
type statusBucket struct {
|
|
Status order.Status `json:"status"`
|
|
Count int `json:"count"`
|
|
Total int64 `json:"total"`
|
|
}
|
|
|
|
type revenueSnapshot struct {
|
|
Total int64 `json:"total"`
|
|
Captured int64 `json:"captured"`
|
|
Refunded int64 `json:"refunded"`
|
|
Pending int64 `json:"pending"`
|
|
}
|
|
|
|
type dailyRevenue struct {
|
|
Date string `json:"date"` // YYYY-MM-DD
|
|
Total int64 `json:"total"`
|
|
}
|
|
|
|
type orderStatsResponse struct {
|
|
OrdersTotal int `json:"ordersTotal"`
|
|
Revenue revenueSnapshot `json:"revenue"`
|
|
ByStatus []statusBucket `json:"byStatus"`
|
|
RecentOrders []orderSummary `json:"recentOrders"`
|
|
DailyRevenue30 []dailyRevenue `json:"dailyRevenue30,omitempty"`
|
|
Alerts []orderAlert `json:"alerts"`
|
|
}
|
|
|
|
// statsCache holds an in-memory aggregate of all orders, rebuilt on mutation
|
|
// instead of scanning *.events.log on every request.
|
|
type statsCache struct {
|
|
mu sync.RWMutex
|
|
data orderStatsResponse
|
|
pool *actor.SimpleGrainPool[order.OrderGrain]
|
|
dataDir string
|
|
}
|
|
|
|
func newStatsCache(dataDir string, pool *actor.SimpleGrainPool[order.OrderGrain]) *statsCache {
|
|
return &statsCache{dataDir: dataDir, pool: pool}
|
|
}
|
|
|
|
// Get returns the cached stats. Call Rebuild() on startup to warm the cache;
|
|
// if never built, returns an empty response.
|
|
func (c *statsCache) Get() orderStatsResponse {
|
|
c.mu.RLock()
|
|
defer c.mu.RUnlock()
|
|
return c.data
|
|
}
|
|
|
|
// Rebuild rebuilds the stats cache from disk. Safe to call concurrently;
|
|
// intended as a goroutine after each mutation so it never blocks an HTTP
|
|
// response.
|
|
func (c *statsCache) Rebuild() {
|
|
c.rebuild()
|
|
}
|
|
|
|
func (c *statsCache) rebuild() {
|
|
if c.pool == nil {
|
|
return
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
defer cancel()
|
|
|
|
matches, _ := filepath.Glob(filepath.Join(c.dataDir, "*.events.log"))
|
|
now := time.Now()
|
|
var totalOrders int
|
|
var totalRevenue, capturedRevenue, refundedRevenue int64
|
|
byStatus := map[order.Status]int{}
|
|
statusTotal := map[order.Status]int64{}
|
|
|
|
dailyByDay := map[string]int64{}
|
|
for i := 0; i < 30; i++ {
|
|
day := now.AddDate(0, 0, -i).Format("2006-01-02")
|
|
dailyByDay[day] = 0
|
|
}
|
|
|
|
var allOrders []orderSummary
|
|
var alerts []orderAlert
|
|
|
|
for _, m := range matches {
|
|
base := strings.TrimSuffix(filepath.Base(m), ".events.log")
|
|
raw, err := strconv.ParseUint(base, 10, 64)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
g, err := c.pool.Get(ctx, raw)
|
|
if err != nil || g.Status == order.StatusNew {
|
|
continue
|
|
}
|
|
totalOrders++
|
|
byStatus[g.Status]++
|
|
statusTotal[g.Status] += g.TotalAmount.Int64()
|
|
|
|
totalRevenue += g.TotalAmount.Int64()
|
|
capturedRevenue += g.CapturedAmount.Int64()
|
|
refundedRevenue += g.RefundedAmount.Int64()
|
|
|
|
if g.PlacedAt != "" {
|
|
if placed, parseErr := time.Parse(time.RFC3339, g.PlacedAt); parseErr == nil {
|
|
dayKey := placed.Format("2006-01-02")
|
|
if _, ok := dailyByDay[dayKey]; ok {
|
|
dailyByDay[dayKey] += g.TotalAmount.Int64()
|
|
}
|
|
}
|
|
}
|
|
|
|
allOrders = append(allOrders, orderSummary{
|
|
OrderId: order.OrderId(raw).String(),
|
|
Reference: g.OrderReference,
|
|
Status: g.Status,
|
|
TotalAmount: g.TotalAmount.Int64(),
|
|
Currency: g.Currency,
|
|
PlacedAt: g.PlacedAt,
|
|
})
|
|
|
|
if g.Status == order.StatusPending {
|
|
alerts = append(alerts, orderAlert{
|
|
Severity: "warn",
|
|
Message: "Payment still pending",
|
|
OrderId: order.OrderId(raw).String(),
|
|
})
|
|
}
|
|
if g.Status == order.StatusCancelled {
|
|
alerts = append(alerts, orderAlert{
|
|
Severity: "warn",
|
|
Message: "Order was cancelled",
|
|
OrderId: order.OrderId(raw).String(),
|
|
})
|
|
}
|
|
}
|
|
|
|
sort.Slice(allOrders, func(i, j int) bool {
|
|
return allOrders[i].PlacedAt > allOrders[j].PlacedAt
|
|
})
|
|
recent := allOrders
|
|
if len(recent) > 10 {
|
|
recent = recent[:10]
|
|
}
|
|
|
|
if len(alerts) > 50 {
|
|
alerts = alerts[:50]
|
|
}
|
|
|
|
buckets := make([]statusBucket, 0, len(byStatus))
|
|
for st, cnt := range byStatus {
|
|
buckets = append(buckets, statusBucket{Status: st, Count: cnt, Total: statusTotal[st]})
|
|
}
|
|
sort.Slice(buckets, func(i, j int) bool {
|
|
return buckets[i].Count > buckets[j].Count
|
|
})
|
|
|
|
var dailyRev []dailyRevenue
|
|
for i := 29; i >= 0; i-- {
|
|
day := now.AddDate(0, 0, -i).Format("2006-01-02")
|
|
dailyRev = append(dailyRev, dailyRevenue{Date: day, Total: dailyByDay[day]})
|
|
}
|
|
|
|
c.mu.Lock()
|
|
c.data = orderStatsResponse{
|
|
OrdersTotal: totalOrders,
|
|
Revenue: revenueSnapshot{
|
|
Total: totalRevenue,
|
|
Captured: capturedRevenue,
|
|
Refunded: refundedRevenue,
|
|
Pending: totalRevenue - capturedRevenue - refundedRevenue,
|
|
},
|
|
ByStatus: buckets,
|
|
RecentOrders: recent,
|
|
DailyRevenue30: dailyRev,
|
|
Alerts: alerts,
|
|
}
|
|
c.mu.Unlock()
|
|
}
|
|
|
|
// handleStats returns dashboard stats from the in-memory cache (O(1) read,
|
|
// rebuilt after each mutation) instead of scanning all *.events.log on every
|
|
// request.
|
|
func (s *server) handleStats(w http.ResponseWriter, r *http.Request) {
|
|
writeJSON(w, http.StatusOK, s.statsCache.Get())
|
|
}
|