diff --git a/.github/workflows/linters.yml b/.github/workflows/linters.yml new file mode 100644 index 0000000..665a8ba --- /dev/null +++ b/.github/workflows/linters.yml @@ -0,0 +1,35 @@ +name: Linters + +on: + pull_request: + merge_group: + +concurrency: + group: ${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }} + cancel-in-progress: true + +permissions: + contents: read + +jobs: + nph: + name: Check NPH formatting + runs-on: ubuntu-latest + steps: + - name: Checkout repository + uses: actions/checkout@v6 + with: + fetch-depth: 2 # In PR, has extra merge commit: ^1 = PR, ^2 = base + persist-credentials: false + + - name: Check NPH formatting + uses: arnetheduck/nph-action@v1 + with: + version: 0.7.0 + # NPH recursively discovers Nim files and applies its default include + # pattern (\.nim, .nims, and .nimble) when given a directory. + # NPH normalizes paths with forward slashes, so this matches a + # nimbledeps directory at every nesting level. + options: "--extend-exclude:/nimbledeps/ ./." + fail: true + suggest: true diff --git a/nim-test-node/connmanager/env.nim b/nim-test-node/connmanager/env.nim index 3f98917..3fa23dc 100644 --- a/nim-test-node/connmanager/env.nim +++ b/nim-test-node/connmanager/env.nim @@ -9,27 +9,30 @@ logScope: type NodeRole* = enum - RoleHub, RolePeer + RoleHub + RolePeer ReconnectMode* = enum - ReconnectNone, ReconnectAggressive, ReconnectBeforeGrace + ReconnectNone + ReconnectAggressive + ReconnectBeforeGrace HubConfig* = object lowWater*: int highWater*: int gracePeriodS*: int silencePeriodS*: int - maxConnections*: int # 0 = no hard cap (Runs A/B/C), >0 adds semaphore (Run D) + maxConnections*: int # 0 = no hard cap (Runs A/B/C), >0 adds semaphore (Run D) protectedPeers*: seq[PeerId] - outboundPeers*: seq[string] # addresses hub dials proactively (Group A in Run A) - numHubs*: int # total hub replicas; >1 triggers hub-to-hub dialing - hubNamespace*: string # k8s namespace used to build peer-hub DNS addresses + outboundPeers*: seq[string] # addresses hub dials proactively (Group A in Run A) + numHubs*: int # total hub replicas; >1 triggers hub-to-hub dialing + hubNamespace*: string # k8s namespace used to build peer-hub DNS addresses PeerConfig* = object - hubAddrs*: seq[string] # one or more hub addresses to connect to (multi-hub support) - dialOut*: bool # true = peer dials hub; false = peer listens, hub dials it + hubAddrs*: seq[string] # one or more hub addresses to connect to (multi-hub support) + dialOut*: bool # true = peer dials hub; false = peer listens, hub dials it reconnect*: ReconnectMode - reconnectIntervalS*: int # for ReconnectBeforeGrace: cycle connection every N seconds + reconnectIntervalS*: int # for ReconnectBeforeGrace: cycle connection every N seconds privateKey*: Opt[PrivateKey] let @@ -49,7 +52,8 @@ proc parseHubConfig*(): HubConfig = let protectedStr = getEnv("PROTECTED_PEERS", "") for entry in protectedStr.split(','): let s = entry.strip() - if s.len == 0: continue + if s.len == 0: + continue let peerId = PeerId.init(s).valueOr: warn "Skipping invalid peer ID in PROTECTED_PEERS", raw = s continue @@ -95,8 +99,15 @@ proc parsePeerConfig*(): PeerConfig = let keys = privKeysStr.split(',') let hostname = getHostname() let parts = hostname.split('-') - let idx = try: parseInt(parts[^1]) except CatchableError: 0 - if idx < keys.len: keys[idx].strip() else: "" + let idx = + try: + parseInt(parts[^1]) + except CatchableError: + 0 + if idx < keys.len: + keys[idx].strip() + else: + "" else: getEnv("PRIVATE_KEY", "") diff --git a/nim-test-node/connmanager/main.nim b/nim-test-node/connmanager/main.nim index 800812e..3f19f14 100644 --- a/nim-test-node/connmanager/main.nim +++ b/nim-test-node/connmanager/main.nim @@ -10,12 +10,13 @@ logScope: proc resolveAndConnect(switch: Switch, address: string): Future[void] {.async.} = var backoff = 1.seconds - for attempt in 1..10: + for attempt in 1 .. 10: try: let addrs = resolveTAddress(address).mapIt(MultiAddress.init(it).tryGet()) if addrs.len == 0: raise newException(CatchableError, "no addresses resolved for " & address) - let peerId = await switch.connect(addrs[0], allowUnknownPeerId = true).wait(10.seconds) + let peerId = + await switch.connect(addrs[0], allowUnknownPeerId = true).wait(10.seconds) info "Connected", address = address, peerId = peerId return except CatchableError as exc: @@ -31,7 +32,8 @@ proc startMetrics() = warn "Failed to start metrics server", error = $res.error return let srv = res.get() - try: waitFor srv.start() + try: + waitFor srv.start() except CatchableError as exc: warn "Metrics server start failed", error = exc.msg @@ -44,10 +46,7 @@ proc runHub(cfg: HubConfig) {.async.} = .withNoise() .withYamux() .withWatermark( - cfg.lowWater, - cfg.highWater, - cfg.gracePeriodS.seconds, - cfg.silencePeriodS.seconds, + cfg.lowWater, cfg.highWater, cfg.gracePeriodS.seconds, cfg.silencePeriodS.seconds ) if cfg.maxConnections > 0: @@ -81,14 +80,20 @@ proc runHub(cfg: HubConfig) {.async.} = if cfg.numHubs > 1: let hostname = getHostname() let parts = hostname.split('-') - let myIdx = try: parseInt(parts[^1]) except CatchableError: 0 - for i in 0..= 2.2.0", - "nimcrypto 0.6.4", - "libp2p#7cc4280e2efd5e6c2ebd732ce33d309376a9627e" + "nimcrypto 0.6.4", "libp2p#7cc4280e2efd5e6c2ebd732ce33d309376a9627e" diff --git a/nim-test-node/gossipsub-queues/env.nim b/nim-test-node/gossipsub-queues/env.nim index b2c455c..ccb4b48 100644 --- a/nim-test-node/gossipsub-queues/env.nim +++ b/nim-test-node/gossipsub-queues/env.nim @@ -6,15 +6,15 @@ logScope: topics = "dst" let - inShadow* = getEnv("SHADOWENV").cmpIgnoreCase("true") == 0 #If Running for shadow simulator + inShadow* = getEnv("SHADOWENV").cmpIgnoreCase("true") == 0 + #If Running for shadow simulator httpPublishPort* = Port(8645) prometheusPort* = Port(8008) myPort* = Port(5000) - chunks* = parseInt(getEnv("FRAGMENTS", "1")) #No. of fragments for each message - + chunks* = parseInt(getEnv("FRAGMENTS", "1")) #No. of fragments for each message proc getPeerDetails*(): Result[(int, int, int, string, string, string), string] = - let + let hostname = getHostname() ordinal = parseInt(hostname.split('-')[^1]) peerIdOffset = parseInt(getEnv("PEER_ID_OFFSET", "0")) @@ -22,19 +22,30 @@ proc getPeerDetails*(): Result[(int, int, int, string, string, string), string] networkSize = parseInt(getEnv("PEERS", "100")) connectTo = parseInt(getEnv("CONNECTTO", "10")) muxer = getEnv("MUXER", "yamux") - filePath = if inShadow: "../" else: getEnv("FILEPATH", "./") - address = if muxer.toLowerAscii() == "quic": - "/ip4/0.0.0.0/udp/" & $myPort & "/quic-v1" - else: - "/ip4/0.0.0.0/tcp/" & $myPort - + filePath = + if inShadow: + "../" + else: + getEnv("FILEPATH", "./") + address = + if muxer.toLowerAscii() == "quic": + "/ip4/0.0.0.0/udp/" & $myPort & "/quic-v1" + else: + "/ip4/0.0.0.0/tcp/" & $myPort + if muxer.toLowerAscii() notin ["quic", "yamux", "mplex"]: return err("Unknown muxer type : " & muxer) if connectTo >= networkSize: - return err("Not enough peers to make target connections. Network size : " & $networkSize) - - info "Host info ", hostname = hostname, peer = myId, muxer = muxer, inShadow = inShadow, address = address + return + err("Not enough peers to make target connections. Network size : " & $networkSize) + + info "Host info ", + hostname = hostname, + peer = myId, + muxer = muxer, + inShadow = inShadow, + address = address return ok((myId, networkSize, connectTo, muxer, filePath, address)) @@ -59,12 +70,13 @@ proc startMetricsServer*( #log metrics if needed (useful for shadow simulations) proc storeMetrics*(myId: int) {.async.} = - await sleepAsync((myId*60).milliseconds) + await sleepAsync((myId * 60).milliseconds) while true: try: - let cmd = "curl -s --connect-timeout 5 --max-time 5 http://localhost:" & - $prometheusPort & "/metrics >> metrics_pod-" & $myId & ".txt" - + let cmd = + "curl -s --connect-timeout 5 --max-time 5 http://localhost:" & $prometheusPort & + "/metrics >> metrics_pod-" & $myId & ".txt" + let exitCode = execCmd(cmd) if exitCode == 0: info "Metrics saved for peer ", pod = myId diff --git a/nim-test-node/gossipsub-queues/main.nim b/nim-test-node/gossipsub-queues/main.nim index 71c7b67..92c2e18 100644 --- a/nim-test-node/gossipsub-queues/main.nim +++ b/nim-test-node/gossipsub-queues/main.nim @@ -3,7 +3,8 @@ import chronos, chronos/apps/http/httpserver import chronicles import env import std/[strformat, random, hashes] -import libp2p, libp2p/[muxers/mplex/lpchannel, stream/connection, crypto/secp, multiaddress] +import + libp2p, libp2p/[muxers/mplex/lpchannel, stream/connection, crypto/secp, multiaddress] import libp2p/protocols/[pubsub/pubsubpeer, pubsub/rpc/messages, ping] import sequtils, math, metrics, metrics/chronos_httpserver @@ -14,7 +15,6 @@ from nativesockets import getHostname logScope: topics = "dst" - template toUnixNanoseconds(t: times.Time): int64 = (t.toUnixFloat() * 1_000_000_000).int64 @@ -29,56 +29,57 @@ var declareCounter( dst_testnode_publish_requests_total, "number of /publish requests accepted by the test node", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) declareCounter( dst_testnode_publish_failures_total, "number of failed local publish attempts", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) declareCounter( dst_testnode_received_chunks_total, "number of application-level message chunks received", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) declareCounter( dst_testnode_completed_messages_total, "number of application-level messages fully received", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) declareCounter( dst_testnode_message_delay_ms_sum, "sum of message delays in milliseconds (use with rate)", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) declareHistogram( dst_testnode_message_delay_ms, "message delay histogram for percentile analysis", labels = ["muxer", "peer_id"], - buckets = [1.0, 5.0, 10.0, 25.0, 50.0, 100.0, 250.0, 500.0, 1000.0, 2500.0, 5000.0, 10000.0] + buckets = + [1.0, 5.0, 10.0, 25.0, 50.0, 100.0, 250.0, 500.0, 1000.0, 2500.0, 5000.0, 10000.0], ) declareGauge( dst_testnode_last_message_delay_ms, "last observed message delay in milliseconds (real-time)", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) declareGauge( dst_testnode_mesh_size, "current GossipSub mesh size for the test topic", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) declareGauge( dst_testnode_topic_peers, "current number of GossipSub peers for the test topic", - labels = ["muxer", "peer_id"] + labels = ["muxer", "peer_id"], ) proc getEnvInt(name: string, defaultValue: int): int = let value = getEnv(name, "") @@ -89,12 +90,9 @@ proc getEnvInt(name: string, defaultValue: int): int = return parseInt(value) except ValueError: warn "Invalid integer ENV value, using default", - name = name, - value = value, - defaultValue = defaultValue + name = name, value = value, defaultValue = defaultValue return defaultValue - proc getEnvFloat(name: string, defaultValue: float): float = let value = getEnv(name, "") if value.len == 0: @@ -104,12 +102,9 @@ proc getEnvFloat(name: string, defaultValue: float): float = return parseFloat(value) except ValueError: warn "Invalid float ENV value, using default", - name = name, - value = value, - defaultValue = defaultValue + name = name, value = value, defaultValue = defaultValue return defaultValue - proc getEnvBool(name: string, defaultValue: bool): bool = let value = getEnv(name, "") if value.len == 0: @@ -119,9 +114,7 @@ proc getEnvBool(name: string, defaultValue: bool): bool = return parseBool(value) except ValueError: warn "Invalid bool ENV value, using default", - name = name, - value = value, - defaultValue = defaultValue + name = name, value = value, defaultValue = defaultValue return defaultValue proc msgIdProvider(m: Message): Result[MessageId, ValidationResult] = @@ -139,7 +132,8 @@ proc createMessageHandler(): proc(topic: string, data: seq[byte]) {.async, gcsaf delay = recvTime - sendTime # warm-up - if timestampNs < 1000000: return + if timestampNs < 1000000: + return # Log received message info "Received message", @@ -148,40 +142,47 @@ proc createMessageHandler(): proc(topic: string, data: seq[byte]) {.async, gcsaf current = recvTime.toUnixNanoseconds(), delayMs = delay.inMilliseconds() - messagesChunks.inc(msgId) # Use msgId instead of timestamp for tracking - if messagesChunks[msgId] < chunks: return + messagesChunks.inc(msgId) # Use msgId instead of timestamp for tracking + if messagesChunks[msgId] < chunks: + return echo msgId, " milliseconds: ", delay.inMilliseconds() dst_testnode_completed_messages_total.inc(labelValues = [gMuxer, gPeerId]) - dst_testnode_message_delay_ms_sum.inc(delay.inMilliseconds().int64, labelValues = [gMuxer, gPeerId]) - dst_testnode_message_delay_ms.observe(delay.inMilliseconds().float64, labelValues = [gMuxer, gPeerId]) - dst_testnode_last_message_delay_ms.set(delay.inMilliseconds().int64, labelValues = [gMuxer, gPeerId]) + dst_testnode_message_delay_ms_sum.inc( + delay.inMilliseconds().int64, labelValues = [gMuxer, gPeerId] + ) + dst_testnode_message_delay_ms.observe( + delay.inMilliseconds().float64, labelValues = [gMuxer, gPeerId] + ) + dst_testnode_last_message_delay_ms.set( + delay.inMilliseconds().int64, labelValues = [gMuxer, gPeerId] + ) proc messageValidator(topic: string, msg: Message): Future[ValidationResult] {.async.} = return ValidationResult.Accept - -proc publishNewMessage(gossipSub: GossipSub, msgSize: int, topic: string): Future[(Time, int)] {.async.} = +proc publishNewMessage( + gossipSub: GossipSub, msgSize: int, topic: string +): Future[(Time, int)] {.async.} = dst_testnode_publish_requests_total.inc(labelValues = [gMuxer, gPeerId]) let now = getTime() - nowInt = now.toUnixFloat() * 1_000_000_000.0 # seconds + nanoseconds as float - msgId = uint64(rand(high(int64))) # Safe 0..<2^63 range + nowInt = now.toUnixFloat() * 1_000_000_000.0 # seconds + nanoseconds as float + msgId = uint64(rand(high(int64))) # Safe 0..<2^63 range var res = 0 - nowBytes = @(toBytesLE(uint64(nowInt))) & @(toBytesLE(msgId)) & - newSeq[byte](msgSize div chunks - 16) + nowBytes = + @(toBytesLE(uint64(nowInt))) & @(toBytesLE(msgId)) & + newSeq[byte](msgSize div chunks - 16) - info "Sent message", - msgId = msgId, - timestamp = getTime().toUnixNanoseconds() + info "Sent message", msgId = msgId, timestamp = getTime().toUnixNanoseconds() #To support message fragmentation, we add fragment #. Each fragment (chunk) differs by one byte - for chunk in 0.. 0: - let responseJson = """{"status":"success","message":"Message published at time """ & $publishTime & "}" - return await req.respond(Http200, responseJson, HttpTable.init([("Content-Type", "application/json")])) + let responseJson = + """{"status":"success","message":"Message published at time """ & + $publishTime & "}" + return await req.respond( + Http200, + responseJson, + HttpTable.init([("Content-Type", "application/json")]), + ) else: - let responseJson = """{"status":"error","message":"Failed to publist at time """ & $publishTime & "}" - return await req.respond(Http500, responseJson, HttpTable.init([("Content-Type", "application/json")])) + let responseJson = + """{"status":"error","message":"Failed to publist at time """ & + $publishTime & "}" + return await req.respond( + Http500, + responseJson, + HttpTable.init([("Content-Type", "application/json")]), + ) else: return await req.respond(Http404, "Not Found") else: return await req.respond(Http405, "Method Not Supported") - except CatchableError as e: info "Error handling http request: ", error = e.msg - let responseJson = """{"status":"error","message":"""" & e.msg.replace("\"", "\\\"") & """"}""" - return await req.respond(Http400, responseJson, HttpTable.init([("Content-Type", "application/json")])) + let responseJson = + """{"status":"error","message":"""" & e.msg.replace("\"", "\\\"") & """"}""" + return await req.respond( + Http400, responseJson, HttpTable.init([("Content-Type", "application/json")]) + ) # http endpoint for publish controller info "starting http server", httpPort = $httpPublishPort @@ -236,7 +253,8 @@ proc startHttpServer(gossipSub: GossipSub, myId: int): Future[HttpServerRef] {.a let serverRes = HttpServerRef.new(serverAddress, processRequests) if serverRes.isErr(): - raise newException(CatchableError, "Failed to create HTTP server: " & $serverRes.error) + raise + newException(CatchableError, "Failed to create HTTP server: " & $serverRes.error) let server = serverRes.get() server.start() @@ -245,13 +263,13 @@ proc startHttpServer(gossipSub: GossipSub, myId: int): Future[HttpServerRef] {.a proc initializeGossipsub(switch: Switch, anonymize: bool): GossipSub = return GossipSub.init( - switch = switch, - triggerSelf = parseBool(getEnv("SELFTRIGGER", "true")), - msgIdProvider = msgIdProvider, - verifySignature = false, - anonymize = anonymize, - rng = libp2p.newRng(), - ) + switch = switch, + triggerSelf = parseBool(getEnv("SELFTRIGGER", "true")), + msgIdProvider = msgIdProvider, + verifySignature = false, + anonymize = anonymize, + rng = libp2p.newRng(), + ) proc configureGossipsubParams(gossipSub: GossipSub) = let @@ -266,7 +284,8 @@ proc configureGossipsubParams(gossipSub: GossipSub) = pruneBackoffSec = getEnvInt("GOSSIPSUB_PRUNE_BACKOFF_SEC", 60) maxHighPriorityQueueLen = getEnvInt("GOSSIPSUB_MAX_HIGH_PRIORITY_QUEUE_LEN", 256) - maxMediumPriorityQueueLen = getEnvInt("GOSSIPSUB_MAX_MEDIUM_PRIORITY_QUEUE_LEN", 512) + maxMediumPriorityQueueLen = + getEnvInt("GOSSIPSUB_MAX_MEDIUM_PRIORITY_QUEUE_LEN", 512) maxLowPriorityQueueLen = getEnvInt("GOSSIPSUB_MAX_LOW_PRIORITY_QUEUE_LEN", 1024) slowPeerPenaltyWeight = getEnvFloat("GOSSIPSUB_SLOW_PEER_PENALTY_WEIGHT", 0.0) @@ -279,9 +298,10 @@ proc configureGossipsubParams(gossipSub: GossipSub) = #gossipThreshold = getEnvFloat("GOSSIPSUB_GOSSIP_THRESHOLD", -100.0) #publishThreshold = getEnvFloat("GOSSIPSUB_PUBLISH_THRESHOLD", -1000.0) #graylistThreshold = getEnvFloat("GOSSIPSUB_GRAYLIST_THRESHOLD", -10000.0) - + gossipSub.parameters.floodPublish = getEnvBool("GOSSIPSUB_FLOOD_PUBLISH", true) - gossipSub.parameters.opportunisticGraftThreshold = getEnvFloat("GOSSIPSUB_OPPORTUNISTIC_GRAFT_THRESHOLD", -10000) + gossipSub.parameters.opportunisticGraftThreshold = + getEnvFloat("GOSSIPSUB_OPPORTUNISTIC_GRAFT_THRESHOLD", -10000) gossipSub.parameters.heartbeatInterval = heartbeatMs.milliseconds gossipSub.parameters.pruneBackoff = pruneBackoffSec.seconds @@ -340,22 +360,22 @@ proc subscribGossipsubTopic(gossipSub: GossipSub, topic: string) = topicWeight: 1, firstMessageDeliveriesWeight: 1, firstMessageDeliveriesCap: 30, - firstMessageDeliveriesDecay: 0.9 + firstMessageDeliveriesDecay: 0.9, ) gossipSub.subscribe(topic, createMessageHandler()) gossipSub.addValidator([topic], messageValidator) - -proc resolveAddress(muxer: string, tAddress: string): Future[Result[seq[MultiAddress], string]] {.async.} = +proc resolveAddress( + muxer: string, tAddress: string +): Future[Result[seq[MultiAddress], string]] {.async.} = while true: try: let resolvedAddrs = if muxer.toLowerAscii() == "quic": let quicV1 = MultiAddress.init("/quic-v1").tryGet() resolveTAddress(tAddress).mapIt( - MultiAddress.init(it, IPPROTO_UDP).tryGet() - .concat(quicV1).tryGet() + MultiAddress.init(it, IPPROTO_UDP).tryGet().concat(quicV1).tryGet() ) else: resolveTAddress(tAddress).mapIt(MultiAddress.init(it).tryGet()) @@ -369,7 +389,7 @@ proc resolveAddress(muxer: string, tAddress: string): Future[Result[seq[MultiAdd await sleepAsync(15.seconds) proc connectGossipsubPeers( - switch: Switch, muxer: string, networkSize: int, myId: int, connectTo: int + switch: Switch, muxer: string, networkSize: int, myId: int, connectTo: int ): Future[Result[int, string]] {.async.} = let rng = libp2p.newRng() var @@ -378,10 +398,11 @@ proc connectGossipsubPeers( connected = 0 if inShadow: - var peers = toSeq(0..= connectTo: break + if connected >= connectTo: + break try: - discard await switch.connect(peer, allowUnknownPeerId=true).wait(5.seconds) + discard await switch.connect(peer, allowUnknownPeerId = true).wait(5.seconds) connected.inc() - info "Connected!: current connections ", connected = $connected, target = connectTo + info "Connected!: current connections ", + connected = $connected, target = connectTo except CatchableError as exc: warn "Failed to dial ", theirAddress = peer, message = exc.msg await sleepAsync(15.seconds) @@ -409,21 +432,21 @@ proc connectGossipsubPeers( if connected == 0: return err("Failed to connect any peer") elif connected < connectTo: - warn "Connected to fewer peers than target", connected = connected, target = connectTo + warn "Connected to fewer peers than target", + connected = connected, target = connectTo return ok(connected) - -proc main {.async.} = +proc main() {.async.} = randomize() let rng = libp2p.newRng() (myId, networkSize, connectTo, muxer, filePath, address) = getPeerDetails().valueOr: error "Node configuration is invalid", error = error return - + # Set global metric labels gMuxer = muxer - + var gossipSub: GossipSub builder = SwitchBuilder @@ -438,14 +461,12 @@ proc main {.async.} = of "quic": builder = builder.withQuicTransport() of "yamux": - builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) - .withYamux() + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}).withYamux() of "mplex": - builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) - .withMplex() + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}).withMplex() let switch = builder.build() - + # Set peerId for metric labels gPeerId = $switch.peerInfo.peerId gossipSub = initializeGossipsub(switch, true) @@ -480,15 +501,13 @@ proc main {.async.} = error "Failed to establish any connections", error = error return - await sleepAsync(15.seconds) # Allow multiple heartbeats to build mesh + await sleepAsync(15.seconds) # Allow multiple heartbeats to build mesh let meshSize = gossipSub.mesh.getOrDefault("test").len let peersConnected = gossipSub.gossipsub.getOrDefault("test").len dst_testnode_mesh_size.set(meshSize.int64, labelValues = [gMuxer, gPeerId]) dst_testnode_topic_peers.set(peersConnected.int64, labelValues = [gMuxer, gPeerId]) - info "Mesh details ", - meshSize = meshSize, - peersConnected = peersConnected + info "Mesh details ", meshSize = meshSize, peersConnected = peersConnected info "Starting listening endpoint for publish controller" discard gossipSub.startHttpServer(myId) diff --git a/nim-test-node/gossipsub-queues/test_node.nimble b/nim-test-node/gossipsub-queues/test_node.nimble index acbce1d..93ad166 100644 --- a/nim-test-node/gossipsub-queues/test_node.nimble +++ b/nim-test-node/gossipsub-queues/test_node.nimble @@ -2,14 +2,15 @@ mode = ScriptMode.Verbose bin = @["main"] -packageName = "test_node" -version = "0.1.0" -author = "Status Research & Development GmbH" -description = "A test node for gossipsub" -license = "MIT" -skipDirs = @[] +packageName = "test_node" +version = "0.1.0" +author = "Status Research & Development GmbH" +description = "A test node for gossipsub" +license = "MIT" +skipDirs = @[] requires "nim >= 2.2.4", - "nimcrypto >= 0.6.0", - "https://github.com/vacp2p/nim-libp2p#9067f2a5b004fc54a70f53ab02f13f59befa8460" # fix(gossip): make slow peer penalty opt-in by default (#2429) - #"ggplotnim" \ No newline at end of file + "nimcrypto >= 0.6.0", + "https://github.com/vacp2p/nim-libp2p#9067f2a5b004fc54a70f53ab02f13f59befa8460" + # fix(gossip): make slow peer penalty opt-in by default (#2429) + #"ggplotnim" diff --git a/nim-test-node/kad-dht/core.nim b/nim-test-node/kad-dht/core.nim index a679862..005c29b 100644 --- a/nim-test-node/kad-dht/core.nim +++ b/nim-test-node/kad-dht/core.nim @@ -1,4 +1,5 @@ -import libp2p, libp2p/[muxers/mplex/lpchannel, stream/connection, crypto/secp, multiaddress] +import + libp2p, libp2p/[muxers/mplex/lpchannel, stream/connection, crypto/secp, multiaddress] import libp2p/protocols/[pubsub/pubsubpeer, pubsub/rpc/messages, ping] import libp2p/protocols/[kademlia, kad_disco] import sequtils, math, metrics, metrics/chronos_httpserver @@ -17,7 +18,7 @@ proc runWarmup*(kad: KadDHT, selfId: PeerId) {.async.} = info "Starting warmup phase" # 5x FIND_NODE(self) - for i in 1..5: + for i in 1 .. 5: debug "Warmup: Finding self", iteration = i let peers = await kad.findNode(selfId.toKey()) var rtPeers = 0 @@ -29,7 +30,7 @@ proc runWarmup*(kad: KadDHT, selfId: PeerId) {.async.} = await sleepAsync(1.seconds) # 15x FIND_NODE(random) - for i in 1..15: + for i in 1 .. 15: let target = getRandomPeerId() debug "Warmup: Finding random node", iteration = i, target = target @@ -50,9 +51,6 @@ proc runProbe*(kad: KadDHT) {.async.} = info "Probe: Finding node", target = $targetPeer let peers = await kad.findNode(targetKey).wait(30.seconds) except CatchableError as exc: - warn "Probe Failed", - target = $targetPeer, - success = false, - error = exc.msg + warn "Probe Failed", target = $targetPeer, success = false, error = exc.msg await sleepAsync(5.seconds) diff --git a/nim-test-node/kad-dht/env.nim b/nim-test-node/kad-dht/env.nim index 206c70a..59ffafe 100644 --- a/nim-test-node/kad-dht/env.nim +++ b/nim-test-node/kad-dht/env.nim @@ -8,9 +8,10 @@ logScope: topics = "dst" # --- Configuration & Types --- -type - NodeType* = enum - RoleBootstrap, RoleNormal, RoleProbe +type NodeType* = enum + RoleBootstrap + RoleNormal + RoleProbe let httpPublishPort* = Port(8645) @@ -18,21 +19,25 @@ let myPort* = Port(parseInt(getEnv("PORT", "5000"))) proc getPeerDetails*(): Result[(int, string, string, NodeType, string), string] = - let + let hostname = getHostname() - myId = try: parseInt(hostname.split('-')[^1]) - except ValueError: 0 + myId = + try: + parseInt(hostname.split('-')[^1]) + except ValueError: + 0 muxer = getEnv("MUXER", "yamux") - address = if muxer.toLowerAscii() == "quic": - "/ip4/0.0.0.0/udp/" & $myPort & "/quic-v1" - else: - "/ip4/0.0.0.0/tcp/" & $myPort + address = + if muxer.toLowerAscii() == "quic": + "/ip4/0.0.0.0/udp/" & $myPort & "/quic-v1" + else: + "/ip4/0.0.0.0/tcp/" & $myPort nodeRole = parseEnum[NodeType](getEnv("NODE_ROLE", "RoleBootstrap")) discovery = getEnv("DISCOVERY", "kad-dht") if muxer.toLowerAscii() notin ["quic", "yamux", "mplex"]: return err("Unknown muxer type : " & muxer) - + info "Host info ", hostname = hostname, peer = myId, muxer = muxer, address = address return ok((myId, muxer, address, nodeRole, discovery)) @@ -56,15 +61,16 @@ proc startMetricsServer*( info "Metrics HTTP server started", serverIp = $serverIp, serverPort = $serverPort ok(metricsServerRes.value) -proc resolveAddress*(muxer: string, tAddress: string): Future[Result[seq[MultiAddress], string]] {.async.} = +proc resolveAddress*( + muxer: string, tAddress: string +): Future[Result[seq[MultiAddress], string]] {.async.} = while true: try: let resolvedAddrs = if muxer.toLowerAscii() == "quic": let quicV1 = MultiAddress.init("/quic-v1").tryGet() resolveTAddress(tAddress).mapIt( - MultiAddress.init(it, IPPROTO_UDP).tryGet() - .concat(quicV1).tryGet() + MultiAddress.init(it, IPPROTO_UDP).tryGet().concat(quicV1).tryGet() ) else: resolveTAddress(tAddress).mapIt(MultiAddress.init(it).tryGet()) @@ -74,9 +80,15 @@ proc resolveAddress*(muxer: string, tAddress: string): Future[Result[seq[MultiAd warn "Failed to resolve address", address = tAddress, error = exc.msg await sleepAsync(15.seconds) -proc resolveService*(muxer: string, service: string): Future[Result[seq[MultiAddress], string]] {.async.} = +proc resolveService*( + muxer: string, service: string +): Future[Result[seq[MultiAddress], string]] {.async.} = # If the service doesn't contain a port, append the default one - let tAddress = if ":" in service: service else: service & ":" & $myPort + let tAddress = + if ":" in service: + service + else: + service & ":" & $myPort # Call the existing resolveAddress helper let resolvedAddrs = (await resolveAddress(muxer, tAddress)).valueOr: diff --git a/nim-test-node/kad-dht/helpers.nim b/nim-test-node/kad-dht/helpers.nim index 7005018..8c36cad 100644 --- a/nim-test-node/kad-dht/helpers.nim +++ b/nim-test-node/kad-dht/helpers.nim @@ -17,34 +17,30 @@ proc getRandomPeerId*(): PeerId = proc buildSwitch*(muxer: string, address: string): Switch = var builder = SwitchBuilder - .new() - .withNoise() - .withRng(crypto.newRng()) - .withAddresses(@[MultiAddress.init(address).tryGet()]) - .withTcpTransport(flags = {ServerFlags.TcpNoDelay}) - .withMaxConnections(200) + .new() + .withNoise() + .withRng(crypto.newRng()) + .withAddresses(@[MultiAddress.init(address).tryGet()]) + .withTcpTransport(flags = {ServerFlags.TcpNoDelay}) + .withMaxConnections(200) case muxer.toLowerAscii() of "quic": builder = builder.withQuicTransport() of "yamux": - builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) - .withYamux() + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}).withYamux() of "mplex": - builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) - .withMplex() + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}).withMplex() builder.build() -proc mountDiscovery*(switch: Switch, discovery: string, addresses: seq[(PeerId, seq[MultiAddress])] - ): Future[KadDHT] {.async.} = +proc mountDiscovery*( + switch: Switch, discovery: string, addresses: seq[(PeerId, seq[MultiAddress])] +): Future[KadDHT] {.async.} = # Discovery selection if discovery == "kad-dht": - let kad = KadDHT.new( - switch, - bootstrapNodes = addresses, - config = KadDHTConfig.new(), - ) + let kad = + KadDHT.new(switch, bootstrapNodes = addresses, config = KadDHTConfig.new()) await kad.start() switch.mount(kad) return kad @@ -61,9 +57,9 @@ proc mountDiscovery*(switch: Switch, discovery: string, addresses: seq[(PeerId, raise newException(ValueError, "Unknown DISCOVERY: " & discovery) - -proc connectToBootstraps*(switch: Switch, muxer: string, service: string - ): Future[Result[seq[(PeerId, seq[MultiAddress])], string]] {.async.} = +proc connectToBootstraps*( + switch: Switch, muxer: string, service: string +): Future[Result[seq[(PeerId, seq[MultiAddress])], string]] {.async.} = let addrsRes = await resolveService(muxer, service) let addrs = addrsRes.valueOr: return err("Failed to resolve bootstrap service '" & service & "': " & error) @@ -74,9 +70,10 @@ proc connectToBootstraps*(switch: Switch, muxer: string, service: string for addr in addrs: var backoff = 1.seconds - for attempt in 1..10: + for attempt in 1 .. 10: try: - let remotePeerId: PeerId = await switch.connect(addr, allowUnknownPeerId = true).wait(10.seconds) + let remotePeerId: PeerId = + await switch.connect(addr, allowUnknownPeerId = true).wait(10.seconds) info "Connected to bootstrap", address = addr, peerId = remotePeerId bootstraps.add((remotePeerId, @[addr])) break @@ -88,12 +85,13 @@ proc connectToBootstraps*(switch: Switch, muxer: string, service: string backoff = min(backoff * 2, 30.seconds) if bootstraps.len == 0: - return err("Could not connect to any bootstrap resolved from '" & service & - "' (candidates=" & $addrs.len & "). Last error: " & lastErr) + return err( + "Could not connect to any bootstrap resolved from '" & service & "' (candidates=" & + $addrs.len & "). Last error: " & lastErr + ) ok(bootstraps) - proc startHealthServer*(port: Port): Future[HttpServerRef] {.async.} = proc handler(request: RequestFence): Future[HttpResponseRef] {.async.} = if request.isErr(): @@ -103,9 +101,7 @@ proc startHealthServer*(port: Port): Future[HttpServerRef] {.async.} = if req.meth == MethodGet and (req.uri.path == "/health" or req.uri.path == "/ready"): return await req.respond( - Http200, - "ok", - HttpTable.init([("Content-Type", "text/plain")]) + Http200, "ok", HttpTable.init([("Content-Type", "text/plain")]) ) return await req.respond(Http404, "Not Found") @@ -113,7 +109,9 @@ proc startHealthServer*(port: Port): Future[HttpServerRef] {.async.} = let addrs = initTAddress("0.0.0.0:" & $port) let serverRes = HttpServerRef.new(addrs, handler) if serverRes.isErr(): - raise newException(CatchableError, "Failed to create health HTTP server: " & $serverRes.error) + raise newException( + CatchableError, "Failed to create health HTTP server: " & $serverRes.error + ) let server = serverRes.get() server.start() diff --git a/nim-test-node/kad-dht/main.nim b/nim-test-node/kad-dht/main.nim index 6d5c252..a145b61 100644 --- a/nim-test-node/kad-dht/main.nim +++ b/nim-test-node/kad-dht/main.nim @@ -3,7 +3,8 @@ import chronos, chronos/apps/http/httpserver import chronicles import env import std/[strformat, random, hashes] -import libp2p, libp2p/[muxers/mplex/lpchannel, stream/connection, crypto/secp, multiaddress] +import + libp2p, libp2p/[muxers/mplex/lpchannel, stream/connection, crypto/secp, multiaddress] import libp2p/protocols/[pubsub/pubsubpeer, pubsub/rpc/messages, ping] import libp2p/protocols/[kademlia, kad_disco] @@ -16,7 +17,7 @@ import core logScope: topics = "dst" -proc main {.async.} = +proc main() {.async.} = randomize() var service = getEnv("SERVICE", "kad-service:5000") @@ -39,8 +40,8 @@ proc main {.async.} = var kad = await mountDiscovery(switch, discovery, @[]) discard await startHealthServer(prometheusPort) # Just stay alive and serve queries - while true: await sleepAsync(1.hours) - + while true: + await sleepAsync(1.hours) of RoleNormal: let jitter = myId * 200 if jitter > 0: @@ -56,8 +57,8 @@ proc main {.async.} = await runWarmup(kad, selfId) discard await startHealthServer(prometheusPort) # Keep node alive for steady state refresh - while true: await sleepAsync(1.hours) - + while true: + await sleepAsync(1.hours) of RoleProbe: let jitter = myId * 200 if jitter > 0: @@ -71,6 +72,7 @@ proc main {.async.} = var kad = await mountDiscovery(switch, discovery, bootAddresses) await runProbe(kad) - while true: await sleepAsync(1.hours) + while true: + await sleepAsync(1.hours) waitFor(main()) diff --git a/nim-test-node/kad-dht/test_node.nimble b/nim-test-node/kad-dht/test_node.nimble index 37549dd..ca8d3d2 100644 --- a/nim-test-node/kad-dht/test_node.nimble +++ b/nim-test-node/kad-dht/test_node.nimble @@ -1,13 +1,11 @@ mode = ScriptMode.Verbose -packageName = "test_node" -version = "0.1.0" -author = "Status Research & Development GmbH" -description = "A test node for gossipsub" -license = "MIT" -skipDirs = @[] +packageName = "test_node" +version = "0.1.0" +author = "Status Research & Development GmbH" +description = "A test node for gossipsub" +license = "MIT" +skipDirs = @[] requires "nim >= 2.2.0", - "nimcrypto 0.6.4", - "libp2p#e653fb093a7ced6f53947ad956f113280ae129a2", - "ggplotnim" + "nimcrypto 0.6.4", "libp2p#e653fb093a7ced6f53947ad956f113280ae129a2", "ggplotnim" diff --git a/nim-test-node/regression/bootstrap/main.nim b/nim-test-node/regression/bootstrap/main.nim index 7bbff81..7895a3c 100644 --- a/nim-test-node/regression/bootstrap/main.nim +++ b/nim-test-node/regression/bootstrap/main.nim @@ -12,13 +12,12 @@ import ../shutdown_utils logScope: topics = "dst" -proc main {.async.} = +proc main() {.async.} = let rng = libp2p.newRng() - (myId, muxer, _, address) = - getPeerDetails().valueOr: - error "Node configuration is invalid", error = error - quit(1) + (myId, muxer, _, address) = getPeerDetails().valueOr: + error "Node configuration is invalid", error = error + quit(1) let switch = buildSwitch(muxer, address) discard mountBaseProtocols(switch, rng) @@ -26,7 +25,8 @@ proc main {.async.} = await switch.start() info "Starting metrics server" - let metricsServer = await startMetricsServer(parseIpAddress("0.0.0.0"), prometheusPort) + let metricsServer = + await startMetricsServer(parseIpAddress("0.0.0.0"), prometheusPort) if metricsServer.isErr: warn "Failed to initialize metrics server", error = metricsServer.error elif inShadow: diff --git a/nim-test-node/regression/env.nim b/nim-test-node/regression/env.nim index 42666c7..6727f98 100644 --- a/nim-test-node/regression/env.nim +++ b/nim-test-node/regression/env.nim @@ -7,11 +7,12 @@ logScope: topics = "dst" let - inShadow* = getEnv("SHADOWENV").cmpIgnoreCase("true") == 0 #If Running for shadow simulator + inShadow* = getEnv("SHADOWENV").cmpIgnoreCase("true") == 0 + #If Running for shadow simulator httpPublishPort* = Port(8645) prometheusPort* = Port(8008) myPort* = Port(5000) - chunks* = parseInt(getEnv("FRAGMENTS", "1")) #No. of fragments for each message + chunks* = parseInt(getEnv("FRAGMENTS", "1")) #No. of fragments for each message # Per-pod startup jitter (pod index * this ms) to spread bootstrap dials and avoid # simultaneous-dial collisions. Mainly for Shadow, which starts all hosts at the same # simulated instant; real deployments get this spread naturally. Node availability @@ -21,8 +22,8 @@ let # that stays fine at 1000 nodes here, since each dials a sparse fixed set so inbound sits at # ~CONNECTTO on average, well under the connection cap. startupJitterStepMs* = parseInt(getEnv("STARTUP_JITTER_STEP_MS", "50")) - metricsIntervalS* = parseInt(getEnv("METRICS_INTERVAL_S", "300")) #storeMetrics scrape interval (s); short for shadow - + metricsIntervalS* = parseInt(getEnv("METRICS_INTERVAL_S", "300")) + #storeMetrics scrape interval (s); short for shadow proc listenHost*(): string = ## The interface the pod routes out of; 0.0.0.0 would announce loopback too. @@ -36,19 +37,33 @@ proc getPeerDetails*(): Result[(int, string, string, string), string] = let hostname = getHostname() listenIp = listenHost() - myId = try: parseInt(hostname.split('-')[^1]) - except ValueError: 0 + myId = + try: + parseInt(hostname.split('-')[^1]) + except ValueError: + 0 muxer = getEnv("MUXER", "yamux") - filePath = if inShadow: "../" else: getEnv("FILEPATH", "./") - address = if muxer.toLowerAscii() == "quic": - "/ip4/" & listenIp & "/udp/" & $myPort & "/quic-v1" - else: - "/ip4/" & listenIp & "/tcp/" & $myPort + filePath = + if inShadow: + "../" + else: + getEnv("FILEPATH", "./") + address = + if muxer.toLowerAscii() == "quic": + "/ip4/" & listenIp & "/udp/" & $myPort & "/quic-v1" + else: + "/ip4/" & listenIp & "/tcp/" & $myPort if muxer.toLowerAscii() notin ["quic", "yamux", "mplex"]: return err("Unknown muxer type : " & muxer) - info "Host info ", hostname = hostname, peer = myId, muxer = muxer, inShadow = inShadow, address = address, jitterStepMs = startupJitterStepMs + info "Host info ", + hostname = hostname, + peer = myId, + muxer = muxer, + inShadow = inShadow, + address = address, + jitterStepMs = startupJitterStepMs return ok((myId, muxer, filePath, address)) @@ -73,12 +88,13 @@ proc startMetricsServer*( #log metrics if needed (useful for shadow simulations) proc storeMetrics*(myId: int) {.async.} = - await sleepAsync((myId*60).milliseconds) + await sleepAsync((myId * 60).milliseconds) while true: try: - let cmd = "curl -s --connect-timeout 5 --max-time 5 http://localhost:" & - $prometheusPort & "/metrics >> metrics_pod-" & $myId & ".txt" - + let cmd = + "curl -s --connect-timeout 5 --max-time 5 http://localhost:" & $prometheusPort & + "/metrics >> metrics_pod-" & $myId & ".txt" + let exitCode = execCmd(cmd) if exitCode == 0: info "Metrics saved for peer ", pod = myId diff --git a/nim-test-node/regression/kad_utils.nim b/nim-test-node/regression/kad_utils.nim index b854d76..77cd230 100644 --- a/nim-test-node/regression/kad_utils.nim +++ b/nim-test-node/regression/kad_utils.nim @@ -13,14 +13,17 @@ logScope: # let FIND_NODE lookups discover the rest of the network. GossipSub then grafts # its mesh from the peers the DHT connected us to. -const - BootstrapDialTimeout = 10.seconds +const BootstrapDialTimeout = 10.seconds proc resolveBootstrapAddrs( muxer: string, service: string ): Future[Result[seq[MultiAddress], string]] {.async.} = # `service` is a k8s DNS name, optionally with a port; default to myPort. - let tAddress = if ':' in service: service else: service & ":" & $myPort + let tAddress = + if ':' in service: + service + else: + service & ":" & $myPort try: let resolved = if muxer.toLowerAscii() == "quic": diff --git a/nim-test-node/regression/node/main.nim b/nim-test-node/regression/node/main.nim index 8a109ba..4cc1a72 100644 --- a/nim-test-node/regression/node/main.nim +++ b/nim-test-node/regression/node/main.nim @@ -6,7 +6,8 @@ import libp2p import libp2p/protocols/[pubsub/pubsubpeer, pubsub/rpc/messages, ping] import math, metrics, metrics/chronos_httpserver -from times import getTime, Time, toUnix, fromUnix, `-`, initTime, `$`, inMilliseconds, toUnixFloat +from times import + getTime, Time, toUnix, fromUnix, `-`, initTime, `$`, inMilliseconds, toUnixFloat from nativesockets import getHostname import ../env @@ -39,7 +40,8 @@ proc createMessageHandler(): proc(topic: string, data: seq[byte]) {.async, gcsaf delay = recvTime - sendTime # warm-up - if timestampNs < 1000000: return + if timestampNs < 1000000: + return notePublishingStarted() @@ -50,30 +52,31 @@ proc createMessageHandler(): proc(topic: string, data: seq[byte]) {.async, gcsaf current = recvTime.toUnixNanoseconds(), delayMs = delay.inMilliseconds() - messagesChunks.inc(msgId) # Use msgId instead of timestamp for tracking - if messagesChunks[msgId] < chunks: return + messagesChunks.inc(msgId) # Use msgId instead of timestamp for tracking + if messagesChunks[msgId] < chunks: + return proc messageValidator(topic: string, msg: Message): Future[ValidationResult] {.async.} = return ValidationResult.Accept - -proc publishNewMessage(gossipSub: GossipSub, msgSize: int, topic: string): Future[(Time, int)] {.async.} = +proc publishNewMessage( + gossipSub: GossipSub, msgSize: int, topic: string +): Future[(Time, int)] {.async.} = let now = getTime() - nowInt = now.toUnixFloat() * 1_000_000_000.0 # seconds + nanoseconds as float - msgId = uint64(rand(high(int64))) # Safe 0..<2^63 range + nowInt = now.toUnixFloat() * 1_000_000_000.0 # seconds + nanoseconds as float + msgId = uint64(rand(high(int64))) # Safe 0..<2^63 range var res = 0 - nowBytes = @(toBytesLE(uint64(nowInt))) & @(toBytesLE(msgId)) & - newSeq[byte](msgSize div chunks - 16) + nowBytes = + @(toBytesLE(uint64(nowInt))) & @(toBytesLE(msgId)) & + newSeq[byte](msgSize div chunks - 16) - info "Sent message", - msgId = msgId, - timestamp = getTime().toUnixNanoseconds() + info "Sent message", msgId = msgId, timestamp = getTime().toUnixNanoseconds() #To support message fragmentation, we add fragment #. Each fragment (chunk) differs by one byte - for chunk in 0.. 0: - let responseJson = """{"status":"success","message":"Message published at time """ & $publishTime & "}" - return await req.respond(Http200, responseJson, HttpTable.init([("Content-Type", "application/json")])) + let responseJson = + """{"status":"success","message":"Message published at time """ & + $publishTime & "}" + return await req.respond( + Http200, + responseJson, + HttpTable.init([("Content-Type", "application/json")]), + ) else: - let responseJson = """{"status":"error","message":"Failed to publist at time """ & $publishTime & "}" - return await req.respond(Http500, responseJson, HttpTable.init([("Content-Type", "application/json")])) + let responseJson = + """{"status":"error","message":"Failed to publist at time """ & + $publishTime & "}" + return await req.respond( + Http500, + responseJson, + HttpTable.init([("Content-Type", "application/json")]), + ) else: return await req.respond(Http404, "Not Found") else: return await req.respond(Http405, "Method Not Supported") - except CatchableError as e: info "Error handling http request: ", error = e.msg - let responseJson = """{"status":"error","message":"""" & e.msg.replace("\"", "\\\"") & """"}""" - return await req.respond(Http400, responseJson, HttpTable.init([("Content-Type", "application/json")])) + let responseJson = + """{"status":"error","message":"""" & e.msg.replace("\"", "\\\"") & """"}""" + return await req.respond( + Http400, responseJson, HttpTable.init([("Content-Type", "application/json")]) + ) # http endpoint for publish controller info "starting http server", httpPort = $httpPublishPort @@ -122,7 +141,8 @@ proc startHttpServer(gossipSub: GossipSub, myId: int): Future[HttpServerRef] {.a let serverRes = HttpServerRef.new(serverAddress, processRequests) if serverRes.isErr(): - raise newException(CatchableError, "Failed to create HTTP server: " & $serverRes.error) + raise + newException(CatchableError, "Failed to create HTTP server: " & $serverRes.error) let server = serverRes.get() server.start() @@ -131,14 +151,14 @@ proc startHttpServer(gossipSub: GossipSub, myId: int): Future[HttpServerRef] {.a proc initializeGossipsub(switch: Switch, anonymize: bool, rng: Rng): GossipSub = return GossipSub.init( - switch = switch, - triggerSelf = parseBool(getEnv("SELFTRIGGER", "true")), - msgIdProvider = msgIdProvider, - verifySignature = false, - anonymize = anonymize, - rng = rng, - customStreamCallbacks = Opt.none(CustomStreamCallbacks) - ) + switch = switch, + triggerSelf = parseBool(getEnv("SELFTRIGGER", "true")), + msgIdProvider = msgIdProvider, + verifySignature = false, + anonymize = anonymize, + rng = rng, + customStreamCallbacks = Opt.none(CustomStreamCallbacks), + ) proc configureGossipsubParams(gossipSub: GossipSub) = gossipSub.parameters.floodPublish = true @@ -158,21 +178,19 @@ proc subscribGossipsubTopic(gossipSub: GossipSub, topic: string) = topicWeight: 1, firstMessageDeliveriesWeight: 1, firstMessageDeliveriesCap: 30, - firstMessageDeliveriesDecay: 0.9 + firstMessageDeliveriesDecay: 0.9, ) gossipSub.subscribe(topic, createMessageHandler()) gossipSub.addValidator([topic], messageValidator) - -proc main {.async.} = +proc main() {.async.} = randomize() let rng = libp2p.newRng() - (myId, muxer, _, address) = - getPeerDetails().valueOr: - error "Node configuration is invalid", error = error - return + (myId, muxer, _, address) = getPeerDetails().valueOr: + error "Node configuration is invalid", error = error + return let switch = buildSwitch(muxer, address) # Mount protocols before starting the switch; switch.start() starts mounted protocols. @@ -196,7 +214,8 @@ proc main {.async.} = discard gossipSub.startHttpServer(myId) info "Starting metrics server" - let metricsServer = await startMetricsServer(parseIpAddress("0.0.0.0"), prometheusPort) + let metricsServer = + await startMetricsServer(parseIpAddress("0.0.0.0"), prometheusPort) if metricsServer.isErr: warn "Failed to initialize metrics server", error = metricsServer.error elif inShadow: @@ -217,7 +236,8 @@ proc main {.async.} = info "kad-dht discovery active", bootstraps = bootstraps.len await sleepAsync(5.seconds) - info "Mesh details ", meshSize = gossipSub.mesh.getOrDefault("test").len, + info "Mesh details ", + meshSize = gossipSub.mesh.getOrDefault("test").len, peersConnected = gossipSub.gossipsub.getOrDefault("test").len # Hold connections open until gossipsub traffic takes over diff --git a/nim-test-node/regression/node_setup.nim b/nim-test-node/regression/node_setup.nim index 2396481..5fdd7c9 100644 --- a/nim-test-node/regression/node_setup.nim +++ b/nim-test-node/regression/node_setup.nim @@ -18,11 +18,9 @@ proc buildSwitch*(muxer: string, address: string): Switch = of "quic": builder = builder.withQuicTransport() of "yamux": - builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) - .withYamux() + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}).withYamux() of "mplex": - builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) - .withMplex() + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}).withMplex() else: raiseAssert("Unknown muxer type: " & muxer) diff --git a/nim-test-node/regression/ping_utils.nim b/nim-test-node/regression/ping_utils.nim index dffb8c6..d1ab6bd 100644 --- a/nim-test-node/regression/ping_utils.nim +++ b/nim-test-node/regression/ping_utils.nim @@ -9,19 +9,19 @@ logScope: topics = "dst" const - PingInterval = 12.seconds + PingInterval = 12.seconds ## A full sweep in batches takes a few seconds, so the gap between two pings of the ## same peer is the interval plus a sweep. Keep that comfortably under quic's 30s ## idle timeout. - PingBatch = 16 + PingBatch = 16 ## Dials in flight at once. All ~250 together stalls the node's own publish ## endpoint, which cost 34 injections in the first minute of a 1000-node run. - PingBatchGap = 200.milliseconds - PingTimeout = 4.seconds - DialTimeout = 4.seconds - CloseTimeout = 2.seconds - SlowDialLog = 500.milliseconds - SlowCloseLog = 500.milliseconds + PingBatchGap = 200.milliseconds + PingTimeout = 4.seconds + DialTimeout = 4.seconds + CloseTimeout = 2.seconds + SlowDialLog = 500.milliseconds + SlowCloseLog = 500.milliseconds var messagesStarted = false @@ -68,7 +68,8 @@ proc pingPeer*(switch: Switch, pingProtocol: Ping, peerId: PeerId) {.async.} = raise exc except CatchableError as exc: let dialDur = Moment.now() - dialStart - warn "keepalive ping failed", peerId = peerId, error = exc.msg, dialMs = dialDur.milliseconds + warn "keepalive ping failed", + peerId = peerId, error = exc.msg, dialMs = dialDur.milliseconds finally: if not stream.isNil and not stream.closed: let closeStart = Moment.now() @@ -82,7 +83,8 @@ proc pingPeer*(switch: Switch, pingProtocol: Ping, peerId: PeerId) {.async.} = finally: let closeDur = Moment.now() - closeStart if closeDur >= SlowCloseLog: - warn "keepalive ping: slow close", peerId = peerId, closeMs = closeDur.milliseconds + warn "keepalive ping: slow close", + peerId = peerId, closeMs = closeDur.milliseconds proc pingAllOnce*(switch: Switch, pingProtocol: Ping) {.async.} = var peers: seq[PeerId] = @[] diff --git a/nim-test-node/regression/test_node.nimble b/nim-test-node/regression/test_node.nimble index 27c85c9..6438a1b 100644 --- a/nim-test-node/regression/test_node.nimble +++ b/nim-test-node/regression/test_node.nimble @@ -1,14 +1,14 @@ mode = ScriptMode.Verbose bin = @["node/main", "bootstrap/main"] -namedBin = {"node/main": "regression-node", "bootstrap/main": "regression-bootstrap"}.toTable() +namedBin = + {"node/main": "regression-node", "bootstrap/main": "regression-bootstrap"}.toTable() -packageName = "test_node" -version = "0.1.0" -author = "Status Research & Development GmbH" -description = "A test node for gossipsub" -license = "MIT" -skipDirs = @[] +packageName = "test_node" +version = "0.1.0" +author = "Status Research & Development GmbH" +description = "A test node for gossipsub" +license = "MIT" +skipDirs = @[] -requires "nim >= 2.2.0", - "libp2p == 2.4.0" \ No newline at end of file +requires "nim >= 2.2.0", "libp2p == 2.4.0" diff --git a/nim-test-node/service-discovery/core.nim b/nim-test-node/service-discovery/core.nim index e852513..93c0286 100644 --- a/nim-test-node/service-discovery/core.nim +++ b/nim-test-node/service-discovery/core.nim @@ -7,22 +7,16 @@ import libp2p/extended_peer_record logScope: topics = "dst" -proc startAdvertisingServices*( - disco: ServiceDiscovery, services: seq[ServiceInfo] -) = +proc startAdvertisingServices*(disco: ServiceDiscovery, services: seq[ServiceInfo]) = if services.len == 0: warn "No services configured for advertising" return for service in services: disco.startAdvertising(service) - info "Advertising service", - service = service.id, - dataLen = service.data.len + info "Advertising service", service = service.id, dataLen = service.data.len -proc startDiscoveringServicesLog*( - disco: ServiceDiscovery, serviceIds: seq[string] -) = +proc startDiscoveringServicesLog*(disco: ServiceDiscovery, serviceIds: seq[string]) = if serviceIds.len == 0: warn "No services configured for discovery" return @@ -50,8 +44,6 @@ proc runLookupLoop*( addrs = ad.data.addresses.mapIt($it.address) info "Lookup completed", - service = serviceId, - advertisements = ads.len, - uniquePeers = uniquePeers.len + service = serviceId, advertisements = ads.len, uniquePeers = uniquePeers.len await sleepAsync(lookupInterval) diff --git a/nim-test-node/service-discovery/env.nim b/nim-test-node/service-discovery/env.nim index a7aa9ea..f595f60 100644 --- a/nim-test-node/service-discovery/env.nim +++ b/nim-test-node/service-discovery/env.nim @@ -8,7 +8,10 @@ logScope: type NodeRole* = enum - RoleBootstrap, RoleAdvertiser, RoleDiscoverer, RoleHybrid + RoleBootstrap + RoleAdvertiser + RoleDiscoverer + RoleHybrid NodeConfig* = object nodeIndex*: int @@ -31,26 +34,15 @@ type maxBootstraps*: int proc `$`(c: NodeConfig): string = - "NodeConfig(" & - "role=" & $c.role & - ", nodeIndex=" & $c.nodeIndex & - ", muxer=" & c.muxer & - ", listenAddress=" & c.listenAddress & - ", bootstrapService=" & c.bootstrapService & - ", advertiseServices=" & $c.advertiseServices & - ", discoverServices=" & $c.discoverServices & - ", listenPort=" & $c.listenPort & - ", healthPort=" & $c.healthPort & - ", lookupInterval=" & $c.lookupInterval & - ", startupJitterMs=" & $c.startupJitterMs & - ", safetyParam=" & $c.safetyParam & - ", ipSimCoefficient=" & $c.ipSimCoefficient & - ", advertExpiry=" & $c.advertExpiry & - ", xprPublishing=" & $c.xprPublishing & - ", maxConnections=" & $c.maxConnections & - ", maxBootstraps=" & $c.maxBootstraps & - ")" - + "NodeConfig(" & "role=" & $c.role & ", nodeIndex=" & $c.nodeIndex & ", muxer=" & + c.muxer & ", listenAddress=" & c.listenAddress & ", bootstrapService=" & + c.bootstrapService & ", advertiseServices=" & $c.advertiseServices & + ", discoverServices=" & $c.discoverServices & ", listenPort=" & $c.listenPort & + ", healthPort=" & $c.healthPort & ", lookupInterval=" & $c.lookupInterval & + ", startupJitterMs=" & $c.startupJitterMs & ", safetyParam=" & $c.safetyParam & + ", ipSimCoefficient=" & $c.ipSimCoefficient & ", advertExpiry=" & $c.advertExpiry & + ", xprPublishing=" & $c.xprPublishing & ", maxConnections=" & $c.maxConnections & + ", maxBootstraps=" & $c.maxBootstraps & ")" proc parseIntEnv(name: string, defaultValue: string): Result[int, string] = let raw = getEnv(name, defaultValue) @@ -97,12 +89,13 @@ proc getNodeConfig*(): Result[NodeConfig, string] = if listenPort <= 0 or listenPort > 65535: return err("PORT out of range: " & $listenPort) - let role = try: - parseEnum[NodeRole](getEnv("NODE_ROLE", "RoleBootstrap")) - except ValueError: - return err( - "Unknown NODE_ROLE. Expected one of: RoleBootstrap, RoleAdvertiser, RoleDiscoverer, RoleHybrid" - ) + let role = + try: + parseEnum[NodeRole](getEnv("NODE_ROLE", "RoleBootstrap")) + except ValueError: + return err( + "Unknown NODE_ROLE. Expected one of: RoleBootstrap, RoleAdvertiser, RoleDiscoverer, RoleHybrid" + ) let healthPort = parseIntEnv("HEALTH_PORT", "8645").valueOr: return err(error) @@ -183,7 +176,7 @@ proc getNodeConfig*(): Result[NodeConfig, string] = advertExpiry: advertExpirySeconds.seconds, xprPublishing: xprPublishing, maxConnections: maxConnections, - maxBootstraps: maxBootstraps + maxBootstraps: maxBootstraps, ) info "Node config loaded", cfg = $cfg diff --git a/nim-test-node/service-discovery/helpers.nim b/nim-test-node/service-discovery/helpers.nim index 8950909..7790e77 100644 --- a/nim-test-node/service-discovery/helpers.nim +++ b/nim-test-node/service-discovery/helpers.nim @@ -70,7 +70,7 @@ proc connectToBootstraps*( muxer: string, service: string, defaultPort: Port, - maxConnections: int + maxConnections: int, ): Future[Result[seq[(PeerId, seq[MultiAddress])], string]] {.async.} = if maxConnections <= 0: return err("maxConnections must be greater than 0") diff --git a/nim-test-node/service-discovery/main.nim b/nim-test-node/service-discovery/main.nim index 4e2037b..00f9862 100644 --- a/nim-test-node/service-discovery/main.nim +++ b/nim-test-node/service-discovery/main.nim @@ -22,16 +22,16 @@ proc main() {.async.} = await switch.start() let selfId = switch.peerInfo.peerId - info "Node started", - peerId = $selfId, nodeType = cfg.role, listen = cfg.listenAddress + info "Node started", peerId = $selfId, nodeType = cfg.role, listen = cfg.listenAddress if cfg.role != RoleBootstrap: if cfg.startupJitterMs > 0: info "Applying startup jitter", delayMs = cfg.startupJitterMs await sleepAsync(cfg.startupJitterMs.milliseconds) - let connectedBootstraps = - await connectToBootstraps(switch, cfg.muxer, cfg.bootstrapService, cfg.listenPort, cfg.maxBootstraps) + let connectedBootstraps = await connectToBootstraps( + switch, cfg.muxer, cfg.bootstrapService, cfg.listenPort, cfg.maxBootstraps + ) let bootstrapNodes = connectedBootstraps.valueOr: error "Failed to connect to bootstrap nodes", service = cfg.bootstrapService, error diff --git a/nim-test-node/service-discovery/test_node.nimble b/nim-test-node/service-discovery/test_node.nimble index f641f1f..70ff1fe 100644 --- a/nim-test-node/service-discovery/test_node.nimble +++ b/nim-test-node/service-discovery/test_node.nimble @@ -1,12 +1,12 @@ mode = ScriptMode.Verbose -packageName = "test_node" -version = "0.1.0" -author = "Status Research & Development GmbH" -description = "A test node for libp2p service discovery" -license = "MIT" -skipDirs = @[] +packageName = "test_node" +version = "0.1.0" +author = "Status Research & Development GmbH" +description = "A test node for libp2p service discovery" +license = "MIT" +skipDirs = @[] requires "nim >= 2.2.4", - "nimcrypto 0.6.4", - "https://github.com/vacp2p/nim-libp2p#26e181e4dd65188051ab4783bcc538ed9579644f" # 2.0.0 + "nimcrypto 0.6.4", + "https://github.com/vacp2p/nim-libp2p#26e181e4dd65188051ab4783bcc538ed9579644f" # 2.0.0