package repository import ( "context" "encoding/json" "errors" "strings" "time" "evobgp/internal/store" "github.com/google/uuid" "github.com/jackc/pgx/v5" ) func (p *Postgres) ListPeerDiscoveries(tenantID, status string) ([]*store.BGPPeerDiscovery, error) { ctx := context.Background() status = strings.TrimSpace(strings.ToLower(status)) q := ` SELECT id::text, tenant_id::text, COALESCE(speaker_id::text,''), COALESCE(neighbor_id,''), neighbor::text, remote_asn, COALESCE(protocol_name,''), COALESCE(session_state,''), status, first_seen_at, last_seen_at, COALESCE(approved_peer_id::text,'') FROM bgp_peer_discovery WHERE tenant_id=$1` args := []any{tenantID} if status != "" { q += ` AND status=$2` args = append(args, status) } q += ` ORDER BY last_seen_at DESC` rows, err := p.pool.Query(ctx, q, args...) if err != nil { return nil, err } defer rows.Close() var out []*store.BGPPeerDiscovery for rows.Next() { d, err := scanPeerDiscovery(rows) if err != nil { continue } out = append(out, d) } return out, nil } func (p *Postgres) GetPeerDiscovery(tenantID, id string) (*store.BGPPeerDiscovery, error) { ctx := context.Background() row := p.pool.QueryRow(ctx, ` SELECT id::text, tenant_id::text, COALESCE(speaker_id::text,''), COALESCE(neighbor_id,''), neighbor::text, remote_asn, COALESCE(protocol_name,''), COALESCE(session_state,''), status, first_seen_at, last_seen_at, COALESCE(approved_peer_id::text,'') FROM bgp_peer_discovery WHERE id=$1 AND tenant_id=$2`, id, tenantID) d, err := scanPeerDiscovery(row) if err != nil { if errors.Is(err, pgx.ErrNoRows) { return nil, store.ErrNotFound } return nil, err } return d, nil } type peerDiscoveryScanner interface { Scan(dest ...any) error } func scanPeerDiscovery(row peerDiscoveryScanner) (*store.BGPPeerDiscovery, error) { var d store.BGPPeerDiscovery var first, last time.Time err := row.Scan( &d.ID, &d.TenantID, &d.SpeakerID, &d.NeighborID, &d.Neighbor, &d.RemoteASN, &d.ProtocolName, &d.SessionState, &d.Status, &first, &last, &d.ApprovedPeerID, ) if err != nil { return nil, err } d.FirstSeenAt = first.UTC() d.LastSeenAt = last.UTC() return &d, nil } func (p *Postgres) UpsertPeerDiscovery(tenantID string, in *store.PeerDiscoveryUpsert) (*store.BGPPeerDiscovery, error) { if in == nil { return nil, store.ErrInvalidInput } neighbor, ok := store.NormalizePeerNeighborString(in.Neighbor) if !ok { return nil, store.ErrInvalidInput } neighborID := strings.TrimSpace(in.NeighborID) seenAt := in.SeenAt if seenAt.IsZero() { seenAt = time.Now().UTC() } ctx := context.Background() var existingID string if neighborID != "" { _ = p.pool.QueryRow(ctx, ` SELECT id::text FROM bgp_peer_discovery WHERE tenant_id=$1 AND neighbor_id=$2 LIMIT 1`, tenantID, neighborID).Scan(&existingID) } if existingID == "" { _ = p.pool.QueryRow(ctx, ` SELECT id::text FROM bgp_peer_discovery WHERE tenant_id=$1 AND neighbor=$2::inet AND remote_asn=$3 AND neighbor_id='' LIMIT 1`, tenantID, neighbor, in.RemoteASN).Scan(&existingID) } if existingID != "" { cur, err := p.GetPeerDiscovery(tenantID, existingID) if err != nil { return nil, err } if cur.Status == store.PeerDiscoveryRejected { return cur, nil } var sp any if s := strings.TrimSpace(in.SpeakerID); s != "" { sp = s } _, err = p.pool.Exec(ctx, ` UPDATE bgp_peer_discovery SET neighbor=$3::inet, remote_asn=$4, neighbor_id=CASE WHEN $5 <> '' THEN $5 ELSE neighbor_id END, protocol_name=$6, session_state=$7, last_seen_at=$8, speaker_id=COALESCE($9::uuid, speaker_id), updated_at=now() WHERE id=$1 AND tenant_id=$2 AND status <> 'rejected'`, existingID, tenantID, neighbor, in.RemoteASN, neighborID, strings.TrimSpace(in.ProtocolName), strings.TrimSpace(in.SessionState), seenAt, sp) if err != nil { return nil, err } return p.GetPeerDiscovery(tenantID, existingID) } id := uuid.NewString() var sp any if s := strings.TrimSpace(in.SpeakerID); s != "" { sp = s } _, err := p.pool.Exec(ctx, ` INSERT INTO bgp_peer_discovery ( id, tenant_id, speaker_id, neighbor_id, neighbor, remote_asn, protocol_name, session_state, status, first_seen_at, last_seen_at ) VALUES ($1,$2,$3::uuid,$4,$5::inet,$6,$7,$8,'pending',$9,$9)`, id, tenantID, sp, neighborID, neighbor, in.RemoteASN, strings.TrimSpace(in.ProtocolName), strings.TrimSpace(in.SessionState), seenAt) if err != nil { return nil, err } return p.GetPeerDiscovery(tenantID, id) } func (p *Postgres) ApprovePeerDiscovery(tenantID, id string, in *store.PeerDiscoveryApproveInput) (*store.BGPPeer, *store.BGPPeerDiscovery, error) { ctx := context.Background() tx, err := p.pool.Begin(ctx) if err != nil { return nil, nil, err } defer func() { _ = tx.Rollback(ctx) }() d, err := p.getPeerDiscoveryTx(ctx, tx, tenantID, id) if err != nil { return nil, nil, err } if d.Status != store.PeerDiscoveryPending { return nil, nil, store.ErrInvalidInput } if d.RemoteASN == 0 { return nil, nil, store.ErrInvalidInput } neighbor, ok := store.NormalizePeerNeighborString(d.Neighbor) if !ok { return nil, nil, store.ErrInvalidInput } name := "" enabled := true var speakerID *string if in != nil { name = strings.TrimSpace(in.Name) if in.Enabled != nil { enabled = *in.Enabled } if in.SpeakerID != nil { v := strings.TrimSpace(*in.SpeakerID) if v != "" { speakerID = &v } } } if name == "" { if d.NeighborID != "" { name = "discovered-" + d.NeighborID } else { name = "discovered-" + neighbor } } if speakerID == nil && strings.TrimSpace(d.SpeakerID) != "" { sp := d.SpeakerID speakerID = &sp } peerID := uuid.NewString() meta, _ := json.Marshal(map[string]any{"name": name, "session_state": d.SessionState}) var sp any if speakerID != nil { sp = *speakerID } _, err = tx.Exec(ctx, ` INSERT INTO bgp_peer (id, tenant_id, bgp_speaker_id, neighbor, remote_asn, enabled, policies_json, meta_json) VALUES ($1,$2,$3::uuid,$4::inet,$5,$6,'{}'::jsonb,$7::jsonb)`, peerID, tenantID, sp, neighbor, d.RemoteASN, enabled, string(meta)) if err != nil { return nil, nil, err } _, err = tx.Exec(ctx, ` UPDATE bgp_peer_discovery SET status='approved', approved_peer_id=$3::uuid, last_seen_at=now(), updated_at=now() WHERE id=$1 AND tenant_id=$2`, id, tenantID, peerID) if err != nil { return nil, nil, err } if err := tx.Commit(ctx); err != nil { return nil, nil, err } peer, err := p.GetPeer(tenantID, peerID) if err != nil { return nil, nil, err } disc, err := p.GetPeerDiscovery(tenantID, id) if err != nil { return nil, nil, err } return peer, disc, nil } func (p *Postgres) RejectPeerDiscovery(tenantID, id string) (*store.BGPPeerDiscovery, error) { ctx := context.Background() cur, err := p.GetPeerDiscovery(tenantID, id) if err != nil { return nil, err } if cur.Status == store.PeerDiscoveryApproved { return nil, store.ErrInvalidInput } tag, err := p.pool.Exec(ctx, ` UPDATE bgp_peer_discovery SET status='rejected', last_seen_at=now(), updated_at=now() WHERE id=$1 AND tenant_id=$2`, id, tenantID) if err != nil { return nil, err } if tag.RowsAffected() == 0 { return nil, store.ErrNotFound } return p.GetPeerDiscovery(tenantID, id) } func (p *Postgres) getPeerDiscoveryTx(ctx context.Context, tx pgx.Tx, tenantID, id string) (*store.BGPPeerDiscovery, error) { row := tx.QueryRow(ctx, ` SELECT id::text, tenant_id::text, COALESCE(speaker_id::text,''), COALESCE(neighbor_id,''), neighbor::text, remote_asn, COALESCE(protocol_name,''), COALESCE(session_state,''), status, first_seen_at, last_seen_at, COALESCE(approved_peer_id::text,'') FROM bgp_peer_discovery WHERE id=$1 AND tenant_id=$2 FOR UPDATE`, id, tenantID) d, err := scanPeerDiscovery(row) if err != nil { if errors.Is(err, pgx.ErrNoRows) { return nil, store.ErrNotFound } return nil, err } return d, nil }