Files
rogee b6d0af1a56
management-images / build-and-publish (push) Successful in 9m15s
docs: add one-click management deployment and image workflow
2026-09-16 17:57:05 +08:00

1345 lines
55 KiB
Go

package management
import (
"context"
"database/sql"
_ "embed"
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
"sync"
"time"
dbgen "github.com/rogee/agent-call-management/internal/db"
_ "modernc.org/sqlite"
)
//go:embed migrations/001_init.sql
var migrationSQL string
type Store struct {
db *sql.DB
queries *dbgen.Queries
mu sync.Mutex
}
func NewStore(path string) (*Store, error) {
if path == "" {
path = "./data/management.db"
}
if path != ":memory:" && !strings.HasPrefix(path, "file:") {
if err := os.MkdirAll(filepath.Dir(path), 0o750); err != nil {
return nil, fmt.Errorf("create database directory: %w", err)
}
}
db, err := sql.Open("sqlite", path)
if err != nil {
return nil, fmt.Errorf("open database: %w", err)
}
db.SetMaxOpenConns(1)
db.SetMaxIdleConns(1)
db.SetConnMaxLifetime(0)
s := &Store{db: db, queries: dbgen.New(db)}
if _, err := db.Exec("PRAGMA busy_timeout = 5000; PRAGMA foreign_keys = ON; PRAGMA journal_mode = WAL;"); err != nil {
db.Close()
return nil, fmt.Errorf("configure database: %w", err)
}
if _, err := db.Exec(migrationSQL); err != nil {
db.Close()
return nil, fmt.Errorf("migrate database: %w", err)
}
if err := ensureSchemaCompatibility(db); err != nil {
db.Close()
return nil, err
}
if _, err := db.Exec("INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES(1, ?)", utcString(time.Now())); err != nil {
db.Close()
return nil, fmt.Errorf("record migration: %w", err)
}
return s, nil
}
func ensureSchemaCompatibility(db *sql.DB) error {
rows, err := db.Query("PRAGMA table_info(egress_pools)")
if err != nil {
return fmt.Errorf("inspect database schema: %w", err)
}
defer rows.Close()
hasProvider := false
for rows.Next() {
var cid, name, typ string
var notNull, pk int
var defaultValue any
if err := rows.Scan(&cid, &name, &typ, &notNull, &defaultValue, &pk); err != nil {
return err
}
if name == "provider_id" {
hasProvider = true
}
}
if err := rows.Err(); err != nil {
return err
}
if !hasProvider {
if _, err := db.Exec("ALTER TABLE egress_pools ADD COLUMN provider_id TEXT NOT NULL DEFAULT ''"); err != nil {
return fmt.Errorf("upgrade egress pool schema: %w", err)
}
}
return nil
}
func NewMemoryStore() (*Store, error) {
return NewStore("file:management-" + newID("test") + "?mode=memory&cache=shared")
}
func (s *Store) Close() error {
if s == nil || s.db == nil {
return nil
}
return s.db.Close()
}
func (s *Store) withTx(ctx context.Context, fn func(*sql.Tx) error) error {
s.mu.Lock()
defer s.mu.Unlock()
tx, err := s.db.BeginTx(ctx, nil)
if err != nil {
return err
}
if err := fn(tx); err != nil {
_ = tx.Rollback()
return err
}
return tx.Commit()
}
func (s *Store) withRead(_ context.Context, fn func(*sql.DB) error) error {
s.mu.Lock()
defer s.mu.Unlock()
return fn(s.db)
}
func (s *Store) GetProvider(ctx context.Context, id string) (Provider, error) {
var out Provider
err := s.withRead(ctx, func(_ *sql.DB) error {
row, err := s.queries.GetProviderRow(ctx, id)
if err != nil {
return err
}
trunkCount, err := s.trunkCount(ctx, id)
if err != nil {
return err
}
out = Provider{ProviderID: row.ProviderID, DisplayName: row.DisplayName, Notes: row.Notes, Lifecycle: row.Lifecycle, Revision: row.Revision, TrunkCount: trunkCount, CreatedAt: row.CreatedAt, UpdatedAt: row.UpdatedAt, UpdatedBy: row.UpdatedBy}
return nil
})
return out, err
}
func (s *Store) ListProviders(ctx context.Context) ([]Provider, error) {
var out []Provider
err := s.withRead(ctx, func(_ *sql.DB) error {
rows, err := s.queries.ListProviderRows(ctx)
if err != nil {
return err
}
out = make([]Provider, 0, len(rows))
for _, row := range rows {
out = append(out, Provider{ProviderID: row.ProviderID, DisplayName: row.DisplayName, Notes: row.Notes, Lifecycle: row.Lifecycle, Revision: row.Revision, TrunkCount: int(row.TrunkCount), CreatedAt: row.CreatedAt, UpdatedAt: row.UpdatedAt, UpdatedBy: row.UpdatedBy})
}
return nil
})
return out, err
}
func (s *Store) trunkCount(ctx context.Context, providerID string) (int, error) {
var count int
err := s.db.QueryRowContext(ctx, "SELECT COUNT(*) FROM trunks WHERE provider_id = ?", providerID).Scan(&count)
return count, err
}
func (s *Store) UpsertProvider(ctx context.Context, id string, input ProviderInput, expected int64, actor, requestID string) (Provider, int, error) {
if !validID(id) {
return Provider{}, 0, newAppError(422, "INVALID_ID", "provider_id is invalid", map[string]any{"provider_id": id})
}
if input.Lifecycle == "" {
input.Lifecycle = "active"
}
if input.Lifecycle != "active" && input.Lifecycle != "archived" {
return Provider{}, 0, newAppError(422, "INVALID_LIFECYCLE", "lifecycle must be active or archived", nil)
}
if strings.TrimSpace(input.DisplayName) == "" || len(input.DisplayName) > 256 || strings.ContainsAny(input.DisplayName, "\r\n") {
return Provider{}, 0, newAppError(422, "INVALID_DISPLAY_NAME", "display_name is invalid", nil)
}
if len(input.Notes) > 2000 || strings.ContainsAny(input.Notes, "\r\n") {
return Provider{}, 0, newAppError(422, "INVALID_NOTES", "notes is invalid", nil)
}
var status int
err := s.withTx(ctx, func(tx *sql.Tx) error {
var currentRevision int64
err := tx.QueryRowContext(ctx, "SELECT revision FROM providers WHERE provider_id=?", id).Scan(&currentRevision)
exists := err == nil
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return err
}
if exists && currentRevision != expected {
return newAppError(409, "REVISION_CONFLICT", "provider revision changed", map[string]any{"expected_revision": expected, "current_revision": currentRevision})
}
if !exists && expected != 0 {
return newAppError(409, "REVISION_CONFLICT", "new provider requires If-Match: 0", map[string]any{"expected_revision": expected, "current_revision": 0})
}
if input.Lifecycle == "archived" {
var enabled, occupied int
if err := tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM trunks WHERE provider_id=? AND active_status <> 'disabled'", id).Scan(&enabled); err != nil {
return err
}
if err := tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM attempts WHERE provider_id=? AND (ended_at IS NULL OR ended_at='') AND origin_status IN ('confirmed','uncertain')", id).Scan(&occupied); err != nil {
return err
}
if enabled > 0 || occupied > 0 {
return newAppError(409, "PROVIDER_HAS_DEPENDENCIES", "provider cannot be archived while enabled trunks or active calls remain", map[string]any{"enabled_trunks": enabled, "active_attempts": occupied})
}
}
now := utcString(time.Now())
revision := int64(1)
if exists {
revision = currentRevision + 1
_, err = tx.ExecContext(ctx, "UPDATE providers SET display_name=?, notes=?, lifecycle=?, revision=?, updated_at=?, updated_by=? WHERE provider_id=?", input.DisplayName, input.Notes, input.Lifecycle, revision, now, actor, id)
} else {
_, err = tx.ExecContext(ctx, "INSERT INTO providers(provider_id,display_name,notes,lifecycle,revision,created_at,updated_at,created_by,updated_by) VALUES(?,?,?,?,?,?,?,?,?)", id, input.DisplayName, input.Notes, input.Lifecycle, revision, now, now, actor, actor)
}
if err != nil {
return err
}
details, _ := json.Marshal(map[string]any{"after": input, "provider_revision": revision})
if err := insertAudit(ctx, tx, "provider", id, "upsert", revision, actor, requestID, string(details), now); err != nil {
return err
}
status = 200
if !exists {
status = 201
}
return nil
})
if err != nil {
return Provider{}, 0, err
}
out, err := s.GetProvider(ctx, id)
return out, status, err
}
func (s *Store) GetCell(ctx context.Context, id string) (Cell, error) {
var cell Cell
err := s.withRead(ctx, func(db *sql.DB) error {
var configJSON string
err := db.QueryRowContext(ctx, "SELECT cell_id,revision,config_json,cloud_instance_id,instance_name,region,boot_id,updated_at,updated_by FROM cells WHERE cell_id=?", id).Scan(&cell.CellID, &cell.Revision, &configJSON, &cell.CloudInstanceID, &cell.InstanceName, &cell.Region, &cell.BootID, &cell.UpdatedAt, &cell.UpdatedBy)
if err != nil {
return err
}
if err := json.Unmarshal([]byte(configJSON), &cell.Config); err != nil {
return fmt.Errorf("decode cell config: %w", err)
}
cell.Mode = ModeMock
return nil
})
return cell, err
}
func (s *Store) ListCells(ctx context.Context) ([]Cell, error) {
var out []Cell
err := s.withRead(ctx, func(db *sql.DB) error {
rows, err := db.QueryContext(ctx, "SELECT cell_id,revision,config_json,cloud_instance_id,instance_name,region,boot_id,updated_at,updated_by FROM cells ORDER BY cell_id")
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var cell Cell
var configJSON string
if err := rows.Scan(&cell.CellID, &cell.Revision, &configJSON, &cell.CloudInstanceID, &cell.InstanceName, &cell.Region, &cell.BootID, &cell.UpdatedAt, &cell.UpdatedBy); err != nil {
return err
}
if err := json.Unmarshal([]byte(configJSON), &cell.Config); err != nil {
return err
}
cell.Mode = ModeMock
out = append(out, cell)
}
return rows.Err()
})
return out, err
}
func (s *Store) UpsertCell(ctx context.Context, id string, input CellConfig, expected int64, actor, requestID, cloudInstanceID, instanceName, region string) (Cell, int, error) {
if !validID(id) {
return Cell{}, 0, newAppError(422, "INVALID_ID", "cell_id is invalid", nil)
}
if err := validateCellConfig(input); err != nil {
return Cell{}, 0, err
}
if cloudInstanceID == "" {
cloudInstanceID = id
}
if len(cloudInstanceID) > 256 || len(instanceName) > 256 || len(region) > 128 || strings.ContainsAny(cloudInstanceID+instanceName+region, "\r\n") {
return Cell{}, 0, newAppError(422, "INVALID_CELL_METADATA", "cell metadata is invalid", nil)
}
configJSON, err := json.Marshal(input)
if err != nil {
return Cell{}, 0, err
}
var status int
err = s.withTx(ctx, func(tx *sql.Tx) error {
var currentRevision int64
err := tx.QueryRowContext(ctx, "SELECT revision FROM cells WHERE cell_id=?", id).Scan(&currentRevision)
exists := err == nil
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return err
}
if exists && currentRevision != expected {
return newAppError(409, "REVISION_CONFLICT", "cell revision changed", map[string]any{"expected_revision": expected, "current_revision": currentRevision})
}
if !exists && expected != 0 {
return newAppError(409, "REVISION_CONFLICT", "new cell requires If-Match: 0", nil)
}
now := utcString(time.Now())
revision := int64(1)
if exists {
revision = currentRevision + 1
_, err = tx.ExecContext(ctx, "UPDATE cells SET revision=?, config_json=?, cloud_instance_id=?, instance_name=?, region=?, updated_at=?, updated_by=? WHERE cell_id=?", revision, string(configJSON), cloudInstanceID, instanceName, region, now, actor, id)
} else {
_, err = tx.ExecContext(ctx, "INSERT INTO cells(cell_id,revision,config_json,cloud_instance_id,instance_name,region,boot_id,created_at,updated_at,updated_by) VALUES(?,?,?,?,?,?,?,?,?,?)", id, revision, string(configJSON), cloudInstanceID, instanceName, region, "", now, now, actor)
}
if err != nil {
return err
}
details, _ := json.Marshal(map[string]any{"after": input, "cell_revision": revision})
if err := insertAudit(ctx, tx, "cell", id, "upsert", revision, actor, requestID, string(details), now); err != nil {
return err
}
status = 200
if !exists {
status = 201
}
return nil
})
if err != nil {
return Cell{}, 0, err
}
out, err := s.GetCell(ctx, id)
return out, status, err
}
func validateCellConfig(input CellConfig) error {
if !validID(input.EgressPoolID) || input.EgressPoolID == "" {
return newAppError(422, "INVALID_EGRESS_POOL", "egress_pool_id is invalid", nil)
}
if len(input.CodecCapabilities) == 0 || input.MaxConcurrency < 1 || input.MaxConcurrency > 1000000 {
return newAppError(422, "INVALID_CELL_CAPACITY", "codec_capabilities and max_concurrency are invalid", nil)
}
for _, codec := range input.CodecCapabilities {
if codec != "PCMA" && codec != "PCMU" {
return newAppError(422, "UNSUPPORTED_CODEC", "codec_capabilities contains an unsupported codec", nil)
}
}
if input.Status != "healthy" && input.Status != "draining" && input.Status != "disabled" {
return newAppError(422, "INVALID_CELL_STATUS", "status must be healthy, draining, or disabled", nil)
}
if strings.ContainsAny(input.ManagementURL, "\r\n") {
return newAppError(422, "INVALID_MANAGEMENT_URL", "management_url is invalid", nil)
}
return nil
}
func (s *Store) UpsertEgressPool(ctx context.Context, id, name string, ips []string, whitelistStatus, source string, checkedAt *time.Time) error {
if !validID(id) || name == "" {
return fmt.Errorf("invalid egress pool")
}
b, err := json.Marshal(ips)
if err != nil {
return err
}
var checked any
if checkedAt != nil {
checked = utcString(*checkedAt)
}
_, err = s.db.ExecContext(ctx, "INSERT INTO egress_pools(egress_pool_id,provider_id,display_name,fixed_ips_json,whitelist_status,whitelist_checked_at,source,updated_at) VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(egress_pool_id) DO UPDATE SET display_name=excluded.display_name,fixed_ips_json=excluded.fixed_ips_json,whitelist_status=excluded.whitelist_status,whitelist_checked_at=excluded.whitelist_checked_at,source=excluded.source,updated_at=excluded.updated_at", id, "", name, string(b), whitelistStatus, checked, source, utcString(time.Now()))
return err
}
func (s *Store) SetEgressPoolProvider(ctx context.Context, poolID, providerID string) error {
if !validID(poolID) || !validID(providerID) {
return fmt.Errorf("invalid egress/provider id")
}
if _, err := s.GetProvider(ctx, providerID); err != nil {
return err
}
_, err := s.db.ExecContext(ctx, "UPDATE egress_pools SET provider_id=?,updated_at=? WHERE egress_pool_id=?", providerID, utcString(time.Now()), poolID)
return err
}
func (s *Store) ListEgressPools(ctx context.Context) ([]map[string]any, error) {
var out []map[string]any
err := s.withRead(ctx, func(db *sql.DB) error {
rows, err := db.QueryContext(ctx, "SELECT egress_pool_id,provider_id,display_name,fixed_ips_json,whitelist_status,whitelist_checked_at,source,updated_at FROM egress_pools ORDER BY egress_pool_id")
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var id, providerID, name, ipsJSON, status, source, updated string
var checked sql.NullString
if err := rows.Scan(&id, &providerID, &name, &ipsJSON, &status, &checked, &source, &updated); err != nil {
return err
}
var ips []string
if err := json.Unmarshal([]byte(ipsJSON), &ips); err != nil {
return err
}
item := map[string]any{"egress_pool_id": id, "provider_id": providerID, "display_name": name, "fixed_ips": ips, "whitelist_status": status, "source": source, "updated_at": updated}
if checked.Valid {
item["whitelist_checked_at"] = checked.String
} else {
item["whitelist_checked_at"] = nil
}
out = append(out, item)
}
return rows.Err()
})
return out, err
}
func (s *Store) compatibleCellIDs(ctx context.Context, cfg TrunkConfig) ([]string, error) {
cells, err := s.ListCells(ctx)
if err != nil {
return nil, err
}
var ids []string
for _, cell := range cells {
if cell.Config.Status == "disabled" || cell.Config.EgressPoolID != cfg.EgressPoolID || len(codecIntersection(cfg.CodecProfile.Allowed, cell.Config.CodecCapabilities)) == 0 {
continue
}
ids = append(ids, cell.CellID)
}
return sortedStrings(ids), nil
}
func (s *Store) GetTrunk(ctx context.Context, id, mode string) (AdminTrunk, error) {
return s.buildAdminTrunk(ctx, id, mode)
}
func (s *Store) ListTrunks(ctx context.Context, mode, providerID string) ([]AdminTrunk, error) {
var ids []string
err := s.withRead(ctx, func(db *sql.DB) error {
query := "SELECT trunk_id FROM trunks"
args := []any{}
if providerID != "" {
query += " WHERE provider_id=?"
args = append(args, providerID)
}
query += " ORDER BY updated_at DESC, trunk_id ASC"
rows, err := db.QueryContext(ctx, query, args...)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var id string
if err := rows.Scan(&id); err != nil {
return err
}
ids = append(ids, id)
}
return rows.Err()
})
if err != nil {
return nil, err
}
out := make([]AdminTrunk, 0, len(ids))
for _, id := range ids {
trunk, err := s.GetTrunk(ctx, id, mode)
if err != nil {
return nil, err
}
out = append(out, trunk)
}
return out, nil
}
func (s *Store) PutTrunk(ctx context.Context, id string, cfg TrunkConfig, expected int64, actor, requestID string) (AdminTrunk, int, error) {
if !validID(id) {
return AdminTrunk{}, 0, newAppError(422, "INVALID_ID", "trunk_id is invalid", nil)
}
if err := s.validateTrunkConfig(ctx, cfg, false); err != nil {
return AdminTrunk{}, 0, err
}
configJSON, err := json.Marshal(cfg)
if err != nil {
return AdminTrunk{}, 0, err
}
digest := hashBytes(configJSON)
var status int
err = s.withTx(ctx, func(tx *sql.Tx) error {
var currentRevision int64
var currentProvider string
err := tx.QueryRowContext(ctx, "SELECT latest_revision,provider_id FROM trunks WHERE trunk_id=?", id).Scan(&currentRevision, &currentProvider)
exists := err == nil
if err != nil && !errors.Is(err, sql.ErrNoRows) {
return err
}
if exists && currentRevision != expected {
return newAppError(409, "REVISION_CONFLICT", "trunk revision changed", map[string]any{"expected_revision": expected, "current_revision": currentRevision})
}
if !exists && expected != 0 {
return newAppError(409, "REVISION_CONFLICT", "new trunk requires If-Match: 0", nil)
}
now := utcString(time.Now())
revision := int64(1)
if exists {
revision = currentRevision + 1
_, err = tx.ExecContext(ctx, "UPDATE trunks SET provider_id=?,latest_revision=?,updated_at=?,updated_by=? WHERE trunk_id=?", cfg.ProviderID, revision, now, actor, id)
} else {
_, err = tx.ExecContext(ctx, "INSERT INTO trunks(trunk_id,provider_id,latest_revision,active_revision,active_status,updated_at,updated_by) VALUES(?,?,?,?,?,?,?)", id, cfg.ProviderID, revision, 0, "draft", now, actor)
}
if err != nil {
return err
}
_, err = tx.ExecContext(ctx, "INSERT INTO trunk_versions(trunk_id,revision,state,config_json,config_sha256,created_at,created_by) VALUES(?,?,?,?,?,?,?)", id, revision, "draft", string(configJSON), digest, now, actor)
if err != nil {
return err
}
details, _ := json.Marshal(map[string]any{"before_provider_id": currentProvider, "after": redactedConfig(cfg), "config_sha256": digest})
if err := insertAudit(ctx, tx, "trunk", id, "upsert", revision, actor, requestID, string(details), now); err != nil {
return err
}
status = 200
if !exists {
status = 201
}
return nil
})
if err != nil {
return AdminTrunk{}, 0, err
}
trunk, err := s.GetTrunk(ctx, id, ModeMock)
return trunk, status, err
}
func (s *Store) validateTrunkConfig(ctx context.Context, cfg TrunkConfig, forReal bool) error {
if !validID(cfg.ProviderID) || cfg.ProviderID == "" {
return newAppError(422, "PROVIDER_REQUIRED", "provider_id is required", nil)
}
provider, err := s.GetProvider(ctx, cfg.ProviderID)
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return newAppError(422, "PROVIDER_NOT_FOUND", "provider does not exist", nil)
}
return err
}
if provider.Lifecycle != "active" {
return newAppError(422, "PROVIDER_ARCHIVED", "trunk must reference an active provider", nil)
}
if strings.TrimSpace(cfg.DisplayName) == "" || len(cfg.DisplayName) > 256 || strings.ContainsAny(cfg.DisplayName, "\r\n") {
return newAppError(422, "INVALID_DISPLAY_NAME", "display_name is invalid", nil)
}
if err := validateHost(cfg.Sip.Host); err != nil {
return newAppError(422, "INVALID_SIP_HOST", err.Error(), nil)
}
if cfg.Sip.Port < 1 || cfg.Sip.Port > 65535 || (cfg.Sip.Transport != "udp" && cfg.Sip.Transport != "tcp" && cfg.Sip.Transport != "tls") || (cfg.Sip.AuthMode != "ip" && cfg.Sip.AuthMode != "digest") {
return newAppError(422, "INVALID_SIP_CONFIG", "sip transport, port, or auth_mode is invalid", nil)
}
if cfg.Sip.CredentialRef != "" && (len(cfg.Sip.CredentialRef) > 256 || strings.ContainsAny(cfg.Sip.CredentialRef, "\r\n=\\\"'{}[]")) {
return newAppError(422, "INVALID_CREDENTIAL_REF", "credential_ref must be a secret-store reference, not a secret or config fragment", nil)
}
if cfg.Sip.AuthMode == "digest" && cfg.Sip.CredentialRef == "" {
return newAppError(422, "CREDENTIAL_REQUIRED", "digest authentication requires a credential reference", nil)
}
if forReal && strings.HasPrefix(cfg.Sip.CredentialRef, "missing:") {
return newAppError(409, "CREDENTIAL_UNAVAILABLE", "the referenced SIP credential is not available", nil)
}
if len(cfg.CodecProfile.Allowed) == 0 || !containsString(cfg.CodecProfile.Allowed, cfg.CodecProfile.Preferred) {
return newAppError(422, "INVALID_CODEC_PROFILE", "preferred codec must be in allowed", nil)
}
for _, codec := range cfg.CodecProfile.Allowed {
if codec != "PCMA" && codec != "PCMU" {
return newAppError(422, "UNSUPPORTED_CODEC", "only PCMA and PCMU are supported", nil)
}
}
if len(cfg.CallerIDs) == 0 || len(cfg.CallerIDs) > 100 {
return newAppError(422, "INVALID_CALLER_IDS", "at least one caller ID is required", nil)
}
for _, callerID := range cfg.CallerIDs {
if callerID == "" || len(callerID) > 128 || strings.ContainsAny(callerID, "\r\n,;{}[]") {
return newAppError(422, "INVALID_CALLER_ID", "caller_ids contains an invalid value", nil)
}
}
if err := safePrefix(cfg.DialPrefix); err != nil {
return newAppError(422, "INVALID_DIAL_PREFIX", err.Error(), nil)
}
if !validID(cfg.EgressPoolID) || cfg.EgressPoolID == "" || cfg.MaxConcurrency < 1 || cfg.MaxCPS < 1 || cfg.MaxConcurrency > 1000000 || cfg.MaxCPS > 1000000 {
return newAppError(422, "INVALID_CAPACITY", "egress_pool_id, max_concurrency, or max_cps is invalid", nil)
}
var fixedIPsJSON, poolProviderID string
err = s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT fixed_ips_json,provider_id FROM egress_pools WHERE egress_pool_id=?", cfg.EgressPoolID).Scan(&fixedIPsJSON, &poolProviderID)
})
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return newAppError(422, "EGRESS_POOL_NOT_FOUND", "egress pool does not exist", nil)
}
return err
}
if poolProviderID != "" && poolProviderID != cfg.ProviderID {
return newAppError(403, "EGRESS_PROVIDER_MISMATCH", "egress pool is owned by another provider", map[string]any{"egress_pool_id": cfg.EgressPoolID})
}
var ips []string
if err := json.Unmarshal([]byte(fixedIPsJSON), &ips); err != nil {
return err
}
for _, ip := range ips {
if ip == "123.56.71.98" && cfg.Sip.Host == ip {
return newAppError(422, "EGRESS_IP_USED_AS_SIP_HOST", "the fixed egress allowlist IP cannot be used as the SIP host", nil)
}
}
if forReal && cfg.Sip.AuthMode == "digest" {
return newAppError(409, "REAL_CAPABILITY_UNSUPPORTED", "the real Cell adapter is not approved for Digest authentication", nil)
}
return nil
}
func (s *Store) BuildValidation(ctx context.Context, id string, revision int64, mode string) (map[string]any, error) {
cfg, err := s.getTrunkConfig(ctx, id, revision)
if err != nil {
return nil, err
}
issues := []map[string]any{}
if validationErr := s.validateTrunkConfig(ctx, cfg, mode == ModeReal); validationErr != nil {
appErr := asAppError(validationErr)
issues = append(issues, map[string]any{"code": appErr.Code, "field": validationField(appErr.Code), "message": appErr.Message})
}
cells, err := s.compatibleCellIDs(ctx, cfg)
if err != nil {
return nil, err
}
if len(cells) == 0 {
issues = append(issues, map[string]any{"code": "NO_COMPATIBLE_CELL", "field": "egress_pool_id", "message": "no registered Cell matches the egress pool and codec"})
}
return map[string]any{"mode": mode, "trunk_id": id, "revision": revision, "valid": len(issues) == 0, "issues": issues, "target_cell_ids": cells, "impact": map[string]any{"reload_required": true, "active_calls": s.activeAttemptCount(ctx, id)}}, nil
}
func validationField(code string) string {
return map[string]string{
"PROVIDER_REQUIRED": "provider_id", "PROVIDER_NOT_FOUND": "provider_id", "PROVIDER_ARCHIVED": "provider_id",
"INVALID_DISPLAY_NAME": "display_name", "INVALID_SIP_HOST": "sip.host", "INVALID_SIP_CONFIG": "sip",
"INVALID_CREDENTIAL_REF": "sip.credential_ref", "CREDENTIAL_REQUIRED": "sip.credential_ref", "CREDENTIAL_UNAVAILABLE": "sip.credential_ref",
"INVALID_CODEC_PROFILE": "codec_profile", "UNSUPPORTED_CODEC": "codec_profile", "INVALID_CALLER_IDS": "caller_ids", "INVALID_CALLER_ID": "caller_ids",
"INVALID_DIAL_PREFIX": "dial_prefix", "INVALID_CAPACITY": "egress_pool_id", "EGRESS_POOL_NOT_FOUND": "egress_pool_id",
"EGRESS_PROVIDER_MISMATCH": "egress_pool_id", "EGRESS_IP_USED_AS_SIP_HOST": "sip.host", "REAL_CAPABILITY_UNSUPPORTED": "sip.auth_mode",
}[code]
}
func (s *Store) getTrunkConfig(ctx context.Context, id string, revision int64) (TrunkConfig, error) {
var raw string
err := s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT config_json FROM trunk_versions WHERE trunk_id=? AND revision=?", id, revision).Scan(&raw)
})
if err != nil {
return TrunkConfig{}, err
}
var cfg TrunkConfig
if err := json.Unmarshal([]byte(raw), &cfg); err != nil {
return TrunkConfig{}, err
}
return cfg, nil
}
func (s *Store) activeAttemptCount(ctx context.Context, trunkID string) int {
var count int
_ = s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT COUNT(*) FROM attempts WHERE trunk_id=? AND origin_status='confirmed' AND answered_at IS NOT NULL AND ended_at IS NULL", trunkID).Scan(&count)
})
return count
}
func (s *Store) buildAdminTrunk(ctx context.Context, id, mode string) (AdminTrunk, error) {
var meta dbgen.Trunk
err := s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT trunk_id,provider_id,latest_revision,active_revision,active_status,updated_at,updated_by FROM trunks WHERE trunk_id=?", id).Scan(&meta.TrunkID, &meta.ProviderID, &meta.LatestRevision, &meta.ActiveRevision, &meta.ActiveStatus, &meta.UpdatedAt, &meta.UpdatedBy)
})
if err != nil {
return AdminTrunk{}, err
}
latestCfg, err := s.getTrunkConfig(ctx, id, meta.LatestRevision)
if err != nil {
return AdminTrunk{}, err
}
activeCfg := TrunkConfig{}
if meta.ActiveRevision > 0 {
activeCfg, err = s.getTrunkConfig(ctx, id, meta.ActiveRevision)
if err != nil {
return AdminTrunk{}, err
}
}
compatible, err := s.compatibleCellIDs(ctx, latestCfg)
if err != nil {
return AdminTrunk{}, err
}
versions, err := s.listVersions(ctx, id)
if err != nil {
return AdminTrunk{}, err
}
status := "draft"
if meta.ActiveStatus == "disabled" {
status = "disabled"
} else if meta.ActiveRevision > 0 && meta.ActiveRevision == meta.LatestRevision && meta.ActiveStatus == "published" {
status = "published"
}
return AdminTrunk{Mode: mode, TrunkID: id, ProviderID: meta.ProviderID, LatestRevision: meta.LatestRevision, ActiveRevision: meta.ActiveRevision, ActiveStatus: meta.ActiveStatus, Status: status, UpdatedAt: meta.UpdatedAt, CompatibleCellIDs: compatible, Latest: ptrTrunkView(id, latestCfg, versionDigest(versions, meta.LatestRevision)), Active: nilIfRevision(meta.ActiveRevision, id, activeCfg, versionDigest(versions, meta.ActiveRevision)), Versions: versions}, nil
}
func ptrTrunkView(id string, cfg TrunkConfig, digest string) *TrunkView {
configured := cfg.Sip.CredentialRef != ""
cfg.Sip.CredentialRef = ""
return &TrunkView{TrunkConfig: cfg, TrunkID: id, CredentialConfigured: configured, AsteriskAllow: asteriskCodecs(cfg.CodecProfile.Allowed), ConfigSHA256: digest}
}
func nilIfRevision(revision int64, id string, cfg TrunkConfig, digest string) *TrunkView {
if revision == 0 {
return nil
}
return ptrTrunkView(id, cfg, digest)
}
func asteriskCodecs(codecs []string) []string {
var out []string
for _, codec := range codecs {
if codec == "PCMA" {
out = append(out, "alaw")
}
if codec == "PCMU" {
out = append(out, "ulaw")
}
}
return out
}
func versionDigest(versions []RevisionInfo, revision int64) string {
for _, version := range versions {
if version.Revision == revision {
return version.ConfigSHA256
}
}
return ""
}
func (s *Store) listVersions(ctx context.Context, id string) ([]RevisionInfo, error) {
var out []RevisionInfo
err := s.withRead(ctx, func(db *sql.DB) error {
rows, err := db.QueryContext(ctx, "SELECT revision,state,config_sha256,created_at,created_by FROM trunk_versions WHERE trunk_id=? ORDER BY revision DESC", id)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var v RevisionInfo
if err := rows.Scan(&v.Revision, &v.State, &v.ConfigSHA256, &v.CreatedAt, &v.CreatedBy); err != nil {
return err
}
out = append(out, v)
}
return rows.Err()
})
return out, err
}
func (s *Store) GetPublicationRows(ctx context.Context, trunkID string, revision *int64) ([]Publication, error) {
var out []Publication
err := s.withRead(ctx, func(db *sql.DB) error {
query := "SELECT trunk_id,revision,cell_id,status,error_code,target_digest,local_revision,local_digest,applied_at,operation_id,updated_at FROM publications WHERE trunk_id=?"
args := []any{trunkID}
if revision != nil {
query += " AND revision=?"
args = append(args, *revision)
}
query += " ORDER BY revision DESC, cell_id ASC"
rows, err := db.QueryContext(ctx, query, args...)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var p Publication
var ec, applied sql.NullString
if err := rows.Scan(&p.TrunkID, &p.Revision, &p.CellID, &p.Status, &ec, &p.TargetDigest, &p.LocalRevision, &p.LocalDigest, &applied, &p.OperationID, &p.UpdatedAt); err != nil {
return err
}
if ec.Valid {
p.ErrorCode = &ec.String
}
if applied.Valid {
p.AppliedAt = &applied.String
}
out = append(out, p)
}
return rows.Err()
})
return out, err
}
func (s *Store) StagePublication(ctx context.Context, trunkID string, revision int64, cfg TrunkConfig, cellIDs []string, operationID, reason string) error {
digest, raw, err := hashJSON(cfg)
if err != nil {
return err
}
_ = raw
now := utcString(time.Now())
return s.withTx(ctx, func(tx *sql.Tx) error {
if _, err := tx.ExecContext(ctx, "INSERT INTO admission_barriers(trunk_id,operation_id,target_revision,reason,state,created_at,updated_at) VALUES(?,?,?,?,?,?,?) ON CONFLICT(trunk_id) DO UPDATE SET operation_id=excluded.operation_id,target_revision=excluded.target_revision,reason=excluded.reason,state='active',updated_at=excluded.updated_at", trunkID, operationID, revision, reason, "active", now, now); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, "UPDATE trunk_versions SET state='publishing' WHERE trunk_id=? AND revision=?", trunkID, revision); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, "DELETE FROM publications WHERE trunk_id=? AND revision=?", trunkID, revision); err != nil {
return err
}
for _, cellID := range cellIDs {
if _, err := tx.ExecContext(ctx, "INSERT INTO publications(trunk_id,revision,cell_id,status,error_code,target_digest,local_revision,local_digest,operation_id,updated_at) VALUES(?,?,?,?,?,?,?,?,?,?)", trunkID, revision, cellID, "pending", nil, digest, 0, "", operationID, now); err != nil {
return err
}
}
return nil
})
}
func (s *Store) UpdatePublication(ctx context.Context, trunkID string, revision int64, cellID, status, errorCode string, localRevision int64, localDigest, operationID string) error {
if status != "pending" && status != "applied" && status != "failed" {
return fmt.Errorf("invalid publication status")
}
now := utcString(time.Now())
var applied any
if status == "applied" {
applied = now
}
return s.withTx(ctx, func(tx *sql.Tx) error {
_, err := tx.ExecContext(ctx, "UPDATE publications SET status=?,error_code=?,local_revision=?,local_digest=?,applied_at=?,operation_id=?,updated_at=? WHERE trunk_id=? AND revision=? AND cell_id=?", status, nullableString(errorCode), localRevision, localDigest, applied, operationID, now, trunkID, revision, cellID)
return err
})
}
func nullableString(v string) any {
if v == "" {
return nil
}
return v
}
func (s *Store) FinalizePublication(ctx context.Context, trunkID string, revision int64, cfg TrunkConfig, operationID, actor, requestID, action string) (bool, error) {
var pubs []Publication
var allApplied bool
err := s.withRead(ctx, func(db *sql.DB) error {
rows, err := db.QueryContext(ctx, "SELECT trunk_id,revision,cell_id,status,error_code,target_digest,local_revision,local_digest,applied_at,operation_id,updated_at FROM publications WHERE trunk_id=? AND revision=?", trunkID, revision)
if err != nil {
return err
}
defer rows.Close()
allApplied = true
for rows.Next() {
var p Publication
var ec, applied sql.NullString
if err := rows.Scan(&p.TrunkID, &p.Revision, &p.CellID, &p.Status, &ec, &p.TargetDigest, &p.LocalRevision, &p.LocalDigest, &applied, &p.OperationID, &p.UpdatedAt); err != nil {
return err
}
if ec.Valid {
p.ErrorCode = &ec.String
}
pubs = append(pubs, p)
if p.Status != "applied" {
allApplied = false
}
}
if len(pubs) == 0 {
allApplied = false
}
return rows.Err()
})
if err != nil || !allApplied {
return false, err
}
now := utcString(time.Now())
return true, s.withTx(ctx, func(tx *sql.Tx) error {
activeStatus := "published"
if !cfg.Enabled {
activeStatus = "disabled"
}
if _, err := tx.ExecContext(ctx, "UPDATE trunks SET active_revision=?,active_status=?,updated_at=?,updated_by=? WHERE trunk_id=?", revision, activeStatus, now, actor, trunkID); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, "UPDATE trunk_versions SET state=CASE WHEN revision=? THEN 'published' ELSE state END WHERE trunk_id=? AND revision=?", revision, trunkID, revision); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, "UPDATE admission_barriers SET state='released',updated_at=? WHERE trunk_id=? AND operation_id=?", now, trunkID, operationID); err != nil {
return err
}
cellResults := make([]map[string]any, 0, len(pubs))
for _, pub := range pubs {
item := map[string]any{"cell_id": pub.CellID, "status": pub.Status, "local_revision": pub.LocalRevision}
if pub.ErrorCode != nil {
item["error_code"] = *pub.ErrorCode
}
cellResults = append(cellResults, item)
}
details, _ := json.Marshal(map[string]any{"action": action, "revision": revision, "cells": cellResults})
return insertAudit(ctx, tx, "trunk", trunkID, action, revision, actor, requestID, string(details), now)
})
}
func (s *Store) ActiveConfig(ctx context.Context, trunkID string) (TrunkConfig, int64, error) {
var rev int64
err := s.withRead(ctx, func(db *sql.DB) error {
return db.QueryRowContext(ctx, "SELECT active_revision FROM trunks WHERE trunk_id=?", trunkID).Scan(&rev)
})
if err != nil {
return TrunkConfig{}, 0, err
}
if rev == 0 {
return TrunkConfig{}, 0, sql.ErrNoRows
}
cfg, err := s.getTrunkConfig(ctx, trunkID, rev)
return cfg, rev, err
}
func (s *Store) MarkPublicationFailure(ctx context.Context, trunkID string, revision int64, operationID string, actor, requestID, action string, code string) error {
now := utcString(time.Now())
return s.withTx(ctx, func(tx *sql.Tx) error {
if _, err := tx.ExecContext(ctx, "UPDATE admission_barriers SET state='pending',updated_at=? WHERE trunk_id=? AND operation_id=?", now, trunkID, operationID); err != nil {
return err
}
details, _ := json.Marshal(map[string]any{"action": action, "revision": revision, "error_code": code, "safety": "admission remains blocked until reconciliation"})
return insertAudit(ctx, tx, "trunk", trunkID, action+"_failed", revision, actor, requestID, string(details), now)
})
}
func insertAudit(ctx context.Context, tx *sql.Tx, resourceType, resourceID, action string, revision int64, actor, requestID, details, now string) error {
_, err := tx.ExecContext(ctx, "INSERT INTO audit_entries(audit_id,resource_type,resource_id,action,revision,actor,request_id,details_json,created_at) VALUES(?,?,?,?,?,?,?,?,?)", newID("audit"), resourceType, resourceID, action, revision, actor, nullableString(requestID), details, now)
return err
}
func (s *Store) ListAudit(ctx context.Context, resourceID, requestID, actor string, limit int) ([]AuditEntry, error) {
if limit < 1 || limit > 200 {
limit = 50
}
var out []AuditEntry
err := s.withRead(ctx, func(db *sql.DB) error {
query := "SELECT audit_id,resource_type,resource_id,action,revision,actor,request_id,details_json,created_at FROM audit_entries WHERE 1=1"
args := []any{}
if resourceID != "" {
query += " AND resource_id=?"
args = append(args, resourceID)
}
if requestID != "" {
query += " AND request_id=?"
args = append(args, requestID)
}
if actor != "" {
query += " AND actor=?"
args = append(args, actor)
}
query += " ORDER BY created_at DESC,audit_id DESC LIMIT ?"
args = append(args, limit)
rows, err := db.QueryContext(ctx, query, args...)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var a AuditEntry
var rid sql.NullString
if err := rows.Scan(&a.AuditID, &a.ResourceType, &a.ResourceID, &a.Action, &a.Revision, &a.Actor, &rid, &a.DetailsJSON, &a.CreatedAt); err != nil {
return err
}
if rid.Valid {
a.RequestID = &rid.String
}
out = append(out, a)
}
return rows.Err()
})
return out, err
}
func (s *Store) FindOperation(ctx context.Context, actor, action, resourceType, resourceID, requestID string) (Operation, error) {
var op Operation
err := s.withRead(ctx, func(db *sql.DB) error {
var response, failure sql.NullString
err := db.QueryRowContext(ctx, "SELECT operation_id,actor,action,resource_type,resource_id,request_id,request_digest,expected_revision,state,http_status,response_json,error_json,created_at,updated_at,expires_at FROM operations WHERE actor=? AND action=? AND resource_type=? AND resource_id=? AND request_id=?", actor, action, resourceType, resourceID, requestID).Scan(&op.OperationID, &op.Actor, &op.Action, &op.ResourceType, &op.ResourceID, &op.RequestID, &op.RequestDigest, &op.ExpectedRevision, &op.State, &op.HTTPStatus, &response, &failure, &op.CreatedAt, &op.UpdatedAt, &op.ExpiresAt)
if response.Valid {
op.ResponseJSON = response.String
}
if failure.Valid {
op.ErrorJSON = failure.String
}
return err
})
if err == nil {
err = s.withRead(ctx, func(db *sql.DB) error {
var response, failure sql.NullString
if err := db.QueryRowContext(ctx, "SELECT response_json,error_json FROM operations WHERE operation_id=?", op.OperationID).Scan(&response, &failure); err != nil {
return err
}
if response.Valid {
op.ResponseJSON = response.String
}
if failure.Valid {
op.ErrorJSON = failure.String
}
return nil
})
}
return op, err
}
func (s *Store) CreateOperation(ctx context.Context, actor, action, resourceType, resourceID, requestID, digest string, expected int64) (Operation, bool, error) {
op := Operation{OperationID: newID("op"), Actor: actor, Action: action, ResourceType: resourceType, ResourceID: resourceID, RequestID: requestID, RequestDigest: digest, ExpectedRevision: expected, State: "in_progress", HTTPStatus: 200, CreatedAt: utcString(time.Now()), UpdatedAt: utcString(time.Now()), ExpiresAt: utcString(time.Now().Add(30 * 24 * time.Hour))}
var existing Operation
err := s.withTx(ctx, func(tx *sql.Tx) error {
var response, failure sql.NullString
err := tx.QueryRowContext(ctx, "SELECT operation_id,actor,action,resource_type,resource_id,request_id,request_digest,expected_revision,state,http_status,response_json,error_json,created_at,updated_at,expires_at FROM operations WHERE actor=? AND action=? AND resource_type=? AND resource_id=? AND request_id=?", actor, action, resourceType, resourceID, requestID).Scan(&existing.OperationID, &existing.Actor, &existing.Action, &existing.ResourceType, &existing.ResourceID, &existing.RequestID, &existing.RequestDigest, &existing.ExpectedRevision, &existing.State, &existing.HTTPStatus, &response, &failure, &existing.CreatedAt, &existing.UpdatedAt, &existing.ExpiresAt)
if err == nil {
if expiry, parseErr := parseStoredTime(existing.ExpiresAt); parseErr == nil && time.Now().After(expiry) {
return newAppError(409, "IDEMPOTENCY_EXPIRED", "the durable idempotency record has expired", map[string]any{"operation_id": existing.OperationID})
}
if response.Valid {
existing.ResponseJSON = response.String
}
if failure.Valid {
existing.ErrorJSON = failure.String
}
return nil
}
if !errors.Is(err, sql.ErrNoRows) {
return err
}
_, err = tx.ExecContext(ctx, "INSERT INTO operations(operation_id,actor,action,resource_type,resource_id,request_id,request_digest,expected_revision,state,http_status,created_at,updated_at,expires_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)", op.OperationID, op.Actor, op.Action, op.ResourceType, op.ResourceID, op.RequestID, op.RequestDigest, op.ExpectedRevision, op.State, op.HTTPStatus, op.CreatedAt, op.UpdatedAt, op.ExpiresAt)
return err
})
if err != nil {
return Operation{}, false, err
}
if existing.OperationID != "" {
return existing, true, nil
}
return op, false, nil
}
func (s *Store) FinishOperation(ctx context.Context, id, state string, status int, responseJSON, errorJSON string) error {
return s.withTx(ctx, func(tx *sql.Tx) error {
_, err := tx.ExecContext(ctx, "UPDATE operations SET state=?,http_status=?,response_json=?,error_json=?,updated_at=? WHERE operation_id=?", state, status, nullableString(responseJSON), nullableString(errorJSON), utcString(time.Now()), id)
return err
})
}
func (s *Store) GetOperation(ctx context.Context, id string) (Operation, error) {
var op Operation
err := s.withRead(ctx, func(db *sql.DB) error {
var r, e sql.NullString
err := db.QueryRowContext(ctx, "SELECT operation_id,actor,action,resource_type,resource_id,request_id,request_digest,expected_revision,state,http_status,response_json,error_json,created_at,updated_at,expires_at FROM operations WHERE operation_id=?", id).Scan(&op.OperationID, &op.Actor, &op.Action, &op.ResourceType, &op.ResourceID, &op.RequestID, &op.RequestDigest, &op.ExpectedRevision, &op.State, &op.HTTPStatus, &r, &e, &op.CreatedAt, &op.UpdatedAt, &op.ExpiresAt)
if r.Valid {
op.ResponseJSON = r.String
}
if e.Valid {
op.ErrorJSON = e.String
}
return err
})
return op, err
}
func (s *Store) AdvanceCellBoot(ctx context.Context, cellID, bootID string) error {
if !validID(cellID) || bootID == "" || len(bootID) > 128 || strings.ContainsAny(bootID, "\r\n") {
return newAppError(422, "INVALID_CELL_BOOT", "cell_id or boot_id is invalid", nil)
}
return s.withTx(ctx, func(tx *sql.Tx) error {
result, err := tx.ExecContext(ctx, "UPDATE cells SET boot_id=?,updated_at=? WHERE cell_id=?", bootID, utcString(time.Now()), cellID)
if err != nil {
return err
}
if affected, _ := result.RowsAffected(); affected == 0 {
return newAppError(404, "CELL_NOT_FOUND", "cell does not exist", nil)
}
return nil
})
}
func (s *Store) InsertObservation(ctx context.Context, input ObservationInput) error {
if !validID(input.ObservationID) {
input.ObservationID = newID("obs")
}
if !validID(input.CellID) || input.BootID == "" || len(input.BootID) > 128 || strings.ContainsAny(input.BootID, "\r\n") || input.Sequence < 1 {
return newAppError(422, "INVALID_OBSERVATION", "cell_id, boot_id, or sequence is invalid", nil)
}
if input.ObservedAt.IsZero() {
input.ObservedAt = time.Now()
}
if input.ReceivedAt.IsZero() {
input.ReceivedAt = time.Now()
}
states, _ := json.Marshal(input.States)
occupancy, _ := json.Marshal(input.Occupancy)
if string(states) == "null" {
states = []byte("{}")
}
if string(occupancy) == "null" {
occupancy = []byte("{}")
}
return s.withTx(ctx, func(tx *sql.Tx) error {
var currentBoot string
if err := tx.QueryRowContext(ctx, "SELECT boot_id FROM cells WHERE cell_id=?", input.CellID).Scan(&currentBoot); err != nil {
if errors.Is(err, sql.ErrNoRows) {
return newAppError(404, "CELL_NOT_FOUND", "cell does not exist", nil)
}
return err
}
if currentBoot != "" && currentBoot != input.BootID {
return newAppError(409, "STALE_BOOT_ID", "observation belongs to a retired Cell boot", map[string]any{"current_boot_id": currentBoot})
}
if currentBoot == "" {
if _, err := tx.ExecContext(ctx, "UPDATE cells SET boot_id=? WHERE cell_id=?", input.BootID, input.CellID); err != nil {
return err
}
}
var last int64
err := tx.QueryRowContext(ctx, "SELECT COALESCE(MAX(sequence),-1) FROM observations WHERE cell_id=? AND boot_id=?", input.CellID, input.BootID).Scan(&last)
if err != nil {
return err
}
if input.Sequence <= last {
return newAppError(409, "STALE_OBSERVATION", "observation sequence is not newer", map[string]any{"last_sequence": last})
}
_, err = tx.ExecContext(ctx, "INSERT INTO observations(observation_id,cell_id,trunk_id,boot_id,sequence,observed_at,received_at,source,config_revision,states_json,occupancy_json,sample_id) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)", input.ObservationID, input.CellID, nullableString(input.TrunkID), input.BootID, input.Sequence, utcString(input.ObservedAt), utcString(input.ReceivedAt), input.Source, input.ConfigRevision, string(states), string(occupancy), input.SampleID)
return err
})
}
func (s *Store) LatestObservations(ctx context.Context, cellID, trunkID string) ([]ObservationInput, error) {
var out []ObservationInput
err := s.withRead(ctx, func(db *sql.DB) error {
query := "SELECT o.observation_id,o.cell_id,COALESCE(o.trunk_id,''),o.boot_id,o.sequence,o.observed_at,o.received_at,o.source,o.config_revision,o.states_json,o.occupancy_json,o.sample_id FROM observations o JOIN cells c ON c.cell_id=o.cell_id WHERE c.boot_id=o.boot_id"
args := []any{}
if cellID != "" {
query += " AND o.cell_id=?"
args = append(args, cellID)
}
if trunkID != "" {
query += " AND o.trunk_id=?"
args = append(args, trunkID)
}
query += " AND o.sequence=(SELECT MAX(o2.sequence) FROM observations o2 WHERE o2.cell_id=o.cell_id AND o2.boot_id=o.boot_id AND COALESCE(o2.trunk_id,'')=COALESCE(o.trunk_id,'') ) ORDER BY o.cell_id,o.trunk_id"
rows, err := db.QueryContext(ctx, query, args...)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var o ObservationInput
var obs, rec string
var states, occ string
if err := rows.Scan(&o.ObservationID, &o.CellID, &o.TrunkID, &o.BootID, &o.Sequence, &obs, &rec, &o.Source, &o.ConfigRevision, &states, &occ, &o.SampleID); err != nil {
return err
}
o.ObservedAt, _ = parseStoredTime(obs)
o.ReceivedAt, _ = parseStoredTime(rec)
_ = json.Unmarshal([]byte(states), &o.States)
_ = json.Unmarshal([]byte(occ), &o.Occupancy)
out = append(out, o)
}
return rows.Err()
})
return out, err
}
func (s *Store) UpsertAttempt(ctx context.Context, a Attempt) error {
if !validID(a.AttemptID) || (a.Mode != ModeMock && a.Mode != ModeReal) {
return newAppError(422, "INVALID_ATTEMPT", "attempt_id or mode is invalid", nil)
}
if a.OriginStatus != "confirmed" && a.OriginStatus != "uncertain" && a.OriginStatus != "not_started" && a.OriginStatus != "rejected" {
return newAppError(422, "INVALID_ORIGIN_STATUS", "origin_status is invalid", nil)
}
if a.AttemptStartedAt != nil && a.AnsweredAt != nil && a.AnsweredAt.Before(*a.AttemptStartedAt) {
return newAppError(422, "INVALID_ATTEMPT_FACT", "answered_at precedes attempt_started_at", nil)
}
if a.AttemptStartedAt != nil && a.EndedAt != nil && a.EndedAt.Before(*a.AttemptStartedAt) {
return newAppError(422, "INVALID_ATTEMPT_FACT", "ended_at precedes attempt_started_at", nil)
}
if a.AnsweredAt != nil && a.EndedAt != nil && a.EndedAt.Before(*a.AnsweredAt) {
return newAppError(422, "INVALID_ATTEMPT_FACT", "ended_at precedes answered_at", nil)
}
if a.UpdatedAt.IsZero() {
a.UpdatedAt = time.Now()
}
return s.withTx(ctx, func(tx *sql.Tx) error {
_, err := tx.ExecContext(ctx, "INSERT INTO attempts(attempt_id,execution_id,call_id,tenant_key,provider_id,trunk_id,cell_id,egress_pool_id,config_revision,mode,dialed_number,attempt_started_at,origin_status,answered_at,ended_at,termination_reason,sip_code,q850,ai_status,source,fact_version,updated_at) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(attempt_id) DO UPDATE SET execution_id=excluded.execution_id,call_id=excluded.call_id,tenant_key=excluded.tenant_key,provider_id=excluded.provider_id,trunk_id=excluded.trunk_id,cell_id=excluded.cell_id,egress_pool_id=excluded.egress_pool_id,config_revision=excluded.config_revision,mode=excluded.mode,dialed_number=excluded.dialed_number,attempt_started_at=excluded.attempt_started_at,origin_status=excluded.origin_status,answered_at=excluded.answered_at,ended_at=excluded.ended_at,termination_reason=excluded.termination_reason,sip_code=excluded.sip_code,q850=excluded.q850,ai_status=excluded.ai_status,source=excluded.source,fact_version=MAX(attempts.fact_version,excluded.fact_version),updated_at=excluded.updated_at WHERE attempts.fact_version<=excluded.fact_version", a.AttemptID, a.ExecutionID, a.CallID, a.TenantKey, a.ProviderID, a.TrunkID, a.CellID, a.EgressPoolID, a.ConfigRevision, a.Mode, a.DialedNumber, nullableTime(a.AttemptStartedAt), a.OriginStatus, nullableTime(a.AnsweredAt), nullableTime(a.EndedAt), a.TerminationReason, a.SIPCode, a.Q850, a.AIStatus, a.Source, a.FactVersion, utcString(a.UpdatedAt))
return err
})
}
func nullableTime(v *time.Time) any {
if v == nil {
return nil
}
return utcString(*v)
}
func (s *Store) RecordAttemptEvent(ctx context.Context, eventID, attemptID, eventType string, eventTime, receivedAt time.Time, source string, version int64, payload map[string]any) error {
if eventID == "" {
eventID = newID("event")
}
if !validID(eventID) || !validID(attemptID) || version < 1 || source == "" {
return newAppError(422, "INVALID_ATTEMPT_EVENT", "event identity, source, or fact_version is invalid", nil)
}
if eventType != "originated" && eventType != "answered" && eventType != "ended" && eventType != "ai_failed" {
return newAppError(422, "INVALID_ATTEMPT_EVENT", "event_type is invalid", nil)
}
if receivedAt.IsZero() {
receivedAt = time.Now()
}
if eventTime.IsZero() {
eventTime = receivedAt
}
b, err := json.Marshal(payload)
if err != nil {
return newAppError(422, "INVALID_ATTEMPT_EVENT", "event payload is not JSON serializable", nil)
}
return s.withTx(ctx, func(tx *sql.Tx) error {
res, err := tx.ExecContext(ctx, "INSERT OR IGNORE INTO attempt_events(event_id,attempt_id,event_type,event_time,received_at,source,fact_version,payload_json) VALUES(?,?,?,?,?,?,?,?)", eventID, attemptID, eventType, utcString(eventTime), utcString(receivedAt), source, version, string(b))
if err != nil {
return err
}
n, _ := res.RowsAffected()
if n == 0 {
return nil
}
where := "WHERE attempt_id=? AND fact_version<=?"
switch eventType {
case "originated":
_, err = tx.ExecContext(ctx, "UPDATE attempts SET attempt_started_at=?,origin_status='confirmed',fact_version=MAX(fact_version,?),updated_at=? "+where+" AND origin_status<>'rejected'", utcString(eventTime), version, utcString(receivedAt), attemptID, version)
case "answered":
_, err = tx.ExecContext(ctx, "UPDATE attempts SET answered_at=?,fact_version=MAX(fact_version,?),updated_at=? "+where, utcString(eventTime), version, utcString(receivedAt), attemptID, version)
case "ended":
reason, _ := payload["termination_reason"].(string)
if reason == "" {
reason = "unknown"
}
_, err = tx.ExecContext(ctx, "UPDATE attempts SET ended_at=?,termination_reason=?,fact_version=MAX(fact_version,?),updated_at=? "+where, utcString(eventTime), reason, version, utcString(receivedAt), attemptID, version)
case "ai_failed":
_, err = tx.ExecContext(ctx, "UPDATE attempts SET ai_status='failed',fact_version=MAX(fact_version,?),updated_at=? "+where, version, utcString(receivedAt), attemptID, version)
}
return err
})
}
func (s *Store) ListAttempts(ctx context.Context, filter AttemptFilter) ([]Attempt, error) {
var out []Attempt
err := s.withRead(ctx, func(db *sql.DB) error {
query := "SELECT attempt_id,execution_id,call_id,tenant_key,provider_id,trunk_id,cell_id,egress_pool_id,config_revision,mode,dialed_number,attempt_started_at,origin_status,answered_at,ended_at,termination_reason,sip_code,q850,ai_status,source,fact_version,updated_at FROM attempts WHERE 1=1"
args := []any{}
if filter.Timeline && filter.From != nil && filter.To != nil {
query += " AND answered_at IS NOT NULL AND answered_at<? AND (ended_at IS NULL OR ended_at>?)"
args = append(args, utcString(*filter.To), utcString(*filter.From))
} else if filter.From != nil && filter.To != nil {
query += " AND ((attempt_started_at>=? AND attempt_started_at<?) OR (attempt_started_at IS NULL AND origin_status IN ('rejected','not_started') AND updated_at>=? AND updated_at<?))"
args = append(args, utcString(*filter.From), utcString(*filter.To), utcString(*filter.From), utcString(*filter.To))
} else if filter.From != nil {
query += " AND (attempt_started_at>=? OR (attempt_started_at IS NULL AND updated_at>=?))"
args = append(args, utcString(*filter.From), utcString(*filter.From))
} else if filter.To != nil {
query += " AND (attempt_started_at<? OR (attempt_started_at IS NULL AND updated_at<?))"
args = append(args, utcString(*filter.To), utcString(*filter.To))
}
if filter.ProviderID != "" {
query += " AND provider_id=?"
args = append(args, filter.ProviderID)
}
if filter.TrunkID != "" {
query += " AND trunk_id=?"
args = append(args, filter.TrunkID)
}
if filter.CellID != "" {
query += " AND cell_id=?"
args = append(args, filter.CellID)
}
if filter.EgressPoolID != "" {
query += " AND egress_pool_id=?"
args = append(args, filter.EgressPoolID)
}
if filter.Mode != "" {
query += " AND mode=?"
args = append(args, filter.Mode)
}
query += " ORDER BY CASE WHEN attempt_started_at IS NULL THEN updated_at ELSE attempt_started_at END ASC,attempt_id ASC"
if filter.Limit > 0 {
query += " LIMIT ?"
args = append(args, filter.Limit)
}
rows, err := db.QueryContext(ctx, query, args...)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
a, err := scanAttempt(rows)
if err != nil {
return err
}
out = append(out, a)
}
return rows.Err()
})
return out, err
}
type AttemptFilter struct {
From, To *time.Time
ProviderID, TrunkID, CellID, EgressPoolID string
Mode string
Limit int
Timeline bool
}
func scanAttempt(scanner interface{ Scan(...any) error }) (Attempt, error) {
var a Attempt
var started, answered, ended, updated sql.NullString
var sip sql.NullInt64
err := scanner.Scan(&a.AttemptID, &a.ExecutionID, &a.CallID, &a.TenantKey, &a.ProviderID, &a.TrunkID, &a.CellID, &a.EgressPoolID, &a.ConfigRevision, &a.Mode, &a.DialedNumber, &started, &a.OriginStatus, &answered, &ended, &a.TerminationReason, &sip, &a.Q850, &a.AIStatus, &a.Source, &a.FactVersion, &updated)
if err != nil {
return a, err
}
if started.Valid {
t, _ := parseStoredTime(started.String)
a.AttemptStartedAt = &t
}
if answered.Valid {
t, _ := parseStoredTime(answered.String)
a.AnsweredAt = &t
}
if ended.Valid {
t, _ := parseStoredTime(ended.String)
a.EndedAt = &t
}
if sip.Valid {
x := int(sip.Int64)
a.SIPCode = &x
}
a.UpdatedAt, _ = parseStoredTime(updated.String)
return a, nil
}
func (s *Store) GetAttempt(ctx context.Context, id string) (Attempt, error) {
var out Attempt
err := s.withRead(ctx, func(db *sql.DB) error {
row := db.QueryRowContext(ctx, "SELECT attempt_id,execution_id,call_id,tenant_key,provider_id,trunk_id,cell_id,egress_pool_id,config_revision,mode,dialed_number,attempt_started_at,origin_status,answered_at,ended_at,termination_reason,sip_code,q850,ai_status,source,fact_version,updated_at FROM attempts WHERE attempt_id=?", id)
var e error
out, e = scanAttempt(row)
return e
})
return out, err
}