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)) }