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, ¬Null, &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(¤tRevision) 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(¤tRevision) 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(¤tRevision, ¤tProvider) 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(¤tBoot); 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?)" 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=? AND updated_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 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 }