diff --git a/.github/workflows/docker.yml b/.github/workflows/docker.yml index 0930624..5f0af9f 100644 --- a/.github/workflows/docker.yml +++ b/.github/workflows/docker.yml @@ -4,6 +4,8 @@ on: push: tags: - "*" + branches: + - refactor/v2 permissions: write-all @@ -29,7 +31,8 @@ jobs: images: ghcr.io/inc4/gonka-exporter-go tags: | type=ref,event=tag - type=raw,value=latest + type=raw,value=latest,enable=${{ github.ref_type == 'tag' }} + type=raw,value=dev,enable=${{ github.ref_type == 'branch' }} - name: Build and push uses: docker/build-push-action@v5 diff --git a/.gitignore b/.gitignore index 1f0d526..8f516c3 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,6 @@ # Compiled binaries -exporter -gonka-exporter +/exporter +/gonka-exporter *.exe # Environment files diff --git a/cmd/exporter/main.go b/cmd/exporter/main.go index 8d73b3c..aced1e5 100644 --- a/cmd/exporter/main.go +++ b/cmd/exporter/main.go @@ -10,10 +10,12 @@ import ( "syscall" "time" + "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promhttp" "github.com/gonka/exporter/internal/collector" "github.com/gonka/exporter/internal/config" + "github.com/gonka/exporter/internal/fetcher" ) func main() { @@ -37,7 +39,7 @@ func main() { slog.Warn("PARTICIPANT_ADDRESS is not set — participant metrics will be skipped") } - c := collector.New(cfg) + c := collector.New(cfg, fetcher.NewHTTPFetcher(), prometheus.DefaultRegisterer) mux := http.NewServeMux() mux.Handle("/metrics", promhttp.Handler()) diff --git a/internal/collector/collector.go b/internal/collector/collector.go index ad0645d..8aea63f 100644 --- a/internal/collector/collector.go +++ b/internal/collector/collector.go @@ -4,10 +4,14 @@ import ( "fmt" "log/slog" "math" + "sort" "strconv" "strings" + "sync" "time" + "github.com/prometheus/client_golang/prometheus" + "github.com/gonka/exporter/internal/config" "github.com/gonka/exporter/internal/fetcher" "github.com/gonka/exporter/internal/metrics" @@ -31,20 +35,26 @@ var ( // Collector holds all mutable state for one collection cycle. type Collector struct { - cfg config.Config - st state.EpochState - history state.History - epochMax map[int64]*state.EpochMaxValues - epochNode state.EpochNodeState - estimatedReward float64 // last calculated estimated reward for current epoch + cfg config.Config + f fetcher.Fetcher + m *metrics.Metrics + st state.EpochState + history state.History + epochMax map[int64]*state.EpochMaxValues + epochNode state.EpochNodeState + estimatedReward float64 // last calculated estimated reward for current epoch } // New creates a Collector, restoring persisted state and history. -func New(cfg config.Config) *Collector { - st := state.LoadState(cfg.StateFile) +// reg is the Prometheus registerer to use; pass prometheus.DefaultRegisterer in production, +// prometheus.NewRegistry() in tests. +func New(cfg config.Config, f fetcher.Fetcher, reg prometheus.Registerer) *Collector { + st := state.LoadState(cfg.StateFile) history := state.LoadHistory(cfg.HistoryFile) c := &Collector{ cfg: cfg, + f: f, + m: metrics.NewMetrics(reg), st: st, history: history, epochMax: make(map[int64]*state.EpochMaxValues), @@ -77,49 +87,51 @@ func (c *Collector) collectChain() { } go func() { - h, t := fetcher.FetchMaxBlockHeightFromNodes(c.cfg.BlockHeightNodes) + h, t := c.f.FetchMaxBlockHeightFromNodes(c.cfg.BlockHeightNodes) if h > 0 { - metrics.BlockHeightMax.WithLabelValues(p).Set(float64(h)) + c.m.BlockHeightMax.WithLabelValues(p).Set(float64(h)) } if t != "" { ts := strings.TrimSuffix(t, "Z") if parsed, err := time.Parse("2006-01-02T15:04:05.999999999", ts); err == nil { - metrics.BlockTime.WithLabelValues(p).Set(float64(parsed.Unix())) + c.m.BlockTimeNetwork.WithLabelValues(p).Set(float64(parsed.Unix())) } } }() - status, err := fetcher.FetchTendermintStatus(c.cfg.NodeRPCURL) + status, err := c.f.FetchTendermintStatus(c.cfg.NodeRPCURL) if err != nil { slog.Warn("tendermint status", "err", err) return } si := status.Result.SyncInfo if h, err := strconv.ParseInt(si.LatestBlockHeight, 10, 64); err == nil { - metrics.BlockHeight.WithLabelValues(p).Set(float64(h)) + c.m.BlockHeight.WithLabelValues(p).Set(float64(h)) } if si.LatestBlockTime != "" { ts := strings.TrimSuffix(si.LatestBlockTime, "Z") if parsed, err := time.Parse("2006-01-02T15:04:05.999999999", ts); err == nil { - metrics.BlockTime.WithLabelValues(p).Set(float64(parsed.Unix())) + c.m.BlockTimeLocal.WithLabelValues(p).Set(float64(parsed.Unix())) } } catching := 0.0 if si.CatchingUp { catching = 1.0 } - metrics.CatchingUp.WithLabelValues(p).Set(catching) + c.m.CatchingUp.WithLabelValues(p).Set(catching) } // --- Network-wide participants --- func (c *Collector) collectNetworkParticipants() { - participants, err := fetcher.FetchNetworkParticipants(c.cfg.APIURL) + participants, err := c.f.FetchNetworkParticipants(c.cfg.APIURL) if err != nil { slog.Warn("network participants", "err", err) return } - metrics.NetTotalParticipantCount.Set(float64(len(participants))) + c.m.NetTotalParticipantCount.Set(float64(len(participants))) + c.m.NetParticipantWeight.Reset() + c.m.NetNodePocWeight.Reset() active := 0 for _, p := range participants { addr := p.Seed.Participant @@ -128,56 +140,56 @@ func (c *Collector) collectNetworkParticipants() { } active++ if p.Weight != nil { - metrics.NetParticipantWeight.WithLabelValues(addr).Set(*p.Weight) + c.m.NetParticipantWeight.WithLabelValues(addr).Set(*p.Weight) } for _, group := range p.MLNodes { for _, node := range group.MLNodes { if node.NodeID != "" && node.PocWeight != nil { - metrics.NetNodePocWeight.WithLabelValues(addr, node.NodeID).Set(*node.PocWeight) + c.m.NetNodePocWeight.WithLabelValues(addr, node.NodeID).Set(*node.PocWeight) } } } } - metrics.NetActiveParticipantCount.Set(float64(active)) + c.m.NetActiveParticipantCount.Set(float64(active)) } // --- Pricing and models --- func (c *Collector) collectPricingAndModels() { - pricing, err := fetcher.FetchPricing(c.cfg.APIURL) + pricing, err := c.f.FetchPricing(c.cfg.APIURL) if err != nil { slog.Warn("pricing", "err", err) } else { if pricing.UnitOfComputePrice != nil { - metrics.PricingUoC.Set(*pricing.UnitOfComputePrice) + c.m.PricingUoC.Set(*pricing.UnitOfComputePrice) } if pricing.DynamicPricingEnabled != nil { v := 0.0 if *pricing.DynamicPricingEnabled { v = 1.0 } - metrics.PricingDynamic.Set(v) + c.m.PricingDynamic.Set(v) } for _, m := range pricing.Models { if m.ID == "" { continue } if m.PricePerToken != nil { - metrics.ModelPrice.WithLabelValues(m.ID).Set(*m.PricePerToken) + c.m.ModelPrice.WithLabelValues(m.ID).Set(*m.PricePerToken) } if m.UnitsOfComputePerToken != nil { - metrics.ModelUnits.WithLabelValues(m.ID).Set(*m.UnitsOfComputePerToken) + c.m.ModelUnits.WithLabelValues(m.ID).Set(*m.UnitsOfComputePerToken) } if m.Utilization != nil { - metrics.ModelUtilization.WithLabelValues(m.ID).Set(*m.Utilization * 100) + c.m.ModelUtilization.WithLabelValues(m.ID).Set(*m.Utilization * 100) } if m.Capacity != nil { - metrics.ModelCapacity.WithLabelValues(m.ID).Set(float64(*m.Capacity)) + c.m.ModelCapacity.WithLabelValues(m.ID).Set(float64(*m.Capacity)) } } } - models, err := fetcher.FetchModels(c.cfg.APIURL) + models, err := c.f.FetchModels(c.cfg.APIURL) if err != nil { slog.Warn("models", "err", err) return @@ -187,13 +199,13 @@ func (c *Collector) collectPricingAndModels() { continue } if m.VRAM != nil { - metrics.ModelVRAM.WithLabelValues(m.ID).Set(*m.VRAM) + c.m.ModelVRAM.WithLabelValues(m.ID).Set(*m.VRAM) } if m.ThroughputPerNonce != nil { - metrics.ModelThroughput.WithLabelValues(m.ID).Set(*m.ThroughputPerNonce) + c.m.ModelThroughput.WithLabelValues(m.ID).Set(*m.ThroughputPerNonce) } if m.ValidationThreshold != nil { - metrics.ModelValThresh.WithLabelValues(m.ID).Set( + c.m.ModelValThresh.WithLabelValues(m.ID).Set( m.ValidationThreshold.Value * math.Pow10(m.ValidationThreshold.Exponent), ) } @@ -209,182 +221,215 @@ func (c *Collector) collectParticipant() { } now := float64(time.Now().UnixNano()) / 1e9 + prevEpochStartTime := c.st.EpochStartTime + chainEpoch := c.fetchAndSetChainEpoch(addr) - // 1. Chain epoch - chainEpoch, err := fetcher.FetchCurrentEpoch(c.cfg.NodeRESTURL) - if err != nil { - slog.Warn("fetch current epoch", "err", err) - } else { - metrics.ChainEpoch.WithLabelValues(addr).Set(float64(chainEpoch)) - } + epochInfo := c.fetchAndSetEpochInfo(chainEpoch) - // 2. Save prev start time before updating - prevEpochStartTime := c.st.EpochStartTime + c.fetchAndSetGroupData(addr, chainEpoch) - // 3. Epoch info — real start time from blockchain - var epochInfo *fetcher.EpochInfo - epochInfo, err = fetcher.FetchEpochInfo(c.cfg.NodeRESTURL) - if err != nil { - slog.Warn("fetch epoch info", "err", err) - } else { - c.st.EpochLength = epochInfo.EpochLength - if epochInfo.PocStartBlockHeight != c.st.PocStartBlockHeight { - if t, err := fetcher.FetchBlockTimeAtHeight(c.cfg.NodeRPCURL, epochInfo.PocStartBlockHeight); err != nil { - slog.Warn("fetch block time", "height", epochInfo.PocStartBlockHeight, "err", err) - } else { - c.st.PocStartBlockHeight = epochInfo.PocStartBlockHeight - c.st.EpochStartTime = t - slog.Info("epoch start block time", - "epoch", chainEpoch, - "block", epochInfo.PocStartBlockHeight, - "time", time.Unix(int64(t), 0).UTC().Format(time.RFC3339)) - } - } + c.fetchAndSetBLS(addr, chainEpoch) + + stats, ok := c.fetchAndSetParticipantStats(addr) + if !ok { + return } - // 4. Epoch group data — total network weight + estimated reward + reputation + network inference count - groupData, err := fetcher.FetchEpochGroupData(c.cfg.NodeRESTURL) - if err != nil { - slog.Warn("fetch epoch group data", "err", err) - } else { - if groupData.NumberOfRequests > 0 { - metrics.NetEpochInferenceCount.WithLabelValues(addr).Set(float64(groupData.NumberOfRequests)) + wallet, walletOK := c.fetchAndSetWallet(addr) + + c.updateEpochMaxAndGauges(addr, chainEpoch, stats, wallet, walletOK, epochInfo, now) + + c.detectAndRecordEpochBoundary(chainEpoch, prevEpochStartTime, wallet, walletOK) +} + +// --- collectParticipant helpers --- + +func (c *Collector) detectAndRecordEpochBoundary(chainEpoch int64, prevEpochStartTime float64, wallet float64, walletOK bool) { + prevEpoch := c.st.ChainEpoch + if c.st.Valid && prevEpoch > 0 && chainEpoch > prevEpoch { + if chainEpoch > prevEpoch+1 { + slog.Warn("skipped epochs", "from", prevEpoch, "to", chainEpoch) } - if groupData.TotalWeight > 0 { - // emission(N) = 323000 × exp(-0.000475 × (N-1)) - emission := 323000.0 * math.Exp(-0.000475*float64(groupData.EpochIndex-1)) - rewardPerWeight := emission / float64(groupData.TotalWeight) - metrics.NetTotalWeight.WithLabelValues(addr).Set(float64(groupData.TotalWeight)) - metrics.NetRewardPerWeight.WithLabelValues(addr).Set(rewardPerWeight) - for _, vw := range groupData.ValidationWeights { - if vw.MemberAddress == addr { - myWeight, _ := strconv.ParseInt(vw.Weight, 10, 64) - c.estimatedReward = float64(myWeight) * rewardPerWeight - metrics.ParticipantReputation.WithLabelValues(addr).Set(float64(vw.Reputation)) - if chainEpoch > 0 { - ce := strconv.FormatInt(chainEpoch, 10) - metrics.EpochEstimated.WithLabelValues(addr, ce).Set(c.estimatedReward) - } - break - } - } + c.recordSnapshot(prevEpoch, prevEpochStartTime, c.st.EpochStartTime, wallet, walletOK) + if walletOK { + c.st.WalletBalanceAtEpochStart = wallet + } + c.pruneEpochMax() + } else if !c.st.Valid { + if walletOK { + c.st.WalletBalanceAtEpochStart = wallet } } + changed := c.st.ChainEpoch != chainEpoch + c.st.ChainEpoch = chainEpoch + c.st.Valid = true + if changed { + state.SaveState(c.cfg.StateFile, c.st) + } +} - // 4b. BLS DKG phase - if chainEpoch > 0 { - bls, blsErr := fetcher.FetchBLSEpoch(c.cfg.APIURL, chainEpoch) - if blsErr != nil { - slog.Warn("fetch bls epoch", "err", blsErr) - } else { - metrics.BLSDKGPhase.WithLabelValues(addr).Set(float64(bls.DKGPhase)) - if bls.DealingPhaseDeadlineBlock > 0 { - metrics.BLSDealingDeadline.WithLabelValues(addr).Set(float64(bls.DealingPhaseDeadlineBlock)) - } - if bls.VerifyingPhaseDeadlineBlock > 0 { - metrics.BLSVerifyingDeadline.WithLabelValues(addr).Set(float64(bls.VerifyingPhaseDeadlineBlock)) - } +func (c *Collector) updateEpochMaxAndGauges(addr string, chainEpoch int64, stats *fetcher.ParticipantStats, wallet float64, walletOK bool, epochInfo *fetcher.EpochInfo, now float64) { + if chainEpoch == 0 { + return + } + if _, ok := c.epochMax[chainEpoch]; !ok { + c.epochMax[chainEpoch] = &state.EpochMaxValues{} + } + em := c.epochMax[chainEpoch] + em.InferenceCount = max(em.InferenceCount, stats.InferenceCount) + em.MissedRequests = max(em.MissedRequests, stats.MissedRequests) + em.EarnedCoins = max(em.EarnedCoins, stats.EarnedCoins) + em.ValidatedInferences = max(em.ValidatedInferences, stats.ValidatedInferences) + em.InvalidatedInferences = max(em.InvalidatedInferences, stats.InvalidatedInferences) + em.CoinBalance = max(em.CoinBalance, stats.CoinBalance) + em.EpochsCompleted = max(em.EpochsCompleted, stats.EpochsCompleted) + + ce := strconv.FormatInt(chainEpoch, 10) + c.m.EpochInferences.WithLabelValues(addr, ce).Set(float64(em.InferenceCount)) + c.m.EpochMissed.WithLabelValues(addr, ce).Set(float64(em.MissedRequests)) + c.m.EpochEarnedCoins.WithLabelValues(addr, ce).Set(float64(em.EarnedCoins)) + c.m.EpochValidated.WithLabelValues(addr, ce).Set(float64(em.ValidatedInferences)) + c.m.EpochInvalidated.WithLabelValues(addr, ce).Set(float64(em.InvalidatedInferences)) + c.m.EpochCoinBalance.WithLabelValues(addr, ce).Set(float64(em.CoinBalance)) + c.m.EpochDone.WithLabelValues(addr, ce).Set(float64(em.EpochsCompleted)) + c.m.EpochMissRate.WithLabelValues(addr, ce).Set(state.MissRate(em.InferenceCount, em.MissedRequests)) + + if walletOK { + c.m.EpochEarnedGNK.WithLabelValues(addr, ce).Set(wallet - c.st.WalletBalanceAtEpochStart) + } + if c.st.EpochStartTime > 0 { + c.m.EpochStartTime.WithLabelValues(addr, ce).Set(c.st.EpochStartTime) + } + if epochInfo != nil && c.st.EpochStartTime > 0 && c.st.PocStartBlockHeight > 0 && c.st.EpochLength > 0 { + blocksElapsed := epochInfo.BlockHeight - c.st.PocStartBlockHeight + timeElapsed := now - c.st.EpochStartTime + if blocksElapsed > 0 && timeElapsed > 0 { + avgBlockTime := timeElapsed / float64(blocksElapsed) + blocksRemaining := c.st.EpochLength - blocksElapsed + c.m.EpochEndTime.WithLabelValues(addr, ce).Set(now + float64(blocksRemaining)*avgBlockTime) } } +} - // 5. Participant stats - stats, err := fetcher.FetchParticipantStats(c.cfg.NodeRESTURL, addr) +func (c *Collector) fetchAndSetWallet(addr string) (float64, bool) { + wallet, err := c.f.FetchWalletBalance(c.cfg.NodeRESTURL, addr) if err != nil { - slog.Warn("fetch participant stats", "err", err) - return + slog.Warn("fetch wallet balance", "err", err) + return 0, false + } + c.m.ParticipantWallet.WithLabelValues(addr).Set(wallet) + if c.st.WalletBalanceAtEpochStart == 0 { + c.st.WalletBalanceAtEpochStart = wallet } - metrics.ParticipantEpochsDone.WithLabelValues(addr).Set(float64(stats.EpochsCompleted)) - metrics.ParticipantCoinBalance.WithLabelValues(addr).Set(float64(stats.CoinBalance)) - metrics.ParticipantInferences.WithLabelValues(addr).Set(float64(stats.InferenceCount)) - metrics.ParticipantMissed.WithLabelValues(addr).Set(float64(stats.MissedRequests)) - metrics.ParticipantEarnedCoins.WithLabelValues(addr).Set(float64(stats.EarnedCoins)) - metrics.ParticipantValidated.WithLabelValues(addr).Set(float64(stats.ValidatedInferences)) - metrics.ParticipantInvalidated.WithLabelValues(addr).Set(float64(stats.InvalidatedInferences)) - // Extended participant health + return wallet, true +} + +func (c *Collector) fetchAndSetParticipantStats(addr string) (*fetcher.ParticipantStats, bool) { + stats, err := c.f.FetchParticipantStats(c.cfg.NodeRESTURL, addr) + if err != nil { + slog.Warn("fetch participant stats", "err", err) + return nil, false + } + c.m.ParticipantEpochsDone.WithLabelValues(addr).Set(float64(stats.EpochsCompleted)) + c.m.ParticipantCoinBalance.WithLabelValues(addr).Set(float64(stats.CoinBalance)) + c.m.ParticipantInferences.WithLabelValues(addr).Set(float64(stats.InferenceCount)) + c.m.ParticipantMissed.WithLabelValues(addr).Set(float64(stats.MissedRequests)) + c.m.ParticipantEarnedCoins.WithLabelValues(addr).Set(float64(stats.EarnedCoins)) + c.m.ParticipantValidated.WithLabelValues(addr).Set(float64(stats.ValidatedInferences)) + c.m.ParticipantInvalidated.WithLabelValues(addr).Set(float64(stats.InvalidatedInferences)) if sv, ok := participantStatusMap[stats.Status]; ok { - metrics.ParticipantStatus.WithLabelValues(addr).Set(sv) + c.m.ParticipantStatus.WithLabelValues(addr).Set(sv) } - metrics.ParticipantConsecutiveInv.WithLabelValues(addr).Set(float64(stats.ConsecutiveInvalidInferences)) - metrics.ParticipantBurnedCoins.WithLabelValues(addr).Set(float64(stats.BurnedCoins)) - metrics.ParticipantRewardedCoins.WithLabelValues(addr).Set(float64(stats.RewardedCoins)) + c.m.ParticipantConsecutiveInv.WithLabelValues(addr).Set(float64(stats.ConsecutiveInvalidInferences)) + c.m.ParticipantBurnedCoins.WithLabelValues(addr).Set(float64(stats.BurnedCoins)) + c.m.ParticipantRewardedCoins.WithLabelValues(addr).Set(float64(stats.RewardedCoins)) + return stats, true +} - // 6. Wallet balance - wallet, walletErr := fetcher.FetchWalletBalance(c.cfg.NodeRESTURL, addr) - if walletErr != nil { - slog.Warn("fetch wallet balance", "err", walletErr) - } else { - metrics.ParticipantWallet.WithLabelValues(addr).Set(wallet) - if c.st.WalletBalanceAtEpochStart == 0 { - c.st.WalletBalanceAtEpochStart = wallet - } +func (c *Collector) fetchAndSetBLS(addr string, chainEpoch int64) { + if chainEpoch == 0 { + return } + bls, err := c.f.FetchBLSEpoch(c.cfg.APIURL, chainEpoch) + if err != nil { + slog.Warn("fetch bls epoch", "err", err) + return + } + c.m.BLSDKGPhase.WithLabelValues(addr).Set(float64(bls.DKGPhase)) + if bls.DealingPhaseDeadlineBlock > 0 { + c.m.BLSDealingDeadline.WithLabelValues(addr).Set(float64(bls.DealingPhaseDeadlineBlock)) + } + if bls.VerifyingPhaseDeadlineBlock > 0 { + c.m.BLSVerifyingDeadline.WithLabelValues(addr).Set(float64(bls.VerifyingPhaseDeadlineBlock)) + } +} - // 7. Per-epoch max tracking + live gauges - if chainEpoch > 0 { - if _, ok := c.epochMax[chainEpoch]; !ok { - c.epochMax[chainEpoch] = &state.EpochMaxValues{} - } - em := c.epochMax[chainEpoch] - em.InferenceCount = max64(em.InferenceCount, stats.InferenceCount) - em.MissedRequests = max64(em.MissedRequests, stats.MissedRequests) - em.EarnedCoins = max64(em.EarnedCoins, stats.EarnedCoins) - em.ValidatedInferences = max64(em.ValidatedInferences, stats.ValidatedInferences) - em.InvalidatedInferences = max64(em.InvalidatedInferences, stats.InvalidatedInferences) - em.CoinBalance = max64(em.CoinBalance, stats.CoinBalance) - em.EpochsCompleted = max64(em.EpochsCompleted, stats.EpochsCompleted) - - ce := strconv.FormatInt(chainEpoch, 10) - metrics.EpochInferences.WithLabelValues(addr, ce).Set(float64(em.InferenceCount)) - metrics.EpochMissed.WithLabelValues(addr, ce).Set(float64(em.MissedRequests)) - metrics.EpochEarnedCoins.WithLabelValues(addr, ce).Set(float64(em.EarnedCoins)) - metrics.EpochValidated.WithLabelValues(addr, ce).Set(float64(em.ValidatedInferences)) - metrics.EpochInvalidated.WithLabelValues(addr, ce).Set(float64(em.InvalidatedInferences)) - metrics.EpochCoinBalance.WithLabelValues(addr, ce).Set(float64(em.CoinBalance)) - metrics.EpochDone.WithLabelValues(addr, ce).Set(float64(em.EpochsCompleted)) - metrics.EpochMissRate.WithLabelValues(addr, ce).Set(state.MissRate(em.InferenceCount, em.MissedRequests)) - - if walletErr == nil { - metrics.EpochEarnedGNK.WithLabelValues(addr, ce).Set(wallet - c.st.WalletBalanceAtEpochStart) - } - if c.st.EpochStartTime > 0 { - metrics.EpochStartTime.WithLabelValues(addr, ce).Set(c.st.EpochStartTime) - } - if epochInfo != nil && c.st.EpochStartTime > 0 && c.st.PocStartBlockHeight > 0 && c.st.EpochLength > 0 { - blocksElapsed := epochInfo.BlockHeight - c.st.PocStartBlockHeight - timeElapsed := now - c.st.EpochStartTime - if blocksElapsed > 0 && timeElapsed > 0 { - avgBlockTime := timeElapsed / float64(blocksElapsed) - blocksRemaining := c.st.EpochLength - blocksElapsed - metrics.EpochEndTime.WithLabelValues(addr, ce).Set(now + float64(blocksRemaining)*avgBlockTime) +func (c *Collector) fetchAndSetGroupData(addr string, chainEpoch int64) { + groupData, err := c.f.FetchEpochGroupData(c.cfg.NodeRESTURL) + if err != nil { + slog.Warn("fetch epoch group data", "err", err) + return + } + if groupData.NumberOfRequests > 0 { + c.m.NetEpochInferenceCount.WithLabelValues(addr).Set(float64(groupData.NumberOfRequests)) + } + if groupData.TotalWeight > 0 { + // emission(N) = 323000 × exp(-0.000475 × (N-1)) + emission := 323000.0 * math.Exp(-0.000475*float64(groupData.EpochIndex-1)) + rewardPerWeight := emission / float64(groupData.TotalWeight) + c.m.NetTotalWeight.WithLabelValues(addr).Set(float64(groupData.TotalWeight)) + c.m.NetRewardPerWeight.WithLabelValues(addr).Set(rewardPerWeight) + for _, vw := range groupData.ValidationWeights { + if vw.MemberAddress == addr { + myWeight, _ := strconv.ParseInt(vw.Weight, 10, 64) + c.estimatedReward = float64(myWeight) * rewardPerWeight + c.m.ParticipantReputation.WithLabelValues(addr).Set(float64(vw.Reputation)) + if chainEpoch > 0 { + ce := strconv.FormatInt(chainEpoch, 10) + c.m.EpochEstimated.WithLabelValues(addr, ce).Set(c.estimatedReward) + } + break } } } +} - // 8. Epoch boundary detection - prevEpoch := c.st.ChainEpoch - if c.st.Valid && prevEpoch > 0 && chainEpoch > prevEpoch { - if chainEpoch > prevEpoch+1 { - slog.Warn("skipped epochs", "from", prevEpoch, "to", chainEpoch) +func (c *Collector) fetchAndSetEpochInfo(chainEpoch int64) *fetcher.EpochInfo { + epochInfo, err := c.f.FetchEpochInfo(c.cfg.NodeRESTURL) + if err != nil { + slog.Warn("fetch epoch info", "err", err) + return nil + } + c.st.EpochLength = epochInfo.EpochLength + if epochInfo.PocStartBlockHeight != c.st.PocStartBlockHeight { + t, err := c.f.FetchBlockTimeAtHeight(c.cfg.NodeRPCURL, epochInfo.PocStartBlockHeight) + if err != nil { + slog.Warn("fetch block time", "height", epochInfo.PocStartBlockHeight, "err", err) + } else { + c.st.PocStartBlockHeight = epochInfo.PocStartBlockHeight + c.st.EpochStartTime = t + slog.Info("epoch start block time", + "epoch", chainEpoch, + "block", epochInfo.PocStartBlockHeight, + "time", time.Unix(int64(t), 0).UTC().Format(time.RFC3339)) } - c.recordSnapshot(prevEpoch, prevEpochStartTime, c.st.EpochStartTime, wallet, walletErr == nil) - c.st.WalletBalanceAtEpochStart = wallet - c.pruneEpochMax() - } else if !c.st.Valid { - c.st.WalletBalanceAtEpochStart = wallet } + return epochInfo +} - // 9. Persist state on epoch change - changed := c.st.ChainEpoch != chainEpoch - c.st.ChainEpoch = chainEpoch - c.st.Valid = true - if changed { - state.SaveState(c.cfg.StateFile, c.st) +func (c *Collector) fetchAndSetChainEpoch(addr string) int64 { + chainEpoch, err := c.f.FetchCurrentEpoch(c.cfg.NodeRESTURL) + if err != nil { + slog.Warn("fetch current epoch", "err", err) + return 0 } + c.m.ChainEpoch.WithLabelValues(addr).Set(float64(chainEpoch)) + return chainEpoch } func (c *Collector) recordSnapshot(epoch int64, startTime, endTime float64, wallet float64, walletOK bool) { addr := c.cfg.Participant - em := c.epochMax[epoch] + em := c.epochMax[epoch] if em == nil { em = &state.EpochMaxValues{} } @@ -416,12 +461,12 @@ func (c *Collector) recordSnapshot(epoch int64, startTime, endTime float64, wall } // On-chain performance summary - perf, err := fetcher.FetchEpochPerfSummary(c.cfg.NodeRESTURL, addr, epoch) + perf, err := c.f.FetchEpochPerfSummary(c.cfg.NodeRESTURL, addr, epoch) if err != nil { slog.Warn("epoch performance summary", "epoch", epoch, "err", err) } else { snap.RewardedGNK = &perf.RewardedGNK - snap.Claimed = &perf.Claimed + snap.Claimed = &perf.Claimed } epochStr := fmt.Sprintf("%d", epoch) @@ -432,32 +477,32 @@ func (c *Collector) recordSnapshot(epoch int64, startTime, endTime float64, wall state.SaveHistory(c.cfg.HistoryFile, c.history, c.cfg.MaxHistory) ce := strconv.FormatInt(epoch, 10) - metrics.EpochInferences.WithLabelValues(addr, ce).Set(float64(snap.InferenceCount)) - metrics.EpochMissed.WithLabelValues(addr, ce).Set(float64(snap.MissedRequests)) - metrics.EpochEarnedCoins.WithLabelValues(addr, ce).Set(float64(snap.EarnedCoins)) - metrics.EpochValidated.WithLabelValues(addr, ce).Set(float64(snap.ValidatedInferences)) - metrics.EpochInvalidated.WithLabelValues(addr, ce).Set(float64(snap.InvalidatedInferences)) - metrics.EpochCoinBalance.WithLabelValues(addr, ce).Set(float64(snap.CoinBalance)) - metrics.EpochDone.WithLabelValues(addr, ce).Set(float64(snap.EpochsCompleted)) - metrics.EpochMissRate.WithLabelValues(addr, ce).Set(snap.MissRatePercent) - metrics.EpochTimeslot.WithLabelValues(addr, ce).Set(float64(snap.TimeslotAssigned)) + c.m.EpochInferences.WithLabelValues(addr, ce).Set(float64(snap.InferenceCount)) + c.m.EpochMissed.WithLabelValues(addr, ce).Set(float64(snap.MissedRequests)) + c.m.EpochEarnedCoins.WithLabelValues(addr, ce).Set(float64(snap.EarnedCoins)) + c.m.EpochValidated.WithLabelValues(addr, ce).Set(float64(snap.ValidatedInferences)) + c.m.EpochInvalidated.WithLabelValues(addr, ce).Set(float64(snap.InvalidatedInferences)) + c.m.EpochCoinBalance.WithLabelValues(addr, ce).Set(float64(snap.CoinBalance)) + c.m.EpochDone.WithLabelValues(addr, ce).Set(float64(snap.EpochsCompleted)) + c.m.EpochMissRate.WithLabelValues(addr, ce).Set(snap.MissRatePercent) + c.m.EpochTimeslot.WithLabelValues(addr, ce).Set(float64(snap.TimeslotAssigned)) for nodeID, pw := range snap.PocWeights { - metrics.EpochPocWeight.WithLabelValues(addr, ce, nodeID).Set(float64(pw)) + c.m.EpochPocWeight.WithLabelValues(addr, ce, nodeID).Set(float64(pw)) } - metrics.EpochStartTime.WithLabelValues(addr, ce).Set(snap.StartTime) - metrics.EpochEndTime.WithLabelValues(addr, ce).Set(snap.EndTime) - metrics.EpochDuration.WithLabelValues(addr, ce).Set(snap.DurationSeconds) + c.m.EpochStartTime.WithLabelValues(addr, ce).Set(snap.StartTime) + c.m.EpochEndTime.WithLabelValues(addr, ce).Set(snap.EndTime) + c.m.EpochDuration.WithLabelValues(addr, ce).Set(snap.DurationSeconds) if snap.EarnedGNK != nil { - metrics.EpochEarnedGNK.WithLabelValues(addr, ce).Set(*snap.EarnedGNK) + c.m.EpochEarnedGNK.WithLabelValues(addr, ce).Set(*snap.EarnedGNK) } if snap.RewardedGNK != nil { - metrics.EpochRewardedGNK.WithLabelValues(addr, ce).Set(*snap.RewardedGNK) + c.m.EpochRewardedGNK.WithLabelValues(addr, ce).Set(*snap.RewardedGNK) } if snap.EstimatedGNK != nil { - metrics.EpochEstimated.WithLabelValues(addr, ce).Set(*snap.EstimatedGNK) + c.m.EpochEstimated.WithLabelValues(addr, ce).Set(*snap.EstimatedGNK) } if snap.Claimed != nil { - metrics.EpochClaimed.WithLabelValues(addr, ce).Set(float64(*snap.Claimed)) + c.m.EpochClaimed.WithLabelValues(addr, ce).Set(float64(*snap.Claimed)) } slog.Info("epoch snapshot saved", @@ -483,112 +528,153 @@ func (c *Collector) collectNodes() { addr = "unknown" } - nodes, err := fetcher.FetchNodes(c.cfg.AdminAPIURL) + nodes, err := c.f.FetchNodes(c.cfg.AdminAPIURL) if err != nil { slog.Warn("fetch nodes", "err", err) return } + c.m.NodeHardwareInfo.Reset() + c.m.NodeStatus.Reset() + c.m.NodeIntended.Reset() + c.m.PocCurrent.Reset() + c.m.PocIntended.Reset() + c.m.NodePocWeight.Reset() + c.m.NodeTimeslot.Reset() + c.m.NodeGPUCount.Reset() + c.m.NodeGPUUtil.Reset() + c.m.NodeGPUDeviceUtil.Reset() + c.m.NodeGPUDeviceTemp.Reset() + c.m.NodeGPUDeviceMemTotal.Reset() + c.m.NodeGPUDeviceMemFree.Reset() + c.m.NodeGPUDeviceMemUsed.Reset() + c.m.NodeGPUDeviceAvail.Reset() + c.m.NodeServiceState.Reset() + c.m.NodeDiskAvailableGB.Reset() + c.m.NodeGPUDriverInfo.Reset() + c.m.NodeManagerRunning.Reset() + c.m.NodeManagerHealthy.Reset() + + var wg sync.WaitGroup + var mu sync.Mutex for _, entry := range nodes { - ni := entry.Node - nodeID := ni.ID - host := ni.Host - st := entry.State + wg.Add(1) + go func(entry fetcher.NodeEntry) { + defer wg.Done() + c.collectOneNode(addr, entry, &mu) + }(entry) + } + wg.Wait() +} - for _, hw := range ni.Hardware { - metrics.NodeHardwareInfo.WithLabelValues(addr, nodeID, host, hw.Type, strconv.Itoa(hw.Count)).Set(1) - } +// collectOneNode collects all metrics for a single node. +// mu protects writes to c.epochNode which is shared across parallel node goroutines. +func (c *Collector) collectOneNode(addr string, entry fetcher.NodeEntry, mu *sync.Mutex) { + ni := entry.Node + nodeID := ni.ID + host := ni.Host + st := entry.State - metrics.NodeStatus.WithLabelValues(addr, nodeID, host).Set(nodeStatusVal(st.CurrentStatus)) - metrics.NodeIntended.WithLabelValues(addr, nodeID, host).Set(nodeStatusVal(st.IntendedStatus)) - metrics.PocCurrent.WithLabelValues(addr, nodeID, host).Set(pocStatusVal(st.PocCurrentStatus)) - metrics.PocIntended.WithLabelValues(addr, nodeID, host).Set(pocStatusVal(st.PocIntendedStatus)) - - for model, md := range st.EpochMLNodes { - if md.PocWeight != nil { - pw := *md.PocWeight - metrics.NodePocWeight.WithLabelValues(addr, nodeID, host, model).Set(float64(pw)) - if pw > c.epochNode.PocWeights[nodeID] { - c.epochNode.PocWeights[nodeID] = pw - } - if c.st.ChainEpoch > 0 { - ce := strconv.FormatInt(c.st.ChainEpoch, 10) - metrics.EpochPocWeight.WithLabelValues(addr, ce, nodeID).Set(float64(c.epochNode.PocWeights[nodeID])) - } + for _, hw := range ni.Hardware { + c.m.NodeHardwareInfo.WithLabelValues(addr, nodeID, host, hw.Type, strconv.Itoa(hw.Count)).Set(1) + } + + c.m.NodeStatus.WithLabelValues(addr, nodeID, host).Set(nodeStatusVal(st.CurrentStatus)) + c.m.NodeIntended.WithLabelValues(addr, nodeID, host).Set(nodeStatusVal(st.IntendedStatus)) + c.m.PocCurrent.WithLabelValues(addr, nodeID, host).Set(pocStatusVal(st.PocCurrentStatus)) + c.m.PocIntended.WithLabelValues(addr, nodeID, host).Set(pocStatusVal(st.PocIntendedStatus)) + + for model, md := range st.EpochMLNodes { + if md.PocWeight != nil { + pw := *md.PocWeight + c.m.NodePocWeight.WithLabelValues(addr, nodeID, host, model).Set(float64(pw)) + + mu.Lock() + if pw > c.epochNode.PocWeights[nodeID] { + c.epochNode.PocWeights[nodeID] = pw } - if len(md.TimeslotAllocation) > 0 && md.TimeslotAllocation[0] { - metrics.NodeTimeslot.WithLabelValues(addr, nodeID, host, model).Set(1) - c.epochNode.TimeslotAssigned = 1 - if c.st.ChainEpoch > 0 { - ce := strconv.FormatInt(c.st.ChainEpoch, 10) - metrics.EpochTimeslot.WithLabelValues(addr, ce).Set(1) - } - } else { - metrics.NodeTimeslot.WithLabelValues(addr, nodeID, host, model).Set(0) + curMax := c.epochNode.PocWeights[nodeID] + mu.Unlock() + + if c.st.ChainEpoch > 0 { + ce := strconv.FormatInt(c.st.ChainEpoch, 10) + c.m.EpochPocWeight.WithLabelValues(addr, ce, nodeID).Set(float64(curMax)) } } - - if ni.PocPort > 0 && host != "" { - gpu := fetcher.FetchGPUStats(host, ni.PocPort) - metrics.NodeGPUCount.WithLabelValues(addr, nodeID, host).Set(float64(gpu.Count)) - metrics.NodeGPUUtil.WithLabelValues(addr, nodeID, host).Set(gpu.AvgUtil) - for _, dev := range gpu.Devices { - di := strconv.Itoa(dev.Index) - metrics.NodeGPUDeviceUtil.WithLabelValues(addr, nodeID, host, di).Set(dev.UtilizationPercent) - avail := 0.0 - if dev.IsAvailable { - avail = 1.0 - } - metrics.NodeGPUDeviceAvail.WithLabelValues(addr, nodeID, host, di).Set(avail) - if dev.TemperatureC != nil { - metrics.NodeGPUDeviceTemp.WithLabelValues(addr, nodeID, host, di).Set(*dev.TemperatureC) - } - if dev.TotalMemoryMB != nil { - metrics.NodeGPUDeviceMemTotal.WithLabelValues(addr, nodeID, host, di).Set(float64(*dev.TotalMemoryMB)) - } - if dev.FreeMemoryMB != nil { - metrics.NodeGPUDeviceMemFree.WithLabelValues(addr, nodeID, host, di).Set(float64(*dev.FreeMemoryMB)) - } - if dev.UsedMemoryMB != nil { - metrics.NodeGPUDeviceMemUsed.WithLabelValues(addr, nodeID, host, di).Set(float64(*dev.UsedMemoryMB)) - } + if len(md.TimeslotAllocation) > 0 && md.TimeslotAllocation[0] { + c.m.NodeTimeslot.WithLabelValues(addr, nodeID, host, model).Set(1) + mu.Lock() + c.epochNode.TimeslotAssigned = 1 + mu.Unlock() + if c.st.ChainEpoch > 0 { + ce := strconv.FormatInt(c.st.ChainEpoch, 10) + c.m.EpochTimeslot.WithLabelValues(addr, ce).Set(1) } + } else { + c.m.NodeTimeslot.WithLabelValues(addr, nodeID, host, model).Set(0) + } + } - // ML node service state - if svcState, svcErr := fetcher.FetchMLNodeState(host, ni.PocPort); svcErr != nil { - slog.Debug("ml node state unavailable", "node", nodeID, "err", svcErr) - } else { - sv := mlNodeServiceStateMap[svcState] - metrics.NodeServiceState.WithLabelValues(addr, nodeID, host).Set(sv) + if ni.PocPort > 0 && host != "" { + gpu := c.f.FetchGPUStats(host, ni.PocPort) + c.m.NodeGPUCount.WithLabelValues(addr, nodeID, host).Set(float64(gpu.Count)) + c.m.NodeGPUUtil.WithLabelValues(addr, nodeID, host).Set(gpu.AvgUtil) + for _, dev := range gpu.Devices { + di := strconv.Itoa(dev.Index) + c.m.NodeGPUDeviceUtil.WithLabelValues(addr, nodeID, host, di).Set(dev.UtilizationPercent) + avail := 0.0 + if dev.IsAvailable { + avail = 1.0 } - - // ML node disk space - if diskGB, diskErr := fetcher.FetchMLNodeDiskSpaceGB(host, ni.PocPort); diskErr != nil { - slog.Debug("ml node disk space unavailable", "node", nodeID, "err", diskErr) - } else { - metrics.NodeDiskAvailableGB.WithLabelValues(addr, nodeID, host).Set(diskGB) + c.m.NodeGPUDeviceAvail.WithLabelValues(addr, nodeID, host, di).Set(avail) + if dev.TemperatureC != nil { + c.m.NodeGPUDeviceTemp.WithLabelValues(addr, nodeID, host, di).Set(*dev.TemperatureC) } - - // GPU driver info - if drvInfo, drvErr := fetcher.FetchGPUDriverInfo(host, ni.PocPort); drvErr != nil { - slog.Debug("gpu driver info unavailable", "node", nodeID, "err", drvErr) - } else { - metrics.NodeGPUDriverInfo.WithLabelValues(addr, nodeID, host, drvInfo.DriverVersion, drvInfo.CudaDriverVersion).Set(1) + if dev.TotalMemoryMB != nil { + c.m.NodeGPUDeviceMemTotal.WithLabelValues(addr, nodeID, host, di).Set(float64(*dev.TotalMemoryMB)) } - - // ML node manager health - if health, healthErr := fetcher.FetchMLNodeHealth(host, ni.PocPort); healthErr != nil { - slog.Debug("ml node health unavailable", "node", nodeID, "err", healthErr) - } else { - setManagerMetric(addr, nodeID, host, "pow", health.ManagerPow) - setManagerMetric(addr, nodeID, host, "inference", health.ManagerInference) - setManagerMetric(addr, nodeID, host, "train", health.ManagerTrain) + if dev.FreeMemoryMB != nil { + c.m.NodeGPUDeviceMemFree.WithLabelValues(addr, nodeID, host, di).Set(float64(*dev.FreeMemoryMB)) } + if dev.UsedMemoryMB != nil { + c.m.NodeGPUDeviceMemUsed.WithLabelValues(addr, nodeID, host, di).Set(float64(*dev.UsedMemoryMB)) + } + } + + // ML node service state + if svcState, svcErr := c.f.FetchMLNodeState(host, ni.PocPort); svcErr != nil { + slog.Debug("ml node state unavailable", "node", nodeID, "err", svcErr) + } else { + sv := mlNodeServiceStateMap[svcState] + c.m.NodeServiceState.WithLabelValues(addr, nodeID, host).Set(sv) + } + + // ML node disk space + if diskGB, diskErr := c.f.FetchMLNodeDiskSpaceGB(host, ni.PocPort); diskErr != nil { + slog.Debug("ml node disk space unavailable", "node", nodeID, "err", diskErr) + } else { + c.m.NodeDiskAvailableGB.WithLabelValues(addr, nodeID, host).Set(diskGB) + } + + // GPU driver info + if drvInfo, drvErr := c.f.FetchGPUDriverInfo(host, ni.PocPort); drvErr != nil { + slog.Debug("gpu driver info unavailable", "node", nodeID, "err", drvErr) + } else { + c.m.NodeGPUDriverInfo.WithLabelValues(addr, nodeID, host, drvInfo.DriverVersion, drvInfo.CudaDriverVersion).Set(1) + } + + // ML node manager health + if health, healthErr := c.f.FetchMLNodeHealth(host, ni.PocPort); healthErr != nil { + slog.Debug("ml node health unavailable", "node", nodeID, "err", healthErr) + } else { + c.setManagerMetric(addr, nodeID, host, "pow", health.ManagerPow) + c.setManagerMetric(addr, nodeID, host, "inference", health.ManagerInference) + c.setManagerMetric(addr, nodeID, host, "train", health.ManagerTrain) } } } -func setManagerMetric(addr, nodeID, host, manager string, s fetcher.MLNodeManagerStatus) { +func (c *Collector) setManagerMetric(addr, nodeID, host, manager string, s fetcher.MLNodeManagerStatus) { running := 0.0 if s.Running { running = 1.0 @@ -597,8 +683,8 @@ func setManagerMetric(addr, nodeID, host, manager string, s fetcher.MLNodeManage if s.Healthy { healthy = 1.0 } - metrics.NodeManagerRunning.WithLabelValues(addr, nodeID, host, manager).Set(running) - metrics.NodeManagerHealthy.WithLabelValues(addr, nodeID, host, manager).Set(healthy) + c.m.NodeManagerRunning.WithLabelValues(addr, nodeID, host, manager).Set(running) + c.m.NodeManagerHealthy.WithLabelValues(addr, nodeID, host, manager).Set(healthy) } // --- Restore metrics from history on startup --- @@ -609,38 +695,38 @@ func (c *Collector) restoreMetrics() { for epochStr, snap := range epochs { p := participant e := epochStr - metrics.EpochInferences.WithLabelValues(p, e).Set(float64(snap.InferenceCount)) - metrics.EpochMissed.WithLabelValues(p, e).Set(float64(snap.MissedRequests)) - metrics.EpochEarnedCoins.WithLabelValues(p, e).Set(float64(snap.EarnedCoins)) - metrics.EpochValidated.WithLabelValues(p, e).Set(float64(snap.ValidatedInferences)) - metrics.EpochInvalidated.WithLabelValues(p, e).Set(float64(snap.InvalidatedInferences)) - metrics.EpochCoinBalance.WithLabelValues(p, e).Set(float64(snap.CoinBalance)) - metrics.EpochDone.WithLabelValues(p, e).Set(float64(snap.EpochsCompleted)) - metrics.EpochMissRate.WithLabelValues(p, e).Set(snap.MissRatePercent) - metrics.EpochTimeslot.WithLabelValues(p, e).Set(float64(snap.TimeslotAssigned)) + c.m.EpochInferences.WithLabelValues(p, e).Set(float64(snap.InferenceCount)) + c.m.EpochMissed.WithLabelValues(p, e).Set(float64(snap.MissedRequests)) + c.m.EpochEarnedCoins.WithLabelValues(p, e).Set(float64(snap.EarnedCoins)) + c.m.EpochValidated.WithLabelValues(p, e).Set(float64(snap.ValidatedInferences)) + c.m.EpochInvalidated.WithLabelValues(p, e).Set(float64(snap.InvalidatedInferences)) + c.m.EpochCoinBalance.WithLabelValues(p, e).Set(float64(snap.CoinBalance)) + c.m.EpochDone.WithLabelValues(p, e).Set(float64(snap.EpochsCompleted)) + c.m.EpochMissRate.WithLabelValues(p, e).Set(snap.MissRatePercent) + c.m.EpochTimeslot.WithLabelValues(p, e).Set(float64(snap.TimeslotAssigned)) for nodeID, pw := range snap.PocWeights { - metrics.EpochPocWeight.WithLabelValues(p, e, nodeID).Set(float64(pw)) + c.m.EpochPocWeight.WithLabelValues(p, e, nodeID).Set(float64(pw)) } if snap.StartTime > 0 { - metrics.EpochStartTime.WithLabelValues(p, e).Set(snap.StartTime) + c.m.EpochStartTime.WithLabelValues(p, e).Set(snap.StartTime) } if snap.EndTime > 0 { - metrics.EpochEndTime.WithLabelValues(p, e).Set(snap.EndTime) + c.m.EpochEndTime.WithLabelValues(p, e).Set(snap.EndTime) } if snap.DurationSeconds > 0 { - metrics.EpochDuration.WithLabelValues(p, e).Set(snap.DurationSeconds) + c.m.EpochDuration.WithLabelValues(p, e).Set(snap.DurationSeconds) } if snap.EarnedGNK != nil { - metrics.EpochEarnedGNK.WithLabelValues(p, e).Set(*snap.EarnedGNK) + c.m.EpochEarnedGNK.WithLabelValues(p, e).Set(*snap.EarnedGNK) } if snap.RewardedGNK != nil { - metrics.EpochRewardedGNK.WithLabelValues(p, e).Set(*snap.RewardedGNK) + c.m.EpochRewardedGNK.WithLabelValues(p, e).Set(*snap.RewardedGNK) } if snap.EstimatedGNK != nil { - metrics.EpochEstimated.WithLabelValues(p, e).Set(*snap.EstimatedGNK) + c.m.EpochEstimated.WithLabelValues(p, e).Set(*snap.EstimatedGNK) } if snap.Claimed != nil { - metrics.EpochClaimed.WithLabelValues(p, e).Set(float64(*snap.Claimed)) + c.m.EpochClaimed.WithLabelValues(p, e).Set(float64(*snap.Claimed)) } total++ } @@ -658,13 +744,7 @@ func (c *Collector) pruneEpochMax() { if len(keys) <= 5 { return } - for i := 0; i < len(keys)-1; i++ { - for j := i + 1; j < len(keys); j++ { - if keys[i] > keys[j] { - keys[i], keys[j] = keys[j], keys[i] - } - } - } + sort.Slice(keys, func(i, j int) bool { return keys[i] < keys[j] }) for _, k := range keys[:len(keys)-5] { delete(c.epochMax, k) } @@ -686,13 +766,6 @@ func pocStatusVal(s string) float64 { return v } -func max64(a, b int64) int64 { - if a > b { - return a - } - return b -} - func copyMap(m map[string]int64) map[string]int64 { out := make(map[string]int64, len(m)) for k, v := range m { @@ -705,7 +778,7 @@ func copyMap(m map[string]int64) map[string]int64 { func (c *Collector) collectStats() { // Per-model stats - models, err := fetcher.FetchStatsModels(c.cfg.APIURL) + models, err := c.f.FetchStatsModels(c.cfg.APIURL) if err != nil { slog.Debug("stats models unavailable", "err", err) } else { @@ -713,57 +786,57 @@ func (c *Collector) collectStats() { if m.Model == "" { continue } - metrics.StatsModelAiTokens.WithLabelValues(m.Model).Set(float64(m.AiTokens)) - metrics.StatsModelInferences.WithLabelValues(m.Model).Set(float64(m.Inferences)) + c.m.StatsModelAiTokens.WithLabelValues(m.Model).Set(float64(m.AiTokens)) + c.m.StatsModelInferences.WithLabelValues(m.Model).Set(float64(m.Inferences)) } } // Network-wide summary - summary, err := fetcher.FetchStatsSummary(c.cfg.APIURL) + summary, err := c.f.FetchStatsSummary(c.cfg.APIURL) if err != nil { slog.Debug("stats summary unavailable", "err", err) return } - metrics.StatsAiTokens.Set(float64(summary.AiTokens)) - metrics.StatsInferences.Set(float64(summary.Inferences)) - metrics.StatsActualCost.Set(float64(summary.ActualCost)) + c.m.StatsAiTokens.Set(float64(summary.AiTokens)) + c.m.StatsInferences.Set(float64(summary.Inferences)) + c.m.StatsActualCost.Set(float64(summary.ActualCost)) } // --- Bridge --- func (c *Collector) collectBridge() { - status, err := fetcher.FetchBridgeStatus(c.cfg.APIURL) + status, err := c.f.FetchBridgeStatus(c.cfg.APIURL) if err != nil { slog.Debug("bridge status unavailable", "err", err) return } - metrics.BridgePendingBlocks.Set(float64(status.PendingBlocks)) - metrics.BridgePendingReceipts.Set(float64(status.PendingReceipts)) + c.m.BridgePendingBlocks.Set(float64(status.PendingBlocks)) + c.m.BridgePendingReceipts.Set(float64(status.PendingReceipts)) ready := 0.0 if status.ReadyToProcess { ready = 1.0 } - metrics.BridgeReadyToProcess.Set(ready) + c.m.BridgeReadyToProcess.Set(ready) if status.EarliestBlock > 0 { - metrics.BridgeEarliestBlock.Set(float64(status.EarliestBlock)) + c.m.BridgeEarliestBlock.Set(float64(status.EarliestBlock)) } if status.LatestBlock > 0 { - metrics.BridgeLatestBlock.Set(float64(status.LatestBlock)) + c.m.BridgeLatestBlock.Set(float64(status.LatestBlock)) } } // --- Tokenomics --- func (c *Collector) collectTokenomics() { - tok, err := fetcher.FetchTokenomics(c.cfg.NodeRESTURL) + tok, err := c.f.FetchTokenomics(c.cfg.NodeRESTURL) if err != nil { slog.Debug("tokenomics unavailable", "err", err) return } - metrics.TokenomicsTotalFees.Set(float64(tok.TotalFees)) - metrics.TokenomicsTotalSubsidies.Set(float64(tok.TotalSubsidies)) - metrics.TokenomicsTotalRefunded.Set(float64(tok.TotalRefunded)) - metrics.TokenomicsTotalBurned.Set(float64(tok.TotalBurned)) + c.m.TokenomicsTotalFees.Set(float64(tok.TotalFees)) + c.m.TokenomicsTotalSubsidies.Set(float64(tok.TotalSubsidies)) + c.m.TokenomicsTotalRefunded.Set(float64(tok.TotalRefunded)) + c.m.TokenomicsTotalBurned.Set(float64(tok.TotalBurned)) } // --- PoC v2 --- @@ -775,15 +848,15 @@ func (c *Collector) collectPoCv2() { } // Artifact count - commit, err := fetcher.FetchPoCv2Commit(c.cfg.NodeRESTURL, c.st.PocStartBlockHeight, addr) + commit, err := c.f.FetchPoCv2Commit(c.cfg.NodeRESTURL, c.st.PocStartBlockHeight, addr) if err != nil { slog.Debug("poc_v2 commit unavailable", "err", err) } else { - metrics.PoCv2ArtifactCount.WithLabelValues(addr).Set(float64(commit.Count)) + c.m.PoCv2ArtifactCount.WithLabelValues(addr).Set(float64(commit.Count)) } // Per-node weight distribution - weights, err := fetcher.FetchMLNodeWeightDist(c.cfg.NodeRESTURL, c.st.PocStartBlockHeight, addr) + weights, err := c.f.FetchMLNodeWeightDist(c.cfg.NodeRESTURL, c.st.PocStartBlockHeight, addr) if err != nil { slog.Debug("poc_v2 node weight dist unavailable", "err", err) return @@ -792,6 +865,6 @@ func (c *Collector) collectPoCv2() { if w.NodeID == "" { continue } - metrics.PoCv2NodeWeight.WithLabelValues(addr, w.NodeID).Set(float64(w.Weight)) + c.m.PoCv2NodeWeight.WithLabelValues(addr, w.NodeID).Set(float64(w.Weight)) } } diff --git a/internal/config/config.go b/internal/config/config.go index 3e2048f..77df068 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -23,14 +23,11 @@ type Config struct { BlockHeightNodes []string } +// defaultBlockNodes is a minimal fallback list used when BLOCK_HEIGHT_NODES is not set. +// Override via BLOCK_HEIGHT_NODES env var in production (comma-separated URLs). var defaultBlockNodes = []string{ "http://node1.gonka.ai:8000", "http://node2.gonka.ai:8000", - "https://node3.gonka.ai", - "http://36.189.234.237:17241", - "http://47.236.26.199:8000", - "http://47.236.19.22:18000", - "http://gonka.spv.re:8000", } // Load reads all configuration from environment variables. diff --git a/internal/fetcher/fetcher.go b/internal/fetcher/fetcher.go index 5a791ae..c6b0e51 100644 --- a/internal/fetcher/fetcher.go +++ b/internal/fetcher/fetcher.go @@ -4,10 +4,8 @@ import ( "encoding/json" "fmt" "io" - "math/rand/v2" "net/http" "strconv" - "strings" "time" ) @@ -50,729 +48,3 @@ func get(url string, dest any) error { return json.Unmarshal(body, dest) } -// --- Tendermint RPC --- - -type TendermintStatus struct { - Result struct { - SyncInfo struct { - LatestBlockHeight string `json:"latest_block_height"` - LatestBlockTime string `json:"latest_block_time"` - CatchingUp bool `json:"catching_up"` - } `json:"sync_info"` - } `json:"result"` -} - -func FetchTendermintStatus(rpcURL string) (*TendermintStatus, error) { - var s TendermintStatus - err := get(rpcURL+"/status", &s) - return &s, err -} - -func FetchBlockTimeAtHeight(rpcURL string, height int64) (float64, error) { - var resp struct { - Result struct { - Block struct { - Header struct { - Time string `json:"time"` - } `json:"header"` - } `json:"block"` - } `json:"result"` - } - if err := get(fmt.Sprintf("%s/block?height=%d", rpcURL, height), &resp); err != nil { - return 0, err - } - t := resp.Result.Block.Header.Time - if t == "" { - return 0, fmt.Errorf("empty time") - } - t = strings.TrimSuffix(t, "Z") - parsed, err := time.Parse("2006-01-02T15:04:05.999999999", t) - if err != nil { - parsed, err = time.Parse(time.RFC3339Nano, t+"Z") - if err != nil { - return 0, fmt.Errorf("parse time %q: %w", t, err) - } - } - return float64(parsed.UTC().Unix()) + float64(parsed.Nanosecond())/1e9, nil -} - -func FetchMaxBlockHeightFromNodes(nodes []string) (int64, string) { - sample := rand.Perm(len(nodes)) - if len(sample) > 5 { - sample = sample[:5] - } - var maxHeight int64 - var latestTime string - for _, i := range sample { - var resp struct { - Result struct { - SyncInfo struct { - LatestBlockHeight string `json:"latest_block_height"` - LatestBlockTime string `json:"latest_block_time"` - } `json:"sync_info"` - } `json:"result"` - } - if err := get(nodes[i]+"/chain-rpc/status", &resp); err != nil { - continue - } - h, err := strconv.ParseInt(resp.Result.SyncInfo.LatestBlockHeight, 10, 64) - if err == nil && h > maxHeight { - maxHeight = h - latestTime = resp.Result.SyncInfo.LatestBlockTime - } - } - return maxHeight, latestTime -} - -// --- Chain REST --- - -func FetchCurrentEpoch(restURL string) (int64, error) { - var r struct { - Epoch string `json:"epoch"` - } - if err := get(restURL+"/productscience/inference/inference/get_current_epoch", &r); err != nil { - return 0, err - } - return strconv.ParseInt(r.Epoch, 10, 64) -} - -// EpochInfo contains epoch block data from the chain. -type EpochInfo struct { - PocStartBlockHeight int64 - EpochLength int64 - BlockHeight int64 -} - -func FetchEpochInfo(restURL string) (*EpochInfo, error) { - var r struct { - LatestEpoch struct { - PocStartBlockHeight string `json:"poc_start_block_height"` - } `json:"latest_epoch"` - Params struct { - EpochParams struct { - EpochLength string `json:"epoch_length"` - } `json:"epoch_params"` - } `json:"params"` - BlockHeight string `json:"block_height"` - } - if err := get(restURL+"/productscience/inference/inference/epoch_info", &r); err != nil { - return nil, err - } - poc, err1 := strconv.ParseInt(r.LatestEpoch.PocStartBlockHeight, 10, 64) - el, err2 := strconv.ParseInt(r.Params.EpochParams.EpochLength, 10, 64) - bh, err3 := strconv.ParseInt(r.BlockHeight, 10, 64) - if err1 != nil || err2 != nil || err3 != nil { - return nil, fmt.Errorf("parse epoch_info: %v %v %v", err1, err2, err3) - } - return &EpochInfo{ - PocStartBlockHeight: poc, - EpochLength: el, - BlockHeight: bh, - }, nil -} - -// --- Epoch group data --- - -// ValidationWeight is a per-participant weight entry. -type ValidationWeight struct { - MemberAddress string `json:"member_address"` - Weight string `json:"weight"` - Reputation flexInt64 `json:"reputation"` -} - -// EpochGroupData contains network-wide epoch weight data. -type EpochGroupData struct { - TotalWeight int64 - EpochIndex int64 - NumberOfRequests int64 - ValidationWeights []ValidationWeight -} - -func FetchEpochGroupData(restURL string) (*EpochGroupData, error) { - var r struct { - EpochGroupData struct { - TotalWeight string `json:"total_weight"` - EpochIndex string `json:"epoch_index"` - NumberOfRequests flexInt64 `json:"number_of_requests"` - ValidationWeights []ValidationWeight `json:"validation_weights"` - } `json:"epoch_group_data"` - } - if err := get(restURL+"/productscience/inference/inference/current_epoch_group_data", &r); err != nil { - return nil, err - } - tw, _ := strconv.ParseInt(r.EpochGroupData.TotalWeight, 10, 64) - ei, _ := strconv.ParseInt(r.EpochGroupData.EpochIndex, 10, 64) - return &EpochGroupData{ - TotalWeight: tw, - EpochIndex: ei, - NumberOfRequests: int64(r.EpochGroupData.NumberOfRequests), - ValidationWeights: r.EpochGroupData.ValidationWeights, - }, nil -} - -// --- Epoch performance summary --- - -// EpochPerfSummary holds on-chain reward data for a completed epoch. -type EpochPerfSummary struct { - RewardedGNK float64 - Claimed int -} - -func FetchEpochPerfSummary(restURL, address string, epochNum int64) (*EpochPerfSummary, error) { - var r struct { - EpochPerformanceSummary struct { - RewardedCoins string `json:"rewarded_coins"` - Claimed bool `json:"claimed"` - } `json:"epochPerformanceSummary"` - } - url := fmt.Sprintf("%s/productscience/inference/inference/epoch_performance_summary/%d/%s", restURL, epochNum, address) - if err := get(url, &r); err != nil { - return nil, err - } - rc, err := strconv.ParseInt(r.EpochPerformanceSummary.RewardedCoins, 10, 64) - if err != nil { - return nil, fmt.Errorf("parse rewarded_coins: %w", err) - } - claimed := 0 - if r.EpochPerformanceSummary.Claimed { - claimed = 1 - } - return &EpochPerfSummary{RewardedGNK: float64(rc) / 1e9, Claimed: claimed}, nil -} - -// --- Participant stats --- - -// ParticipantStats holds current-epoch stats for one participant. -type ParticipantStats struct { - EpochsCompleted int64 - CoinBalance int64 - InferenceCount int64 - MissedRequests int64 - EarnedCoins int64 - ValidatedInferences int64 - InvalidatedInferences int64 - Status string // "ACTIVE", "INACTIVE", "INVALID", "UNCONFIRMED", "UNSPECIFIED" - ConsecutiveInvalidInferences int64 - BurnedCoins int64 - RewardedCoins int64 -} - -type participantStatsResp struct { - Participant struct { - EpochsCompleted flexInt64 `json:"epochs_completed"` - CoinBalance flexInt64 `json:"coin_balance"` - Status string `json:"status"` - ConsecutiveInvalidInferences flexInt64 `json:"consecutive_invalid_inferences"` - CurrentEpochStats struct { - InferenceCount flexInt64 `json:"inference_count"` - MissedRequests flexInt64 `json:"missed_requests"` - EarnedCoins flexInt64 `json:"earned_coins"` - ValidatedInferences flexInt64 `json:"validated_inferences"` - InvalidatedInferences flexInt64 `json:"invalidated_inferences"` - BurnedCoins flexInt64 `json:"burned_coins"` - RewardedCoins flexInt64 `json:"rewarded_coins"` - } `json:"current_epoch_stats"` - } `json:"participant"` -} - -func FetchParticipantStats(restURL, address string) (*ParticipantStats, error) { - var r participantStatsResp - if err := get(restURL+"/productscience/inference/inference/participant/"+address, &r); err != nil { - return nil, err - } - p := r.Participant - return &ParticipantStats{ - EpochsCompleted: int64(p.EpochsCompleted), - CoinBalance: int64(p.CoinBalance), - InferenceCount: int64(p.CurrentEpochStats.InferenceCount), - MissedRequests: int64(p.CurrentEpochStats.MissedRequests), - EarnedCoins: int64(p.CurrentEpochStats.EarnedCoins), - ValidatedInferences: int64(p.CurrentEpochStats.ValidatedInferences), - InvalidatedInferences: int64(p.CurrentEpochStats.InvalidatedInferences), - Status: p.Status, - ConsecutiveInvalidInferences: int64(p.ConsecutiveInvalidInferences), - BurnedCoins: int64(p.CurrentEpochStats.BurnedCoins), - RewardedCoins: int64(p.CurrentEpochStats.RewardedCoins), - }, nil -} - -// --- Bank balance --- - -func FetchWalletBalance(restURL, address string) (float64, error) { - var r struct { - Balances []struct { - Denom string `json:"denom"` - Amount string `json:"amount"` - } `json:"balances"` - } - if err := get(restURL+"/cosmos/bank/v1beta1/balances/"+address, &r); err != nil { - return 0, err - } - for _, b := range r.Balances { - if b.Denom == "ngonka" { - n, err := strconv.ParseInt(b.Amount, 10, 64) - if err != nil { - return 0, err - } - return float64(n) / 1e9, nil - } - } - return 0, fmt.Errorf("ngonka balance not found") -} - -// --- Node list (admin API) --- - -// NodeEntry represents a single node returned by the admin API. -type NodeEntry struct { - Node struct { - ID string `json:"id"` - Host string `json:"host"` - PocPort int `json:"poc_port"` - Hardware []struct { - Type string `json:"type"` - Count int `json:"count"` - } `json:"hardware"` - } `json:"node"` - State struct { - CurrentStatus string `json:"current_status"` - IntendedStatus string `json:"intended_status"` - PocCurrentStatus string `json:"poc_current_status"` - PocIntendedStatus string `json:"poc_intended_status"` - EpochMLNodes map[string]struct { - PocWeight *int64 `json:"poc_weight"` - TimeslotAllocation []bool `json:"timeslot_allocation"` - } `json:"epoch_ml_nodes"` - } `json:"state"` -} - -func FetchNodes(adminURL string) ([]NodeEntry, error) { - var nodes []NodeEntry - err := get(adminURL+"/admin/v1/nodes", &nodes) - return nodes, err -} - -// --- GPU stats --- - -// GPUDevice holds per-device GPU metrics. -type GPUDevice struct { - Index int - UtilizationPercent float64 - TemperatureC *float64 - TotalMemoryMB *int64 - FreeMemoryMB *int64 - UsedMemoryMB *int64 - IsAvailable bool -} - -// GPUStats contains aggregated GPU info for a node. -type GPUStats struct { - Count int - AvgUtil float64 - Devices []GPUDevice -} - -func FetchGPUStats(host string, port int) GPUStats { - var r struct { - Devices []struct { - Index int `json:"index"` - UtilizationPercent float64 `json:"utilization_percent"` - TemperatureC *float64 `json:"temperature_c"` - TotalMemoryMB *int64 `json:"total_memory_mb"` - FreeMemoryMB *int64 `json:"free_memory_mb"` - UsedMemoryMB *int64 `json:"used_memory_mb"` - IsAvailable bool `json:"is_available"` - } `json:"devices"` - } - url := fmt.Sprintf("http://%s:%d/v3.0.8/api/v1/gpu/devices", host, port) - if err := get(url, &r); err != nil || len(r.Devices) == 0 { - return GPUStats{} - } - var total float64 - devices := make([]GPUDevice, len(r.Devices)) - for i, d := range r.Devices { - total += d.UtilizationPercent - devices[i] = GPUDevice{ - Index: d.Index, - UtilizationPercent: d.UtilizationPercent, - TemperatureC: d.TemperatureC, - TotalMemoryMB: d.TotalMemoryMB, - FreeMemoryMB: d.FreeMemoryMB, - UsedMemoryMB: d.UsedMemoryMB, - IsAvailable: d.IsAvailable, - } - } - return GPUStats{Count: len(r.Devices), AvgUtil: total / float64(len(r.Devices)), Devices: devices} -} - -// FetchMLNodeState returns the current service state string ("POW", "INFERENCE", "TRAIN", "STOPPED"). -func FetchMLNodeState(host string, port int) (string, error) { - var r struct { - State string `json:"state"` - } - url := fmt.Sprintf("http://%s:%d/v3.0.8/api/v1/state", host, port) - if err := get(url, &r); err != nil { - return "", err - } - return r.State, nil -} - -// FetchMLNodeDiskSpaceGB returns available disk space in GB for the model cache. -func FetchMLNodeDiskSpaceGB(host string, port int) (float64, error) { - var r struct { - AvailableGB float64 `json:"available_gb"` - } - url := fmt.Sprintf("http://%s:%d/v3.0.8/api/v1/models/space", host, port) - if err := get(url, &r); err != nil { - return 0, err - } - return r.AvailableGB, nil -} - -// --- Participants (network-wide weights) --- - -// ParticipantEntry is one entry in the active participants list. -type ParticipantEntry struct { - Seed struct { - Participant string `json:"participant"` - } `json:"seed"` - Weight *float64 `json:"weight"` - MLNodes []struct { - MLNodes []struct { - NodeID string `json:"node_id"` - PocWeight *float64 `json:"poc_weight"` - } `json:"ml_nodes"` - } `json:"ml_nodes"` -} - -func FetchNetworkParticipants(apiURL string) ([]ParticipantEntry, error) { - var r struct { - ActiveParticipants struct { - Participants []ParticipantEntry `json:"participants"` - } `json:"active_participants"` - } - if err := get(apiURL+"/v1/epochs/current/participants", &r); err != nil { - return nil, err - } - return r.ActiveParticipants.Participants, nil -} - -// --- Pricing --- - -// PricingData holds current pricing configuration. -type PricingData struct { - UnitOfComputePrice *float64 `json:"unit_of_compute_price"` - DynamicPricingEnabled *bool `json:"dynamic_pricing_enabled"` - Models []struct { - ID string `json:"id"` - PricePerToken *float64 `json:"price_per_token"` - UnitsOfComputePerToken *float64 `json:"units_of_compute_per_token"` - Utilization *float64 `json:"utilization"` - Capacity *int64 `json:"capacity"` - } `json:"models"` -} - -func FetchPricing(apiURL string) (*PricingData, error) { - var r PricingData - err := get(apiURL+"/v1/pricing", &r) - return &r, err -} - -// --- Models --- - -// ModelData holds model definitions from the API. -type ModelData struct { - Models []struct { - ID string `json:"id"` - VRAM *float64 `json:"v_ram"` - ThroughputPerNonce *float64 `json:"throughput_per_nonce"` - ValidationThreshold *struct { - Value float64 `json:"value"` - Exponent int `json:"exponent"` - } `json:"validation_threshold"` - } `json:"models"` -} - -func FetchModels(apiURL string) (*ModelData, error) { - var r ModelData - err := get(apiURL+"/v1/models", &r) - return &r, err -} - -// BLSEpochData holds BLS DKG phase information for an epoch. -type BLSEpochData struct { - DKGPhase int32 - DealingPhaseDeadlineBlock int64 - VerifyingPhaseDeadlineBlock int64 -} - -// dkgPhaseNames maps DKG phase string names to numeric values. -var dkgPhaseNames = map[string]int32{ - "DKG_PHASE_UNDEFINED": 0, - "DKG_PHASE_DEALING": 1, - "DKG_PHASE_VERIFYING": 2, - "DKG_PHASE_COMPLETED": 3, - "DKG_PHASE_FAILED": 4, - "DKG_PHASE_SIGNED": 5, -} - -// --- Stats API --- - -// StatsSummaryData holds network-wide inference statistics. -type StatsSummaryData struct { - AiTokens int64 - Inferences int32 - ActualCost int64 -} - -// FetchStatsSummary returns aggregated network stats from /v1/stats/summary/time. -func FetchStatsSummary(apiURL string) (*StatsSummaryData, error) { - var r struct { - AiTokens int64 `json:"ai_tokens"` - Inferences int32 `json:"inferences"` - ActualInferencesCost int64 `json:"actual_inferences_cost"` - } - if err := get(apiURL+"/v1/stats/summary/time", &r); err != nil { - return nil, err - } - return &StatsSummaryData{AiTokens: r.AiTokens, Inferences: r.Inferences, ActualCost: r.ActualInferencesCost}, nil -} - -// StatsModelEntry holds per-model inference statistics. -type StatsModelEntry struct { - Model string - AiTokens int64 - Inferences int32 -} - -// FetchStatsModels returns per-model stats from /v1/stats/models. -func FetchStatsModels(apiURL string) ([]StatsModelEntry, error) { - var r struct { - StatsModels []struct { - Model string `json:"model"` - AiTokens int64 `json:"ai_tokens"` - Inferences int32 `json:"inferences"` - } `json:"stats_models"` - } - if err := get(apiURL+"/v1/stats/models", &r); err != nil { - return nil, err - } - out := make([]StatsModelEntry, len(r.StatsModels)) - for i, m := range r.StatsModels { - out[i] = StatsModelEntry{Model: m.Model, AiTokens: m.AiTokens, Inferences: m.Inferences} - } - return out, nil -} - -// --- Bridge --- - -// BridgeStatusData holds bridge queue status. -type BridgeStatusData struct { - PendingBlocks int - PendingReceipts int - ReadyToProcess bool - EarliestBlock uint64 - LatestBlock uint64 -} - -// FetchBridgeStatus returns bridge queue status from /v1/bridge/status. -func FetchBridgeStatus(apiURL string) (*BridgeStatusData, error) { - var r struct { - PendingBlocksCount int `json:"pendingBlocksCount"` - PendingReceiptsCount int `json:"pendingReceiptsCount"` - EarliestBlockNumber uint64 `json:"earliestBlockNumber"` - LatestBlockNumber uint64 `json:"latestBlockNumber"` - ReadyToProcess bool `json:"readyToProcess"` - } - if err := get(apiURL+"/v1/bridge/status", &r); err != nil { - return nil, err - } - return &BridgeStatusData{ - PendingBlocks: r.PendingBlocksCount, - PendingReceipts: r.PendingReceiptsCount, - ReadyToProcess: r.ReadyToProcess, - EarliestBlock: r.EarliestBlockNumber, - LatestBlock: r.LatestBlockNumber, - }, nil -} - -// --- ML node health --- - -// MLNodeManagerStatus holds running/healthy status for one manager. -type MLNodeManagerStatus struct { - Running bool - Healthy bool -} - -// MLNodeHealthData holds manager health for all ML node managers. -type MLNodeHealthData struct { - ManagerPow MLNodeManagerStatus - ManagerInference MLNodeManagerStatus - ManagerTrain MLNodeManagerStatus -} - -// FetchMLNodeHealth returns manager health. Tries /health first (root endpoint), -// then /v3.0.8/api/v1/health as fallback. -func FetchMLNodeHealth(host string, port int) (*MLNodeHealthData, error) { - var r struct { - Managers struct { - Pow struct { - Running bool `json:"running"` - Healthy bool `json:"healthy"` - } `json:"pow"` - Inference struct { - Running bool `json:"running"` - Healthy bool `json:"healthy"` - } `json:"inference"` - Train struct { - Running bool `json:"running"` - Healthy bool `json:"healthy"` - } `json:"train"` - } `json:"managers"` - } - url := fmt.Sprintf("http://%s:%d/health", host, port) - if err := get(url, &r); err != nil { - // Fallback: versioned path - url = fmt.Sprintf("http://%s:%d/v3.0.8/api/v1/health", host, port) - if err2 := get(url, &r); err2 != nil { - return nil, err2 - } - } - return &MLNodeHealthData{ - ManagerPow: MLNodeManagerStatus{Running: r.Managers.Pow.Running, Healthy: r.Managers.Pow.Healthy}, - ManagerInference: MLNodeManagerStatus{Running: r.Managers.Inference.Running, Healthy: r.Managers.Inference.Healthy}, - ManagerTrain: MLNodeManagerStatus{Running: r.Managers.Train.Running, Healthy: r.Managers.Train.Healthy}, - }, nil -} - -// --- GPU driver info --- - -// GPUDriverData holds GPU driver and CUDA version strings. -type GPUDriverData struct { - DriverVersion string - CudaDriverVersion string -} - -// FetchGPUDriverInfo returns GPU driver info from /v3.0.8/api/v1/gpu/driver. -func FetchGPUDriverInfo(host string, port int) (*GPUDriverData, error) { - var r struct { - DriverVersion string `json:"driver_version"` - CudaDriverVersion string `json:"cuda_driver_version"` - } - url := fmt.Sprintf("http://%s:%d/v3.0.8/api/v1/gpu/driver", host, port) - if err := get(url, &r); err != nil { - return nil, err - } - if r.DriverVersion == "" { - return nil, fmt.Errorf("empty driver version") - } - return &GPUDriverData{DriverVersion: r.DriverVersion, CudaDriverVersion: r.CudaDriverVersion}, nil -} - -// --- Tokenomics --- - -// TokenomicsData holds chain-wide tokenomics counters. -type TokenomicsData struct { - TotalFees uint64 - TotalSubsidies uint64 - TotalRefunded uint64 - TotalBurned uint64 -} - -// FetchTokenomics returns tokenomics data from the chain REST endpoint. -func FetchTokenomics(restURL string) (*TokenomicsData, error) { - var r struct { - TokenomicsData struct { - TotalFees flexInt64 `json:"total_fees"` - TotalSubsidies flexInt64 `json:"total_subsidies"` - TotalRefunded flexInt64 `json:"total_refunded"` - TotalBurned flexInt64 `json:"total_burned"` - } `json:"tokenomics_data"` - } - if err := get(restURL+"/productscience/inference/inference/tokenomics_data", &r); err != nil { - return nil, err - } - td := r.TokenomicsData - return &TokenomicsData{ - TotalFees: uint64(td.TotalFees), - TotalSubsidies: uint64(td.TotalSubsidies), - TotalRefunded: uint64(td.TotalRefunded), - TotalBurned: uint64(td.TotalBurned), - }, nil -} - -// --- PoC v2 --- - -// PoCv2CommitData holds the artifact commit count for a participant. -type PoCv2CommitData struct { - Count uint32 -} - -// FetchPoCv2Commit returns the PoC v2 store commit for a participant at a given PoC stage start block. -func FetchPoCv2Commit(restURL string, pocStartBlock int64, address string) (*PoCv2CommitData, error) { - var r struct { - PoCV2StoreCommit struct { - Count flexInt64 `json:"count"` - } `json:"poc_v2_store_commit"` - } - url := fmt.Sprintf("%s/productscience/inference/inference/poc_v2_store_commit/%d/%s", restURL, pocStartBlock, address) - if err := get(url, &r); err != nil { - return nil, err - } - return &PoCv2CommitData{Count: uint32(r.PoCV2StoreCommit.Count)}, nil -} - -// MLNodeWeightEntry holds per-node weight from MLNode weight distribution. -type MLNodeWeightEntry struct { - NodeID string - Weight uint32 -} - -// FetchMLNodeWeightDist returns the per-node weight distribution for a participant at a given PoC stage start block. -func FetchMLNodeWeightDist(restURL string, pocStartBlock int64, address string) ([]MLNodeWeightEntry, error) { - var r struct { - MLNodeWeightDistribution struct { - Weights []struct { - NodeID string `json:"node_id"` - Weight flexInt64 `json:"weight"` - } `json:"weights"` - } `json:"mlnode_weight_distribution"` - } - url := fmt.Sprintf("%s/productscience/inference/inference/mlnode_weight_distribution/%d/%s", restURL, pocStartBlock, address) - if err := get(url, &r); err != nil { - return nil, err - } - out := make([]MLNodeWeightEntry, len(r.MLNodeWeightDistribution.Weights)) - for i, w := range r.MLNodeWeightDistribution.Weights { - out[i] = MLNodeWeightEntry{NodeID: w.NodeID, Weight: uint32(w.Weight)} - } - return out, nil -} - -// FetchBLSEpoch fetches BLS DKG phase data for the given epoch ID from the public API. -func FetchBLSEpoch(apiURL string, epochID int64) (*BLSEpochData, error) { - var r struct { - EpochData struct { - DKGPhase interface{} `json:"dkg_phase"` - DealingPhaseDeadlineBlock flexInt64 `json:"dealing_phase_deadline_block"` - VerifyingPhaseDeadlineBlock flexInt64 `json:"verifying_phase_deadline_block"` - } `json:"epoch_data"` - } - url := fmt.Sprintf("%s/v1/bls/epoch/%d", apiURL, epochID) - if err := get(url, &r); err != nil { - return nil, err - } - var phase int32 - switch v := r.EpochData.DKGPhase.(type) { - case float64: - phase = int32(v) - case string: - if n, ok := dkgPhaseNames[v]; ok { - phase = n - } - } - return &BLSEpochData{ - DKGPhase: phase, - DealingPhaseDeadlineBlock: int64(r.EpochData.DealingPhaseDeadlineBlock), - VerifyingPhaseDeadlineBlock: int64(r.EpochData.VerifyingPhaseDeadlineBlock), - }, nil -} diff --git a/internal/fetcher/fetcher_admin.go b/internal/fetcher/fetcher_admin.go new file mode 100644 index 0000000..8262632 --- /dev/null +++ b/internal/fetcher/fetcher_admin.go @@ -0,0 +1,30 @@ +package fetcher + +// NodeEntry represents a single node returned by the admin API. +type NodeEntry struct { + Node struct { + ID string `json:"id"` + Host string `json:"host"` + PocPort int `json:"poc_port"` + Hardware []struct { + Type string `json:"type"` + Count int `json:"count"` + } `json:"hardware"` + } `json:"node"` + State struct { + CurrentStatus string `json:"current_status"` + IntendedStatus string `json:"intended_status"` + PocCurrentStatus string `json:"poc_current_status"` + PocIntendedStatus string `json:"poc_intended_status"` + EpochMLNodes map[string]struct { + PocWeight *int64 `json:"poc_weight"` + TimeslotAllocation []bool `json:"timeslot_allocation"` + } `json:"epoch_ml_nodes"` + } `json:"state"` +} + +func (h *HTTPFetcher) FetchNodes(adminURL string) ([]NodeEntry, error) { + var nodes []NodeEntry + err := get(adminURL+"/admin/v1/nodes", &nodes) + return nodes, err +} diff --git a/internal/fetcher/fetcher_api.go b/internal/fetcher/fetcher_api.go new file mode 100644 index 0000000..69626aa --- /dev/null +++ b/internal/fetcher/fetcher_api.go @@ -0,0 +1,189 @@ +package fetcher + +import "fmt" + +// ParticipantEntry is one entry in the active participants list. +type ParticipantEntry struct { + Seed struct { + Participant string `json:"participant"` + } `json:"seed"` + Weight *float64 `json:"weight"` + MLNodes []struct { + MLNodes []struct { + NodeID string `json:"node_id"` + PocWeight *float64 `json:"poc_weight"` + } `json:"ml_nodes"` + } `json:"ml_nodes"` +} + +func (h *HTTPFetcher) FetchNetworkParticipants(apiURL string) ([]ParticipantEntry, error) { + var r struct { + ActiveParticipants struct { + Participants []ParticipantEntry `json:"participants"` + } `json:"active_participants"` + } + if err := get(apiURL+"/v1/epochs/current/participants", &r); err != nil { + return nil, err + } + return r.ActiveParticipants.Participants, nil +} + +// PricingData holds current pricing configuration. +type PricingData struct { + UnitOfComputePrice *float64 `json:"unit_of_compute_price"` + DynamicPricingEnabled *bool `json:"dynamic_pricing_enabled"` + Models []struct { + ID string `json:"id"` + PricePerToken *float64 `json:"price_per_token"` + UnitsOfComputePerToken *float64 `json:"units_of_compute_per_token"` + Utilization *float64 `json:"utilization"` + Capacity *int64 `json:"capacity"` + } `json:"models"` +} + +func (h *HTTPFetcher) FetchPricing(apiURL string) (*PricingData, error) { + var r PricingData + err := get(apiURL+"/v1/pricing", &r) + return &r, err +} + +// ModelData holds model definitions from the API. +type ModelData struct { + Models []struct { + ID string `json:"id"` + VRAM *float64 `json:"v_ram"` + ThroughputPerNonce *float64 `json:"throughput_per_nonce"` + ValidationThreshold *struct { + Value float64 `json:"value"` + Exponent int `json:"exponent"` + } `json:"validation_threshold"` + } `json:"models"` +} + +func (h *HTTPFetcher) FetchModels(apiURL string) (*ModelData, error) { + var r ModelData + err := get(apiURL+"/v1/models", &r) + return &r, err +} + +// BLSEpochData holds BLS DKG phase information for an epoch. +type BLSEpochData struct { + DKGPhase int32 + DealingPhaseDeadlineBlock int64 + VerifyingPhaseDeadlineBlock int64 +} + +// dkgPhaseNames maps DKG phase string names to numeric values. +var dkgPhaseNames = map[string]int32{ + "DKG_PHASE_UNDEFINED": 0, + "DKG_PHASE_DEALING": 1, + "DKG_PHASE_VERIFYING": 2, + "DKG_PHASE_COMPLETED": 3, + "DKG_PHASE_FAILED": 4, + "DKG_PHASE_SIGNED": 5, +} + +// FetchBLSEpoch fetches BLS DKG phase data for the given epoch ID from the public API. +func (h *HTTPFetcher) FetchBLSEpoch(apiURL string, epochID int64) (*BLSEpochData, error) { + var r struct { + EpochData struct { + DKGPhase any `json:"dkg_phase"` + DealingPhaseDeadlineBlock flexInt64 `json:"dealing_phase_deadline_block"` + VerifyingPhaseDeadlineBlock flexInt64 `json:"verifying_phase_deadline_block"` + } `json:"epoch_data"` + } + url := fmt.Sprintf("%s/v1/bls/epoch/%d", apiURL, epochID) + if err := get(url, &r); err != nil { + return nil, err + } + var phase int32 + switch v := r.EpochData.DKGPhase.(type) { + case float64: + phase = int32(v) + case string: + if n, ok := dkgPhaseNames[v]; ok { + phase = n + } + } + return &BLSEpochData{ + DKGPhase: phase, + DealingPhaseDeadlineBlock: int64(r.EpochData.DealingPhaseDeadlineBlock), + VerifyingPhaseDeadlineBlock: int64(r.EpochData.VerifyingPhaseDeadlineBlock), + }, nil +} + +// StatsSummaryData holds network-wide inference statistics. +type StatsSummaryData struct { + AiTokens int64 + Inferences int32 + ActualCost int64 +} + +// FetchStatsSummary returns aggregated network stats from /v1/stats/summary/time. +func (h *HTTPFetcher) FetchStatsSummary(apiURL string) (*StatsSummaryData, error) { + var r struct { + AiTokens int64 `json:"ai_tokens"` + Inferences int32 `json:"inferences"` + ActualInferencesCost int64 `json:"actual_inferences_cost"` + } + if err := get(apiURL+"/v1/stats/summary/time", &r); err != nil { + return nil, err + } + return &StatsSummaryData{AiTokens: r.AiTokens, Inferences: r.Inferences, ActualCost: r.ActualInferencesCost}, nil +} + +// StatsModelEntry holds per-model inference statistics. +type StatsModelEntry struct { + Model string + AiTokens int64 + Inferences int32 +} + +// FetchStatsModels returns per-model stats from /v1/stats/models. +func (h *HTTPFetcher) FetchStatsModels(apiURL string) ([]StatsModelEntry, error) { + var r struct { + StatsModels []struct { + Model string `json:"model"` + AiTokens int64 `json:"ai_tokens"` + Inferences int32 `json:"inferences"` + } `json:"stats_models"` + } + if err := get(apiURL+"/v1/stats/models", &r); err != nil { + return nil, err + } + out := make([]StatsModelEntry, len(r.StatsModels)) + for i, m := range r.StatsModels { + out[i] = StatsModelEntry{Model: m.Model, AiTokens: m.AiTokens, Inferences: m.Inferences} + } + return out, nil +} + +// BridgeStatusData holds bridge queue status. +type BridgeStatusData struct { + PendingBlocks int + PendingReceipts int + ReadyToProcess bool + EarliestBlock uint64 + LatestBlock uint64 +} + +// FetchBridgeStatus returns bridge queue status from /v1/bridge/status. +func (h *HTTPFetcher) FetchBridgeStatus(apiURL string) (*BridgeStatusData, error) { + var r struct { + PendingBlocksCount int `json:"pendingBlocksCount"` + PendingReceiptsCount int `json:"pendingReceiptsCount"` + EarliestBlockNumber uint64 `json:"earliestBlockNumber"` + LatestBlockNumber uint64 `json:"latestBlockNumber"` + ReadyToProcess bool `json:"readyToProcess"` + } + if err := get(apiURL+"/v1/bridge/status", &r); err != nil { + return nil, err + } + return &BridgeStatusData{ + PendingBlocks: r.PendingBlocksCount, + PendingReceipts: r.PendingReceiptsCount, + ReadyToProcess: r.ReadyToProcess, + EarliestBlock: r.EarliestBlockNumber, + LatestBlock: r.LatestBlockNumber, + }, nil +} diff --git a/internal/fetcher/fetcher_chain.go b/internal/fetcher/fetcher_chain.go new file mode 100644 index 0000000..cc5476f --- /dev/null +++ b/internal/fetcher/fetcher_chain.go @@ -0,0 +1,100 @@ +package fetcher + +import ( + "fmt" + "math/rand/v2" + "strconv" + "strings" + "sync" + "time" +) + +// TendermintStatus is the response from the Tendermint /status RPC endpoint. +type TendermintStatus struct { + Result struct { + SyncInfo struct { + LatestBlockHeight string `json:"latest_block_height"` + LatestBlockTime string `json:"latest_block_time"` + CatchingUp bool `json:"catching_up"` + } `json:"sync_info"` + } `json:"result"` +} + +func (h *HTTPFetcher) FetchTendermintStatus(rpcURL string) (*TendermintStatus, error) { + var s TendermintStatus + err := get(rpcURL+"/status", &s) + return &s, err +} + +func (h *HTTPFetcher) FetchBlockTimeAtHeight(rpcURL string, height int64) (float64, error) { + var resp struct { + Result struct { + Block struct { + Header struct { + Time string `json:"time"` + } `json:"header"` + } `json:"block"` + } `json:"result"` + } + if err := get(fmt.Sprintf("%s/block?height=%d", rpcURL, height), &resp); err != nil { + return 0, err + } + t := resp.Result.Block.Header.Time + if t == "" { + return 0, fmt.Errorf("empty time") + } + t = strings.TrimSuffix(t, "Z") + parsed, err := time.Parse("2006-01-02T15:04:05.999999999", t) + if err != nil { + parsed, err = time.Parse(time.RFC3339Nano, t+"Z") + if err != nil { + return 0, fmt.Errorf("parse time %q: %w", t, err) + } + } + return float64(parsed.UTC().Unix()) + float64(parsed.Nanosecond())/1e9, nil +} + +func (h *HTTPFetcher) FetchMaxBlockHeightFromNodes(nodes []string) (int64, string) { + if len(nodes) == 0 { + return 0, "" + } + sample := rand.Perm(len(nodes)) + if len(sample) > 5 { + sample = sample[:5] + } + + var mu sync.Mutex + var maxHeight int64 + var latestTime string + + var wg sync.WaitGroup + for _, idx := range sample { + wg.Add(1) + go func(nodeURL string) { + defer wg.Done() + var resp struct { + Result struct { + SyncInfo struct { + LatestBlockHeight string `json:"latest_block_height"` + LatestBlockTime string `json:"latest_block_time"` + } `json:"sync_info"` + } `json:"result"` + } + if err := get(nodeURL+"/chain-rpc/status", &resp); err != nil { + return + } + height, err := strconv.ParseInt(resp.Result.SyncInfo.LatestBlockHeight, 10, 64) + if err != nil { + return + } + mu.Lock() + if height > maxHeight { + maxHeight = height + latestTime = resp.Result.SyncInfo.LatestBlockTime + } + mu.Unlock() + }(nodes[idx]) + } + wg.Wait() + return maxHeight, latestTime +} diff --git a/internal/fetcher/fetcher_cosmos.go b/internal/fetcher/fetcher_cosmos.go new file mode 100644 index 0000000..f67f9c6 --- /dev/null +++ b/internal/fetcher/fetcher_cosmos.go @@ -0,0 +1,274 @@ +package fetcher + +import ( + "fmt" + "strconv" +) + +func (h *HTTPFetcher) FetchCurrentEpoch(restURL string) (int64, error) { + var r struct { + Epoch string `json:"epoch"` + } + if err := get(restURL+"/productscience/inference/inference/get_current_epoch", &r); err != nil { + return 0, err + } + return strconv.ParseInt(r.Epoch, 10, 64) +} + +// EpochInfo contains epoch block data from the chain. +type EpochInfo struct { + PocStartBlockHeight int64 + EpochLength int64 + BlockHeight int64 +} + +func (h *HTTPFetcher) FetchEpochInfo(restURL string) (*EpochInfo, error) { + var r struct { + LatestEpoch struct { + PocStartBlockHeight string `json:"poc_start_block_height"` + } `json:"latest_epoch"` + Params struct { + EpochParams struct { + EpochLength string `json:"epoch_length"` + } `json:"epoch_params"` + } `json:"params"` + BlockHeight string `json:"block_height"` + } + if err := get(restURL+"/productscience/inference/inference/epoch_info", &r); err != nil { + return nil, err + } + poc, err1 := strconv.ParseInt(r.LatestEpoch.PocStartBlockHeight, 10, 64) + el, err2 := strconv.ParseInt(r.Params.EpochParams.EpochLength, 10, 64) + bh, err3 := strconv.ParseInt(r.BlockHeight, 10, 64) + if err1 != nil || err2 != nil || err3 != nil { + return nil, fmt.Errorf("parse epoch_info: %v %v %v", err1, err2, err3) + } + return &EpochInfo{ + PocStartBlockHeight: poc, + EpochLength: el, + BlockHeight: bh, + }, nil +} + +// ValidationWeight is a per-participant weight entry. +type ValidationWeight struct { + MemberAddress string `json:"member_address"` + Weight string `json:"weight"` + Reputation flexInt64 `json:"reputation"` +} + +// EpochGroupData contains network-wide epoch weight data. +type EpochGroupData struct { + TotalWeight int64 + EpochIndex int64 + NumberOfRequests int64 + ValidationWeights []ValidationWeight +} + +func (h *HTTPFetcher) FetchEpochGroupData(restURL string) (*EpochGroupData, error) { + var r struct { + EpochGroupData struct { + TotalWeight string `json:"total_weight"` + EpochIndex string `json:"epoch_index"` + NumberOfRequests flexInt64 `json:"number_of_requests"` + ValidationWeights []ValidationWeight `json:"validation_weights"` + } `json:"epoch_group_data"` + } + if err := get(restURL+"/productscience/inference/inference/current_epoch_group_data", &r); err != nil { + return nil, err + } + tw, err := strconv.ParseInt(r.EpochGroupData.TotalWeight, 10, 64) + if err != nil { + return nil, fmt.Errorf("parse current_epoch_group_data total_weight %q: %w", r.EpochGroupData.TotalWeight, err) + } + ei, err := strconv.ParseInt(r.EpochGroupData.EpochIndex, 10, 64) + if err != nil { + return nil, fmt.Errorf("parse current_epoch_group_data epoch_index %q: %w", r.EpochGroupData.EpochIndex, err) + } + return &EpochGroupData{ + TotalWeight: tw, + EpochIndex: ei, + NumberOfRequests: int64(r.EpochGroupData.NumberOfRequests), + ValidationWeights: r.EpochGroupData.ValidationWeights, + }, nil +} + +// EpochPerfSummary holds on-chain reward data for a completed epoch. +type EpochPerfSummary struct { + RewardedGNK float64 + Claimed int +} + +func (h *HTTPFetcher) FetchEpochPerfSummary(restURL, address string, epochNum int64) (*EpochPerfSummary, error) { + var r struct { + EpochPerformanceSummary struct { + RewardedCoins string `json:"rewarded_coins"` + Claimed bool `json:"claimed"` + } `json:"epochPerformanceSummary"` + } + url := fmt.Sprintf("%s/productscience/inference/inference/epoch_performance_summary/%d/%s", restURL, epochNum, address) + if err := get(url, &r); err != nil { + return nil, err + } + rc, err := strconv.ParseInt(r.EpochPerformanceSummary.RewardedCoins, 10, 64) + if err != nil { + return nil, fmt.Errorf("parse rewarded_coins: %w", err) + } + claimed := 0 + if r.EpochPerformanceSummary.Claimed { + claimed = 1 + } + return &EpochPerfSummary{RewardedGNK: float64(rc) / 1e9, Claimed: claimed}, nil +} + +// ParticipantStats holds current-epoch stats for one participant. +type ParticipantStats struct { + EpochsCompleted int64 + CoinBalance int64 + InferenceCount int64 + MissedRequests int64 + EarnedCoins int64 + ValidatedInferences int64 + InvalidatedInferences int64 + Status string // "ACTIVE", "INACTIVE", "INVALID", "UNCONFIRMED", "UNSPECIFIED" + ConsecutiveInvalidInferences int64 + BurnedCoins int64 + RewardedCoins int64 +} + +type participantStatsResp struct { + Participant struct { + EpochsCompleted flexInt64 `json:"epochs_completed"` + CoinBalance flexInt64 `json:"coin_balance"` + Status string `json:"status"` + ConsecutiveInvalidInferences flexInt64 `json:"consecutive_invalid_inferences"` + CurrentEpochStats struct { + InferenceCount flexInt64 `json:"inference_count"` + MissedRequests flexInt64 `json:"missed_requests"` + EarnedCoins flexInt64 `json:"earned_coins"` + ValidatedInferences flexInt64 `json:"validated_inferences"` + InvalidatedInferences flexInt64 `json:"invalidated_inferences"` + BurnedCoins flexInt64 `json:"burned_coins"` + RewardedCoins flexInt64 `json:"rewarded_coins"` + } `json:"current_epoch_stats"` + } `json:"participant"` +} + +func (h *HTTPFetcher) FetchParticipantStats(restURL, address string) (*ParticipantStats, error) { + var r participantStatsResp + if err := get(restURL+"/productscience/inference/inference/participant/"+address, &r); err != nil { + return nil, err + } + p := r.Participant + return &ParticipantStats{ + EpochsCompleted: int64(p.EpochsCompleted), + CoinBalance: int64(p.CoinBalance), + InferenceCount: int64(p.CurrentEpochStats.InferenceCount), + MissedRequests: int64(p.CurrentEpochStats.MissedRequests), + EarnedCoins: int64(p.CurrentEpochStats.EarnedCoins), + ValidatedInferences: int64(p.CurrentEpochStats.ValidatedInferences), + InvalidatedInferences: int64(p.CurrentEpochStats.InvalidatedInferences), + Status: p.Status, + ConsecutiveInvalidInferences: int64(p.ConsecutiveInvalidInferences), + BurnedCoins: int64(p.CurrentEpochStats.BurnedCoins), + RewardedCoins: int64(p.CurrentEpochStats.RewardedCoins), + }, nil +} + +func (h *HTTPFetcher) FetchWalletBalance(restURL, address string) (float64, error) { + var r struct { + Balances []struct { + Denom string `json:"denom"` + Amount string `json:"amount"` + } `json:"balances"` + } + if err := get(restURL+"/cosmos/bank/v1beta1/balances/"+address, &r); err != nil { + return 0, err + } + for _, b := range r.Balances { + if b.Denom == "ngonka" { + n, err := strconv.ParseInt(b.Amount, 10, 64) + if err != nil { + return 0, err + } + return float64(n) / 1e9, nil + } + } + return 0, nil // new wallet with zero balance — valid state, not an error +} + +// TokenomicsData holds chain-wide tokenomics counters. +type TokenomicsData struct { + TotalFees uint64 + TotalSubsidies uint64 + TotalRefunded uint64 + TotalBurned uint64 +} + +// FetchTokenomics returns tokenomics data from the chain REST endpoint. +func (h *HTTPFetcher) FetchTokenomics(restURL string) (*TokenomicsData, error) { + var r struct { + TokenomicsData struct { + TotalFees flexInt64 `json:"total_fees"` + TotalSubsidies flexInt64 `json:"total_subsidies"` + TotalRefunded flexInt64 `json:"total_refunded"` + TotalBurned flexInt64 `json:"total_burned"` + } `json:"tokenomics_data"` + } + if err := get(restURL+"/productscience/inference/inference/tokenomics_data", &r); err != nil { + return nil, err + } + td := r.TokenomicsData + return &TokenomicsData{ + TotalFees: uint64(td.TotalFees), + TotalSubsidies: uint64(td.TotalSubsidies), + TotalRefunded: uint64(td.TotalRefunded), + TotalBurned: uint64(td.TotalBurned), + }, nil +} + +// PoCv2CommitData holds the artifact commit count for a participant. +type PoCv2CommitData struct { + Count uint32 +} + +// FetchPoCv2Commit returns the PoC v2 store commit for a participant at a given PoC stage start block. +func (h *HTTPFetcher) FetchPoCv2Commit(restURL string, pocStartBlock int64, address string) (*PoCv2CommitData, error) { + var r struct { + PoCV2StoreCommit struct { + Count flexInt64 `json:"count"` + } `json:"poc_v2_store_commit"` + } + url := fmt.Sprintf("%s/productscience/inference/inference/poc_v2_store_commit/%d/%s", restURL, pocStartBlock, address) + if err := get(url, &r); err != nil { + return nil, err + } + return &PoCv2CommitData{Count: uint32(r.PoCV2StoreCommit.Count)}, nil +} + +// MLNodeWeightEntry holds per-node weight from MLNode weight distribution. +type MLNodeWeightEntry struct { + NodeID string + Weight uint32 +} + +// FetchMLNodeWeightDist returns the per-node weight distribution for a participant at a given PoC stage start block. +func (h *HTTPFetcher) FetchMLNodeWeightDist(restURL string, pocStartBlock int64, address string) ([]MLNodeWeightEntry, error) { + var r struct { + MLNodeWeightDistribution struct { + Weights []struct { + NodeID string `json:"node_id"` + Weight flexInt64 `json:"weight"` + } `json:"weights"` + } `json:"mlnode_weight_distribution"` + } + url := fmt.Sprintf("%s/productscience/inference/inference/mlnode_weight_distribution/%d/%s", restURL, pocStartBlock, address) + if err := get(url, &r); err != nil { + return nil, err + } + out := make([]MLNodeWeightEntry, len(r.MLNodeWeightDistribution.Weights)) + for i, w := range r.MLNodeWeightDistribution.Weights { + out[i] = MLNodeWeightEntry{NodeID: w.NodeID, Weight: uint32(w.Weight)} + } + return out, nil +} diff --git a/internal/fetcher/fetcher_ml.go b/internal/fetcher/fetcher_ml.go new file mode 100644 index 0000000..c673c2d --- /dev/null +++ b/internal/fetcher/fetcher_ml.go @@ -0,0 +1,149 @@ +package fetcher + +import "fmt" + +const mlNodeAPIBase = "/v3.0.8/api/v1" + +// GPUDevice holds per-device GPU metrics. +type GPUDevice struct { + Index int + UtilizationPercent float64 + TemperatureC *float64 + TotalMemoryMB *int64 + FreeMemoryMB *int64 + UsedMemoryMB *int64 + IsAvailable bool +} + +// GPUStats contains aggregated GPU info for a node. +type GPUStats struct { + Count int + AvgUtil float64 + Devices []GPUDevice +} + +func (h *HTTPFetcher) FetchGPUStats(host string, port int) GPUStats { + var r struct { + Devices []struct { + Index int `json:"index"` + UtilizationPercent float64 `json:"utilization_percent"` + TemperatureC *float64 `json:"temperature_c"` + TotalMemoryMB *int64 `json:"total_memory_mb"` + FreeMemoryMB *int64 `json:"free_memory_mb"` + UsedMemoryMB *int64 `json:"used_memory_mb"` + IsAvailable bool `json:"is_available"` + } `json:"devices"` + } + url := fmt.Sprintf("http://%s:%d%s/gpu/devices", host, port, mlNodeAPIBase) + if err := get(url, &r); err != nil || len(r.Devices) == 0 { + return GPUStats{} + } + var total float64 + devices := make([]GPUDevice, len(r.Devices)) + for i, d := range r.Devices { + total += d.UtilizationPercent + devices[i] = GPUDevice{ + Index: d.Index, + UtilizationPercent: d.UtilizationPercent, + TemperatureC: d.TemperatureC, + TotalMemoryMB: d.TotalMemoryMB, + FreeMemoryMB: d.FreeMemoryMB, + UsedMemoryMB: d.UsedMemoryMB, + IsAvailable: d.IsAvailable, + } + } + return GPUStats{Count: len(r.Devices), AvgUtil: total / float64(len(r.Devices)), Devices: devices} +} + +// FetchMLNodeState returns the current service state string ("POW", "INFERENCE", "TRAIN", "STOPPED"). +func (h *HTTPFetcher) FetchMLNodeState(host string, port int) (string, error) { + var r struct { + State string `json:"state"` + } + url := fmt.Sprintf("http://%s:%d%s/state", host, port, mlNodeAPIBase) + if err := get(url, &r); err != nil { + return "", err + } + return r.State, nil +} + +// FetchMLNodeDiskSpaceGB returns available disk space in GB for the model cache. +func (h *HTTPFetcher) FetchMLNodeDiskSpaceGB(host string, port int) (float64, error) { + var r struct { + AvailableGB float64 `json:"available_gb"` + } + url := fmt.Sprintf("http://%s:%d%s/models/space", host, port, mlNodeAPIBase) + if err := get(url, &r); err != nil { + return 0, err + } + return r.AvailableGB, nil +} + +// MLNodeManagerStatus holds running/healthy status for one manager. +type MLNodeManagerStatus struct { + Running bool + Healthy bool +} + +// MLNodeHealthData holds manager health for all ML node managers. +type MLNodeHealthData struct { + ManagerPow MLNodeManagerStatus + ManagerInference MLNodeManagerStatus + ManagerTrain MLNodeManagerStatus +} + +// FetchMLNodeHealth returns manager health. Tries /health first (root endpoint), +// then /v3.0.8/api/v1/health as fallback. +func (h *HTTPFetcher) FetchMLNodeHealth(host string, port int) (*MLNodeHealthData, error) { + var r struct { + Managers struct { + Pow struct { + Running bool `json:"running"` + Healthy bool `json:"healthy"` + } `json:"pow"` + Inference struct { + Running bool `json:"running"` + Healthy bool `json:"healthy"` + } `json:"inference"` + Train struct { + Running bool `json:"running"` + Healthy bool `json:"healthy"` + } `json:"train"` + } `json:"managers"` + } + url := fmt.Sprintf("http://%s:%d/health", host, port) + if err := get(url, &r); err != nil { + // Fallback: versioned path + url = fmt.Sprintf("http://%s:%d%s/health", host, port, mlNodeAPIBase) + if err2 := get(url, &r); err2 != nil { + return nil, err2 + } + } + return &MLNodeHealthData{ + ManagerPow: MLNodeManagerStatus{Running: r.Managers.Pow.Running, Healthy: r.Managers.Pow.Healthy}, + ManagerInference: MLNodeManagerStatus{Running: r.Managers.Inference.Running, Healthy: r.Managers.Inference.Healthy}, + ManagerTrain: MLNodeManagerStatus{Running: r.Managers.Train.Running, Healthy: r.Managers.Train.Healthy}, + }, nil +} + +// GPUDriverData holds GPU driver and CUDA version strings. +type GPUDriverData struct { + DriverVersion string + CudaDriverVersion string +} + +// FetchGPUDriverInfo returns GPU driver info from /v3.0.8/api/v1/gpu/driver. +func (h *HTTPFetcher) FetchGPUDriverInfo(host string, port int) (*GPUDriverData, error) { + var r struct { + DriverVersion string `json:"driver_version"` + CudaDriverVersion string `json:"cuda_driver_version"` + } + url := fmt.Sprintf("http://%s:%d%s/gpu/driver", host, port, mlNodeAPIBase) + if err := get(url, &r); err != nil { + return nil, err + } + if r.DriverVersion == "" { + return nil, fmt.Errorf("empty driver version") + } + return &GPUDriverData{DriverVersion: r.DriverVersion, CudaDriverVersion: r.CudaDriverVersion}, nil +} diff --git a/internal/fetcher/interface.go b/internal/fetcher/interface.go new file mode 100644 index 0000000..071f7a6 --- /dev/null +++ b/internal/fetcher/interface.go @@ -0,0 +1,42 @@ +package fetcher + +// Fetcher defines all upstream API calls used by the collector. +// HTTPFetcher is the production implementation; tests can inject a mock. +type Fetcher interface { + FetchTendermintStatus(rpcURL string) (*TendermintStatus, error) + FetchBlockTimeAtHeight(rpcURL string, height int64) (float64, error) + FetchMaxBlockHeightFromNodes(nodes []string) (int64, string) + + FetchCurrentEpoch(restURL string) (int64, error) + FetchEpochInfo(restURL string) (*EpochInfo, error) + FetchEpochGroupData(restURL string) (*EpochGroupData, error) + FetchEpochPerfSummary(restURL, address string, epochNum int64) (*EpochPerfSummary, error) + FetchBLSEpoch(apiURL string, epochID int64) (*BLSEpochData, error) + + FetchParticipantStats(restURL, address string) (*ParticipantStats, error) + FetchWalletBalance(restURL, address string) (float64, error) + + FetchNetworkParticipants(apiURL string) ([]ParticipantEntry, error) + FetchPricing(apiURL string) (*PricingData, error) + FetchModels(apiURL string) (*ModelData, error) + FetchStatsSummary(apiURL string) (*StatsSummaryData, error) + FetchStatsModels(apiURL string) ([]StatsModelEntry, error) + FetchBridgeStatus(apiURL string) (*BridgeStatusData, error) + + FetchNodes(adminURL string) ([]NodeEntry, error) + FetchGPUStats(host string, port int) GPUStats + FetchMLNodeState(host string, port int) (string, error) + FetchMLNodeDiskSpaceGB(host string, port int) (float64, error) + FetchGPUDriverInfo(host string, port int) (*GPUDriverData, error) + FetchMLNodeHealth(host string, port int) (*MLNodeHealthData, error) + + FetchTokenomics(restURL string) (*TokenomicsData, error) + FetchPoCv2Commit(restURL string, pocStartBlock int64, address string) (*PoCv2CommitData, error) + FetchMLNodeWeightDist(restURL string, pocStartBlock int64, address string) ([]MLNodeWeightEntry, error) +} + +// HTTPFetcher is the production implementation of Fetcher using real HTTP calls. +type HTTPFetcher struct{} + +// NewHTTPFetcher returns a Fetcher backed by real HTTP calls. +func NewHTTPFetcher() *HTTPFetcher { return &HTTPFetcher{} } diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index 923f00b..c28d603 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -2,177 +2,362 @@ package metrics import "github.com/prometheus/client_golang/prometheus" -// Chain / sync -var ( - BlockHeight = gaugeVec("gonka_block_height", "Latest block height from local node", "participant") - BlockHeightMax = gaugeVec("gonka_block_height_max", "Maximum block height seen across public nodes", "participant") - BlockTime = gaugeVec("gonka_block_time_seconds","Timestamp of latest block (unix)", "participant") - CatchingUp = gaugeVec("gonka_chain_catching_up", "1 = syncing, 0 = synced", "participant") - ChainEpoch = gaugeVec("gonka_chain_epoch", "Current global chain epoch number", "participant") -) - -// Node hardware / status -var ( - NodeStatus = gaugeVec("gonka_node_status", "Node hardware status (0=UNKNOWN 1=INFERENCE 2=POC 3=TRAINING 4=STOPPED 5=FAILED)", "participant", "node_id", "host") - NodeIntended = gaugeVec("gonka_node_intended_status", "Node intended status", "participant", "node_id", "host") - PocCurrent = gaugeVec("gonka_node_poc_current_status", "PoC current status (0=IDLE 1=GENERATING 2=VALIDATING)", "participant", "node_id", "host") - PocIntended = gaugeVec("gonka_node_poc_intended_status", "PoC intended status", "participant", "node_id", "host") - NodePocWeight = gaugeVec("gonka_node_poc_weight", "PoC weight per node per model", "participant", "node_id", "host", "model") - NodeTimeslot = gaugeVec("gonka_node_poc_timeslot_assigned","Timeslot assigned (1/0)", "participant", "node_id", "host", "model") - NodeGPUCount = gaugeVec("gonka_node_gpu_device_count", "GPU device count", "participant", "node_id", "host") - NodeGPUUtil = gaugeVec("gonka_node_gpu_avg_utilization_percent", "Average GPU utilization %", "participant", "node_id", "host") - NodeHardwareInfo = gaugeVec("gonka_node_hardware_info", "Hardware info (value=1, metadata in labels)", "participant", "node_id", "host", "hardware_type", "hardware_count") -) - -// Node GPU — per device -var ( - NodeGPUDeviceUtil = gaugeVec("gonka_node_gpu_device_utilization_percent", "Per-device GPU compute utilization %", "participant", "node_id", "host", "device_index") - NodeGPUDeviceTemp = gaugeVec("gonka_node_gpu_device_temperature_celsius", "Per-device GPU temperature °C", "participant", "node_id", "host", "device_index") - NodeGPUDeviceMemTotal= gaugeVec("gonka_node_gpu_device_memory_total_mb", "Per-device GPU total memory MB", "participant", "node_id", "host", "device_index") - NodeGPUDeviceMemFree = gaugeVec("gonka_node_gpu_device_memory_free_mb", "Per-device GPU free memory MB", "participant", "node_id", "host", "device_index") - NodeGPUDeviceMemUsed = gaugeVec("gonka_node_gpu_device_memory_used_mb", "Per-device GPU used memory MB", "participant", "node_id", "host", "device_index") - NodeGPUDeviceAvail = gaugeVec("gonka_node_gpu_device_available", "Per-device GPU available (1=yes 0=no)", "participant", "node_id", "host", "device_index") -) - -// Node ML service state -var ( - NodeServiceState = gaugeVec("gonka_node_service_state", "ML node service state (0=STOPPED 1=INFERENCE 2=POW 3=TRAIN)", "participant", "node_id", "host") - NodeDiskAvailableGB = gaugeVec("gonka_node_disk_available_gb", "ML node model cache available disk space GB", "participant", "node_id", "host") -) - -// Network-wide -var ( - NetParticipantWeight = gaugeVec("gonka_network_participant_weight", "Per-participant weight in active epoch", "participant") - NetNodePocWeight = gaugeVec("gonka_network_node_poc_weight", "Per-node PoC weight in active epoch", "participant", "node_id") - NetTotalWeight = gaugeVec("gonka_network_total_weight", "Total weight of all participants", "participant") - NetRewardPerWeight = gaugeVec("gonka_network_reward_per_weight", "Estimated GNK per unit of weight", "participant") -) - -// Pricing / models -var ( - PricingUoC = gauge("gonka_pricing_unit_of_compute_price", "Unit of compute price") - PricingDynamic = gauge("gonka_pricing_dynamic_enabled", "Dynamic pricing enabled (1/0)") - ModelPrice = gaugeVec("gonka_pricing_model_price_per_token", "Price per token per model", "model_id") - ModelUnits = gaugeVec("gonka_pricing_model_units_per_token", "Compute units per token", "model_id") - ModelVRAM = gaugeVec("gonka_model_v_ram", "VRAM (GB)", "model_id") - ModelThroughput = gaugeVec("gonka_model_throughput_per_nonce", "Throughput per nonce", "model_id") - ModelValThresh = gaugeVec("gonka_model_validation_threshold", "Validation threshold", "model_id") -) - -// Participant live (current epoch) -var ( - ParticipantEpochsDone = gaugeVec("gonka_participant_epochs_completed", "Personal epochs completed", "participant") - ParticipantCoinBalance = gaugeVec("gonka_participant_coin_balance", "Coin balance (internal points)", "participant") - ParticipantWallet = gaugeVec("gonka_participant_wallet_balance_gonka", "Wallet balance in GNK (ngonka/1e9)", "participant") - ParticipantInferences = gaugeVec("gonka_participant_inference_count", "Inferences in current epoch", "participant") - ParticipantMissed = gaugeVec("gonka_participant_missed_requests", "Missed requests in current epoch", "participant") - ParticipantEarnedCoins = gaugeVec("gonka_participant_earned_coins", "Earned coins in current epoch", "participant") - ParticipantValidated = gaugeVec("gonka_participant_validated_inferences", "Validated inferences in current epoch", "participant") - ParticipantInvalidated = gaugeVec("gonka_participant_invalidated_inferences", "Invalidated inferences in current epoch","participant") -) - -// Participant health (extended) -var ( - ParticipantStatus = gaugeVec("gonka_participant_status", "Participant status (0=UNSPECIFIED 1=ACTIVE 2=INACTIVE 3=INVALID 5=UNCONFIRMED)", "participant") - ParticipantConsecutiveInv = gaugeVec("gonka_participant_consecutive_invalid_inferences", "Consecutive invalid inferences counter", "participant") - ParticipantBurnedCoins = gaugeVec("gonka_participant_burned_coins", "Burned (penalized) coins in current epoch", "participant") - ParticipantRewardedCoins = gaugeVec("gonka_participant_rewarded_coins", "Rewarded coins in current epoch (after distribution)", "participant") - ParticipantReputation = gaugeVec("gonka_participant_reputation", "Participant reputation score from epoch group data", "participant") -) - -// Network — counts and epoch-level -var ( - NetActiveParticipantCount = gauge("gonka_network_active_participant_count", "Number of active participants in current epoch") - NetTotalParticipantCount = gauge("gonka_network_total_participant_count", "Total number of participants in current epoch") - NetEpochInferenceCount = gaugeVec("gonka_network_epoch_inference_count", "Total network inferences in current epoch", "participant") -) - -// BLS DKG phase -var ( - BLSDKGPhase = gaugeVec("gonka_bls_dkg_phase", "BLS DKG phase (0=UNDEFINED 1=DEALING 2=VERIFYING 3=COMPLETED 4=FAILED 5=SIGNED)", "participant") - BLSDealingDeadline = gaugeVec("gonka_bls_dealing_deadline_block", "BLS dealing phase deadline block height", "participant") - BLSVerifyingDeadline = gaugeVec("gonka_bls_verifying_deadline_block", "BLS verifying phase deadline block height", "participant") -) - -// Model utilization / capacity -var ( - ModelUtilization = gaugeVec("gonka_model_utilization_percent", "Model utilization % (inferences/capacity)", "model_id") - ModelCapacity = gaugeVec("gonka_model_capacity", "Model maximum capacity (AI tokens/epoch)", "model_id") -) - -// Epoch history (label epoch = chain epoch number as string) -var ( - EpochInferences = gaugeVec("gonka_epoch_inference_count", "Inferences for epoch N", "participant", "epoch") - EpochMissed = gaugeVec("gonka_epoch_missed_requests", "Missed requests for epoch N", "participant", "epoch") - EpochEarnedCoins = gaugeVec("gonka_epoch_earned_coins", "Earned coins (internal points) for epoch N", "participant", "epoch") - EpochValidated = gaugeVec("gonka_epoch_validated_inferences", "Validated inferences for epoch N", "participant", "epoch") - EpochInvalidated = gaugeVec("gonka_epoch_invalidated_inferences", "Invalidated inferences for epoch N", "participant", "epoch") - EpochCoinBalance = gaugeVec("gonka_epoch_coin_balance", "Coin balance at epoch end", "participant", "epoch") - EpochDone = gaugeVec("gonka_epoch_epochs_completed", "Personal epochs completed at epoch end", "participant", "epoch") - EpochMissRate = gaugeVec("gonka_epoch_miss_rate_percent", "Miss rate %% for epoch N", "participant", "epoch") - EpochPocWeight = gaugeVec("gonka_epoch_poc_weight", "PoC weight for epoch N", "participant", "epoch", "node_id") - EpochTimeslot = gaugeVec("gonka_epoch_timeslot_assigned", "Timeslot assigned in epoch N (1/0)", "participant", "epoch") - EpochStartTime = gaugeVec("gonka_epoch_start_time", "Epoch start unix timestamp", "participant", "epoch") - EpochEndTime = gaugeVec("gonka_epoch_end_time", "Epoch end unix timestamp (estimated for live)", "participant", "epoch") - EpochDuration = gaugeVec("gonka_epoch_duration_seconds", "Real epoch duration (seconds)", "participant", "epoch") - EpochEarnedGNK = gaugeVec("gonka_epoch_earned_gonka", "Wallet balance delta for epoch (GNK)", "participant", "epoch") - EpochRewardedGNK = gaugeVec("gonka_epoch_rewarded_gonka", "On-chain rewarded GNK for epoch N (rewarded_coins/1e9)", "participant", "epoch") - EpochClaimed = gaugeVec("gonka_epoch_claimed", "1 if epoch reward was claimed, 0 if not", "participant", "epoch") - EpochEstimated = gaugeVec("gonka_epoch_estimated_reward_gonka", "Estimated GNK reward for epoch N (weight × emission/total_weight)", "participant", "epoch") -) - -// Stats — network-wide inference statistics -var ( - StatsAiTokens = gauge("gonka_stats_ai_tokens_total", "Total AI tokens processed network-wide (cumulative)") - StatsInferences = gauge("gonka_stats_inferences_total", "Total inferences processed network-wide (cumulative)") - StatsActualCost = gauge("gonka_stats_actual_cost_total", "Total actual inference cost network-wide (coins, cumulative)") - StatsModelAiTokens = gaugeVec("gonka_stats_model_ai_tokens", "AI tokens processed per model (cumulative)", "model_id") - StatsModelInferences = gaugeVec("gonka_stats_model_inferences", "Inferences processed per model (cumulative)", "model_id") -) - -// Bridge — queue status for the Gonka bridge -var ( - BridgePendingBlocks = gauge("gonka_bridge_pending_blocks", "Bridge pending block count") - BridgePendingReceipts = gauge("gonka_bridge_pending_receipts", "Bridge pending receipt count") - BridgeReadyToProcess = gauge("gonka_bridge_ready_to_process", "Bridge ready to process (1=yes 0=no)") - BridgeEarliestBlock = gauge("gonka_bridge_earliest_block_number", "Bridge earliest pending block number") - BridgeLatestBlock = gauge("gonka_bridge_latest_block_number", "Bridge latest pending block number") -) - -// Node managers — running/healthy state per ML node manager -var ( - NodeManagerRunning = gaugeVec("gonka_node_manager_running", "ML node manager running (1=yes 0=no)", "participant", "node_id", "host", "manager") - NodeManagerHealthy = gaugeVec("gonka_node_manager_healthy", "ML node manager healthy (1=yes 0=no)", "participant", "node_id", "host", "manager") -) - -// GPU driver info — version in labels, value always 1 -var ( - NodeGPUDriverInfo = gaugeVec("gonka_node_gpu_driver_info", "GPU driver info (value=1, version in labels)", "participant", "node_id", "host", "driver_version", "cuda_version") -) - -// Tokenomics — chain-wide token flow counters -var ( - TokenomicsTotalFees = gauge("gonka_tokenomics_total_fees", "Total fees collected (ngonka)") - TokenomicsTotalSubsidies = gauge("gonka_tokenomics_total_subsidies", "Total subsidies issued (ngonka)") - TokenomicsTotalRefunded = gauge("gonka_tokenomics_total_refunded", "Total amount refunded (ngonka)") - TokenomicsTotalBurned = gauge("gonka_tokenomics_total_burned", "Total amount burned (ngonka)") -) - - -// PoC v2 — proof-of-compute artifact and weight data -var ( - PoCv2ArtifactCount = gaugeVec("gonka_poc_v2_artifact_count", "PoC v2 artifact count from store commit", "participant") - PoCv2NodeWeight = gaugeVec("gonka_poc_v2_node_weight", "PoC v2 per-node weight distribution", "participant", "node_id") -) - -func gauge(name, help string) prometheus.Gauge { - g := prometheus.NewGauge(prometheus.GaugeOpts{Name: name, Help: help}) - prometheus.MustRegister(g) - return g +// Metrics holds all Prometheus gauge definitions for the exporter. +// Create with NewMetrics(reg) — no global registration. +type Metrics struct { + // Chain / sync + BlockHeight *prometheus.GaugeVec + BlockHeightMax *prometheus.GaugeVec + BlockTimeLocal *prometheus.GaugeVec // timestamp from local node (/status) + BlockTimeNetwork *prometheus.GaugeVec // timestamp from public network nodes + CatchingUp *prometheus.GaugeVec + ChainEpoch *prometheus.GaugeVec + + // Node hardware / status + NodeStatus *prometheus.GaugeVec + NodeIntended *prometheus.GaugeVec + PocCurrent *prometheus.GaugeVec + PocIntended *prometheus.GaugeVec + NodePocWeight *prometheus.GaugeVec + NodeTimeslot *prometheus.GaugeVec + NodeGPUCount *prometheus.GaugeVec + NodeGPUUtil *prometheus.GaugeVec + NodeHardwareInfo *prometheus.GaugeVec + + // Node GPU — per device + NodeGPUDeviceUtil *prometheus.GaugeVec + NodeGPUDeviceTemp *prometheus.GaugeVec + NodeGPUDeviceMemTotal *prometheus.GaugeVec + NodeGPUDeviceMemFree *prometheus.GaugeVec + NodeGPUDeviceMemUsed *prometheus.GaugeVec + NodeGPUDeviceAvail *prometheus.GaugeVec + + // Node ML service state + NodeServiceState *prometheus.GaugeVec + NodeDiskAvailableGB *prometheus.GaugeVec + + // Network-wide + NetParticipantWeight *prometheus.GaugeVec + NetNodePocWeight *prometheus.GaugeVec + NetTotalWeight *prometheus.GaugeVec + NetRewardPerWeight *prometheus.GaugeVec + + // Pricing / models + PricingUoC prometheus.Gauge + PricingDynamic prometheus.Gauge + ModelPrice *prometheus.GaugeVec + ModelUnits *prometheus.GaugeVec + ModelVRAM *prometheus.GaugeVec + ModelThroughput *prometheus.GaugeVec + ModelValThresh *prometheus.GaugeVec + + // Participant live (current epoch) + ParticipantEpochsDone *prometheus.GaugeVec + ParticipantCoinBalance *prometheus.GaugeVec + ParticipantWallet *prometheus.GaugeVec + ParticipantInferences *prometheus.GaugeVec + ParticipantMissed *prometheus.GaugeVec + ParticipantEarnedCoins *prometheus.GaugeVec + ParticipantValidated *prometheus.GaugeVec + ParticipantInvalidated *prometheus.GaugeVec + + // Participant health (extended) + ParticipantStatus *prometheus.GaugeVec + ParticipantConsecutiveInv *prometheus.GaugeVec + ParticipantBurnedCoins *prometheus.GaugeVec + ParticipantRewardedCoins *prometheus.GaugeVec + ParticipantReputation *prometheus.GaugeVec + + // Network — counts and epoch-level + NetActiveParticipantCount prometheus.Gauge + NetTotalParticipantCount prometheus.Gauge + NetEpochInferenceCount *prometheus.GaugeVec + + // BLS DKG phase + BLSDKGPhase *prometheus.GaugeVec + BLSDealingDeadline *prometheus.GaugeVec + BLSVerifyingDeadline *prometheus.GaugeVec + + // Model utilization / capacity + ModelUtilization *prometheus.GaugeVec + ModelCapacity *prometheus.GaugeVec + + // Epoch history (label epoch = chain epoch number as string) + EpochInferences *prometheus.GaugeVec + EpochMissed *prometheus.GaugeVec + EpochEarnedCoins *prometheus.GaugeVec + EpochValidated *prometheus.GaugeVec + EpochInvalidated *prometheus.GaugeVec + EpochCoinBalance *prometheus.GaugeVec + EpochDone *prometheus.GaugeVec + EpochMissRate *prometheus.GaugeVec + EpochPocWeight *prometheus.GaugeVec + EpochTimeslot *prometheus.GaugeVec + EpochStartTime *prometheus.GaugeVec + EpochEndTime *prometheus.GaugeVec + EpochDuration *prometheus.GaugeVec + EpochEarnedGNK *prometheus.GaugeVec + EpochRewardedGNK *prometheus.GaugeVec + EpochClaimed *prometheus.GaugeVec + EpochEstimated *prometheus.GaugeVec + + // Stats — network-wide inference statistics + StatsAiTokens prometheus.Gauge + StatsInferences prometheus.Gauge + StatsActualCost prometheus.Gauge + StatsModelAiTokens *prometheus.GaugeVec + StatsModelInferences *prometheus.GaugeVec + + // Bridge — queue status for the Gonka bridge + BridgePendingBlocks prometheus.Gauge + BridgePendingReceipts prometheus.Gauge + BridgeReadyToProcess prometheus.Gauge + BridgeEarliestBlock prometheus.Gauge + BridgeLatestBlock prometheus.Gauge + + // Node managers — running/healthy state per ML node manager + NodeManagerRunning *prometheus.GaugeVec + NodeManagerHealthy *prometheus.GaugeVec + + // GPU driver info — version in labels, value always 1 + NodeGPUDriverInfo *prometheus.GaugeVec + + // Tokenomics — chain-wide token flow counters + TokenomicsTotalFees prometheus.Gauge + TokenomicsTotalSubsidies prometheus.Gauge + TokenomicsTotalRefunded prometheus.Gauge + TokenomicsTotalBurned prometheus.Gauge + + // PoC v2 — proof-of-compute artifact and weight data + PoCv2ArtifactCount *prometheus.GaugeVec + PoCv2NodeWeight *prometheus.GaugeVec } -func gaugeVec(name, help string, labels ...string) *prometheus.GaugeVec { - g := prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: name, Help: help}, labels) - prometheus.MustRegister(g) - return g +// NewMetrics creates all Prometheus metrics and registers them with reg. +// Pass prometheus.DefaultRegisterer in production, prometheus.NewRegistry() in tests. +func NewMetrics(reg prometheus.Registerer) *Metrics { + m := &Metrics{ + // Chain / sync + BlockHeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_block_height", Help: "Latest block height from local node"}, []string{"participant"}), + BlockHeightMax: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_block_height_max", Help: "Maximum block height seen across public nodes"}, []string{"participant"}), + BlockTimeLocal: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_block_time_local_seconds", Help: "Timestamp of latest block from local node (unix)"}, []string{"participant"}), + BlockTimeNetwork: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_block_time_network_seconds", Help: "Timestamp of latest block from public network nodes (unix)"}, []string{"participant"}), + CatchingUp: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_chain_catching_up", Help: "1 = syncing, 0 = synced"}, []string{"participant"}), + ChainEpoch: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_chain_epoch", Help: "Current global chain epoch number"}, []string{"participant"}), + + // Node hardware / status + NodeStatus: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_status", Help: "Node hardware status (0=UNKNOWN 1=INFERENCE 2=POC 3=TRAINING 4=STOPPED 5=FAILED)"}, []string{"participant", "node_id", "host"}), + NodeIntended: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_intended_status", Help: "Node intended status"}, []string{"participant", "node_id", "host"}), + PocCurrent: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_poc_current_status", Help: "PoC current status (0=IDLE 1=GENERATING 2=VALIDATING)"}, []string{"participant", "node_id", "host"}), + PocIntended: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_poc_intended_status", Help: "PoC intended status"}, []string{"participant", "node_id", "host"}), + NodePocWeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_poc_weight", Help: "PoC weight per node per model"}, []string{"participant", "node_id", "host", "model"}), + NodeTimeslot: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_poc_timeslot_assigned", Help: "Timeslot assigned (1/0)"}, []string{"participant", "node_id", "host", "model"}), + NodeGPUCount: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_device_count", Help: "GPU device count"}, []string{"participant", "node_id", "host"}), + NodeGPUUtil: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_avg_utilization_percent", Help: "Average GPU utilization %"}, []string{"participant", "node_id", "host"}), + NodeHardwareInfo: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_hardware_info", Help: "Hardware info (value=1, metadata in labels)"}, []string{"participant", "node_id", "host", "hardware_type", "hardware_count"}), + + // Node GPU — per device + NodeGPUDeviceUtil: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_device_utilization_percent", Help: "Per-device GPU compute utilization %"}, []string{"participant", "node_id", "host", "device_index"}), + NodeGPUDeviceTemp: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_device_temperature_celsius", Help: "Per-device GPU temperature °C"}, []string{"participant", "node_id", "host", "device_index"}), + NodeGPUDeviceMemTotal: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_device_memory_total_mb", Help: "Per-device GPU total memory MB"}, []string{"participant", "node_id", "host", "device_index"}), + NodeGPUDeviceMemFree: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_device_memory_free_mb", Help: "Per-device GPU free memory MB"}, []string{"participant", "node_id", "host", "device_index"}), + NodeGPUDeviceMemUsed: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_device_memory_used_mb", Help: "Per-device GPU used memory MB"}, []string{"participant", "node_id", "host", "device_index"}), + NodeGPUDeviceAvail: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_device_available", Help: "Per-device GPU available (1=yes 0=no)"}, []string{"participant", "node_id", "host", "device_index"}), + + // Node ML service state + NodeServiceState: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_service_state", Help: "ML node service state (0=STOPPED 1=INFERENCE 2=POW 3=TRAIN)"}, []string{"participant", "node_id", "host"}), + NodeDiskAvailableGB: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_disk_available_gb", Help: "ML node model cache available disk space GB"}, []string{"participant", "node_id", "host"}), + + // Network-wide + NetParticipantWeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_network_participant_weight", Help: "Per-participant weight in active epoch"}, []string{"participant"}), + NetNodePocWeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_network_node_poc_weight", Help: "Per-node PoC weight in active epoch"}, []string{"participant", "node_id"}), + NetTotalWeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_network_total_weight", Help: "Total weight of all participants"}, []string{"participant"}), + NetRewardPerWeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_network_reward_per_weight", Help: "Estimated GNK per unit of weight"}, []string{"participant"}), + + // Pricing / models + PricingUoC: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_pricing_unit_of_compute_price", Help: "Unit of compute price"}), + PricingDynamic: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_pricing_dynamic_enabled", Help: "Dynamic pricing enabled (1/0)"}), + ModelPrice: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_pricing_model_price_per_token", Help: "Price per token per model"}, []string{"model_id"}), + ModelUnits: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_pricing_model_units_per_token", Help: "Compute units per token"}, []string{"model_id"}), + ModelVRAM: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_model_v_ram", Help: "VRAM (GB)"}, []string{"model_id"}), + ModelThroughput: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_model_throughput_per_nonce", Help: "Throughput per nonce"}, []string{"model_id"}), + ModelValThresh: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_model_validation_threshold", Help: "Validation threshold"}, []string{"model_id"}), + + // Participant live (current epoch) + ParticipantEpochsDone: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_epochs_completed", Help: "Personal epochs completed"}, []string{"participant"}), + ParticipantCoinBalance: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_coin_balance", Help: "Coin balance (internal points)"}, []string{"participant"}), + ParticipantWallet: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_wallet_balance_gonka", Help: "Wallet balance in GNK (ngonka/1e9)"}, []string{"participant"}), + ParticipantInferences: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_inference_count", Help: "Inferences in current epoch"}, []string{"participant"}), + ParticipantMissed: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_missed_requests", Help: "Missed requests in current epoch"}, []string{"participant"}), + ParticipantEarnedCoins: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_earned_coins", Help: "Earned coins in current epoch"}, []string{"participant"}), + ParticipantValidated: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_validated_inferences", Help: "Validated inferences in current epoch"}, []string{"participant"}), + ParticipantInvalidated: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_invalidated_inferences", Help: "Invalidated inferences in current epoch"}, []string{"participant"}), + + // Participant health (extended) + ParticipantStatus: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_status", Help: "Participant status (0=UNSPECIFIED 1=ACTIVE 2=INACTIVE 3=INVALID 5=UNCONFIRMED)"}, []string{"participant"}), + ParticipantConsecutiveInv: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_consecutive_invalid_inferences", Help: "Consecutive invalid inferences counter"}, []string{"participant"}), + ParticipantBurnedCoins: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_burned_coins", Help: "Burned (penalized) coins in current epoch"}, []string{"participant"}), + ParticipantRewardedCoins: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_rewarded_coins", Help: "Rewarded coins in current epoch (after distribution)"}, []string{"participant"}), + ParticipantReputation: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_participant_reputation", Help: "Participant reputation score from epoch group data"}, []string{"participant"}), + + // Network — counts and epoch-level + NetActiveParticipantCount: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_network_active_participant_count", Help: "Number of participants with a non-empty address in the current epoch response"}), + NetTotalParticipantCount: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_network_total_participant_count", Help: "Total number of participants in current epoch"}), + NetEpochInferenceCount: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_network_epoch_inference_count", Help: "Total network inferences in current epoch"}, []string{"participant"}), + + // BLS DKG phase + BLSDKGPhase: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_bls_dkg_phase", Help: "BLS DKG phase (0=UNDEFINED 1=DEALING 2=VERIFYING 3=COMPLETED 4=FAILED 5=SIGNED)"}, []string{"participant"}), + BLSDealingDeadline: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_bls_dealing_deadline_block", Help: "BLS dealing phase deadline block height"}, []string{"participant"}), + BLSVerifyingDeadline: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_bls_verifying_deadline_block", Help: "BLS verifying phase deadline block height"}, []string{"participant"}), + + // Model utilization / capacity + ModelUtilization: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_model_utilization_percent", Help: "Model utilization % (inferences/capacity)"}, []string{"model_id"}), + ModelCapacity: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_model_capacity", Help: "Model maximum capacity (AI tokens/epoch)"}, []string{"model_id"}), + + // Epoch history (label epoch = chain epoch number as string) + EpochInferences: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_inference_count", Help: "Inferences for epoch N"}, []string{"participant", "epoch"}), + EpochMissed: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_missed_requests", Help: "Missed requests for epoch N"}, []string{"participant", "epoch"}), + EpochEarnedCoins: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_earned_coins", Help: "Earned coins (internal points) for epoch N"}, []string{"participant", "epoch"}), + EpochValidated: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_validated_inferences", Help: "Validated inferences for epoch N"}, []string{"participant", "epoch"}), + EpochInvalidated: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_invalidated_inferences", Help: "Invalidated inferences for epoch N"}, []string{"participant", "epoch"}), + EpochCoinBalance: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_coin_balance", Help: "Coin balance at epoch end"}, []string{"participant", "epoch"}), + EpochDone: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_epochs_completed", Help: "Personal epochs completed at epoch end"}, []string{"participant", "epoch"}), + EpochMissRate: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_miss_rate_percent", Help: "Miss rate %% for epoch N"}, []string{"participant", "epoch"}), + EpochPocWeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_poc_weight", Help: "PoC weight for epoch N"}, []string{"participant", "epoch", "node_id"}), + EpochTimeslot: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_timeslot_assigned", Help: "Timeslot assigned in epoch N (1/0)"}, []string{"participant", "epoch"}), + EpochStartTime: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_start_time", Help: "Epoch start unix timestamp"}, []string{"participant", "epoch"}), + EpochEndTime: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_end_time", Help: "Epoch end unix timestamp (estimated for live)"}, []string{"participant", "epoch"}), + EpochDuration: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_duration_seconds", Help: "Real epoch duration (seconds)"}, []string{"participant", "epoch"}), + EpochEarnedGNK: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_earned_gonka", Help: "Wallet balance delta for epoch (GNK)"}, []string{"participant", "epoch"}), + EpochRewardedGNK: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_rewarded_gonka", Help: "On-chain rewarded GNK for epoch N (rewarded_coins/1e9)"}, []string{"participant", "epoch"}), + EpochClaimed: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_claimed", Help: "1 if epoch reward was claimed, 0 if not"}, []string{"participant", "epoch"}), + EpochEstimated: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_epoch_estimated_reward_gonka", Help: "Estimated GNK reward for epoch N (weight × emission/total_weight)"}, []string{"participant", "epoch"}), + + // Stats — network-wide inference statistics + StatsAiTokens: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_stats_ai_tokens_total", Help: "Total AI tokens processed network-wide (cumulative)"}), + StatsInferences: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_stats_inferences_total", Help: "Total inferences processed network-wide (cumulative)"}), + StatsActualCost: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_stats_actual_cost_total", Help: "Total actual inference cost network-wide (coins, cumulative)"}), + StatsModelAiTokens: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_stats_model_ai_tokens", Help: "AI tokens processed per model (cumulative)"}, []string{"model_id"}), + StatsModelInferences: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_stats_model_inferences", Help: "Inferences processed per model (cumulative)"}, []string{"model_id"}), + + // Bridge — queue status for the Gonka bridge + BridgePendingBlocks: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_bridge_pending_blocks", Help: "Bridge pending block count"}), + BridgePendingReceipts: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_bridge_pending_receipts", Help: "Bridge pending receipt count"}), + BridgeReadyToProcess: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_bridge_ready_to_process", Help: "Bridge ready to process (1=yes 0=no)"}), + BridgeEarliestBlock: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_bridge_earliest_block_number", Help: "Bridge earliest pending block number"}), + BridgeLatestBlock: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_bridge_latest_block_number", Help: "Bridge latest pending block number"}), + + // Node managers — running/healthy state per ML node manager + NodeManagerRunning: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_manager_running", Help: "ML node manager running (1=yes 0=no)"}, []string{"participant", "node_id", "host", "manager"}), + NodeManagerHealthy: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_manager_healthy", Help: "ML node manager healthy (1=yes 0=no)"}, []string{"participant", "node_id", "host", "manager"}), + + // GPU driver info — version in labels, value always 1 + NodeGPUDriverInfo: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_node_gpu_driver_info", Help: "GPU driver info (value=1, version in labels)"}, []string{"participant", "node_id", "host", "driver_version", "cuda_version"}), + + // Tokenomics — chain-wide token flow counters + TokenomicsTotalFees: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_tokenomics_total_fees", Help: "Total fees collected (ngonka)"}), + TokenomicsTotalSubsidies: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_tokenomics_total_subsidies", Help: "Total subsidies issued (ngonka)"}), + TokenomicsTotalRefunded: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_tokenomics_total_refunded", Help: "Total amount refunded (ngonka)"}), + TokenomicsTotalBurned: prometheus.NewGauge(prometheus.GaugeOpts{Name: "gonka_tokenomics_total_burned", Help: "Total amount burned (ngonka)"}), + + // PoC v2 — proof-of-compute artifact and weight data + PoCv2ArtifactCount: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_poc_v2_artifact_count", Help: "PoC v2 artifact count from store commit"}, []string{"participant"}), + PoCv2NodeWeight: prometheus.NewGaugeVec(prometheus.GaugeOpts{Name: "gonka_poc_v2_node_weight", Help: "PoC v2 per-node weight distribution"}, []string{"participant", "node_id"}), + } + + reg.MustRegister( + m.BlockHeight, + m.BlockHeightMax, + m.BlockTimeLocal, + m.BlockTimeNetwork, + m.CatchingUp, + m.ChainEpoch, + m.NodeStatus, + m.NodeIntended, + m.PocCurrent, + m.PocIntended, + m.NodePocWeight, + m.NodeTimeslot, + m.NodeGPUCount, + m.NodeGPUUtil, + m.NodeHardwareInfo, + m.NodeGPUDeviceUtil, + m.NodeGPUDeviceTemp, + m.NodeGPUDeviceMemTotal, + m.NodeGPUDeviceMemFree, + m.NodeGPUDeviceMemUsed, + m.NodeGPUDeviceAvail, + m.NodeServiceState, + m.NodeDiskAvailableGB, + m.NetParticipantWeight, + m.NetNodePocWeight, + m.NetTotalWeight, + m.NetRewardPerWeight, + m.PricingUoC, + m.PricingDynamic, + m.ModelPrice, + m.ModelUnits, + m.ModelVRAM, + m.ModelThroughput, + m.ModelValThresh, + m.ParticipantEpochsDone, + m.ParticipantCoinBalance, + m.ParticipantWallet, + m.ParticipantInferences, + m.ParticipantMissed, + m.ParticipantEarnedCoins, + m.ParticipantValidated, + m.ParticipantInvalidated, + m.ParticipantStatus, + m.ParticipantConsecutiveInv, + m.ParticipantBurnedCoins, + m.ParticipantRewardedCoins, + m.ParticipantReputation, + m.NetActiveParticipantCount, + m.NetTotalParticipantCount, + m.NetEpochInferenceCount, + m.BLSDKGPhase, + m.BLSDealingDeadline, + m.BLSVerifyingDeadline, + m.ModelUtilization, + m.ModelCapacity, + m.EpochInferences, + m.EpochMissed, + m.EpochEarnedCoins, + m.EpochValidated, + m.EpochInvalidated, + m.EpochCoinBalance, + m.EpochDone, + m.EpochMissRate, + m.EpochPocWeight, + m.EpochTimeslot, + m.EpochStartTime, + m.EpochEndTime, + m.EpochDuration, + m.EpochEarnedGNK, + m.EpochRewardedGNK, + m.EpochClaimed, + m.EpochEstimated, + m.StatsAiTokens, + m.StatsInferences, + m.StatsActualCost, + m.StatsModelAiTokens, + m.StatsModelInferences, + m.BridgePendingBlocks, + m.BridgePendingReceipts, + m.BridgeReadyToProcess, + m.BridgeEarliestBlock, + m.BridgeLatestBlock, + m.NodeManagerRunning, + m.NodeManagerHealthy, + m.NodeGPUDriverInfo, + m.TokenomicsTotalFees, + m.TokenomicsTotalSubsidies, + m.TokenomicsTotalRefunded, + m.TokenomicsTotalBurned, + m.PoCv2ArtifactCount, + m.PoCv2NodeWeight, + ) + + return m } diff --git a/internal/metrics/metrics_test.go b/internal/metrics/metrics_test.go new file mode 100644 index 0000000..6d9b042 --- /dev/null +++ b/internal/metrics/metrics_test.go @@ -0,0 +1,106 @@ +package metrics + +import ( + "testing" + + "github.com/prometheus/client_golang/prometheus" +) + +// TestNewMetrics_NoRegistrationPanic verifies that NewMetrics registers all +// metrics without panicking and returns a fully-initialised *Metrics. +// If any metric name is duplicated or invalid, MustRegister panics. +func TestNewMetrics_NoRegistrationPanic(t *testing.T) { + reg := prometheus.NewRegistry() + defer func() { + if r := recover(); r != nil { + t.Fatalf("NewMetrics panicked: %v", r) + } + }() + m := NewMetrics(reg) + if m == nil { + t.Fatal("NewMetrics returned nil") + } +} + +// TestNewMetrics_IndependentRegistries verifies that two separate registries +// can each host a full Metrics set without collision — no shared global state. +func TestNewMetrics_IndependentRegistries(t *testing.T) { + reg1 := prometheus.NewRegistry() + reg2 := prometheus.NewRegistry() + + // Both must succeed without panic. + m1 := NewMetrics(reg1) + m2 := NewMetrics(reg2) + + if m1 == nil || m2 == nil { + t.Fatal("one of the Metrics instances is nil") + } + + // Write different values into each registry and confirm Gather succeeds. + m1.ChainEpoch.WithLabelValues("addr1").Set(42) + m2.ChainEpoch.WithLabelValues("addr1").Set(7) + + if _, err := reg1.Gather(); err != nil { + t.Fatalf("reg1.Gather() error: %v", err) + } + if _, err := reg2.Gather(); err != nil { + t.Fatalf("reg2.Gather() error: %v", err) + } +} + +// TestNewMetrics_NonVecGaugesExist verifies that plain (non-vec) Gauge metrics +// are immediately present in Gather output after registration — no values need +// to be set for them to appear. +func TestNewMetrics_NonVecGaugesExist(t *testing.T) { + reg := prometheus.NewRegistry() + NewMetrics(reg) + + mfs, err := reg.Gather() + if err != nil { + t.Fatalf("Gather() error: %v", err) + } + + names := make(map[string]bool, len(mfs)) + for _, mf := range mfs { + names[mf.GetName()] = true + } + + // Spot-check a few well-known plain-Gauge metric names (non-vec). + required := []string{ + "gonka_network_active_participant_count", + "gonka_network_total_participant_count", + "gonka_stats_ai_tokens_total", + "gonka_bridge_pending_blocks", + "gonka_tokenomics_total_fees", + } + for _, name := range required { + if !names[name] { + t.Errorf("metric %q not found in Gather output", name) + } + } +} + +// TestNewMetrics_VecGaugeRecordsValue verifies that GaugeVec metrics work +// end-to-end: set a value, gather, confirm it's present with correct value. +func TestNewMetrics_VecGaugeRecordsValue(t *testing.T) { + reg := prometheus.NewRegistry() + m := NewMetrics(reg) + + m.ChainEpoch.WithLabelValues("test_addr").Set(123) + + mfs, err := reg.Gather() + if err != nil { + t.Fatalf("Gather() error: %v", err) + } + + for _, mf := range mfs { + if mf.GetName() == "gonka_chain_epoch" { + for _, mm := range mf.GetMetric() { + if mm.GetGauge().GetValue() == 123 { + return // found + } + } + } + } + t.Fatal("gonka_chain_epoch{participant=\"test_addr\"} = 123 not found in Gather output") +} diff --git a/internal/state/state.go b/internal/state/state.go index 686b588..a54ec71 100644 --- a/internal/state/state.go +++ b/internal/state/state.go @@ -132,6 +132,8 @@ func LoadHistory(path string) History { } } + // TODO: remove after all deployments migrate to format_version=2 + // Safe to delete after 2026-Q3 (one cycle after this version was deployed). // Old flat format: map[epoch]*EpochSnapshot — migrate by grouping on snap.Participant. var oldH map[string]*EpochSnapshot if err := json.Unmarshal(data, &oldH); err != nil { diff --git a/internal/state/state_test.go b/internal/state/state_test.go new file mode 100644 index 0000000..a9611c4 --- /dev/null +++ b/internal/state/state_test.go @@ -0,0 +1,167 @@ +package state + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "testing" +) + +// --- MissRate --- + +func TestMissRate_ZeroTotal(t *testing.T) { + if got := MissRate(0, 0); got != 0 { + t.Fatalf("MissRate(0,0) = %v, want 0", got) + } +} + +func TestMissRate_NoMisses(t *testing.T) { + if got := MissRate(100, 0); got != 0 { + t.Fatalf("MissRate(100,0) = %v, want 0", got) + } +} + +func TestMissRate_AllMissed(t *testing.T) { + if got := MissRate(0, 50); got != 100 { + t.Fatalf("MissRate(0,50) = %v, want 100", got) + } +} + +func TestMissRate_Half(t *testing.T) { + if got := MissRate(50, 50); got != 50 { + t.Fatalf("MissRate(50,50) = %v, want 50", got) + } +} + +func TestMissRate_Rounded(t *testing.T) { + // 1/3 ≈ 33.33% + got := MissRate(2, 1) + if got != 33.33 { + t.Fatalf("MissRate(2,1) = %v, want 33.33", got) + } +} + +// --- LoadHistory / SaveHistory --- + +func TestLoadHistory_EmptyFile(t *testing.T) { + dir := t.TempDir() + h := LoadHistory(filepath.Join(dir, "nonexistent.json")) + if h == nil { + t.Fatal("expected non-nil History map") + } + if len(h) != 0 { + t.Fatalf("expected empty map, got %d entries", len(h)) + } +} + +func TestLoadHistory_NewFormat(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "history.json") + + h := History{ + "addr1": { + "10": &EpochSnapshot{Participant: "addr1", InferenceCount: 100}, + "11": &EpochSnapshot{Participant: "addr1", InferenceCount: 200}, + }, + } + data, err := json.Marshal(h) + if err != nil { + t.Fatalf("json.Marshal failed: %v", err) + } + if err := os.WriteFile(path, data, 0644); err != nil { + t.Fatalf("os.WriteFile failed: %v", err) + } + + loaded := LoadHistory(path) + if len(loaded) != 1 { + t.Fatalf("expected 1 participant, got %d", len(loaded)) + } + if len(loaded["addr1"]) != 2 { + t.Fatalf("expected 2 epochs, got %d", len(loaded["addr1"])) + } + if loaded["addr1"]["10"].InferenceCount != 100 { + t.Fatalf("wrong inference count: %d", loaded["addr1"]["10"].InferenceCount) + } +} + +func TestLoadHistory_OldFlatFormat_Migration(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "history.json") + + // Old format: map[epoch_string]*EpochSnapshot (keys are numeric strings) + old := map[string]*EpochSnapshot{ + "5": {Participant: "addr_old", InferenceCount: 42}, + "6": {Participant: "addr_old", InferenceCount: 55}, + } + data, err := json.Marshal(old) + if err != nil { + t.Fatalf("json.Marshal failed: %v", err) + } + if err := os.WriteFile(path, data, 0644); err != nil { + t.Fatalf("os.WriteFile failed: %v", err) + } + + loaded := LoadHistory(path) + if len(loaded) != 1 { + t.Fatalf("migration: expected 1 participant, got %d", len(loaded)) + } + epochs := loaded["addr_old"] + if len(epochs) != 2 { + t.Fatalf("migration: expected 2 epochs, got %d", len(epochs)) + } + if epochs["5"].InferenceCount != 42 { + t.Fatalf("migration: wrong inference count for epoch 5: %d", epochs["5"].InferenceCount) + } +} + +func TestSaveHistory_Pruning(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "history.json") + + // Build 10 epochs for one participant, max=3 + h := History{"addr1": make(map[string]*EpochSnapshot)} + for i := 1; i <= 10; i++ { + key := fmt.Sprintf("%d", i) + h["addr1"][key] = &EpochSnapshot{Participant: "addr1", InferenceCount: int64(i)} + } + + SaveHistory(path, h, 3) + + loaded := LoadHistory(path) + if len(loaded["addr1"]) != 3 { + t.Fatalf("pruning: expected 3 epochs after maxEntries=3, got %d", len(loaded["addr1"])) + } + // Must keep the 3 highest epoch numbers (8, 9, 10) + for _, epoch := range []string{"8", "9", "10"} { + if loaded["addr1"][epoch] == nil { + t.Fatalf("pruning: expected epoch %s to be retained", epoch) + } + } +} + +func TestSaveHistory_Roundtrip(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "history.json") + + snap := &EpochSnapshot{ + Participant: "addr_test", + InferenceCount: 999, + MissRatePercent: 12.5, + } + h := History{"addr_test": {"42": snap}} + + SaveHistory(path, h, 500) + loaded := LoadHistory(path) + + got := loaded["addr_test"]["42"] + if got == nil { + t.Fatal("roundtrip: epoch 42 not found") + } + if got.InferenceCount != 999 { + t.Fatalf("roundtrip: InferenceCount = %d, want 999", got.InferenceCount) + } + if got.MissRatePercent != 12.5 { + t.Fatalf("roundtrip: MissRatePercent = %v, want 12.5", got.MissRatePercent) + } +}