CI / changes (push) Successful in 8s
CI / commitlint (push) Has been skipped
CI / openapi (push) Successful in 25s
CI / web (push) Successful in 32s
CI / go (push) Successful in 57s
CI / bird2 (push) Successful in 15s
CI / release (push) Successful in 3m18s
Implemented new endpoints for estimating and pruning revisions, including detailed schemas for requests and responses. The `RevisionPruneEstimate` and `RevisionPruneResult` components were added to the OpenAPI documentation, enhancing the API's functionality for managing revision retention. Updated the backend to support these operations and integrated them into the tenant settings UI for improved user interaction.
1354 lines
37 KiB
Go
1354 lines
37 KiB
Go
// Package repository implements SQL-backed store.Backend (PostgreSQL).
|
|
package repository
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"runtime"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"evobgp/internal/store"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
)
|
|
|
|
// #region agent log
|
|
func agentDebugNDJSON3214(hypothesisID, location, message string, data map[string]any) {
|
|
if os.Getenv("EVOBGP_DEBUG_LOG") != "1" {
|
|
return
|
|
}
|
|
f, err := os.OpenFile("debug-3214dc.log", os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
|
|
if err != nil {
|
|
return
|
|
}
|
|
defer func() { _ = f.Close() }()
|
|
var ms runtime.MemStats
|
|
runtime.ReadMemStats(&ms)
|
|
payload := map[string]any{
|
|
"sessionId": "3214dc",
|
|
"hypothesisId": hypothesisID,
|
|
"location": location,
|
|
"message": message,
|
|
"data": data,
|
|
"timestamp": time.Now().UnixMilli(),
|
|
"allocBytes": ms.Alloc,
|
|
}
|
|
b, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return
|
|
}
|
|
_, _ = f.Write(append(b, '\n'))
|
|
}
|
|
|
|
// #endregion
|
|
|
|
// Postgres implements store.Backend using pgxpool.
|
|
type Postgres struct {
|
|
pool *pgxpool.Pool
|
|
// demo IDs after seed
|
|
demoTenant, demoCDN, demoIP, demoRev, demoSpk string
|
|
}
|
|
|
|
// NewPostgres opens migrations-applied pool is assumed; seedDemo inserts demo tenant graph.
|
|
func NewPostgres(ctx context.Context, pool *pgxpool.Pool, seedDemo bool) (*Postgres, error) {
|
|
p := &Postgres{pool: pool}
|
|
if seedDemo {
|
|
if err := p.seedDemo(ctx); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
return p, nil
|
|
}
|
|
|
|
func (p *Postgres) DemoIDs() (tenant, moduleCDN, moduleIP, revision, speaker string) {
|
|
return p.demoTenant, p.demoCDN, p.demoIP, p.demoRev, p.demoSpk
|
|
}
|
|
|
|
// Ping checks PostgreSQL connectivity.
|
|
func (p *Postgres) Ping(ctx context.Context) error {
|
|
return p.pool.Ping(ctx)
|
|
}
|
|
|
|
func (p *Postgres) MaterializedPrefixStats() (max int, sum int) {
|
|
ctx := context.Background()
|
|
// Агрегация в БД — не тащим все строки config_revision в память.
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT
|
|
COALESCE(MAX((meta_json->>'materialized_prefix_count')::int), 0),
|
|
COALESCE(SUM((meta_json->>'materialized_prefix_count')::int), 0)
|
|
FROM config_revision`).Scan(&max, &sum)
|
|
if err != nil {
|
|
return 0, 0
|
|
}
|
|
return max, sum
|
|
}
|
|
|
|
func (p *Postgres) PeerCount() int {
|
|
ctx := context.Background()
|
|
var n int
|
|
_ = p.pool.QueryRow(ctx, `SELECT COUNT(*) FROM bgp_peer`).Scan(&n)
|
|
return n
|
|
}
|
|
|
|
func (p *Postgres) PeerSessionCountsByState() map[string]int {
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `SELECT COALESCE(meta_json->>'session_state','unknown'), COUNT(*) FROM bgp_peer GROUP BY 1`)
|
|
if err != nil {
|
|
return map[string]int{}
|
|
}
|
|
defer rows.Close()
|
|
out := make(map[string]int)
|
|
for rows.Next() {
|
|
var st string
|
|
var c int
|
|
if rows.Scan(&st, &c) == nil {
|
|
out[st] = c
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (p *Postgres) ListModules(tenantID string) []*store.Module {
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `
|
|
SELECT id, type, name, enabled, priority, doh_profile_id::text, doh_resolver_policy,
|
|
refresh_interval_sec, cron_expr, default_community_id::text, last_refreshed_at
|
|
FROM module WHERE tenant_id = $1 AND deleted_at IS NULL ORDER BY priority, name`, tenantID)
|
|
if err != nil {
|
|
return nil
|
|
}
|
|
defer rows.Close()
|
|
var out []*store.Module
|
|
moduleByID := make(map[string]*store.Module)
|
|
for rows.Next() {
|
|
var m store.Module
|
|
m.TenantID = tenantID
|
|
var doh, dc, cron *string
|
|
var refresh *int32
|
|
var last *time.Time
|
|
if err := rows.Scan(&m.ID, &m.Type, &m.Name, &m.Enabled, &m.Priority, &doh, &m.DohResolverPolicy, &refresh, &cron, &dc, &last); err != nil {
|
|
continue
|
|
}
|
|
m.DohResolverPolicy = store.NormalizeDohResolverPolicy(m.DohResolverPolicy)
|
|
if refresh != nil {
|
|
m.RefreshIntervalSec = int(*refresh)
|
|
}
|
|
if cron != nil {
|
|
m.CronExpr = *cron
|
|
}
|
|
if doh != nil && *doh != "" {
|
|
m.DohProfileID = doh
|
|
}
|
|
if dc != nil && *dc != "" {
|
|
m.DefaultCommunityID = dc
|
|
}
|
|
if last != nil {
|
|
t := last.UTC()
|
|
m.LastRefreshedAt = &t
|
|
}
|
|
out = append(out, &m)
|
|
moduleByID[m.ID] = &m
|
|
}
|
|
if err := p.batchFillModuleDohFields(ctx, moduleByID); err != nil {
|
|
return nil
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (p *Postgres) ListModulesPage(tenantID, cursor string, limit int) ([]*store.Module, string, bool) {
|
|
if limit <= 0 {
|
|
limit = 50
|
|
}
|
|
off := 0
|
|
if cursor != "" {
|
|
if n, err := strconv.Atoi(cursor); err == nil && n >= 0 {
|
|
off = n
|
|
}
|
|
}
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `
|
|
SELECT id, type, name, enabled, priority, doh_profile_id::text, doh_resolver_policy,
|
|
refresh_interval_sec, cron_expr, default_community_id::text, last_refreshed_at
|
|
FROM module WHERE tenant_id = $1 AND deleted_at IS NULL
|
|
ORDER BY priority, name
|
|
LIMIT $2 OFFSET $3`, tenantID, limit+1, off)
|
|
if err != nil {
|
|
return nil, "", false
|
|
}
|
|
defer rows.Close()
|
|
var out []*store.Module
|
|
moduleByID := make(map[string]*store.Module)
|
|
for rows.Next() {
|
|
var m store.Module
|
|
m.TenantID = tenantID
|
|
var doh, dc, cron *string
|
|
var refresh *int32
|
|
var last *time.Time
|
|
if err := rows.Scan(&m.ID, &m.Type, &m.Name, &m.Enabled, &m.Priority, &doh, &m.DohResolverPolicy, &refresh, &cron, &dc, &last); err != nil {
|
|
continue
|
|
}
|
|
m.DohResolverPolicy = store.NormalizeDohResolverPolicy(m.DohResolverPolicy)
|
|
if refresh != nil {
|
|
m.RefreshIntervalSec = int(*refresh)
|
|
}
|
|
if cron != nil {
|
|
m.CronExpr = *cron
|
|
}
|
|
if doh != nil && *doh != "" {
|
|
m.DohProfileID = doh
|
|
}
|
|
if dc != nil && *dc != "" {
|
|
m.DefaultCommunityID = dc
|
|
}
|
|
if last != nil {
|
|
t := last.UTC()
|
|
m.LastRefreshedAt = &t
|
|
}
|
|
out = append(out, &m)
|
|
moduleByID[m.ID] = &m
|
|
}
|
|
if err := p.batchFillModuleDohFields(ctx, moduleByID); err != nil {
|
|
return nil, "", false
|
|
}
|
|
more := len(out) > limit
|
|
if more {
|
|
out = out[:limit]
|
|
}
|
|
next := ""
|
|
if more {
|
|
next = fmt.Sprintf("%d", off+limit)
|
|
}
|
|
if len(out) == 0 {
|
|
return nil, "", false
|
|
}
|
|
return out, next, more
|
|
}
|
|
|
|
func (p *Postgres) GetModule(tenantID, moduleID string) (*store.Module, error) {
|
|
ctx := context.Background()
|
|
var m store.Module
|
|
m.TenantID = tenantID
|
|
var doh, dc, cron *string
|
|
var refresh *int32
|
|
var last *time.Time
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT id, type, name, enabled, priority, doh_profile_id::text, doh_resolver_policy,
|
|
refresh_interval_sec, cron_expr, default_community_id::text, last_refreshed_at
|
|
FROM module WHERE id = $1 AND tenant_id = $2 AND deleted_at IS NULL`, moduleID, tenantID).Scan(
|
|
&m.ID, &m.Type, &m.Name, &m.Enabled, &m.Priority, &doh, &m.DohResolverPolicy, &refresh, &cron, &dc, &last)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, store.ErrNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
if refresh != nil {
|
|
m.RefreshIntervalSec = int(*refresh)
|
|
}
|
|
if cron != nil {
|
|
m.CronExpr = *cron
|
|
}
|
|
if doh != nil && *doh != "" {
|
|
m.DohProfileID = doh
|
|
}
|
|
if dc != nil && *dc != "" {
|
|
m.DefaultCommunityID = dc
|
|
}
|
|
if last != nil {
|
|
t := last.UTC()
|
|
m.LastRefreshedAt = &t
|
|
}
|
|
m.DohResolverPolicy = store.NormalizeDohResolverPolicy(m.DohResolverPolicy)
|
|
if err := p.fillModuleDohFields(ctx, &m); err != nil {
|
|
return nil, err
|
|
}
|
|
return &m, nil
|
|
}
|
|
|
|
func (p *Postgres) CreateModule(tenantID string, in *store.Module) (*store.Module, error) {
|
|
if in == nil {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
store.NormalizeModuleDoh(in)
|
|
ctx := context.Background()
|
|
id := uuid.NewString()
|
|
var doh, dc any
|
|
if in.DohProfileID != nil && strings.TrimSpace(*in.DohProfileID) != "" {
|
|
doh = strings.TrimSpace(*in.DohProfileID)
|
|
}
|
|
if in.DefaultCommunityID != nil && strings.TrimSpace(*in.DefaultCommunityID) != "" {
|
|
dc = strings.TrimSpace(*in.DefaultCommunityID)
|
|
}
|
|
var ri any
|
|
if in.RefreshIntervalSec != 0 {
|
|
ri = in.RefreshIntervalSec
|
|
}
|
|
var cronArg any
|
|
if strings.TrimSpace(in.CronExpr) != "" {
|
|
cronArg = strings.TrimSpace(in.CronExpr)
|
|
}
|
|
var lastArg any
|
|
if in.LastRefreshedAt != nil {
|
|
lastArg = in.LastRefreshedAt.UTC()
|
|
}
|
|
policy := store.NormalizeDohResolverPolicy(in.DohResolverPolicy)
|
|
_, err := p.pool.Exec(ctx, `
|
|
INSERT INTO module (id, tenant_id, type, name, enabled, priority, doh_profile_id, doh_resolver_policy, refresh_interval_sec, cron_expr, default_community_id, last_refreshed_at)
|
|
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12)`,
|
|
id, tenantID, in.Type, in.Name, in.Enabled, in.Priority, doh, policy, ri, cronArg, dc, lastArg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := p.setModuleDohProfiles(ctx, id, in.DohProfileIDs); err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetModule(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) UpdateModule(tenantID, moduleID string, patch *store.ModulePatch) (*store.Module, error) {
|
|
if patch == nil {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
ctx := context.Background()
|
|
base, err := p.GetModule(tenantID, moduleID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
work := *base
|
|
if patch.Name != nil {
|
|
work.Name = strings.TrimSpace(*patch.Name)
|
|
}
|
|
if patch.Enabled != nil {
|
|
work.Enabled = *patch.Enabled
|
|
}
|
|
if patch.Priority != nil {
|
|
work.Priority = *patch.Priority
|
|
}
|
|
if patch.RefreshIntervalSec != nil {
|
|
work.RefreshIntervalSec = *patch.RefreshIntervalSec
|
|
}
|
|
if patch.CronExpr != nil {
|
|
work.CronExpr = *patch.CronExpr
|
|
}
|
|
if patch.DefaultCommunityID != nil {
|
|
v := strings.TrimSpace(*patch.DefaultCommunityID)
|
|
if v == "" {
|
|
work.DefaultCommunityID = nil
|
|
} else {
|
|
work.DefaultCommunityID = &v
|
|
}
|
|
}
|
|
store.ApplyModuleDohPatch(&work, patch)
|
|
if patch.LastRefreshedAt != nil {
|
|
t := patch.LastRefreshedAt.UTC()
|
|
work.LastRefreshedAt = &t
|
|
}
|
|
var dcArg, dohArg any
|
|
if work.DefaultCommunityID != nil {
|
|
dcArg = *work.DefaultCommunityID
|
|
}
|
|
if work.DohProfileID != nil {
|
|
dohArg = *work.DohProfileID
|
|
}
|
|
var riArg any
|
|
if work.RefreshIntervalSec != 0 {
|
|
riArg = work.RefreshIntervalSec
|
|
}
|
|
var cronArg any
|
|
if strings.TrimSpace(work.CronExpr) != "" {
|
|
cronArg = strings.TrimSpace(work.CronExpr)
|
|
}
|
|
var lastArg any
|
|
if work.LastRefreshedAt != nil {
|
|
lastArg = work.LastRefreshedAt.UTC()
|
|
}
|
|
policy := store.NormalizeDohResolverPolicy(work.DohResolverPolicy)
|
|
_, err = p.pool.Exec(ctx, `
|
|
UPDATE module SET name=$3, enabled=$4, priority=$5, refresh_interval_sec=$6, cron_expr=$7,
|
|
default_community_id=$8, doh_profile_id=$9, doh_resolver_policy=$10, last_refreshed_at=$11, updated_at=now()
|
|
WHERE id=$1 AND tenant_id=$2 AND deleted_at IS NULL`,
|
|
moduleID, tenantID, work.Name, work.Enabled, work.Priority, riArg, cronArg, dcArg, dohArg, policy, lastArg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if patch.DohProfileIDs != nil || patch.DohProfileID != nil {
|
|
if err := p.setModuleDohProfiles(ctx, moduleID, work.DohProfileIDs); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
return p.GetModule(tenantID, moduleID)
|
|
}
|
|
|
|
func (p *Postgres) SoftDeleteModule(tenantID, moduleID string) error {
|
|
ctx := context.Background()
|
|
tag, err := p.pool.Exec(ctx, `UPDATE module SET deleted_at=now(), updated_at=now() WHERE id=$1 AND tenant_id=$2 AND deleted_at IS NULL`, moduleID, tenantID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return store.ErrNotFound
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *Postgres) ListPeers(tenantID string) []*store.BGPPeer {
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `
|
|
SELECT id::text, tenant_id::text, bgp_speaker_id::text, neighbor::text, remote_asn, enabled,
|
|
COALESCE(meta_json->>'name',''), COALESCE(meta_json->>'session_state',''), COALESCE(policies_json::text,'{}')
|
|
FROM bgp_peer WHERE tenant_id=$1 ORDER BY neighbor`, tenantID)
|
|
if err != nil {
|
|
return nil
|
|
}
|
|
defer rows.Close()
|
|
var out []*store.BGPPeer
|
|
for rows.Next() {
|
|
var peer store.BGPPeer
|
|
var sp *string
|
|
if err := rows.Scan(&peer.ID, &peer.TenantID, &sp, &peer.Neighbor, &peer.RemoteASN, &peer.Enabled, &peer.Name, &peer.SessionState, &peer.PoliciesJSON); err != nil {
|
|
continue
|
|
}
|
|
peer.SpeakerID = sp
|
|
out = append(out, &peer)
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (p *Postgres) GetPeer(tenantID, id string) (*store.BGPPeer, error) {
|
|
ctx := context.Background()
|
|
var peer store.BGPPeer
|
|
var sp *string
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT id::text, tenant_id::text, bgp_speaker_id::text, neighbor::text, remote_asn, enabled,
|
|
COALESCE(meta_json->>'name',''), COALESCE(meta_json->>'session_state',''), COALESCE(policies_json::text,'{}')
|
|
FROM bgp_peer WHERE id=$1 AND tenant_id=$2`, id, tenantID).Scan(
|
|
&peer.ID, &peer.TenantID, &sp, &peer.Neighbor, &peer.RemoteASN, &peer.Enabled, &peer.Name, &peer.SessionState, &peer.PoliciesJSON)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, store.ErrNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
peer.SpeakerID = sp
|
|
return &peer, nil
|
|
}
|
|
|
|
func (p *Postgres) CreatePeer(tenantID string, in *store.BGPPeer) (*store.BGPPeer, error) {
|
|
if in == nil || in.RemoteASN == 0 {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
neighbor, ok := store.NormalizePeerNeighborString(in.Neighbor)
|
|
if !ok {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
ctx := context.Background()
|
|
id := uuid.NewString()
|
|
meta := map[string]any{"name": in.Name, "session_state": in.SessionState}
|
|
mb, _ := json.Marshal(meta)
|
|
pol := "{}"
|
|
if strings.TrimSpace(in.PoliciesJSON) != "" {
|
|
pol = in.PoliciesJSON
|
|
}
|
|
var sp any
|
|
if in.SpeakerID != nil && strings.TrimSpace(*in.SpeakerID) != "" {
|
|
sp = strings.TrimSpace(*in.SpeakerID)
|
|
}
|
|
enabled := store.EffectivePeerEnabledOnCreate(in.Enabled, in.SessionState)
|
|
_, err := p.pool.Exec(ctx, `
|
|
INSERT INTO bgp_peer (id, tenant_id, bgp_speaker_id, neighbor, remote_asn, enabled, policies_json, meta_json)
|
|
VALUES ($1,$2,$3,$4::inet, $5, $6, $7::jsonb, $8::jsonb)`,
|
|
id, tenantID, sp, neighbor, in.RemoteASN, enabled, pol, string(mb))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetPeer(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) UpdatePeer(tenantID, id string, patch *store.PeerPatch) (*store.BGPPeer, error) {
|
|
cur, err := p.GetPeer(tenantID, id)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if patch.Neighbor != nil {
|
|
n, ok := store.NormalizePeerNeighborString(*patch.Neighbor)
|
|
if !ok {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
cur.Neighbor = n
|
|
}
|
|
if patch.RemoteASN != nil {
|
|
cur.RemoteASN = *patch.RemoteASN
|
|
}
|
|
if patch.Enabled != nil {
|
|
cur.Enabled = *patch.Enabled
|
|
}
|
|
if patch.Name != nil {
|
|
cur.Name = *patch.Name
|
|
}
|
|
if patch.SessionState != nil {
|
|
cur.SessionState = *patch.SessionState
|
|
}
|
|
if patch.PoliciesJSON != nil {
|
|
cur.PoliciesJSON = *patch.PoliciesJSON
|
|
}
|
|
if patch.SpeakerID != nil {
|
|
v := strings.TrimSpace(*patch.SpeakerID)
|
|
if v == "" {
|
|
cur.SpeakerID = nil
|
|
} else {
|
|
cur.SpeakerID = &v
|
|
}
|
|
}
|
|
ctx := context.Background()
|
|
meta := map[string]any{"name": cur.Name, "session_state": cur.SessionState}
|
|
mb, _ := json.Marshal(meta)
|
|
pol := "{}"
|
|
if strings.TrimSpace(cur.PoliciesJSON) != "" {
|
|
pol = cur.PoliciesJSON
|
|
}
|
|
var sp any
|
|
if cur.SpeakerID != nil && strings.TrimSpace(*cur.SpeakerID) != "" {
|
|
sp = strings.TrimSpace(*cur.SpeakerID)
|
|
}
|
|
_, err = p.pool.Exec(ctx, `
|
|
UPDATE bgp_peer SET neighbor=$3::inet, remote_asn=$4, enabled=$5, policies_json=$6::jsonb, meta_json=$7::jsonb,
|
|
bgp_speaker_id=$8, updated_at=now()
|
|
WHERE id=$1 AND tenant_id=$2`, id, tenantID, cur.Neighbor, cur.RemoteASN, cur.Enabled, pol, string(mb), sp)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetPeer(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) DeletePeer(tenantID, id string) error {
|
|
ctx := context.Background()
|
|
tag, err := p.pool.Exec(ctx, `DELETE FROM bgp_peer WHERE id=$1 AND tenant_id=$2`, id, tenantID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return store.ErrNotFound
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *Postgres) ListSpeakersForTenant(tenantID string) []*store.Speaker {
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `
|
|
SELECT id::text, role, COALESCE(endpoint,''), last_applied_revision_id::text, COALESCE(meta_json::text,'{}')
|
|
FROM bgp_speaker WHERE tenant_id=$1 ORDER BY id`, tenantID)
|
|
if err != nil {
|
|
return nil
|
|
}
|
|
defer rows.Close()
|
|
var out []*store.Speaker
|
|
for rows.Next() {
|
|
var s store.Speaker
|
|
s.TenantID = tenantID
|
|
var lap *string
|
|
if err := rows.Scan(&s.ID, &s.Role, &s.Endpoint, &lap, &s.MetaJSON); err != nil {
|
|
continue
|
|
}
|
|
s.LastAppliedRevisionID = strOrNil(lap)
|
|
out = append(out, &s)
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (p *Postgres) GetSpeaker(tenantID, speakerID string) (*store.Speaker, error) {
|
|
sp, err := p.getSpeakerRow(context.Background(), speakerID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if sp.TenantID != tenantID {
|
|
return nil, store.ErrTenantScope
|
|
}
|
|
return sp, nil
|
|
}
|
|
|
|
func (p *Postgres) GetSpeakerAnyTenant(speakerID string) (*store.Speaker, error) {
|
|
return p.getSpeakerRow(context.Background(), speakerID)
|
|
}
|
|
|
|
func (p *Postgres) getSpeakerRow(ctx context.Context, speakerID string) (*store.Speaker, error) {
|
|
var s store.Speaker
|
|
var lap *string
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT id::text, tenant_id::text, role, COALESCE(endpoint,''), last_applied_revision_id::text, COALESCE(meta_json::text,'{}')
|
|
FROM bgp_speaker WHERE id=$1`, speakerID).Scan(&s.ID, &s.TenantID, &s.Role, &s.Endpoint, &lap, &s.MetaJSON)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, store.ErrNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
s.LastAppliedRevisionID = strOrNil(lap)
|
|
return &s, nil
|
|
}
|
|
|
|
func (p *Postgres) CreateSpeaker(tenantID string, in *store.Speaker) (*store.Speaker, error) {
|
|
if in == nil {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
ctx := context.Background()
|
|
id := uuid.NewString()
|
|
meta := in.MetaJSON
|
|
if meta == "" {
|
|
meta = "{}"
|
|
}
|
|
var ep any
|
|
if strings.TrimSpace(in.Endpoint) != "" {
|
|
ep = strings.TrimSpace(in.Endpoint)
|
|
}
|
|
_, err := p.pool.Exec(ctx, `INSERT INTO bgp_speaker (id, tenant_id, role, endpoint, meta_json) VALUES ($1,$2,$3,$4,$5::jsonb)`,
|
|
id, tenantID, in.Role, ep, meta)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetSpeaker(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) UpdateSpeaker(tenantID, id string, patch *store.SpeakerPatch) (*store.Speaker, error) {
|
|
cur, err := p.GetSpeaker(tenantID, id)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if patch.Role != nil {
|
|
cur.Role = strings.TrimSpace(*patch.Role)
|
|
}
|
|
if patch.Endpoint != nil {
|
|
cur.Endpoint = *patch.Endpoint
|
|
}
|
|
if patch.MetaJSON != nil {
|
|
cur.MetaJSON = *patch.MetaJSON
|
|
}
|
|
ctx := context.Background()
|
|
meta := cur.MetaJSON
|
|
if strings.TrimSpace(meta) == "" {
|
|
meta = "{}"
|
|
}
|
|
var ep any
|
|
if strings.TrimSpace(cur.Endpoint) != "" {
|
|
ep = strings.TrimSpace(cur.Endpoint)
|
|
}
|
|
_, err = p.pool.Exec(ctx, `UPDATE bgp_speaker SET role=$3, endpoint=$4, meta_json=$5::jsonb, updated_at=now() WHERE id=$1 AND tenant_id=$2`,
|
|
id, tenantID, cur.Role, ep, meta)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetSpeaker(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) DeleteSpeaker(tenantID, id string) error {
|
|
ctx := context.Background()
|
|
tag, err := p.pool.Exec(ctx, `DELETE FROM bgp_speaker WHERE id=$1 AND tenant_id=$2`, id, tenantID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return store.ErrNotFound
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *Postgres) GetRevision(tenantID, revisionID string) (*store.Revision, error) {
|
|
ctx, cancel := boundedRepoCtx(context.Background())
|
|
defer cancel()
|
|
var r store.Revision
|
|
var mod *string
|
|
var parent *string
|
|
var meta []byte
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT id::text, tenant_id::text, module_id::text, content_hash, parent_revision_id::text, meta_json, created_at
|
|
FROM config_revision WHERE id=$1 AND tenant_id=$2`, revisionID, tenantID).Scan(
|
|
&r.ID, &r.TenantID, &mod, &r.ContentHash, &parent, &meta, &r.CreatedAt)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, store.ErrNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
if mod != nil {
|
|
r.ModuleID = *mod
|
|
}
|
|
r.ParentRevisionID = strOrNil(parent)
|
|
var mj struct {
|
|
PreviewFragments map[string]string `json:"preview_fragments"`
|
|
MaterializedPrefixCount int `json:"materialized_prefix_count"`
|
|
}
|
|
_ = json.Unmarshal(meta, &mj)
|
|
if mj.PreviewFragments == nil {
|
|
mj.PreviewFragments = map[string]string{}
|
|
}
|
|
r.PreviewFragments = loadRevisionPreview(ctx, p.pool, revisionID, mj.PreviewFragments)
|
|
r.MaterializedPrefixCount = mj.MaterializedPrefixCount
|
|
return &r, nil
|
|
}
|
|
|
|
func (p *Postgres) GetRevisionSummary(tenantID, revisionID string) (*store.Revision, error) {
|
|
ctx, cancel := boundedRepoCtx(context.Background())
|
|
defer cancel()
|
|
var r store.Revision
|
|
var mod *string
|
|
var parent *string
|
|
var prefixCount int
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT id::text, tenant_id::text, module_id::text, content_hash, parent_revision_id::text,
|
|
COALESCE((meta_json->>'materialized_prefix_count')::int, 0), created_at
|
|
FROM config_revision WHERE id=$1 AND tenant_id=$2`, revisionID, tenantID).Scan(
|
|
&r.ID, &r.TenantID, &mod, &r.ContentHash, &parent, &prefixCount, &r.CreatedAt)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, store.ErrNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
if mod != nil {
|
|
r.ModuleID = *mod
|
|
}
|
|
r.ParentRevisionID = strOrNil(parent)
|
|
r.MaterializedPrefixCount = prefixCount
|
|
r.PreviewFragments = map[string]string{}
|
|
return &r, nil
|
|
}
|
|
|
|
func (p *Postgres) ListRevisions(tenantID, moduleID string, cursor string, limit int) ([]*store.Revision, string, bool) {
|
|
if limit <= 0 {
|
|
limit = 50
|
|
}
|
|
off := 0
|
|
if cursor != "" {
|
|
if n, err := strconv.Atoi(cursor); err == nil && n >= 0 {
|
|
off = n
|
|
}
|
|
}
|
|
ctx := context.Background()
|
|
// List endpoint only needs materialized_prefix_count from meta_json — not full preview_fragments blobs.
|
|
// LIMIT/OFFSET in SQL avoids loading every revision for the tenant into memory (was O(N) per request).
|
|
q := `SELECT id::text, module_id::text, content_hash, parent_revision_id::text,
|
|
COALESCE((meta_json->>'materialized_prefix_count')::int, 0), created_at
|
|
FROM config_revision WHERE tenant_id=$1`
|
|
args := []any{tenantID}
|
|
n := 2
|
|
if moduleID != "" {
|
|
q += fmt.Sprintf(` AND module_id=$%d`, n)
|
|
args = append(args, moduleID)
|
|
n++
|
|
}
|
|
q += ` ORDER BY created_at DESC`
|
|
// Fetch limit+1 rows to compute has_more without COUNT(*).
|
|
q += fmt.Sprintf(` LIMIT $%d OFFSET $%d`, n, n+1)
|
|
args = append(args, limit+1, off)
|
|
rows, err := p.pool.Query(ctx, q, args...)
|
|
if err != nil {
|
|
return nil, "", false
|
|
}
|
|
defer rows.Close()
|
|
var all []*store.Revision
|
|
for rows.Next() {
|
|
var r store.Revision
|
|
r.TenantID = tenantID
|
|
var mod, parent *string
|
|
var mpc int
|
|
if err := rows.Scan(&r.ID, &mod, &r.ContentHash, &parent, &mpc, &r.CreatedAt); err != nil {
|
|
continue
|
|
}
|
|
if mod != nil {
|
|
r.ModuleID = *mod
|
|
}
|
|
r.ParentRevisionID = strOrNil(parent)
|
|
r.MaterializedPrefixCount = mpc
|
|
r.PreviewFragments = map[string]string{}
|
|
all = append(all, &r)
|
|
}
|
|
agentDebugNDJSON3214("A", "repository/postgres.go:ListRevisions", "list_revisions_fetched", map[string]any{
|
|
"rows": len(all), "limit": limit, "offset": off,
|
|
})
|
|
hasMore := len(all) > limit
|
|
if hasMore {
|
|
all = all[:limit]
|
|
}
|
|
next := ""
|
|
if hasMore {
|
|
next = fmt.Sprintf("%d", off+limit)
|
|
}
|
|
if len(all) == 0 {
|
|
return nil, "", false
|
|
}
|
|
return all, next, hasMore
|
|
}
|
|
|
|
func (p *Postgres) ListRevisionPrefixes(tenantID, revisionID string, cursor string, limit int) ([]store.PrefixRow, string, bool) {
|
|
if limit <= 0 {
|
|
limit = 50
|
|
}
|
|
afterID, off, useOffset := store.ParsePrefixPageCursor(cursor)
|
|
ctx := context.Background()
|
|
var one int
|
|
if err := p.pool.QueryRow(ctx, `
|
|
SELECT 1 FROM config_revision WHERE id = $1::uuid AND tenant_id = $2::uuid`,
|
|
revisionID, tenantID).Scan(&one); err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, "", false
|
|
}
|
|
return nil, "", false
|
|
}
|
|
if snapID, ok := p.revisionPrefixSnapshotID(ctx, revisionID); ok {
|
|
return p.listSnapshotPrefixes(ctx, snapID, cursor, limit)
|
|
}
|
|
var rows pgx.Rows
|
|
var err error
|
|
if useOffset {
|
|
rows, err = p.pool.Query(ctx, `
|
|
SELECT id, prefix::text, community_id::text, source FROM revision_materialized_prefix
|
|
WHERE revision_id=$1::uuid ORDER BY id
|
|
LIMIT $2 OFFSET $3`, revisionID, limit+1, off)
|
|
} else {
|
|
var afterArg any
|
|
if afterID != nil {
|
|
afterArg = *afterID
|
|
}
|
|
rows, err = p.pool.Query(ctx, `
|
|
SELECT id, prefix::text, community_id::text, source FROM revision_materialized_prefix
|
|
WHERE revision_id=$1::uuid AND ($2::bigint IS NULL OR id > $2::bigint)
|
|
ORDER BY id
|
|
LIMIT $3`, revisionID, afterArg, limit+1)
|
|
}
|
|
if err != nil {
|
|
return nil, "", false
|
|
}
|
|
defer rows.Close()
|
|
var all []store.PrefixRow
|
|
var ids []int64
|
|
for rows.Next() {
|
|
var rowID int64
|
|
var pr store.PrefixRow
|
|
var comm *string
|
|
if err := rows.Scan(&rowID, &pr.Prefix, &comm, &pr.Source); err != nil {
|
|
continue
|
|
}
|
|
pr.CommunityID = comm
|
|
ids = append(ids, rowID)
|
|
all = append(all, pr)
|
|
}
|
|
agentDebugNDJSON3214("B", "repository/postgres.go:ListRevisionPrefixes", "list_prefixes_fetched", map[string]any{
|
|
"rows": len(all), "limit": limit, "keyset": !useOffset,
|
|
})
|
|
more := len(all) > limit
|
|
if more {
|
|
all = all[:limit]
|
|
ids = ids[:limit]
|
|
}
|
|
next := ""
|
|
if more && len(ids) > 0 {
|
|
next = store.FormatPrefixPageCursor(ids[len(ids)-1])
|
|
}
|
|
if len(all) == 0 {
|
|
return nil, "", false
|
|
}
|
|
return all, next, more
|
|
}
|
|
|
|
func (p *Postgres) CreateRollbackRevision(tenantID, sourceRevisionID string) (string, error) {
|
|
src, err := p.GetRevision(tenantID, sourceRevisionID)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
ctx := context.Background()
|
|
newID := uuid.NewString()
|
|
parent := sourceRevisionID
|
|
meta, err := revisionMetaWithoutPreview(src.MaterializedPrefixCount)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
var modArg any
|
|
if strings.TrimSpace(src.ModuleID) != "" {
|
|
modArg = src.ModuleID
|
|
}
|
|
tx, err := p.pool.Begin(ctx)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
defer func() { _ = tx.Rollback(ctx) }()
|
|
_, err = tx.Exec(ctx, `
|
|
INSERT INTO config_revision (id, tenant_id, module_id, content_hash, parent_revision_id, meta_json)
|
|
VALUES ($1,$2,$3,$4,$5::uuid,$6::jsonb)`,
|
|
newID, tenantID, modArg, src.ContentHash+":rollback", parent, meta)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if err := copyRevisionPreview(ctx, tx, newID, sourceRevisionID); err != nil {
|
|
return "", err
|
|
}
|
|
if err := p.copyRevisionPrefixSnapshotRef(ctx, tx, newID, sourceRevisionID); err != nil {
|
|
return "", err
|
|
}
|
|
if err := tx.Commit(ctx); err != nil {
|
|
return "", err
|
|
}
|
|
return newID, nil
|
|
}
|
|
|
|
const maxRevisionDiffRows = 5000
|
|
|
|
func (p *Postgres) RevisionDiff(tenantID, aID, bID string) (map[string]any, error) {
|
|
if _, err := p.GetRevision(tenantID, aID); err != nil {
|
|
return nil, err
|
|
}
|
|
if _, err := p.GetRevision(tenantID, bID); err != nil {
|
|
return nil, err
|
|
}
|
|
ctx := context.Background()
|
|
var unchanged int
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT COUNT(*)::int FROM (
|
|
SELECT b.prefix FROM (`+sqlRevisionPrefixes("$2")+`) b
|
|
INNER JOIN (`+sqlRevisionPrefixes("$1")+`) a ON a.prefix = b.prefix
|
|
) t`, aID, bID).Scan(&unchanged)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
rowsAdded, err := p.pool.Query(ctx, `
|
|
SELECT b.prefix::text FROM (`+sqlRevisionPrefixes("$2")+`) b
|
|
LEFT JOIN (`+sqlRevisionPrefixes("$1")+`) a ON a.prefix = b.prefix
|
|
WHERE a.prefix IS NULL
|
|
ORDER BY b.prefix
|
|
LIMIT $3`, aID, bID, maxRevisionDiffRows+1)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rowsAdded.Close()
|
|
var added []string
|
|
for rowsAdded.Next() {
|
|
var s string
|
|
if err := rowsAdded.Scan(&s); err != nil {
|
|
continue
|
|
}
|
|
added = append(added, s)
|
|
if len(added) > maxRevisionDiffRows {
|
|
added = added[:maxRevisionDiffRows]
|
|
break
|
|
}
|
|
}
|
|
addedTruncated := len(added) >= maxRevisionDiffRows
|
|
rowsRem, err := p.pool.Query(ctx, `
|
|
SELECT a.prefix::text FROM (`+sqlRevisionPrefixes("$1")+`) a
|
|
LEFT JOIN (`+sqlRevisionPrefixes("$2")+`) b ON b.prefix = a.prefix
|
|
WHERE b.prefix IS NULL
|
|
ORDER BY a.prefix
|
|
LIMIT $3`, aID, bID, maxRevisionDiffRows+1)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rowsRem.Close()
|
|
var removed []string
|
|
for rowsRem.Next() {
|
|
var s string
|
|
if err := rowsRem.Scan(&s); err != nil {
|
|
continue
|
|
}
|
|
removed = append(removed, s)
|
|
if len(removed) > maxRevisionDiffRows {
|
|
removed = removed[:maxRevisionDiffRows]
|
|
break
|
|
}
|
|
}
|
|
return map[string]any{
|
|
"revision_a": aID,
|
|
"revision_b": bID,
|
|
"prefixes": map[string]any{
|
|
"added": added, "removed": removed, "unchanged_count": unchanged,
|
|
"truncated": addedTruncated || len(removed) >= maxRevisionDiffRows,
|
|
},
|
|
}, nil
|
|
}
|
|
|
|
func (p *Postgres) SetLastAppliedRevision(tenantID, speakerID, revisionID string) error {
|
|
ctx := context.Background()
|
|
tag, err := p.pool.Exec(ctx, `
|
|
UPDATE bgp_speaker SET last_applied_revision_id=$3::uuid, updated_at=now()
|
|
WHERE id=$1 AND tenant_id=$2`, speakerID, tenantID, revisionID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return store.ErrNotFound
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *Postgres) PublishRevisionForSpeaker(speakerID, revisionID string) error {
|
|
ctx := context.Background()
|
|
tag, err := p.pool.Exec(ctx, `
|
|
UPDATE bgp_speaker SET published_revision_id=$2::uuid, published_at=now(), updated_at=now() WHERE id=$1`, speakerID, revisionID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return store.ErrNotFound
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *Postgres) LatestPublishedRevision(speakerID string) (string, time.Time, error) {
|
|
ctx := context.Background()
|
|
var rid string
|
|
var at time.Time
|
|
err := p.pool.QueryRow(ctx, `
|
|
SELECT published_revision_id::text, published_at FROM bgp_speaker
|
|
WHERE id=$1 AND published_revision_id IS NOT NULL`, speakerID).Scan(&rid, &at)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return "", time.Time{}, store.ErrNotFound
|
|
}
|
|
return "", time.Time{}, err
|
|
}
|
|
return rid, at, nil
|
|
}
|
|
|
|
func (p *Postgres) ListTenantIDs() ([]string, error) {
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `SELECT id::text FROM tenant ORDER BY id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []string
|
|
for rows.Next() {
|
|
var id string
|
|
if err := rows.Scan(&id); err != nil {
|
|
continue
|
|
}
|
|
out = append(out, id)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (p *Postgres) CreateRenderRevision(revisionID, tenantID, moduleID string, parentRevisionID *string, contentHash string, previewFragments map[string]string, prefixes []store.PrefixRow) error {
|
|
if strings.TrimSpace(revisionID) == "" {
|
|
return store.ErrInvalidInput
|
|
}
|
|
if _, err := p.GetModule(tenantID, moduleID); err != nil {
|
|
return err
|
|
}
|
|
ctx := context.Background()
|
|
if previewFragments == nil {
|
|
previewFragments = map[string]string{}
|
|
}
|
|
meta, err := revisionMetaWithoutPreview(len(prefixes))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tx, err := p.pool.Begin(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer func() { _ = tx.Rollback(ctx) }()
|
|
|
|
var parent any
|
|
if parentRevisionID != nil && strings.TrimSpace(*parentRevisionID) != "" {
|
|
parent = strings.TrimSpace(*parentRevisionID)
|
|
}
|
|
revID := strings.TrimSpace(revisionID)
|
|
_, err = tx.Exec(ctx, `
|
|
INSERT INTO config_revision (id, tenant_id, module_id, content_hash, parent_revision_id, meta_json)
|
|
VALUES ($1::uuid, $2::uuid, $3::uuid, $4, $5::uuid, $6::jsonb)`,
|
|
revID, tenantID, moduleID, strings.TrimSpace(contentHash), parent, meta)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := insertRevisionPreview(ctx, tx, revID, previewFragments); err != nil {
|
|
return err
|
|
}
|
|
snapID, err := p.ensurePrefixSnapshot(ctx, tx, contentHash, prefixes)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := p.linkRevisionPrefixSnapshot(ctx, tx, revID, snapID); err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Commit(ctx); err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *Postgres) ListDohProfiles(tenantID string) ([]*store.DohProfile, error) {
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `SELECT id::text, name, url, timeout_ms, secret_ref FROM doh_profile WHERE tenant_id=$1`, tenantID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []*store.DohProfile
|
|
for rows.Next() {
|
|
var d store.DohProfile
|
|
d.TenantID = tenantID
|
|
var to *int32
|
|
if err := rows.Scan(&d.ID, &d.Name, &d.URL, &to, &d.SecretRef); err != nil {
|
|
continue
|
|
}
|
|
if to != nil {
|
|
v := int(*to)
|
|
d.TimeoutMs = &v
|
|
}
|
|
out = append(out, &d)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (p *Postgres) GetDohProfile(tenantID, id string) (*store.DohProfile, error) {
|
|
ctx := context.Background()
|
|
var d store.DohProfile
|
|
d.TenantID = tenantID
|
|
var to *int32
|
|
err := p.pool.QueryRow(ctx, `SELECT id::text, name, url, timeout_ms, secret_ref FROM doh_profile WHERE id=$1 AND tenant_id=$2`, id, tenantID).Scan(
|
|
&d.ID, &d.Name, &d.URL, &to, &d.SecretRef)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, store.ErrNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
if to != nil {
|
|
v := int(*to)
|
|
d.TimeoutMs = &v
|
|
}
|
|
return &d, nil
|
|
}
|
|
|
|
func (p *Postgres) CreateDohProfile(tenantID string, in *store.DohProfile) (*store.DohProfile, error) {
|
|
if in == nil {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
ctx := context.Background()
|
|
id := uuid.NewString()
|
|
_, err := p.pool.Exec(ctx, `INSERT INTO doh_profile (id, tenant_id, name, url, timeout_ms, secret_ref) VALUES ($1,$2,$3,$4,$5,$6)`,
|
|
id, tenantID, in.Name, in.URL, nullInt32Ptr(in.TimeoutMs), in.SecretRef)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetDohProfile(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) UpdateDohProfile(tenantID, id string, patch *store.DohProfilePatch) (*store.DohProfile, error) {
|
|
cur, err := p.GetDohProfile(tenantID, id)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if patch.Name != nil {
|
|
cur.Name = *patch.Name
|
|
}
|
|
if patch.URL != nil {
|
|
cur.URL = *patch.URL
|
|
}
|
|
if patch.TimeoutMs != nil {
|
|
cur.TimeoutMs = patch.TimeoutMs
|
|
}
|
|
if patch.SecretRef != nil {
|
|
cur.SecretRef = patch.SecretRef
|
|
}
|
|
ctx := context.Background()
|
|
_, err = p.pool.Exec(ctx, `UPDATE doh_profile SET name=$3, url=$4, timeout_ms=$5, secret_ref=$6, updated_at=now() WHERE id=$1 AND tenant_id=$2`,
|
|
id, tenantID, cur.Name, cur.URL, nullInt32Ptr(cur.TimeoutMs), cur.SecretRef)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetDohProfile(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) DeleteDohProfile(tenantID, id string) error {
|
|
ctx := context.Background()
|
|
inUse, err := p.moduleDohProfileInUse(ctx, id)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if inUse {
|
|
return store.ErrInvalidInput
|
|
}
|
|
var n int
|
|
_ = p.pool.QueryRow(ctx, `SELECT COUNT(*) FROM module WHERE doh_profile_id=$1::uuid AND deleted_at IS NULL`, id).Scan(&n)
|
|
if n > 0 {
|
|
return store.ErrInvalidInput
|
|
}
|
|
tag, err := p.pool.Exec(ctx, `DELETE FROM doh_profile WHERE id=$1 AND tenant_id=$2`, id, tenantID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return store.ErrNotFound
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *Postgres) ListCommunities(tenantID string) ([]*store.Community, error) {
|
|
ctx := context.Background()
|
|
rows, err := p.pool.Query(ctx, `SELECT id::text, community, title, value_json::text FROM bgp_community WHERE tenant_id=$1 ORDER BY COALESCE(NULLIF(trim(title), ''), community)`, tenantID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var out []*store.Community
|
|
for rows.Next() {
|
|
var c store.Community
|
|
c.TenantID = tenantID
|
|
if err := rows.Scan(&c.ID, &c.Community, &c.Title, &c.ValueJSON); err != nil {
|
|
continue
|
|
}
|
|
out = append(out, &c)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (p *Postgres) GetCommunity(tenantID, id string) (*store.Community, error) {
|
|
ctx := context.Background()
|
|
var c store.Community
|
|
c.TenantID = tenantID
|
|
err := p.pool.QueryRow(ctx, `SELECT id::text, community, title, value_json::text FROM bgp_community WHERE id=$1 AND tenant_id=$2`, id, tenantID).Scan(
|
|
&c.ID, &c.Community, &c.Title, &c.ValueJSON)
|
|
if err != nil {
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return nil, store.ErrNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
return &c, nil
|
|
}
|
|
|
|
func (p *Postgres) CreateCommunity(tenantID string, in *store.Community) (*store.Community, error) {
|
|
if in == nil {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
if strings.TrimSpace(in.Community) == "" {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
ctx := context.Background()
|
|
id := uuid.NewString()
|
|
vj := in.ValueJSON
|
|
if strings.TrimSpace(vj) == "" {
|
|
vj = "{}"
|
|
}
|
|
comm := strings.TrimSpace(in.Community)
|
|
title := strings.TrimSpace(in.Title)
|
|
_, err := p.pool.Exec(ctx, `INSERT INTO bgp_community (id, tenant_id, community, title, value_json) VALUES ($1,$2,$3,$4,$5::jsonb)`,
|
|
id, tenantID, comm, title, vj)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetCommunity(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) UpdateCommunity(tenantID, id string, patch *store.CommunityPatch) (*store.Community, error) {
|
|
cur, err := p.GetCommunity(tenantID, id)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if patch.Community != nil {
|
|
cur.Community = strings.TrimSpace(*patch.Community)
|
|
}
|
|
if patch.Title != nil {
|
|
cur.Title = strings.TrimSpace(*patch.Title)
|
|
}
|
|
if patch.ValueJSON != nil {
|
|
cur.ValueJSON = *patch.ValueJSON
|
|
}
|
|
if strings.TrimSpace(cur.Community) == "" {
|
|
return nil, store.ErrInvalidInput
|
|
}
|
|
ctx := context.Background()
|
|
_, err = p.pool.Exec(ctx, `UPDATE bgp_community SET community=$3, title=$4, value_json=$5::jsonb, updated_at=now() WHERE id=$1 AND tenant_id=$2`,
|
|
id, tenantID, cur.Community, cur.Title, cur.ValueJSON)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return p.GetCommunity(tenantID, id)
|
|
}
|
|
|
|
func (p *Postgres) DeleteCommunity(tenantID, id string) error {
|
|
ctx := context.Background()
|
|
var n int
|
|
_ = p.pool.QueryRow(ctx, `SELECT COUNT(*) FROM module WHERE default_community_id=$1::uuid AND deleted_at IS NULL`, id).Scan(&n)
|
|
if n > 0 {
|
|
return store.ErrInvalidInput
|
|
}
|
|
tag, err := p.pool.Exec(ctx, `DELETE FROM bgp_community WHERE id=$1 AND tenant_id=$2`, id, tenantID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tag.RowsAffected() == 0 {
|
|
return store.ErrNotFound
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// CDN / AS / domain / IP / settings: см. postgres_entities.go.
|
|
|
|
func strPtrUUID(s *string) *string {
|
|
if s == nil || strings.TrimSpace(*s) == "" {
|
|
return nil
|
|
}
|
|
v := strings.TrimSpace(*s)
|
|
return &v
|
|
}
|
|
|
|
func strOrNil(s *string) *string {
|
|
if s == nil || *s == "" {
|
|
return nil
|
|
}
|
|
return s
|
|
}
|
|
|
|
func nullStr(s string) *string {
|
|
if strings.TrimSpace(s) == "" {
|
|
return nil
|
|
}
|
|
v := strings.TrimSpace(s)
|
|
return &v
|
|
}
|
|
|
|
func nullInt32(i int) *int32 {
|
|
if i == 0 {
|
|
return nil
|
|
}
|
|
v := int32(i)
|
|
return &v
|
|
}
|
|
|
|
func nullIntOrZero(i int) any {
|
|
if i == 0 {
|
|
return nil
|
|
}
|
|
return i
|
|
}
|
|
|
|
func nullInt32Ptr(i *int) *int32 {
|
|
if i == nil {
|
|
return nil
|
|
}
|
|
v := int32(*i)
|
|
return &v
|
|
}
|
|
|
|
func nullTimePtr(t *time.Time) *time.Time {
|
|
if t == nil {
|
|
return nil
|
|
}
|
|
v := t.UTC()
|
|
return &v
|
|
}
|
|
|
|
func nullJSON(s string) *string {
|
|
if strings.TrimSpace(s) == "" {
|
|
v := "{}"
|
|
return &v
|
|
}
|
|
return &s
|
|
}
|