feat: Prometheus metrics endpoint for knox nodes
Refs #2 - /metrics served on a dedicated port (KNOX_METRICS_ADDR, default localhost:8932) via prometheus/client_golang, with Go runtime + process collectors - DB-derived gauges refreshed per scrape: observations by source, last-24h observations, entries, projects, sessions, pending reflections, threads by status, peers, observations by origin node, knowledge vector (max hcl per node) - live gossip counters (pulls/pushes, observations pulled/pushed, errors) incremented during the anti-entropy sweep; Run accepts an optional metrics handle (nil for one-shot CLI) - knox_node_info{node_id,name} for scrape identification - internal/metrics package + db MetricsSnapshot; tests for snapshot, scrape output, and counter increments
This commit is contained in:
@@ -61,7 +61,7 @@ func NewGossipCmd(kdb *db.KnoxDB) *cobra.Command {
|
||||
if len(peers) == 0 {
|
||||
return fmt.Errorf("no peers configured (set KNOX_PEERS)")
|
||||
}
|
||||
watch.Run(kdb, peers)
|
||||
watch.Run(kdb, nil, peers)
|
||||
created, linked, err := watch.Reconcile(kdb)
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -0,0 +1,108 @@
|
||||
package db
|
||||
|
||||
// MetricsSnapshot holds the scrape-time gauges derived from the database.
|
||||
type MetricsSnapshot struct {
|
||||
Observations int
|
||||
ObservationsLast24h int
|
||||
BySource map[string]int
|
||||
Entries int
|
||||
Projects int
|
||||
Sessions int
|
||||
PendingReflections int
|
||||
ThreadsByStatus map[string]int
|
||||
Peers int
|
||||
ByOriginNode map[string]int // node_id → observation count
|
||||
KnowledgeVector map[string]int64 // node_id → max hcl
|
||||
EarliestObservation string
|
||||
}
|
||||
|
||||
// MetricsSnapshot computes database-derived gauges for Prometheus scraping.
|
||||
// All queries are cheap aggregations; nothing is stored or mutated.
|
||||
func (k *KnoxDB) MetricsSnapshot() (*MetricsSnapshot, error) {
|
||||
s := &MetricsSnapshot{
|
||||
BySource: make(map[string]int),
|
||||
ThreadsByStatus: make(map[string]int),
|
||||
ByOriginNode: make(map[string]int),
|
||||
}
|
||||
|
||||
if err := k.db.QueryRow("SELECT COUNT(*) FROM observations").Scan(&s.Observations); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := k.db.QueryRow("SELECT COUNT(*) FROM observations WHERE collected_at > datetime('now', '-1 day')").Scan(&s.ObservationsLast24h); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := k.db.QueryRow("SELECT COUNT(*) FROM entries").Scan(&s.Entries); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := k.db.QueryRow("SELECT COUNT(DISTINCT project) FROM entries").Scan(&s.Projects); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := k.db.QueryRow("SELECT COUNT(*) FROM sessions").Scan(&s.Sessions); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := k.db.QueryRow("SELECT COUNT(*) FROM sessions WHERE indexed=0").Scan(&s.PendingReflections); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := k.db.QueryRow("SELECT COUNT(*) FROM peers").Scan(&s.Peers); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := k.db.QueryRow("SELECT COALESCE(MIN(collected_at),'') FROM observations").Scan(&s.EarliestObservation); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
rows, err := k.db.Query("SELECT COALESCE(source_id,'unknown'), COUNT(*) FROM observations GROUP BY source_id")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
for rows.Next() {
|
||||
var src string
|
||||
var n int
|
||||
if err := rows.Scan(&src, &n); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s.BySource[src] = n
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
rows2, err := k.db.Query("SELECT COALESCE(status,'active'), COUNT(*) FROM threads GROUP BY status")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows2.Close()
|
||||
for rows2.Next() {
|
||||
var st string
|
||||
var n int
|
||||
if err := rows2.Scan(&st, &n); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s.ThreadsByStatus[st] = n
|
||||
}
|
||||
if err := rows2.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
rows3, err := k.db.Query("SELECT node_id, COUNT(*) FROM observations WHERE node_id<>'' GROUP BY node_id")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows3.Close()
|
||||
for rows3.Next() {
|
||||
var nid string
|
||||
var n int
|
||||
if err := rows3.Scan(&nid, &n); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s.ByOriginNode[nid] = n
|
||||
}
|
||||
if err := rows3.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if s.KnowledgeVector, err = k.KnowledgeVector(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return s, nil
|
||||
}
|
||||
@@ -0,0 +1,164 @@
|
||||
// Package metrics exposes Prometheus-format metrics for a knox node.
|
||||
//
|
||||
// Gauges are recomputed from the database on each scrape (cheap aggregates);
|
||||
// gossip counters are in-memory and incremented as the daemon exchanges data
|
||||
// with peers. The registry also gains the standard Go runtime and process
|
||||
// collectors from prometheus/client_golang.
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
|
||||
"github.com/david/knox/internal/db"
|
||||
)
|
||||
|
||||
// Metrics holds the gossip event counters (incremented by watch) and the
|
||||
// scrape-time gauges derived from the database (refreshed on each scrape).
|
||||
type Metrics struct {
|
||||
// Gossip counters (live).
|
||||
pullsTotal prometheus.Counter
|
||||
pushesTotal prometheus.Counter
|
||||
obsPulledTotal prometheus.Counter
|
||||
obsPushedTotal prometheus.Counter
|
||||
errorsTotal prometheus.Counter
|
||||
|
||||
// Snapshot gauges (updated per scrape).
|
||||
observationsGauge *prometheus.GaugeVec
|
||||
entriesGauge prometheus.Gauge
|
||||
projectsGauge prometheus.Gauge
|
||||
sessionsGauge prometheus.Gauge
|
||||
pendingReflections prometheus.Gauge
|
||||
peersGauge prometheus.Gauge
|
||||
threadsByStatus *prometheus.GaugeVec
|
||||
byOriginNode *prometheus.GaugeVec
|
||||
knowledgeVector *prometheus.GaugeVec
|
||||
observationsLast24h prometheus.Gauge
|
||||
|
||||
registry *prometheus.Registry
|
||||
kdb *db.KnoxDB
|
||||
}
|
||||
|
||||
// New builds the metrics registry bound to a knowledge index.
|
||||
func New(kdb *db.KnoxDB, nodeName string) *Metrics {
|
||||
reg := prometheus.NewRegistry()
|
||||
|
||||
m := &Metrics{
|
||||
registry: reg,
|
||||
kdb: kdb,
|
||||
}
|
||||
|
||||
// Node identity aids scraping: which node produced this output.
|
||||
m.nodeInfo(kdb.NodeID(), nodeName)
|
||||
|
||||
m.pullsTotal = newCounter(reg, "knox_gossip_pulls_total", "Peer pull round-trips completed.")
|
||||
m.pushesTotal = newCounter(reg, "knox_gossip_pushes_total", "Peer push round-trips completed.")
|
||||
m.obsPulledTotal = newCounter(reg, "knox_gossip_observations_pulled_total", "Observations received from peers.")
|
||||
m.obsPushedTotal = newCounter(reg, "knox_gossip_observations_pushed_total", "Observations sent to peers.")
|
||||
m.errorsTotal = newCounter(reg, "knox_gossip_errors_total", "Gossip errors (ping/pull/push failures).")
|
||||
|
||||
m.observationsGauge = newGaugeVec(reg, "knox_observations_total", "Observation log size.", "source_id")
|
||||
m.observationsLast24h = newGauge(reg, "knox_observations_last_24h", "Observations collected in the last 24h.")
|
||||
m.entriesGauge = newGauge(reg, "knox_entries_total", "Materialized entry cache size.")
|
||||
m.projectsGauge = newGauge(reg, "knox_projects_total", "Distinct projects in the entry cache.")
|
||||
m.sessionsGauge = newGauge(reg, "knox_sessions_total", "Sessions tracked.")
|
||||
m.pendingReflections = newGauge(reg, "knox_pending_reflections", "Sessions awaiting reflection.")
|
||||
m.peersGauge = newGauge(reg, "knox_peers_total", "Known peer nodes.")
|
||||
m.threadsByStatus = newGaugeVec(reg, "knox_threads_total", "Threads by status.", "status")
|
||||
m.byOriginNode = newGaugeVec(reg, "knox_observations_by_node", "Observations per originating node.", "node_id")
|
||||
m.knowledgeVector = newGaugeVec(reg, "knox_knowledge_max_hcl", "Highest HCL seen per originating node.", "node_id")
|
||||
|
||||
// Go runtime + process collectors come from the official library.
|
||||
reg.MustRegister(prometheus.NewGoCollector())
|
||||
reg.MustRegister(prometheus.NewProcessCollector(prometheus.ProcessCollectorOpts{}))
|
||||
|
||||
return m
|
||||
}
|
||||
|
||||
func newCounter(reg *prometheus.Registry, name, help string) prometheus.Counter {
|
||||
c := prometheus.NewCounter(prometheus.CounterOpts{Name: name, Help: help})
|
||||
reg.MustRegister(c)
|
||||
return c
|
||||
}
|
||||
|
||||
func newGauge(reg *prometheus.Registry, name, help string) prometheus.Gauge {
|
||||
g := prometheus.NewGauge(prometheus.GaugeOpts{Name: name, Help: help})
|
||||
reg.MustRegister(g)
|
||||
return g
|
||||
}
|
||||
|
||||
func newGaugeVec(reg *prometheus.Registry, name, help string, labels ...string) *prometheus.GaugeVec {
|
||||
g := prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: name, Help: help}, labels)
|
||||
reg.MustRegister(g)
|
||||
return g
|
||||
}
|
||||
|
||||
func (m *Metrics) nodeInfo(nodeID, name string) {
|
||||
info := prometheus.NewGauge(prometheus.GaugeOpts{
|
||||
Name: "knox_node_info",
|
||||
Help: "Node identity (always 1).",
|
||||
ConstLabels: prometheus.Labels{
|
||||
"node_id": nodeID,
|
||||
"name": name,
|
||||
},
|
||||
})
|
||||
m.registry.MustRegister(info)
|
||||
info.Set(1)
|
||||
}
|
||||
|
||||
// Capture refresh the DB-derived gauges from a fresh snapshot.
|
||||
func (m *Metrics) Capture(s *db.MetricsSnapshot) {
|
||||
m.observationsGauge.Reset()
|
||||
for src, n := range s.BySource {
|
||||
m.observationsGauge.WithLabelValues(src).Set(float64(n))
|
||||
}
|
||||
m.observationsLast24h.Set(float64(s.ObservationsLast24h))
|
||||
m.entriesGauge.Set(float64(s.Entries))
|
||||
m.projectsGauge.Set(float64(s.Projects))
|
||||
m.sessionsGauge.Set(float64(s.Sessions))
|
||||
m.pendingReflections.Set(float64(s.PendingReflections))
|
||||
m.peersGauge.Set(float64(s.Peers))
|
||||
|
||||
m.threadsByStatus.Reset()
|
||||
for st, n := range s.ThreadsByStatus {
|
||||
m.threadsByStatus.WithLabelValues(st).Set(float64(n))
|
||||
}
|
||||
|
||||
m.byOriginNode.Reset()
|
||||
for nid, n := range s.ByOriginNode {
|
||||
m.byOriginNode.WithLabelValues(nid).Set(float64(n))
|
||||
}
|
||||
m.knowledgeVector.Reset()
|
||||
for nid, hcl := range s.KnowledgeVector {
|
||||
m.knowledgeVector.WithLabelValues(nid).Set(float64(hcl))
|
||||
}
|
||||
}
|
||||
|
||||
// IncrementPull records a completed pull and its accepted observation count.
|
||||
func (m *Metrics) IncrementPull(newObs int) {
|
||||
m.pullsTotal.Inc()
|
||||
m.obsPulledTotal.Add(float64(newObs))
|
||||
}
|
||||
|
||||
// IncrementPush records a completed push and its accepted observation count.
|
||||
func (m *Metrics) IncrementPush(newObs int) {
|
||||
m.pushesTotal.Inc()
|
||||
m.obsPushedTotal.Add(float64(newObs))
|
||||
}
|
||||
|
||||
// IncrementErrors counts a failed gossip attempt.
|
||||
func (m *Metrics) IncrementErrors() { m.errorsTotal.Inc() }
|
||||
|
||||
// Handler returns the /metrics scrape handler. On each scrape it refreshes
|
||||
// DB-derived gauges before rendering (cost is a few cheap aggregates).
|
||||
func (m *Metrics) Handler() http.Handler {
|
||||
h := promhttp.HandlerFor(m.registry, promhttp.HandlerOpts{})
|
||||
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if s, err := m.kdb.MetricsSnapshot(); err == nil {
|
||||
m.Capture(s)
|
||||
}
|
||||
h.ServeHTTP(w, r)
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"io"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/david/knox/internal/db"
|
||||
"github.com/prometheus/client_golang/prometheus/testutil"
|
||||
)
|
||||
|
||||
func tmpKdb(t *testing.T) *db.KnoxDB {
|
||||
t.Helper()
|
||||
k, err := db.Open(filepath.Join(t.TempDir(), "index.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open db: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { k.Close() })
|
||||
return k
|
||||
}
|
||||
|
||||
func seed(t *testing.T, k *db.KnoxDB, src string, n int) {
|
||||
t.Helper()
|
||||
for i := 0; i < n; i++ {
|
||||
_, _, err := k.RecordObservation(db.ObservationRecord{
|
||||
Fingerprint: "fp-" + src + "-" + string(rune('a'+i)),
|
||||
SourceID: src,
|
||||
SourcePath: src,
|
||||
Project: "test",
|
||||
ContentType: "test",
|
||||
Title: src,
|
||||
Summary: "s",
|
||||
CreatedAt: "2026-08-29T00:00:00Z",
|
||||
Confidence: 0.9,
|
||||
IngesterVersion: "itest/v1",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("seed: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMetricsSnapshot(t *testing.T) {
|
||||
k := tmpKdb(t)
|
||||
seed(t, k, "git", 2)
|
||||
seed(t, k, "browser-history", 3)
|
||||
|
||||
s, err := k.MetricsSnapshot()
|
||||
if err != nil {
|
||||
t.Fatalf("snapshot: %v", err)
|
||||
}
|
||||
if s.Observations != 5 {
|
||||
t.Errorf("observations = %d, want 5", s.Observations)
|
||||
}
|
||||
if s.BySource["git"] != 2 || s.BySource["browser-history"] != 3 {
|
||||
t.Errorf("by source = %v", s.BySource)
|
||||
}
|
||||
if s.KnowledgeVector[k.NodeID()] == 0 {
|
||||
t.Errorf("knowledge vector missing own node")
|
||||
}
|
||||
if s.ByOriginNode[k.NodeID()] != 5 {
|
||||
t.Errorf("by origin node = %v", s.ByOriginNode)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMetricsScrape(t *testing.T) {
|
||||
k := tmpKdb(t)
|
||||
seed(t, k, "git", 2)
|
||||
|
||||
m := New(k, "testnode")
|
||||
h := m.Handler()
|
||||
|
||||
rec := httptest.NewRecorder()
|
||||
h.ServeHTTP(rec, httptest.NewRequest("GET", "/metrics", nil))
|
||||
|
||||
body, _ := io.ReadAll(rec.Body)
|
||||
out := string(body)
|
||||
for _, want := range []string{
|
||||
`knox_node_info{name="testnode"`,
|
||||
`knox_observations_total{source_id="git"} 2`,
|
||||
`knox_gossip_pulls_total 0`,
|
||||
"go_goroutines",
|
||||
"process_cpu_seconds_total",
|
||||
} {
|
||||
if !strings.Contains(out, want) {
|
||||
t.Errorf("scrape output missing %q", want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestMetricsCounters(t *testing.T) {
|
||||
k := tmpKdb(t)
|
||||
m := New(k, "t")
|
||||
|
||||
// Manually drive counters through the Metrics API.
|
||||
m.IncrementPull(3)
|
||||
m.IncrementPush(7)
|
||||
m.IncrementErrors()
|
||||
|
||||
if got := testutil.ToFloat64(m.pullsTotal); got != 1 {
|
||||
t.Errorf("pullsTotal = %v, want 1", got)
|
||||
}
|
||||
if got := testutil.ToFloat64(m.obsPulledTotal); got != 3 {
|
||||
t.Errorf("obsPulled = %v, want 3", got)
|
||||
}
|
||||
if got := testutil.ToFloat64(m.obsPushedTotal); got != 7 {
|
||||
t.Errorf("obsPushed = %v, want 7", got)
|
||||
}
|
||||
if got := testutil.ToFloat64(m.errorsTotal); got != 1 {
|
||||
t.Errorf("errorsTotal = %v, want 1", got)
|
||||
}
|
||||
}
|
||||
@@ -19,14 +19,16 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/david/knox/internal/db"
|
||||
"github.com/david/knox/internal/metrics"
|
||||
)
|
||||
|
||||
const defaultPort = "8931"
|
||||
|
||||
type Node struct {
|
||||
Kdb *db.KnoxDB
|
||||
Name string
|
||||
Addr string // advertised base URL, e.g. http://192.168.1.20:8931
|
||||
Kdb *db.KnoxDB
|
||||
Name string
|
||||
Addr string // advertised base URL, e.g. http://192.168.1.20:8931
|
||||
Metrics *metrics.Metrics
|
||||
}
|
||||
|
||||
// pingResponse is the anti-entropy summary returned by /v1/ping.
|
||||
@@ -269,7 +271,9 @@ func (c *Client) Diff() (*DiffSummary, error) {
|
||||
// with any one node reveals who else is in the swarm, and those nodes are then
|
||||
// swept too. There is no relay of observations — only membership is shared; each
|
||||
// node pulls/pushes directly with every other node it learns about.
|
||||
func Run(kdb *db.KnoxDB, static []string) {
|
||||
//
|
||||
// m, when non-nil, receives gossip event counters (nil for one-shot CLI runs).
|
||||
func Run(kdb *db.KnoxDB, m *metrics.Metrics, static []string) {
|
||||
myID := kdb.NodeID()
|
||||
|
||||
// Seed the work queue with static config plus persisted discoveries.
|
||||
@@ -292,6 +296,9 @@ func Run(kdb *db.KnoxDB, static []string) {
|
||||
p, err := c.Ping()
|
||||
if err != nil {
|
||||
log.Printf("[gossip] ping %s: %v", addr, err)
|
||||
if m != nil {
|
||||
m.IncrementErrors()
|
||||
}
|
||||
continue
|
||||
}
|
||||
if p.NodeID == myID {
|
||||
@@ -323,13 +330,22 @@ func Run(kdb *db.KnoxDB, static []string) {
|
||||
n, err := c.Pull(remoteNode, localHCL, kdb)
|
||||
if err != nil {
|
||||
log.Printf("[gossip] pull %s@%s: %v", remoteNode, addr, err)
|
||||
if m != nil {
|
||||
m.IncrementErrors()
|
||||
}
|
||||
continue
|
||||
}
|
||||
pulled += n
|
||||
}
|
||||
}
|
||||
if m != nil {
|
||||
m.IncrementPull(pulled)
|
||||
}
|
||||
|
||||
pushed, _ := c.Push(kdb, p.Vector)
|
||||
if m != nil {
|
||||
m.IncrementPush(pushed)
|
||||
}
|
||||
|
||||
if err := kdb.UpsertPeer(p.NodeID, addr, p.Name, vectorMax(p.Vector)); err != nil {
|
||||
log.Printf("[gossip] peer upsert: %v", err)
|
||||
@@ -368,6 +384,18 @@ func ListenAddr() (addr string) {
|
||||
return "localhost:" + defaultPort
|
||||
}
|
||||
|
||||
const defaultMetricsPort = "8932"
|
||||
|
||||
// MetricsAddr returns the Prometheus scrape address (KNOX_METRICS_ADDR or
|
||||
// default). It is a separate port from gossip so scraping never contends with
|
||||
// the peer protocol.
|
||||
func MetricsAddr() (addr string) {
|
||||
if addr = os.Getenv("KNOX_METRICS_ADDR"); addr != "" {
|
||||
return addr
|
||||
}
|
||||
return "localhost:" + defaultMetricsPort
|
||||
}
|
||||
|
||||
// PeerAddrs returns the configured peer list (KNOX_PEERS, comma-separated).
|
||||
func PeerAddrs() []string {
|
||||
raw := os.Getenv("KNOX_PEERS")
|
||||
|
||||
@@ -60,8 +60,8 @@ func TestGossipConvergence(t *testing.T) {
|
||||
defer sb.Close()
|
||||
|
||||
// A pulls from B, then B pulls from A (bidirectional sweep).
|
||||
Run(a, []string{sb.URL})
|
||||
Run(b, []string{sa.URL})
|
||||
Run(a, nil, []string{sb.URL})
|
||||
Run(b, nil, []string{sa.URL})
|
||||
|
||||
av, err := a.KnowledgeVector()
|
||||
if err != nil {
|
||||
@@ -107,11 +107,11 @@ func TestGossipSwarmDiscovery(t *testing.T) {
|
||||
defer sc.Close()
|
||||
|
||||
// A discovers B (A pings B) so A can advertise B to the swarm.
|
||||
Run(a, []string{sb.URL})
|
||||
Run(a, nil, []string{sb.URL})
|
||||
|
||||
// C only knows A. A single sweep should surface B (membership in ping)
|
||||
// and pull B's observations directly.
|
||||
Run(c, []string{sa.URL})
|
||||
Run(c, nil, []string{sa.URL})
|
||||
|
||||
// C must know B and hold all three origin logs.
|
||||
peers, err := c.ListPeers()
|
||||
@@ -227,12 +227,12 @@ func TestGossipIdempotent(t *testing.T) {
|
||||
sb := httptest.NewServer(nodeB.Handler())
|
||||
defer sb.Close()
|
||||
|
||||
Run(b, []string{sa.URL})
|
||||
Run(b, nil, []string{sa.URL})
|
||||
before, err := b.EntryCount()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
Run(b, []string{sa.URL})
|
||||
Run(b, nil, []string{sa.URL})
|
||||
after, err := b.EntryCount()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
+16
-2
@@ -10,6 +10,7 @@ import (
|
||||
|
||||
"github.com/david/knox/internal/db"
|
||||
"github.com/david/knox/internal/ingest"
|
||||
"github.com/david/knox/internal/metrics"
|
||||
"github.com/fsnotify/fsnotify"
|
||||
)
|
||||
|
||||
@@ -24,6 +25,7 @@ type Watcher struct {
|
||||
vault string
|
||||
debounce time.Duration
|
||||
fileIngesters []ingest.Ingester
|
||||
metrics *metrics.Metrics
|
||||
}
|
||||
|
||||
func New(kdb *db.KnoxDB, dirs []string) *Watcher {
|
||||
@@ -40,6 +42,7 @@ func New(kdb *db.KnoxDB, dirs []string) *Watcher {
|
||||
ingest.NewLogIngester(),
|
||||
ingest.NewSkillsIngester(),
|
||||
},
|
||||
metrics: metrics.New(kdb, "knox"),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -53,7 +56,7 @@ func (w *Watcher) Start() error {
|
||||
// Start the gossip server first so peers can reach us while the initial
|
||||
// seed is still ingesting. Vault/dirs are logged after the seed below.
|
||||
gossipAddr := ListenAddr()
|
||||
node := &Node{Kdb: w.knoxDB, Name: "knox", Addr: gossipAddr}
|
||||
node := &Node{Kdb: w.knoxDB, Name: "knox", Addr: gossipAddr, Metrics: w.metrics}
|
||||
srv := &http.Server{Addr: gossipAddr, Handler: node.Handler()}
|
||||
go func() {
|
||||
log.Printf("[knox] gossip listening on %s", gossipAddr)
|
||||
@@ -61,6 +64,17 @@ func (w *Watcher) Start() error {
|
||||
log.Printf("[knox] gossip server: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
// Prometheus scraping on a dedicated port (KNOX_METRICS_ADDR).
|
||||
metricsAddr := MetricsAddr()
|
||||
metricsSrv := &http.Server{Addr: metricsAddr, Handler: node.Metrics.Handler()}
|
||||
go func() {
|
||||
log.Printf("[knox] metrics listening on %s", metricsAddr)
|
||||
if err := metricsSrv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
log.Printf("[knox] metrics server: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
if peers := PeerAddrs(); len(peers) > 0 {
|
||||
log.Printf("[knox] gossip peers: %v", peers)
|
||||
}
|
||||
@@ -447,7 +461,7 @@ func (w *Watcher) syncGossip() {
|
||||
before = n
|
||||
}
|
||||
|
||||
Run(w.knoxDB, peers)
|
||||
Run(w.knoxDB, w.metrics, peers)
|
||||
|
||||
// If new observations arrived, reconcile to pick up entries/threads they
|
||||
// imply (deterministic log → derived rebuild).
|
||||
|
||||
Reference in New Issue
Block a user