290 lines
8.9 KiB
Go
290 lines
8.9 KiB
Go
package db
|
|
|
|
import (
|
|
"database/sql"
|
|
"fmt"
|
|
)
|
|
|
|
// GossipObservation is the serializable wire form of an observation exchanged
|
|
// between nodes. It carries the locator (NodeID, HCL) so the receiver can
|
|
// idempotently INSERT OR IGNORE on the unique index.
|
|
type GossipObservation struct {
|
|
NodeID string
|
|
HCL int64
|
|
Fingerprint string
|
|
SourceID string
|
|
SourcePath string
|
|
Project string
|
|
ContentType string
|
|
Title string
|
|
Summary string
|
|
CollectedAt string
|
|
CreatedAt string
|
|
LineStart int
|
|
LineEnd int
|
|
Confidence float64
|
|
IngesterVersion string
|
|
Trigger string
|
|
Provenance string
|
|
}
|
|
|
|
// PushObservations inserts gossip rows idempotently. Rows whose (node_id, hcl)
|
|
// already exist are ignored (a node's HCL is locally monotonic, so a given
|
|
// locator is immutable). Returns the number newly inserted.
|
|
func (k *KnoxDB) PushObservations(rows []GossipObservation) (int, error) {
|
|
tx, err := k.db.Begin()
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
defer tx.Rollback()
|
|
|
|
inserted := 0
|
|
for _, o := range rows {
|
|
res, err := tx.Exec(
|
|
`INSERT OR IGNORE INTO observations
|
|
(fingerprint, source_id, source_path, project, content_type, title, summary,
|
|
collected_at, created_at, line_start, line_end, confidence, ingester_version, trigger, provenance,
|
|
node_id, hcl)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
|
o.Fingerprint, o.SourceID, o.SourcePath, o.Project, o.ContentType, o.Title, o.Summary,
|
|
o.CollectedAt, o.CreatedAt, o.LineStart, o.LineEnd, o.Confidence, o.IngesterVersion, o.Trigger, o.Provenance,
|
|
o.NodeID, o.HCL,
|
|
)
|
|
if err != nil {
|
|
return inserted, fmt.Errorf("push observation: %w", err)
|
|
}
|
|
if n, _ := res.RowsAffected(); n > 0 {
|
|
inserted++
|
|
}
|
|
}
|
|
return inserted, tx.Commit()
|
|
}
|
|
|
|
// ObservationsAfter returns gossip rows for one node after a given HCL cursor,
|
|
// in HCL order. after <= 0 means "everything" (initial handshake).
|
|
func (k *KnoxDB) ObservationsAfter(nodeID string, after int64, limit int) ([]GossipObservation, error) {
|
|
if limit <= 0 {
|
|
limit = 500
|
|
}
|
|
query := `SELECT node_id, hcl, fingerprint, source_id, COALESCE(source_path,''), COALESCE(project,''),
|
|
COALESCE(content_type,''), COALESCE(title,''), COALESCE(summary,''),
|
|
COALESCE(collected_at,''), COALESCE(created_at,''),
|
|
COALESCE(line_start,0), COALESCE(line_end,0), COALESCE(confidence,0.5),
|
|
COALESCE(ingester_version,''), COALESCE(trigger,''), COALESCE(provenance,'{}')
|
|
FROM observations WHERE node_id=?`
|
|
args := []any{nodeID}
|
|
if after > 0 {
|
|
query += ` AND hcl > ?`
|
|
args = append(args, after)
|
|
}
|
|
query += ` ORDER BY hcl ASC LIMIT ?`
|
|
args = append(args, limit)
|
|
|
|
rows, err := k.db.Query(query, args...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []GossipObservation
|
|
for rows.Next() {
|
|
var o GossipObservation
|
|
if err := rows.Scan(&o.NodeID, &o.HCL, &o.Fingerprint, &o.SourceID, &o.SourcePath,
|
|
&o.Project, &o.ContentType, &o.Title, &o.Summary, &o.CollectedAt, &o.CreatedAt,
|
|
&o.LineStart, &o.LineEnd, &o.Confidence, &o.IngesterVersion, &o.Trigger, &o.Provenance); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, o)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// KnowledgeVector returns node_id → max HCL held locally. It is the anti-entropy
|
|
// summary exchanged at handshake.
|
|
func (k *KnoxDB) KnowledgeVector() (map[string]int64, error) {
|
|
rows, err := k.db.Query(`SELECT node_id, MAX(hcl) FROM observations WHERE node_id<>'' AND hcl IS NOT NULL GROUP BY node_id`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
kv := make(map[string]int64)
|
|
for rows.Next() {
|
|
var nid string
|
|
var hcl int64
|
|
if err := rows.Scan(&nid, &hcl); err != nil {
|
|
return nil, err
|
|
}
|
|
kv[nid] = hcl
|
|
}
|
|
return kv, rows.Err()
|
|
}
|
|
|
|
// ObservePairs returns all observations from other nodes (for consume/echo
|
|
// suppression checks).
|
|
func (k *KnoxDB) ObserveForeign(nodeID string) (int, error) {
|
|
var n int
|
|
err := k.db.QueryRow(`SELECT COUNT(*) FROM observations WHERE node_id<>? AND node_id<>''`, nodeID).Scan(&n)
|
|
return n, err
|
|
}
|
|
|
|
// UpsertPeer records a peer's last-handshake metadata.
|
|
func (k *KnoxDB) UpsertPeer(peerID, addr, name string, maxHCL int64) error {
|
|
_, err := k.db.Exec(
|
|
`INSERT INTO peers (peer_id, addr, name, last_handshake, cursor) VALUES (?, ?, ?, datetime('now'), ?)
|
|
ON CONFLICT(peer_id) DO UPDATE SET
|
|
addr=?, name=?,
|
|
last_handshake=datetime('now'),
|
|
cursor=?`,
|
|
peerID, addr, name, maxHCL, addr, name, maxHCL,
|
|
)
|
|
return err
|
|
}
|
|
|
|
// MergePeer records a peer discovered indirectly (via another peer's ping).
|
|
// Unlike UpsertPeer it never clobbers the cursor — a freshly learned address
|
|
// has no known knowledge yet; the anti-entropy pull will set it on contact.
|
|
func (k *KnoxDB) MergePeer(peerID, addr, name string) error {
|
|
_, err := k.db.Exec(
|
|
`INSERT INTO peers (peer_id, addr, name) VALUES (?, ?, ?)
|
|
ON CONFLICT(peer_id) DO UPDATE SET addr=?, name=?`,
|
|
peerID, addr, name, addr, name,
|
|
)
|
|
return err
|
|
}
|
|
|
|
// ListPeers returns known peers ordered by first-seen.
|
|
func (k *KnoxDB) ListPeers() ([]Peer, error) {
|
|
rows, err := k.db.Query(`SELECT peer_id, COALESCE(addr,''), COALESCE(name,''), COALESCE(last_handshake,''), COALESCE(cursor,0), COALESCE(created_at,'') FROM peers ORDER BY created_at`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var peers []Peer
|
|
for rows.Next() {
|
|
var p Peer
|
|
if err := rows.Scan(&p.PeerID, &p.Addr, &p.Name, &p.LastHandshake, &p.Cursor, &p.CreatedAt); err != nil {
|
|
return nil, err
|
|
}
|
|
peers = append(peers, p)
|
|
}
|
|
return peers, rows.Err()
|
|
}
|
|
|
|
// GetPeerByAddr finds a peer whose address matches (statically configured).
|
|
func (k *KnoxDB) GetPeerByAddr(addr string) (*Peer, error) {
|
|
row := k.db.QueryRow(`SELECT peer_id, COALESCE(addr,''), COALESCE(name,''), COALESCE(last_handshake,''), COALESCE(cursor,0), COALESCE(created_at,'') FROM peers WHERE addr=?`, addr)
|
|
p := &Peer{}
|
|
err := row.Scan(&p.PeerID, &p.Addr, &p.Name, &p.LastHandshake, &p.Cursor, &p.CreatedAt)
|
|
if err == sql.ErrNoRows {
|
|
return nil, nil
|
|
}
|
|
return p, err
|
|
}
|
|
|
|
// Peer models a row in the peers table.
|
|
type Peer struct {
|
|
PeerID string
|
|
Addr string
|
|
Name string
|
|
LastHandshake string
|
|
Cursor int64
|
|
CreatedAt string
|
|
}
|
|
|
|
// PeerInfo is the shareable (non-secret) subset of a peer that /v1/ping
|
|
// advertises so other nodes can discover the swarm.
|
|
type PeerInfo struct {
|
|
PeerID string `json:"peer_id"`
|
|
Addr string `json:"addr"`
|
|
Name string `json:"name"`
|
|
}
|
|
|
|
// ShareablePeers returns the peers this node knows about, for dissemination in
|
|
// ping responses. Self and peers without an address are excluded.
|
|
func (k *KnoxDB) ShareablePeers() ([]PeerInfo, error) {
|
|
rows, err := k.db.Query(
|
|
`SELECT peer_id, COALESCE(addr,''), COALESCE(name,'') FROM peers WHERE addr<>'' AND peer_id<>? ORDER BY peer_id`,
|
|
k.NodeID(),
|
|
)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []PeerInfo
|
|
for rows.Next() {
|
|
var p PeerInfo
|
|
if err := rows.Scan(&p.PeerID, &p.Addr, &p.Name); err != nil {
|
|
return nil, err
|
|
}
|
|
out = append(out, p)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// MaxHCLForNode returns the maximum HCL this node holds for a given origin
|
|
// node, or 0 if none.
|
|
func (k *KnoxDB) MaxHCLForNode(nodeID string) int64 {
|
|
var m int64
|
|
if err := k.db.QueryRow(`SELECT MAX(hcl) FROM observations WHERE node_id=?`, nodeID).Scan(&m); err != nil {
|
|
return 0
|
|
}
|
|
return m
|
|
}
|
|
|
|
// SwarmPeerAddrs returns the addresses of all known peers (shareable set). It
|
|
// is what the sweep iterates after a bootstrap join.
|
|
func (k *KnoxDB) SwarmPeerAddrs() ([]string, error) {
|
|
peers, err := k.ShareablePeers()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
addrs := make([]string, 0, len(peers))
|
|
for _, p := range peers {
|
|
addrs = append(addrs, p.Addr)
|
|
}
|
|
return addrs, nil
|
|
}
|
|
|
|
// DistinctFingerprints returns the set of all observed fingerprints — the
|
|
// ground-truth index of what this node knows. Used by gossip diff.
|
|
func (k *KnoxDB) DistinctFingerprints() (map[string]bool, error) {
|
|
rows, err := k.db.Query(`SELECT DISTINCT fingerprint FROM observations`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
set := make(map[string]bool)
|
|
for rows.Next() {
|
|
var f string
|
|
if err := rows.Scan(&f); err != nil {
|
|
return nil, err
|
|
}
|
|
set[f] = true
|
|
}
|
|
return set, rows.Err()
|
|
}
|
|
|
|
// ThreadStatusByCluster returns cluster_key → status for auto-threaded threads,
|
|
// excluding human-created threads (no cluster key). Diff uses it to show
|
|
// tombstoned threads: a thread resolved on one node but active on another.
|
|
func (k *KnoxDB) ThreadStatusByCluster() (map[string]string, error) {
|
|
rows, err := k.db.Query(
|
|
`SELECT cluster_key, status FROM threads WHERE cluster_key IS NOT NULL AND cluster_key<>''`)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
out := make(map[string]string)
|
|
for rows.Next() {
|
|
var key, status string
|
|
if err := rows.Scan(&key, &status); err != nil {
|
|
return nil, err
|
|
}
|
|
out[key] = status
|
|
}
|
|
return out, rows.Err()
|
|
}
|