balíček úložiska
import (
"database/sql"
"errors"
"fmt"
"time"
"github.com/company/my-service/internal/event"
"github.com/company/my-service/internal/metric"
)
type Partner struct {
ID string `db:"id" json:"id"`
Balance float32 `db:"balance" json:"balance"`
Hold float32 `db:"hold" json:"hold"`
}
type PayoutRepository struct {
db sql.DB
eventPublisher event.Publisher
metricCollector metric.Collector
totalPayouts int
}
func NewPayoutRepository() *PayoutRepository {
db, err := sql.Open("postgres", "host=localhost user=postgres password=postgres dbname=affiliate")
if err != nil {
panic(err)
}
r := PayoutRepository{
db: db,
eventPublisher: event.NewKafkaPublisher("localhost:9092"),
metricCollector: metric.NewPrometheusCollector("localhost:9090"),
}
go func() {
for range time.Tick(time.Minute) {
_ = r.metricCollector.Send("TotalPayouts", r.totalPayouts)
}
}()
return r
}
func (r *PayoutRepository) Payout(advertiserID, partnerID string, amount float32) error {
advertiser, err := r.find(advertiserID)
if err != nil {
return err
}
partner, err := r.find(partnerID)
if err != nil {
return err
}
advertiser.Balance -= amount
partner.Hold -= amount
partner.Balance += amount
if err := r.save(advertiser); err != nil {
return err
}
if err := r.save(partner); err != nil {
return err
}
r.totalPayouts++
return nil
}
func (r *PayoutRepository) find(id string) (Partner, error) {
p := Partner{}
err := r.db.QueryRow("SELECT * FROM partners WHERE id = " + id).Scan(&p)
if err != nil {
if err == sql.ErrNoRows {
return p, errors.New("partner not found")
}
return p, err
}
return p, nil
}
func (r *PayoutRepository) save(p Partner) error {
_, err := r.db.Exec(
fmt.Sprintf("UPDATE partners SET balance = %s, hold = %s WHERE id = %s",
p.Balance, p.Hold, p.ID),
)
if err != nil {
return err
}
if err := r.eventPublisher.Publish("BalanceChanged", p); err != nil {
return err
}
return nil
}