Files
EvoBGP/internal/repository/module_snapshot.go
T
DenozordecandCursor 3500bd4624 refactor(db): normalize module prefix snapshot rows
Строки префиксов в module_prefix_snapshot_row вместо JSONB blobs.

Co-authored-by: Cursor <[email protected]>
2026-05-25 10:58:56 +07:00

138 lines
3.9 KiB
Go

package repository
import (
"context"
"encoding/json"
"errors"
"strings"
"time"
"evobgp/internal/store"
"github.com/jackc/pgx/v5"
)
func moduleSnapshotRowTableExists(ctx context.Context, q queryRower) bool {
var n int
err := q.QueryRow(ctx, `
SELECT 1 FROM information_schema.tables
WHERE table_schema = 'public' AND table_name = 'module_prefix_snapshot_row'
LIMIT 1`).Scan(&n)
return err == nil
}
func (p *Postgres) GetModulePrefixSnapshot(tenantID, moduleID string) (*store.ModulePrefixSnapshot, bool, error) {
ctx := context.Background()
var inputHash string
var collectedAt time.Time
err := p.pool.QueryRow(ctx, `
SELECT input_hash, collected_at
FROM module_prefix_snapshot
WHERE tenant_id = $1 AND module_id = $2`,
tenantID, moduleID).Scan(&inputHash, &collectedAt)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, false, nil
}
return nil, false, err
}
var prefixes []store.PrefixRow
if moduleSnapshotRowTableExists(ctx, p.pool) {
rows, qerr := p.pool.Query(ctx, `
SELECT prefix::text, community_id::text, source
FROM module_prefix_snapshot_row
WHERE tenant_id = $1::uuid AND module_id = $2::uuid
ORDER BY ord`, tenantID, moduleID)
if qerr != nil {
return nil, false, qerr
}
defer rows.Close()
for rows.Next() {
var pr store.PrefixRow
var comm *string
if err := rows.Scan(&pr.Prefix, &comm, &pr.Source); err != nil {
continue
}
pr.CommunityID = comm
prefixes = append(prefixes, pr)
}
} else {
var raw []byte
if err := p.pool.QueryRow(ctx, `
SELECT prefixes_json FROM module_prefix_snapshot
WHERE tenant_id = $1 AND module_id = $2`, tenantID, moduleID).Scan(&raw); err == nil && len(raw) > 0 {
_ = json.Unmarshal(raw, &prefixes)
}
}
return &store.ModulePrefixSnapshot{
InputHash: inputHash,
CollectedAt: collectedAt.UTC(),
Prefixes: prefixes,
}, true, nil
}
func (p *Postgres) SetModulePrefixSnapshot(tenantID, moduleID, inputHash string, prefixes []store.PrefixRow) error {
if strings.TrimSpace(tenantID) == "" || strings.TrimSpace(moduleID) == "" || strings.TrimSpace(inputHash) == "" {
return store.ErrInvalidInput
}
ctx := context.Background()
tx, err := p.pool.Begin(ctx)
if err != nil {
return err
}
defer func() { _ = tx.Rollback(ctx) }()
_, err = tx.Exec(ctx, `
INSERT INTO module_prefix_snapshot (tenant_id, module_id, input_hash, collected_at)
VALUES ($1::uuid, $2::uuid, $3, now())
ON CONFLICT (tenant_id, module_id) DO UPDATE SET
input_hash = EXCLUDED.input_hash,
collected_at = EXCLUDED.collected_at`,
tenantID, moduleID, inputHash)
if err != nil {
return err
}
if moduleSnapshotRowTableExists(ctx, tx) {
if _, err := tx.Exec(ctx, `
DELETE FROM module_prefix_snapshot_row
WHERE tenant_id = $1::uuid AND module_id = $2::uuid`, tenantID, moduleID); err != nil {
return err
}
for i, pr := range prefixes {
var comm any
if pr.CommunityID != nil && strings.TrimSpace(*pr.CommunityID) != "" {
comm = strings.TrimSpace(*pr.CommunityID)
}
src := pr.Source
if strings.TrimSpace(src) == "" {
src = "render"
}
if _, err := tx.Exec(ctx, `
INSERT INTO module_prefix_snapshot_row (tenant_id, module_id, ord, prefix, community_id, source)
VALUES ($1::uuid, $2::uuid, $3, $4::cidr, $5::uuid, $6)`,
tenantID, moduleID, i, strings.TrimSpace(pr.Prefix), comm, src); err != nil {
return err
}
}
} else {
raw, err := json.Marshal(prefixes)
if err != nil {
return err
}
if _, err := tx.Exec(ctx, `
UPDATE module_prefix_snapshot SET prefixes_json = $3::jsonb
WHERE tenant_id = $1::uuid AND module_id = $2::uuid`,
tenantID, moduleID, string(raw)); err != nil {
return err
}
}
return tx.Commit(ctx)
}
func (p *Postgres) DeleteModulePrefixSnapshot(tenantID, moduleID string) error {
ctx := context.Background()
_, err := p.pool.Exec(ctx, `
DELETE FROM module_prefix_snapshot WHERE tenant_id = $1 AND module_id = $2`,
tenantID, moduleID)
return err
}