Skip to content

Commit c920916

Browse files
authored
Merge pull request #529 from NilFoundation/prom-metrics
[telemetry] Export prometheus metrics (used in libp2p)
2 parents b744b8d + 26403c2 commit c920916

8 files changed

Lines changed: 111 additions & 64 deletions

File tree

go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@ require (
5353
github.com/icza/bitio v1.1.0
5454
github.com/ipfs/go-log/v2 v2.5.1
5555
github.com/klauspost/compress v1.18.0
56+
github.com/prometheus/client_golang v1.21.0
5657
github.com/shopspring/decimal v1.4.0
5758
github.com/spf13/viper v1.19.0
5859
go.dedis.ch/kyber/v3 v3.1.0
@@ -198,7 +199,6 @@ require (
198199
github.com/pion/turn/v4 v4.0.0 // indirect
199200
github.com/pion/webrtc/v4 v4.0.12 // indirect
200201
github.com/polydawn/refmt v0.89.0 // indirect
201-
github.com/prometheus/client_golang v1.21.0 // indirect
202202
github.com/prometheus/client_model v0.6.1 // indirect
203203
github.com/prometheus/common v0.62.0 // indirect
204204
github.com/prometheus/procfs v0.15.1 // indirect

nil/cmd/nild/devnet.go

Lines changed: 77 additions & 57 deletions
Original file line numberDiff line numberDiff line change
@@ -31,11 +31,12 @@ type nodeSpec struct {
3131
ArchiveNodeIndices []int `yaml:"archiveNodeIndices"`
3232
}
3333

34-
type devnetSpec struct {
34+
type clusterSpec struct {
3535
NilServerName string `yaml:"nil_server_name"`
3636
NilCertEmail string `yaml:"nil_cert_email"`
3737
NildConfigDir string `yaml:"nild_config_dir"`
3838
NildCredentialsDir string `yaml:"nild_credentials_dir"`
39+
NildPromBasePort int `yaml:"nild_prom_base_port"`
3940
NildP2PBaseTCPPort int `yaml:"nild_p2p_base_tcp_port"`
4041
PprofBaseTCPPort int `yaml:"pprof_base_tcp_port"`
4142
NilWipeOnUpdate bool `yaml:"nil_wipe_on_update"`
@@ -68,7 +69,8 @@ type server struct {
6869
service string
6970
name string
7071
identity string
71-
port int
72+
p2pPort int
73+
promPort int
7274
pprofPort int
7375
rpcPort int
7476
credsDir string
@@ -78,8 +80,8 @@ type server struct {
7880
vkm *keys.ValidatorKeysManager
7981
}
8082

81-
type devnet struct {
82-
spec *devnetSpec
83+
type cluster struct {
84+
spec *clusterSpec
8385
baseDir string
8486
validators []server
8587
archivers []server
@@ -104,15 +106,15 @@ func validatorKeysFile(credsDir string) string {
104106
return filepath.Join(credsDir, "validator-keys.yaml")
105107
}
106108

107-
func (spec *devnetSpec) ensureValidatorKeys(srv *server) (*keys.ValidatorKeysManager, error) {
109+
func (spec *clusterSpec) ensureValidatorKeys(srv *server) (*keys.ValidatorKeysManager, error) {
108110
vkm := keys.NewValidatorKeyManager(validatorKeysFile(srv.credsDir))
109111
if err := vkm.InitKey(); err != nil {
110112
return nil, err
111113
}
112114
return vkm, nil
113115
}
114116

115-
func (devnet devnet) generateZeroState(nShards uint32, servers []server) (*execution.ZeroStateConfig, error) {
117+
func (c *cluster) generateZeroState(nShards uint32, servers []server) (*execution.ZeroStateConfig, error) {
116118
validators := make([]config.ListValidators, nShards-1)
117119
for _, srv := range servers {
118120
key, err := srv.vkm.GetPublicKey()
@@ -131,7 +133,7 @@ func (devnet devnet) generateZeroState(nShards uint32, servers []server) (*execu
131133
}
132134
}
133135

134-
mainKeyPath := devnet.spec.NildCredentialsDir + "/keys.yaml"
136+
mainKeyPath := c.spec.NildCredentialsDir + "/keys.yaml"
135137
mainPublicKey, err := ensurePublicKey(mainKeyPath)
136138
if err != nil {
137139
return nil, err
@@ -183,7 +185,7 @@ func genDevnet(cmd *cobra.Command, args []string) error {
183185
return fmt.Errorf("can't read devnet spec %s: %w", specFile, err)
184186
}
185187

186-
spec := &devnetSpec{}
188+
spec := &clusterSpec{}
187189
if err := yaml.Unmarshal(specYaml, spec); err != nil {
188190
return fmt.Errorf("can't parse devnet spec %s: %w", specFile, err)
189191
}
@@ -192,23 +194,32 @@ func genDevnet(cmd *cobra.Command, args []string) error {
192194
if spec.EnableRPCOnValidators {
193195
validatorRPCBasePort = spec.NilRPCPort + len(spec.NilRPCConfig)
194196
}
195-
validators, err := spec.makeServers(spec.NilConfig, spec.NildP2PBaseTCPPort, spec.PprofBaseTCPPort, validatorRPCBasePort, "nil", baseDir, false)
197+
validators, err := spec.makeServers(spec.NilConfig,
198+
spec.NildP2PBaseTCPPort, spec.NildPromBasePort, spec.PprofBaseTCPPort, validatorRPCBasePort,
199+
"nil", baseDir, false)
196200
if err != nil {
197201
return fmt.Errorf("failed to setup validator nodes: %w", err)
198202
}
199203

200-
devnet := devnet{spec: spec, baseDir: baseDir, validators: validators}
204+
c := &cluster{spec: spec, baseDir: baseDir, validators: validators}
201205

202206
archiveBaseP2P := spec.NildP2PBaseTCPPort + len(validators)
207+
archiveBaseProm := spec.NildPromBasePort + len(validators)
203208
archiveBasePprof := spec.PprofBaseTCPPort + len(validators)
204209

205-
if devnet.archivers, err = spec.makeServers(spec.NilArchiveConfig, archiveBaseP2P, archiveBasePprof, 0, "nil-archive", baseDir, false); err != nil {
210+
c.archivers, err = spec.makeServers(spec.NilArchiveConfig,
211+
archiveBaseP2P, archiveBaseProm, archiveBasePprof, 0,
212+
"nil-archive", baseDir, false)
213+
if err != nil {
206214
return fmt.Errorf("failed to setup archive nodes: %w", err)
207215
}
208216

209-
rpcBasePprof := spec.PprofBaseTCPPort + len(validators) + len(devnet.archivers)
217+
rpcBasePprof := spec.PprofBaseTCPPort + len(validators) + len(c.archivers)
210218

211-
if devnet.rpcNodes, err = spec.makeServers(spec.NilRPCConfig, 0, rpcBasePprof, spec.NilRPCPort, "nil-rpc", baseDir, true); err != nil {
219+
c.rpcNodes, err = spec.makeServers(spec.NilRPCConfig,
220+
0, 0, rpcBasePprof, spec.NilRPCPort,
221+
"nil-rpc", baseDir, true)
222+
if err != nil {
212223
return fmt.Errorf("failed to setup rpc nodes: %w", err)
213224
}
214225

@@ -217,52 +228,55 @@ func genDevnet(cmd *cobra.Command, args []string) error {
217228
return fmt.Errorf("failed to get only flag: %w", err)
218229
}
219230

220-
if devnet.zeroState, err = devnet.generateZeroState(spec.NShards, devnet.validators); err != nil {
231+
if c.zeroState, err = c.generateZeroState(spec.NShards, c.validators); err != nil {
221232
return err
222233
}
223-
if err := devnet.writeConfigs(devnet.validators, "validator", only); err != nil {
234+
if err := c.writeConfigs(c.validators, "validator", only); err != nil {
224235
return err
225236
}
226-
if err := devnet.writeConfigs(devnet.archivers, "archiver", only); err != nil {
237+
if err := c.writeConfigs(c.archivers, "archiver", only); err != nil {
227238
return err
228239
}
229-
if err := devnet.writeConfigs(devnet.rpcNodes, "RPC node", only); err != nil {
240+
if err := c.writeConfigs(c.rpcNodes, "RPC node", only); err != nil {
230241
return err
231242
}
232243

233244
os.Exit(0)
234245
return nil
235246
}
236247

237-
func (spec *devnetSpec) EnsureValidatorKeys(srv server) (*keys.ValidatorKeysManager, error) {
248+
func (spec *clusterSpec) EnsureValidatorKeys(srv server) (*keys.ValidatorKeysManager, error) {
238249
if err := os.MkdirAll(srv.credsDir, directoryPermissions); err != nil {
239250
return nil, err
240251
}
241252

242253
return spec.ensureValidatorKeys(&srv)
243254
}
244255

245-
func (devnet devnet) writeConfigs(servers []server, name string, only string) error {
256+
func (c *cluster) writeConfigs(servers []server, name string, only string) error {
246257
for i, server := range servers {
247-
if err := devnet.writeServerConfig(i, server, only); err != nil {
258+
if err := c.writeServerConfig(i, server, only); err != nil {
248259
return fmt.Errorf("failed to write %s config: %w", name, err)
249260
}
250261
}
251262
return nil
252263
}
253264

254-
func (spec devnetSpec) makeServers(nodeSpecs []nodeSpec, basePort int, pprofBasePort int, baseHTTPPort int, service string, baseDir string, logClientEvents bool) ([]server, error) {
265+
func (spec *clusterSpec) makeServers(nodeSpecs []nodeSpec, baseP2pPort, basePromPort, basePprofPort, baseHTTPPort int, service string, baseDir string, logClientEvents bool) ([]server, error) {
255266
servers := make([]server, len(nodeSpecs))
256267
for i, nodeSpec := range nodeSpecs {
257268
servers[i].service = service
258269
servers[i].name = fmt.Sprintf("%s-%d", service, i)
259270
servers[i].nodeSpec = nodeSpec
260271
servers[i].logClientEvents = logClientEvents
261-
if basePort != 0 {
262-
servers[i].port = basePort + i
272+
if baseP2pPort != 0 {
273+
servers[i].p2pPort = baseP2pPort + i
263274
}
264-
if pprofBasePort != 0 {
265-
servers[i].pprofPort = pprofBasePort + i
275+
if basePromPort != 0 {
276+
servers[i].promPort = basePromPort + i
277+
}
278+
if basePprofPort != 0 {
279+
servers[i].pprofPort = basePprofPort + i
266280
}
267281
if baseHTTPPort != 0 {
268282
servers[i].rpcPort = baseHTTPPort + i
@@ -292,7 +306,7 @@ const (
292306
filePermissions = 0o644
293307
)
294308

295-
func (spec *devnetSpec) EnsureIdentity(srv server) (string, error) {
309+
func (spec *clusterSpec) EnsureIdentity(srv server) (string, error) {
296310
if err := os.MkdirAll(srv.credsDir, directoryPermissions); err != nil {
297311
return "", err
298312
}
@@ -304,20 +318,43 @@ func (spec *devnetSpec) EnsureIdentity(srv server) (string, error) {
304318
return identity.String(), err
305319
}
306320

307-
func (devnet *devnet) writeServerConfig(instanceId int, srv server, only string) error {
321+
func (c *cluster) writeServerConfig(instanceId int, srv server, only string) error {
308322
if only != "" && srv.name != only {
309323
return nil
310324
}
311325

312-
cfg := nildconfig.Config{}
313-
cfg.Config = &nilservice.Config{
314-
Network: &network.Config{},
315-
Telemetry: &telemetry.Config{},
326+
spec := c.spec
327+
inst := srv.nodeSpec
328+
329+
cfg := nildconfig.Config{
330+
Config: &nilservice.Config{
331+
NShards: spec.NShards,
332+
AllowDbDrop: spec.NilWipeOnUpdate,
333+
LogClientRpcEvents: srv.logClientEvents,
334+
335+
MyShards: inst.Shards,
336+
SplitShards: inst.SplitShards,
337+
BootstrapPeers: getPeers(c.validators, inst.BootstrapPeersIdx),
338+
339+
RPCPort: srv.rpcPort,
340+
PprofPort: srv.pprofPort,
341+
AdminSocketPath: srv.workDir + "/admin_socket",
342+
343+
ZeroState: c.zeroState,
344+
345+
Network: &network.Config{
346+
TcpPort: srv.p2pPort,
347+
348+
DHTEnabled: true,
349+
DHTBootstrapPeers: getPeers(c.validators, inst.DHTBootstrapPeersIdx),
350+
},
351+
Telemetry: &telemetry.Config{
352+
ExportMetrics: true,
353+
PrometheusPort: srv.promPort,
354+
},
355+
},
356+
DB: db.NewDefaultBadgerDBOptions(),
316357
}
317-
cfg.RPCPort = srv.rpcPort
318-
cfg.Network.TcpPort = srv.port
319-
cfg.PprofPort = srv.pprofPort
320-
cfg.ZeroState = devnet.zeroState
321358

322359
var err error
323360
cfg.Network.KeysPath, err = filepath.Abs(srv.NetworkKeysFile())
@@ -330,43 +367,26 @@ func (devnet *devnet) writeServerConfig(instanceId int, srv server, only string)
330367
return fmt.Errorf("failed to get absolute path for validator keys: %w", err)
331368
}
332369

333-
spec := devnet.spec
334-
cfg.NShards = spec.NShards
335-
inst := srv.nodeSpec
336-
cfg.AllowDbDrop = spec.NilWipeOnUpdate
337-
cfg.MyShards = inst.Shards
338-
cfg.SplitShards = inst.SplitShards
339-
cfg.BootstrapPeers = getPeers(devnet.validators, inst.BootstrapPeersIdx)
340-
cfg.AdminSocketPath = srv.workDir + "/admin_socket"
341-
cfg.LogClientRpcEvents = srv.logClientEvents
342-
cfg.DB = db.NewDefaultBadgerDBOptions()
343-
cfg.DB.Path = srv.workDir + "/database"
344-
cfg.Network.DHTEnabled = true
345-
cfg.Network.DHTBootstrapPeers = getPeers(devnet.validators, inst.DHTBootstrapPeersIdx)
346-
347370
if len(inst.ArchiveNodeIndices) > 0 {
348371
cfg.RpcNode = &nilservice.RpcNodeConfig{
349-
ArchiveNodeList: getPeers(devnet.archivers, inst.ArchiveNodeIndices),
372+
ArchiveNodeList: getPeers(c.archivers, inst.ArchiveNodeIndices),
350373
}
351374
}
352375

353-
cfg.Telemetry.ExportMetrics = true
354-
355376
serialized, err := yaml.Marshal(cfg)
356377
if err != nil {
357378
return fmt.Errorf("failed to marshal config: %w", err)
358379
}
359380

360-
err = os.MkdirAll(srv.workDir, directoryPermissions)
361-
if err != nil {
381+
if err := os.MkdirAll(srv.workDir, directoryPermissions); err != nil {
362382
return err
363383
}
364384

365385
configDir := fmt.Sprintf("%s/%s-%d", spec.NildConfigDir, srv.service, instanceId)
366-
err = os.MkdirAll(configDir, directoryPermissions)
367-
if err != nil {
386+
if err := os.MkdirAll(configDir, directoryPermissions); err != nil {
368387
return err
369388
}
389+
370390
return os.WriteFile(configDir+"/nild.yaml", serialized, filePermissions)
371391
}
372392

@@ -376,7 +396,7 @@ func identityToAddress(port int, identity string) string {
376396

377397
func getPeer(srv server) network.AddrInfo {
378398
var peer network.AddrInfo
379-
address := identityToAddress(srv.port, srv.identity)
399+
address := identityToAddress(srv.p2pPort, srv.identity)
380400
check.PanicIfErr(peer.Set(address))
381401
return peer
382402
}

nil/internal/cobrax/cmdflags/telemetry.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,4 +7,5 @@ import (
77

88
func AddTelemetry(fset *pflag.FlagSet, config *telemetry.Config) {
99
fset.BoolVar(&config.ExportMetrics, "metrics", config.ExportMetrics, "export metrics via grpc")
10+
fset.IntVar(&config.PrometheusPort, "prometheus-port", config.PrometheusPort, "port to serve prometheus metrics; 0 to disable")
1011
}

nil/internal/telemetry/internal/config.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,4 +5,6 @@ type Config struct {
55

66
ExportMetrics bool `yaml:"exportMetrics,omitempty"`
77
GrpcEndpoint string `yaml:"grpcEndpoint,omitempty"`
8+
9+
PrometheusPort int `yaml:"prometheusPort,omitempty"`
810
}

nil/internal/telemetry/internal/metric.go

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,8 +13,7 @@ import (
1313
const metricExportInterval = 10 * time.Second
1414

1515
func InitMetrics(ctx context.Context, config *Config) error {
16-
if config == nil || !config.ExportMetrics {
17-
// no metrics
16+
if !config.ExportMetrics {
1817
return nil
1918
}
2019

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
package internal
2+
3+
import (
4+
"net/http"
5+
"strconv"
6+
7+
"github.com/NilFoundation/nil/nil/common/check"
8+
"github.com/prometheus/client_golang/prometheus/promhttp"
9+
)
10+
11+
func StartPrometheusServer(port int) {
12+
if port == 0 {
13+
return
14+
}
15+
16+
server := http.NewServeMux()
17+
server.Handle("/metrics", promhttp.Handler())
18+
go func() {
19+
check.PanicIfErr(http.ListenAndServe(":"+strconv.Itoa(port), server)) //nolint:gosec
20+
}()
21+
}

nil/internal/telemetry/telemetry.go

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -22,10 +22,14 @@ func NewDefaultConfig() *Config {
2222
}
2323

2424
func Init(ctx context.Context, config *Config) error {
25-
if err := internal.InitMetrics(ctx, config); err != nil {
26-
return err
25+
if config == nil {
26+
// no telemetry
27+
return nil
2728
}
28-
return nil
29+
30+
internal.StartPrometheusServer(config.PrometheusPort)
31+
32+
return internal.InitMetrics(ctx, config)
2933
}
3034

3135
func Shutdown(ctx context.Context) {

nix/nil.nix

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ buildGo124Module rec {
4949
];
5050

5151
# to obtain run `nix build` with vendorHash = "";
52-
vendorHash = "sha256-XjIVsjO1duIM/gKmHHHkvlUtkY+ChAQiEqzAMd7in3w=";
52+
vendorHash = "sha256-FHORBSOas4WET3TGJcEw8Ey3gQk2K8FG3elC9tXxo00=";
5353
hardeningDisable = [ "all" ];
5454

5555
postInstall = ''

0 commit comments

Comments
 (0)