// app/billing/billing.go // Модуль учёта использования SQS-операций для биллинга. // Записывает каждую успешную SQS-операцию в PostgreSQL: tenant_id, operation, msg_count, msg_bytes. // Если PostgreSQL не сконфигурирован — billing отключён, SQS работает как раньше. // Created: 2026-04-12 package billing import ( "database/sql" "fmt" "os" // PostgreSQL драйвер — регистрируется в database/sql через init() _ "github.com/lib/pq" log "github.com/sirupsen/logrus" ) // db — подключение к PostgreSQL для записи usage-данных. // nil если billing отключён. var db *sql.DB // Init — подключается к PostgreSQL и создаёт таблицу sqs_usage_records если не существует. // Env переменные: BILLING_PG_HOST, BILLING_PG_PORT, BILLING_PG_DATABASE, BILLING_PG_USER, // BILLING_PG_PASSWORD, BILLING_PG_SSLMODE. // Если BILLING_PG_HOST не задан — billing отключён, сервис работает без него. func Init() { host := os.Getenv("BILLING_PG_HOST") if host == "" { log.Info("billing: BILLING_PG_HOST not set, usage tracking disabled") return } port := os.Getenv("BILLING_PG_PORT") if port == "" { port = "5432" } dbname := os.Getenv("BILLING_PG_DATABASE") user := os.Getenv("BILLING_PG_USER") password := os.Getenv("BILLING_PG_PASSWORD") sslmode := os.Getenv("BILLING_PG_SSLMODE") if sslmode == "" { sslmode = "require" } dsn := fmt.Sprintf("host=%s port=%s dbname=%s user=%s password=%s sslmode=%s", host, port, dbname, user, password, sslmode) var err error db, err = sql.Open("postgres", dsn) if err != nil { log.Errorf("billing: failed to open PostgreSQL: %v", err) return } // Проверяем реальное подключение (Open не подключается) if err = db.Ping(); err != nil { log.Errorf("billing: failed to connect to PostgreSQL: %v", err) db.Close() db = nil return } // Ограничиваем пул — billing не должен отжирать коннекты у основной БД db.SetMaxOpenConns(5) db.SetMaxIdleConns(2) if err = autoMigrate(); err != nil { log.Errorf("billing: failed to create table: %v", err) db.Close() db = nil return } log.Infof("billing: connected to PostgreSQL %s:%s/%s, usage tracking enabled", host, port, dbname) } // autoMigrate — создаёт таблицу и индекс если не существуют. // Идемпотентно — безопасно вызывать при каждом старте. func autoMigrate() error { _, err := db.Exec(` CREATE TABLE IF NOT EXISTS sqs_usage_records ( id BIGSERIAL PRIMARY KEY, tenant_id TEXT NOT NULL, operation TEXT NOT NULL, queue_name TEXT NOT NULL DEFAULT '', msg_count INTEGER DEFAULT 1, msg_bytes BIGINT DEFAULT 0, recorded_at TIMESTAMPTZ DEFAULT NOW() ); CREATE INDEX IF NOT EXISTS idx_sqs_usage_tenant_time ON sqs_usage_records(tenant_id, recorded_at); -- Миграция: добавляем queue_name если таблица уже существует без него DO $$ BEGIN ALTER TABLE sqs_usage_records ADD COLUMN queue_name TEXT NOT NULL DEFAULT ''; EXCEPTION WHEN duplicate_column THEN NULL; END $$; `) return err } // RecordUsage — записывает одну SQS-операцию в таблицу биллинга. // Вызывается асинхронно (горутина) чтобы не добавлять latency к SQS-ответу. // Если billing отключён — no-op. Ошибки логируются, SQS-операция не ломается. func RecordUsage(tenantID, operation, queueName string, msgCount int, msgBytes int64) { if db == nil { return } go func() { _, err := db.Exec( `INSERT INTO sqs_usage_records (tenant_id, operation, queue_name, msg_count, msg_bytes) VALUES ($1, $2, $3, $4, $5)`, tenantID, operation, queueName, msgCount, msgBytes, ) if err != nil { log.Errorf("billing: failed to record usage [%s/%s]: %v", tenantID, operation, err) } }() } // Enabled — возвращает true если billing подключён к PostgreSQL func Enabled() bool { return db != nil } // Close — закрывает подключение к PostgreSQL. Вызывается при graceful shutdown. func Close() { if db != nil { db.Close() db = nil log.Info("billing: PostgreSQL connection closed") } }