Live KPIs
This commit is contained in:
+21
-3
@@ -461,7 +461,7 @@ func (m Model) Update(msg tea.Msg) (tea.Model, tea.Cmd) {
|
||||
|
||||
case snapshotLoadedMsg:
|
||||
snap := viz.Snapshot(msg)
|
||||
m.stats = &snap.Stats
|
||||
m.stats = viz.MergeAuthoritativeStats(m.stats, snap.Stats, viz.MergeStatsOpts{})
|
||||
m.network = &snap.Network
|
||||
items := viz.FeedFromSnapshot(snap, m.cfg.Buffer)
|
||||
m.store.Merge(items)
|
||||
@@ -659,10 +659,11 @@ func (m *Model) handleSSE(evt transport.Event) {
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
nowMS := time.Now().UnixMilli()
|
||||
switch typ {
|
||||
case "snapshot":
|
||||
snap := data.(viz.Snapshot)
|
||||
m.stats = &snap.Stats
|
||||
m.stats = viz.MergeAuthoritativeStats(m.stats, snap.Stats, viz.MergeStatsOpts{})
|
||||
m.network = &snap.Network
|
||||
m.store.Merge(viz.FeedFromSnapshot(snap, m.cfg.Buffer))
|
||||
if m.follow {
|
||||
@@ -676,6 +677,9 @@ func (m *Model) handleSSE(evt transport.Event) {
|
||||
if m.follow {
|
||||
m.cursor = 0
|
||||
}
|
||||
m.stats = viz.BumpStatsOnAttack(m.stats, a, nowMS)
|
||||
m.network = viz.BumpPeerOnAttack(m.network, a, nowMS)
|
||||
m.syncMeshFromNetwork()
|
||||
if m.cfg.Bell && a.TierAfter == "block" {
|
||||
fmt.Print("\a")
|
||||
}
|
||||
@@ -690,6 +694,18 @@ func (m *Model) handleSSE(evt transport.Event) {
|
||||
if m.network == nil {
|
||||
m.network = &viz.Network{}
|
||||
}
|
||||
found, wasOnline := viz.PeerOnlineStatus(m.network, p.ID)
|
||||
alreadyOnline := p.Status == "online" && wasOnline
|
||||
alreadyOffline := p.Status == "offline" && (!found || !wasOnline)
|
||||
delta := 0
|
||||
if p.Status == "online" {
|
||||
delta = 1
|
||||
} else if p.Status == "offline" {
|
||||
delta = -1
|
||||
}
|
||||
if delta != 0 {
|
||||
m.stats = viz.BumpConnectedPeers(m.stats, delta, alreadyOnline, alreadyOffline)
|
||||
}
|
||||
net := viz.ApplyPeerEvent(m.network, p)
|
||||
m.network = net
|
||||
m.syncMeshFromNetwork()
|
||||
@@ -704,6 +720,7 @@ func (m *Model) handleSSE(evt transport.Event) {
|
||||
if m.follow {
|
||||
m.cursor = 0
|
||||
}
|
||||
m.stats = viz.BumpStatsOnBlock(m.stats, b)
|
||||
m.peerMesh.TriggerBlock(m.meshPalette())
|
||||
m.sound.PlayBlock()
|
||||
case "moderation":
|
||||
@@ -717,12 +734,13 @@ func (m *Model) handleSSE(evt transport.Event) {
|
||||
fmt.Print("\a")
|
||||
}
|
||||
if mod.TierAfter == "block" && mod.TierBefore != "block" {
|
||||
// Match web: block KPI bumps come from dedicated block SSE; mesh pulse only here.
|
||||
m.peerMesh.TriggerBlock(m.meshPalette())
|
||||
}
|
||||
m.sound.PlayModeration()
|
||||
case "stats":
|
||||
s := data.(viz.Stats)
|
||||
m.stats = &s
|
||||
m.stats = viz.MergeAuthoritativeStats(m.stats, s, viz.MergeStatsOpts{Absolute: true})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,189 @@
|
||||
package viz
|
||||
|
||||
const hourMS = 60 * 60 * 1000
|
||||
|
||||
func emptyStats() Stats {
|
||||
return Stats{
|
||||
ServiceBreakdown: map[string]int{},
|
||||
GeoBreakdown: map[string]int{},
|
||||
}
|
||||
}
|
||||
|
||||
func cloneIntMap(m map[string]int) map[string]int {
|
||||
out := make(map[string]int, len(m))
|
||||
for k, v := range m {
|
||||
out[k] = v
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func ensureStats(s *Stats) Stats {
|
||||
if s == nil {
|
||||
return emptyStats()
|
||||
}
|
||||
out := *s
|
||||
out.ServiceBreakdown = cloneIntMap(out.ServiceBreakdown)
|
||||
out.GeoBreakdown = cloneIntMap(out.GeoBreakdown)
|
||||
if out.GeoBreakdownLastHour != nil {
|
||||
out.GeoBreakdownLastHour = cloneIntMap(out.GeoBreakdownLastHour)
|
||||
}
|
||||
if out.PeersByGeo != nil {
|
||||
out.PeersByGeo = cloneIntMap(out.PeersByGeo)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// BumpStatsOnAttack applies optimistic KPI updates for a live attack SSE event
|
||||
// (mirrors master/web/app/viz/vizStatsLive.ts bumpVizStatsOnAttack).
|
||||
func BumpStatsOnAttack(s *Stats, a AttackEvent, nowMS int64) *Stats {
|
||||
base := ensureStats(s)
|
||||
base.TotalAttacks++
|
||||
if a.Service != "" {
|
||||
base.ServiceBreakdown[a.Service]++
|
||||
}
|
||||
oneHourAgo := nowMS - hourMS
|
||||
ts := a.Timestamp
|
||||
if ts <= 0 {
|
||||
ts = nowMS
|
||||
}
|
||||
if ts >= oneHourAgo {
|
||||
base.AttacksLastHour++
|
||||
if a.Geo != "" {
|
||||
if base.GeoBreakdownLastHour == nil {
|
||||
base.GeoBreakdownLastHour = map[string]int{}
|
||||
}
|
||||
base.GeoBreakdownLastHour[a.Geo]++
|
||||
}
|
||||
}
|
||||
if a.Geo != "" {
|
||||
base.GeoBreakdown[a.Geo]++
|
||||
}
|
||||
return &base
|
||||
}
|
||||
|
||||
// BumpStatsOnBlock applies optimistic block KPI updates
|
||||
// (mirrors bumpVizStatsOnBlock).
|
||||
func BumpStatsOnBlock(s *Stats, b BlockEvent) *Stats {
|
||||
base := ensureStats(s)
|
||||
if b.SubnetBlock {
|
||||
base.BlockedSubnets++
|
||||
members := b.CoConspiratorCount
|
||||
if members <= 0 {
|
||||
members = 1
|
||||
}
|
||||
base.BlockedIPs += members
|
||||
} else {
|
||||
base.BlockedIPs++
|
||||
}
|
||||
base.RecentBlocksLastHour++
|
||||
return &base
|
||||
}
|
||||
|
||||
// BumpConnectedPeers adjusts connectedPeers with reconnect dedupe
|
||||
// (mirrors bumpVizConnectedPeers).
|
||||
func BumpConnectedPeers(s *Stats, delta int, alreadyOnline, alreadyOffline bool) *Stats {
|
||||
base := ensureStats(s)
|
||||
if delta > 0 && alreadyOnline {
|
||||
return &base
|
||||
}
|
||||
if delta < 0 && alreadyOffline {
|
||||
return &base
|
||||
}
|
||||
next := base.ConnectedPeers + delta
|
||||
if next < 0 {
|
||||
next = 0
|
||||
}
|
||||
base.ConnectedPeers = next
|
||||
return &base
|
||||
}
|
||||
|
||||
// MergeStatsOpts controls how server KPI snapshots merge into local optimistic state.
|
||||
type MergeStatsOpts struct {
|
||||
// Absolute takes server values as-is (live stats SSE / unblock).
|
||||
// When false, counter KPIs never drop below the previous optimistic value.
|
||||
Absolute bool
|
||||
}
|
||||
|
||||
// MergeAuthoritativeStats merges a server KPI snapshot into local state
|
||||
// (mirrors mergeAuthoritativeVizStats).
|
||||
func MergeAuthoritativeStats(prev *Stats, next Stats, opts MergeStatsOpts) *Stats {
|
||||
if prev == nil {
|
||||
out := next
|
||||
return &out
|
||||
}
|
||||
pick := func(a, b int) int {
|
||||
if opts.Absolute {
|
||||
return b
|
||||
}
|
||||
if a > b {
|
||||
return a
|
||||
}
|
||||
return b
|
||||
}
|
||||
out := *prev
|
||||
out.TotalAttacks = pick(prev.TotalAttacks, next.TotalAttacks)
|
||||
if opts.Absolute {
|
||||
out.ConnectedPeers = next.ConnectedPeers
|
||||
} else if next.ConnectedPeers > prev.ConnectedPeers {
|
||||
out.ConnectedPeers = next.ConnectedPeers
|
||||
}
|
||||
out.BlockedIPs = pick(prev.BlockedIPs, next.BlockedIPs)
|
||||
out.BlockedSubnets = pick(prev.BlockedSubnets, next.BlockedSubnets)
|
||||
out.AttacksLastHour = pick(prev.AttacksLastHour, next.AttacksLastHour)
|
||||
|
||||
if next.ServiceBreakdown != nil {
|
||||
out.ServiceBreakdown = cloneIntMap(next.ServiceBreakdown)
|
||||
}
|
||||
if next.GeoBreakdown != nil {
|
||||
out.GeoBreakdown = cloneIntMap(next.GeoBreakdown)
|
||||
}
|
||||
if next.GeoBreakdownLastHour != nil {
|
||||
out.GeoBreakdownLastHour = cloneIntMap(next.GeoBreakdownLastHour)
|
||||
}
|
||||
if next.PeersByGeo != nil {
|
||||
out.PeersByGeo = cloneIntMap(next.PeersByGeo)
|
||||
}
|
||||
if next.RecentBlocksLastHour != 0 || opts.Absolute {
|
||||
out.RecentBlocksLastHour = next.RecentBlocksLastHour
|
||||
}
|
||||
return &out
|
||||
}
|
||||
|
||||
// BumpPeerOnAttack updates per-peer hour/service counters on the live roster.
|
||||
func BumpPeerOnAttack(net *Network, a AttackEvent, nowMS int64) *Network {
|
||||
if net == nil || a.PeerID == "" {
|
||||
return net
|
||||
}
|
||||
ts := a.Timestamp
|
||||
if ts <= 0 {
|
||||
ts = nowMS
|
||||
}
|
||||
inHour := ts >= nowMS-hourMS
|
||||
for i := range net.Peers {
|
||||
if net.Peers[i].ID != a.PeerID {
|
||||
continue
|
||||
}
|
||||
if inHour {
|
||||
net.Peers[i].AttacksLastHour++
|
||||
}
|
||||
if a.Service != "" {
|
||||
net.Peers[i].TopServices = cloneIntMap(net.Peers[i].TopServices)
|
||||
net.Peers[i].TopServices[a.Service]++
|
||||
}
|
||||
return net
|
||||
}
|
||||
return net
|
||||
}
|
||||
|
||||
// PeerOnlineStatus reports whether a peer id is currently online in the roster.
|
||||
func PeerOnlineStatus(net *Network, id string) (found, online bool) {
|
||||
if net == nil || id == "" {
|
||||
return false, false
|
||||
}
|
||||
for _, p := range net.Peers {
|
||||
if p.ID == id {
|
||||
return true, p.Status == "online"
|
||||
}
|
||||
}
|
||||
return false, false
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
package viz_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/honeypeer/cli-viz/internal/viz"
|
||||
)
|
||||
|
||||
func TestBumpStatsOnAttack(t *testing.T) {
|
||||
now := int64(1_000_000_000)
|
||||
s := viz.BumpStatsOnAttack(nil, viz.AttackEvent{
|
||||
Timestamp: now,
|
||||
Service: "SSH",
|
||||
Geo: "CN",
|
||||
}, now)
|
||||
if s.TotalAttacks != 1 || s.AttacksLastHour != 1 {
|
||||
t.Fatalf("totals: %+v", s)
|
||||
}
|
||||
if s.ServiceBreakdown["SSH"] != 1 || s.GeoBreakdown["CN"] != 1 {
|
||||
t.Fatalf("breakdown: %+v", s)
|
||||
}
|
||||
if s.GeoBreakdownLastHour["CN"] != 1 {
|
||||
t.Fatalf("hour geo: %+v", s.GeoBreakdownLastHour)
|
||||
}
|
||||
|
||||
old := viz.BumpStatsOnAttack(s, viz.AttackEvent{
|
||||
Timestamp: now - 2*60*60*1000,
|
||||
Service: "HTTP",
|
||||
Geo: "US",
|
||||
}, now)
|
||||
if old.TotalAttacks != 2 || old.AttacksLastHour != 1 {
|
||||
t.Fatalf("old attack should not bump hour: %+v", old)
|
||||
}
|
||||
if old.ServiceBreakdown["HTTP"] != 1 || old.GeoBreakdown["US"] != 1 {
|
||||
t.Fatalf("all-time still bumps: %+v", old)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBumpStatsOnBlock(t *testing.T) {
|
||||
s := viz.BumpStatsOnBlock(nil, viz.BlockEvent{})
|
||||
if s.BlockedIPs != 1 || s.BlockedSubnets != 0 || s.RecentBlocksLastHour != 1 {
|
||||
t.Fatalf("host block: %+v", s)
|
||||
}
|
||||
s = viz.BumpStatsOnBlock(s, viz.BlockEvent{SubnetBlock: true, CoConspiratorCount: 3})
|
||||
if s.BlockedSubnets != 1 || s.BlockedIPs != 4 {
|
||||
t.Fatalf("subnet block: %+v", s)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBumpConnectedPeers(t *testing.T) {
|
||||
s := viz.BumpConnectedPeers(nil, 1, false, false)
|
||||
if s.ConnectedPeers != 1 {
|
||||
t.Fatalf("online: %+v", s)
|
||||
}
|
||||
same := viz.BumpConnectedPeers(s, 1, true, false)
|
||||
if same.ConnectedPeers != 1 {
|
||||
t.Fatalf("dedupe online: %+v", same)
|
||||
}
|
||||
s = viz.BumpConnectedPeers(s, -1, false, false)
|
||||
if s.ConnectedPeers != 0 {
|
||||
t.Fatalf("offline: %+v", s)
|
||||
}
|
||||
same = viz.BumpConnectedPeers(s, -1, false, true)
|
||||
if same.ConnectedPeers != 0 {
|
||||
t.Fatalf("dedupe offline: %+v", same)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMergeAuthoritativeStats(t *testing.T) {
|
||||
prev := &viz.Stats{TotalAttacks: 100, ConnectedPeers: 5, BlockedIPs: 10, AttacksLastHour: 20}
|
||||
next := viz.Stats{TotalAttacks: 90, ConnectedPeers: 4, BlockedIPs: 12, AttacksLastHour: 15, ServiceBreakdown: map[string]int{"SSH": 1}}
|
||||
|
||||
merged := viz.MergeAuthoritativeStats(prev, next, viz.MergeStatsOpts{})
|
||||
if merged.TotalAttacks != 100 || merged.ConnectedPeers != 5 || merged.BlockedIPs != 12 || merged.AttacksLastHour != 20 {
|
||||
t.Fatalf("non-absolute merge: %+v", merged)
|
||||
}
|
||||
if merged.ServiceBreakdown["SSH"] != 1 {
|
||||
t.Fatalf("breakdown replaced: %+v", merged.ServiceBreakdown)
|
||||
}
|
||||
|
||||
abs := viz.MergeAuthoritativeStats(prev, next, viz.MergeStatsOpts{Absolute: true})
|
||||
if abs.TotalAttacks != 90 || abs.ConnectedPeers != 4 || abs.AttacksLastHour != 15 {
|
||||
t.Fatalf("absolute merge: %+v", abs)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBumpPeerOnAttack(t *testing.T) {
|
||||
now := int64(1_000_000_000)
|
||||
net := &viz.Network{Peers: []viz.Peer{{ID: "p1", Status: "online"}}}
|
||||
net = viz.BumpPeerOnAttack(net, viz.AttackEvent{PeerID: "p1", Service: "SSH", Timestamp: now}, now)
|
||||
if net.Peers[0].AttacksLastHour != 1 || net.Peers[0].TopServices["SSH"] != 1 {
|
||||
t.Fatalf("peer bump: %+v", net.Peers[0])
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user