Files
apskel-pos-backend/internal/service/wallet_reconciliation_job.go
T
efrilmandClaude Opus 5.5 040780cd2d feat(wallet): reconcile balances, ledger and lots on a schedule
Adds the reconciliation of docs/prd-point-coin.md §7.5 (PC-108). One
aggregate query per check, across every wallet:

- wallet balance = SUM(ledger), per currency, including customers with
  ledger rows but no wallet row
- wallet balance = SUM(lot remaining)
- lot original - SUM(allocations) = remaining
- SUM(allocations) = |amount| for every deduction
- lots created = amount for every addition, which the engine keeps and the
  other checks rely on

The check on payments.points_used waits for that column (PC-305).

WalletReconciliationJob runs the checks at startup and every six hours,
alongside the omset scheduler. It is silent while the data is consistent.
Each discrepancy is logged with its check, customer, object and the
expected and actual values, and the organization's admins, owners and
managers get a high-priority notification. An organization is notified
again only when its set of discrepancies changes. Nothing is corrected
automatically. At most 50 discrepancies per check are reported.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-30 10:13:46 +07:00

213 lines
6.3 KiB
Go

package service
import (
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"sort"
"sync"
"time"
"github.com/google/uuid"
"apskel-pos-be/internal/entities"
"apskel-pos-be/internal/logger"
"apskel-pos-be/internal/models"
"apskel-pos-be/internal/repository"
)
const (
defaultWalletReconciliationInterval = 6 * time.Hour
// Per check, so one systematic bug cannot flood the log or the notification.
walletReconciliationLimit = 50
)
type walletDiscrepancyFinder interface {
FindDiscrepancies(ctx context.Context, limit int) ([]repository.WalletDiscrepancy, error)
}
type organizationUserLister interface {
GetByOrganizationID(ctx context.Context, organizationID uuid.UUID) ([]*entities.User, error)
}
type notificationSender interface {
Send(ctx context.Context, req *models.SendNotificationRequest) (*models.NotificationResponse, error)
}
// WalletReconciliationJob periodically runs the §7.5 checks of
// docs/prd-point-coin.md over every wallet (PC-108). It is silent while the data is
// consistent. When it finds a discrepancy it logs each one and notifies the admins,
// owners and managers of the organization concerned.
//
// An organization is notified again only when its set of discrepancies changes, so an
// unfixed problem does not page the same people every run. That memory is in-process:
// a restart notifies once more, and each running instance keeps its own.
type WalletReconciliationJob struct {
finder walletDiscrepancyFinder
users organizationUserLister
notifier notificationSender
mu sync.Mutex
notified map[uuid.UUID]string // organization -> fingerprint last notified
stopCh chan struct{}
stopOnce sync.Once
}
func NewWalletReconciliationJob(finder walletDiscrepancyFinder, users organizationUserLister, notifier notificationSender) *WalletReconciliationJob {
return &WalletReconciliationJob{
finder: finder,
users: users,
notifier: notifier,
notified: make(map[uuid.UUID]string),
stopCh: make(chan struct{}),
}
}
// Start runs the checks once now and then every interval, in the background.
func (j *WalletReconciliationJob) Start(interval time.Duration) {
if interval <= 0 {
interval = defaultWalletReconciliationInterval
}
go func() {
j.runLogged()
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
j.runLogged()
case <-j.stopCh:
return
}
}
}()
logger.NonContext.Infof("Wallet reconciliation job started (interval: %s)", interval)
}
func (j *WalletReconciliationJob) Stop() {
j.stopOnce.Do(func() { close(j.stopCh) })
}
func (j *WalletReconciliationJob) runLogged() {
if _, err := j.RunOnce(context.Background()); err != nil {
logger.NonContext.Error("Wallet reconciliation failed to run", err)
}
}
// RunOnce runs every check, reports what it finds, and returns it.
func (j *WalletReconciliationJob) RunOnce(ctx context.Context) ([]repository.WalletDiscrepancy, error) {
found, err := j.finder.FindDiscrepancies(ctx, walletReconciliationLimit)
if err != nil {
return nil, err
}
byOrg := make(map[uuid.UUID][]repository.WalletDiscrepancy)
for _, d := range found {
fields := map[string]interface{}{
"check": d.Check,
"organization_id": d.OrganizationID.String(),
"customer_id": d.CustomerID.String(),
"currency": d.Currency,
"expected": d.Expected,
"actual": d.Actual,
}
if d.ObjectID != nil {
fields["object_id"] = d.ObjectID.String()
}
logger.NonContext.WarnWithFields("Wallet reconciliation found a discrepancy", fields, nil)
byOrg[d.OrganizationID] = append(byOrg[d.OrganizationID], d)
}
j.mu.Lock()
defer j.mu.Unlock()
// Organizations that are clean again are forgotten, so a later problem notifies.
for org := range j.notified {
if _, still := byOrg[org]; !still {
delete(j.notified, org)
}
}
for org, discrepancies := range byOrg {
fingerprint := walletDiscrepancyFingerprint(discrepancies)
if j.notified[org] == fingerprint {
continue
}
if err := j.notify(ctx, org, discrepancies); err != nil {
logger.NonContext.Error(fmt.Sprintf("Wallet reconciliation could not notify organization %s", org), err)
continue
}
j.notified[org] = fingerprint
}
return found, nil
}
func (j *WalletReconciliationJob) notify(ctx context.Context, organizationID uuid.UUID, discrepancies []repository.WalletDiscrepancy) error {
if organizationID == uuid.Nil {
return fmt.Errorf("discrepancy without an organization")
}
users, err := j.users.GetByOrganizationID(ctx, organizationID)
if err != nil {
return err
}
var receivers []uuid.UUID
for _, u := range users {
switch u.Role {
case entities.RoleAdmin, entities.RoleOwner, entities.RoleManager:
receivers = append(receivers, u.ID)
}
}
if len(receivers) == 0 {
return nil
}
perCheck := map[string]int{}
customers := map[string]bool{}
for _, d := range discrepancies {
perCheck[d.Check]++
customers[d.CustomerID.String()] = true
}
customerIDs := make([]string, 0, len(customers))
for id := range customers {
customerIDs = append(customerIDs, id)
}
sort.Strings(customerIDs)
_, err = j.notifier.Send(ctx, &models.SendNotificationRequest{
Title: "Selisih saldo EnakPoint/EnakCoin terdeteksi",
Body: fmt.Sprintf("Pemeriksaan rutin menemukan %d selisih pada saldo %d customer. Saldo belum dikoreksi otomatis; tim teknis perlu memeriksanya.",
len(discrepancies), len(customerIDs)),
Type: "system",
Category: "wallet_reconciliation",
Priority: entities.NotificationPriorityHigh,
NotifiableType: "organization",
NotifiableID: &organizationID,
ReceiverIDs: receivers,
Data: map[string]interface{}{
"organization_id": organizationID.String(),
"discrepancies": len(discrepancies),
"per_check": perCheck,
"customer_ids": customerIDs,
},
})
return err
}
// walletDiscrepancyFingerprint identifies a set of discrepancies regardless of order.
func walletDiscrepancyFingerprint(discrepancies []repository.WalletDiscrepancy) string {
keys := make([]string, 0, len(discrepancies))
for _, d := range discrepancies {
object := ""
if d.ObjectID != nil {
object = d.ObjectID.String()
}
keys = append(keys, fmt.Sprintf("%s|%s|%s|%s|%d|%d", d.Check, d.CustomerID, d.Currency, object, d.Expected, d.Actual))
}
sort.Strings(keys)
h := sha256.New()
for _, k := range keys {
h.Write([]byte(k))
h.Write([]byte{'\n'})
}
return hex.EncodeToString(h.Sum(nil))
}