Files
EvoBGP/internal/repository/postgres.go
T
Denozordec f39df7c4bf
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
feat(revisions): add pruning estimate and cleanup endpoints
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.
2026-06-12 21:56:36 +07:00

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
}