diff --git a/nim-test-node/gossipsub-queues/.gitignore b/nim-test-node/gossipsub-queues/.gitignore new file mode 100644 index 0000000..f3685e2 --- /dev/null +++ b/nim-test-node/gossipsub-queues/.gitignore @@ -0,0 +1,3 @@ +nimble.develop +nimble.paths +nimbledeps diff --git a/nim-test-node/gossipsub-queues/Dockerfile_amd64 b/nim-test-node/gossipsub-queues/Dockerfile_amd64 new file mode 100644 index 0000000..c4f28ae --- /dev/null +++ b/nim-test-node/gossipsub-queues/Dockerfile_amd64 @@ -0,0 +1,66 @@ +FROM debian:bookworm-slim AS builder + +WORKDIR /node + +COPY . . + +RUN apt-get update && apt-get install -y \ + curl git build-essential ca-certificates \ + libssl-dev \ + && rm -rf /var/lib/apt/lists/* + +RUN apt-get update && apt-get install -y \ + gcc-multilib g++-multilib libc6-dev-i386 \ + && rm -rf /var/lib/apt/lists/* + +RUN apt-get update && apt-get install -y \ + curl git build-essential ca-certificates \ + gcc libc6-dev-i386 \ + libssl3 iproute2 procps \ + && rm -rf /var/lib/apt/lists/* \ + && apt-get clean + +RUN apt-get update && \ + apt-get install -y --no-install-recommends \ + ca-certificates \ + libssl3 \ + iproute2 \ + curl \ + procps \ + && rm -rf /var/lib/apt/lists/* \ + && apt-get clean + +RUN git config --global http.sslVerify false + +RUN curl https://nim-lang.org/download/nim-2.2.10-linux_x64.tar.xz -o /tmp/nim.tar.xz \ + && tar -xf /tmp/nim.tar.xz -C /tmp \ + && cp -r /tmp/nim-2.2.10/* /usr/local \ + && rm -rf /usr/local/nim* && mv /tmp/nim-2.2.10 /usr/local/nim-2.2.10 + +RUN nimble refresh + +RUN ln -sf /usr/local/nim-2.2.10/bin/nim /usr/local/bin/nim +RUN ln -sf /usr/local/nim-2.2.10/bin/nimble /usr/local/bin/nimble +RUN ln -sf /usr/local/nim-2.2.10/bin/nim /usr/local/bin/nimcheck + +RUN echo 'import strutils; echo "OK"' > temp.nim && nim r temp.nim && rm temp.nim + +RUN nim --version +RUN nimble --version +RUN nimble install -dy --verbose | tee /tmp/install.log && grep -A 20 "libp2p" /tmp/install.log | grep -oP '\b\d+\.\d+\.\d+\b' | head -n1 > ./libp2p_version + +RUN nimble c \ + --os:linux --cpu:amd64 --passL:"-static -mmusl" \ + -d:chronicles_colors=None --threads:on --mm:refc \ + -d:metrics -d:libp2p_network_protocols_metrics -d:release -d:pubsubpeer_queue_metrics \ + -d:libp2p_expensive_metrics -d:libp2p_agents_metrics \ + --passL:"-static-libgcc -static-libstdc++" \ + main + +RUN chmod +x /node/main +RUN nim --version > nim_version +RUN nimble --version > nimble_version + +EXPOSE 5000 8008 8645 + +ENTRYPOINT ["/node/main"] \ No newline at end of file diff --git a/nim-test-node/gossipsub-queues/Dockerfile_arm64 b/nim-test-node/gossipsub-queues/Dockerfile_arm64 new file mode 100644 index 0000000..79ee333 --- /dev/null +++ b/nim-test-node/gossipsub-queues/Dockerfile_arm64 @@ -0,0 +1,49 @@ +FROM nimlang/nim:latest AS builder + +WORKDIR /node + +COPY . . + +RUN apt-get update && apt-get install -y \ + curl git build-essential ca-certificates \ + && rm -rf /var/lib/apt/lists/* +RUN git config --global http.sslVerify false +RUN git config --global --add safe.directory /node + +RUN apt-get update && \ + apt-get install -y --no-install-recommends \ + ca-certificates \ + libssl3 \ + iproute2 \ + curl \ + procps \ + && rm -rf /var/lib/apt/lists/* \ + && apt-get clean + +RUN git config --global http.sslVerify false + +RUN nimble refresh +RUN nim --version +RUN nimble --version + +RUN echo 'import strutils; echo "OK"' > temp.nim && nim r temp.nim && rm temp.nim + +RUN nimble install -dy --verbose | tee /tmp/install.log && grep -A 20 "libp2p" /tmp/install.log | grep -oP '\b\d+\.\d+\.\d+\b' | head -n1 > ./libp2p_version + + +RUN nimble c \ + --os:linux --cpu:arm64 --passL:"-static -mmusl" \ + -d:chronicles_colors=None --threads:on --mm:refc \ + -d:metrics -d:libp2p_network_protocols_metrics -d:release -d:pubsubpeer_queue_metrics \ + -d:libp2p_expensive_metrics -d:libp2p_agents_metrics \ + --passL:"-static-libgcc -static-libstdc++" \ + main + + +RUN chmod +x /node/main +RUN nim --version > nim_version +RUN nimble --version > nimble_version + +EXPOSE 5000 8008 8645 + +ENTRYPOINT ["/node/main"] \ No newline at end of file diff --git a/nim-test-node/gossipsub-queues/config.nims b/nim-test-node/gossipsub-queues/config.nims new file mode 100644 index 0000000..8ee48d2 --- /dev/null +++ b/nim-test-node/gossipsub-queues/config.nims @@ -0,0 +1,4 @@ +# begin Nimble config (version 2) +when withDir(thisDir(), system.fileExists("nimble.paths")): + include "nimble.paths" +# end Nimble config diff --git a/nim-test-node/gossipsub-queues/env.nim b/nim-test-node/gossipsub-queues/env.nim new file mode 100644 index 0000000..9ccdc99 --- /dev/null +++ b/nim-test-node/gossipsub-queues/env.nim @@ -0,0 +1,73 @@ +import strutils, os, osproc +import chronos, metrics/chronos_httpserver, chronicles +from nativesockets import getHostname + +let + 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 + + +proc getPeerDetails*(): Result[(int, int, int, string, string, string), string] = + let + hostname = getHostname() + ordinal = parseInt(hostname.split('-')[^1]) + peerIdOffset = parseInt(getEnv("PEER_ID_OFFSET", "0")) + myId = peerIdOffset + ordinal + 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 + + 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 ok((myId, networkSize, connectTo, muxer, filePath, address)) + +#Prometheus metrics +proc startMetricsServer*( + serverIp: IpAddress, serverPort: Port +): Result[MetricsHttpServerRef, string] = + info "Starting metrics HTTP server", serverIp = $serverIp, serverPort = $serverPort + + let metricsServerRes = MetricsHttpServerRef.new($serverIp, serverPort) + if metricsServerRes.isErr(): + return err("metrics HTTP server start failed: " & $metricsServerRes.error) + + let server = metricsServerRes.value + try: + waitFor server.start() + except CatchableError: + return err("metrics HTTP server start failed: " & getCurrentExceptionMsg()) + + info "Metrics HTTP server started", serverIp = $serverIp, serverPort = $serverPort + ok(metricsServerRes.value) + +#log metrics if needed (useful for shadow simulations) +proc storeMetrics*(myId: int) {.async.} = + 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 exitCode = execCmd(cmd) + if exitCode == 0: + info "Metrics saved for peer ", pod = myId + else: + info "Failed to fetch metrics for peer ", pod = myId, curlExitCode = $exitCode + except CatchableError as e: + info "Error storing metrics: ", error = e.msg + return + await sleepAsync(5.minutes) \ No newline at end of file diff --git a/nim-test-node/gossipsub-queues/main.nim b/nim-test-node/gossipsub-queues/main.nim new file mode 100644 index 0000000..bf35253 --- /dev/null +++ b/nim-test-node/gossipsub-queues/main.nim @@ -0,0 +1,487 @@ +import stew/endians2, stew/byteutils, tables, strutils, os, json +import chronos, chronos/apps/http/httpserver +import env +import std/[strformat, random, hashes] +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 +from times import getTime, Time, toUnix, fromUnix, `-`, initTime, `$`, inMilliseconds +from times import getTime, toUnixFloat, `-`, initTime, `$`, inMilliseconds, Time +from nativesockets import getHostname + + +template toUnixNanoseconds(t: times.Time): int64 = + (t.toUnixFloat() * 1_000_000_000).int64 + +template fromUnixNanoseconds(ns: int64): times.Time = + initTime(ns div 1_000_000_000, ns mod 1_000_000_000) + +# Global variables for metric labels (set in main) - thread local for GC safety +var + gMuxer* {.threadvar.}: string + gPeerId* {.threadvar.}: string + +declareCounter( + dst_testnode_publish_requests_total, + "number of /publish requests accepted by the test node", + labels = ["muxer", "peer_id"] +) + +declareCounter( + dst_testnode_publish_failures_total, + "number of failed local publish attempts", + labels = ["muxer", "peer_id"] +) + +declareCounter( + dst_testnode_received_chunks_total, + "number of application-level message chunks received", + labels = ["muxer", "peer_id"] +) + +declareCounter( + dst_testnode_completed_messages_total, + "number of application-level messages fully received", + labels = ["muxer", "peer_id"] +) + +declareCounter( + dst_testnode_message_delay_ms_sum, + "sum of message delays in milliseconds (use with rate)", + 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] +) + +declareGauge( + dst_testnode_last_message_delay_ms, + "last observed message delay in milliseconds (real-time)", + labels = ["muxer", "peer_id"] +) + +declareGauge( + dst_testnode_mesh_size, + "current GossipSub mesh size for the test topic", + labels = ["muxer", "peer_id"] +) + +declareGauge( + dst_testnode_topic_peers, + "current number of GossipSub peers for the test topic", + labels = ["muxer", "peer_id"] +) +proc getEnvInt(name: string, defaultValue: int): int = + let value = getEnv(name, "") + if value.len == 0: + return defaultValue + + try: + return parseInt(value) + except ValueError: + warn "Invalid integer ENV value, using default", + name = name, + value = value, + defaultValue = defaultValue + return defaultValue + + +proc getEnvFloat(name: string, defaultValue: float): float = + let value = getEnv(name, "") + if value.len == 0: + return defaultValue + + try: + return parseFloat(value) + except ValueError: + warn "Invalid float ENV value, using default", + name = name, + value = value, + defaultValue = defaultValue + return defaultValue + + +proc getEnvBool(name: string, defaultValue: bool): bool = + let value = getEnv(name, "") + if value.len == 0: + return defaultValue + + try: + return parseBool(value) + except ValueError: + warn "Invalid bool ENV value, using default", + name = name, + value = value, + defaultValue = defaultValue + return defaultValue + +proc msgIdProvider(m: Message): Result[MessageId, ValidationResult] = + return ok(($m.data.hash).toBytes()) + +proc createMessageHandler(): proc(topic: string, data: seq[byte]) {.async, gcsafe.} = + var messagesChunks: CountTable[uint64] + + return proc(topic: string, data: seq[byte]) {.async, gcsafe.} = + let + timestampNs = uint64.fromBytesLE(data[0 ..< 8]).int64 + sendTime = fromUnixNanoseconds(timestampNs) + msgId = uint64.fromBytesLE(data[8 ..< 16]) + recvTime = getTime() + delay = recvTime - sendTime + + # warm-up + if timestampNs < 1000000: return + + # Log received message + info "Received message", + msgId = msgId, + sentAt = timestampNs, + current = recvTime.toUnixNanoseconds(), + delayMs = delay.inMilliseconds() + + 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]) + +proc messageValidator(topic: string, msg: Message): Future[ValidationResult] {.async.} = + return ValidationResult.Accept + + +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 + + var + res = 0 + nowBytes = @(toBytesLE(uint64(nowInt))) & @(toBytesLE(msgId)) & + newSeq[byte](msgSize div chunks - 16) + + 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")])) + else: + 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")])) + + # http endpoint for publish controller + info "starting http server", httpPort = $httpPublishPort + let serverAddress = initTAddress("0.0.0.0:" & $httpPublishPort) + let serverRes = HttpServerRef.new(serverAddress, processRequests) + + if serverRes.isErr(): + raise newException(CatchableError, "Failed to create HTTP server: " & $serverRes.error) + + let server = serverRes.get() + server.start() + info "http server started ", httpPort = $httpPublishPort + return server + +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(), + ) + +proc configureGossipsubParams(gossipSub: GossipSub) = + let + d = getEnvInt("GOSSIPSUB_D", 6) + dLow = getEnvInt("GOSSIPSUB_D_LOW", 4) + dHigh = getEnvInt("GOSSIPSUB_D_HIGH", 8) + dScore = getEnvInt("GOSSIPSUB_D_SCORE", dLow) + dOut = getEnvInt("GOSSIPSUB_D_OUT", d div 2) + dLazy = getEnvInt("GOSSIPSUB_D_LAZY", d) + + heartbeatMs = getEnvInt("GOSSIPSUB_HEARTBEAT_MS", 1000) + 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) + maxLowPriorityQueueLen = getEnvInt("GOSSIPSUB_MAX_LOW_PRIORITY_QUEUE_LEN", 1024) + + slowPeerPenaltyWeight = getEnvFloat("GOSSIPSUB_SLOW_PEER_PENALTY_WEIGHT", 0.0) + slowPeerPenaltyThreshold = getEnvFloat("GOSSIPSUB_SLOW_PEER_PENALTY_THRESHOLD", 2.0) + slowPeerPenaltyDecay = getEnvFloat("GOSSIPSUB_SLOW_PEER_PENALTY_DECAY", 0.2) + + decayIntervalMs = getEnvInt("GOSSIPSUB_DECAY_INTERVAL_MS", 1000) + decayToZero = getEnvFloat("GOSSIPSUB_DECAY_TO_ZERO", 0.01) + + #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.heartbeatInterval = heartbeatMs.milliseconds + gossipSub.parameters.pruneBackoff = pruneBackoffSec.seconds + gossipSub.parameters.gossipFactor = getEnvFloat("GOSSIPSUB_GOSSIP_FACTOR", 0.25) + + gossipSub.parameters.d = d + gossipSub.parameters.dLow = dLow + gossipSub.parameters.dHigh = dHigh + gossipSub.parameters.dScore = dScore + gossipSub.parameters.dOut = dOut + gossipSub.parameters.dLazy = dLazy + + gossipSub.parameters.maxHighPriorityQueueLen = maxHighPriorityQueueLen + gossipSub.parameters.maxMediumPriorityQueueLen = maxMediumPriorityQueueLen + gossipSub.parameters.maxLowPriorityQueueLen = maxLowPriorityQueueLen + + gossipSub.parameters.slowPeerPenaltyWeight = slowPeerPenaltyWeight + gossipSub.parameters.slowPeerPenaltyThreshold = slowPeerPenaltyThreshold + gossipSub.parameters.slowPeerPenaltyDecay = slowPeerPenaltyDecay + + gossipSub.parameters.decayInterval = decayIntervalMs.milliseconds + gossipSub.parameters.decayToZero = decayToZero + + #gossipSub.parameters.gossipThreshold = gossipThreshold + #gossipSub.parameters.publishThreshold = publishThreshold + #gossipSub.parameters.graylistThreshold = graylistThreshold + + info "Configured GossipSub mesh params", + floodPublish = gossipSub.parameters.floodPublish, + opportunisticGraftThreshold = gossipSub.parameters.opportunisticGraftThreshold, + heartbeatMs = heartbeatMs, + pruneBackoffSec = pruneBackoffSec, + gossipFactor = gossipSub.parameters.gossipFactor, + d = d, + dLow = dLow, + dHigh = dHigh, + dScore = dScore, + dOut = dOut, + dLazy = dLazy + + info "Configured GossipSub queue and scoring params", + maxHighPriorityQueueLen = maxHighPriorityQueueLen, + maxMediumPriorityQueueLen = maxMediumPriorityQueueLen, + maxLowPriorityQueueLen = maxLowPriorityQueueLen, + slowPeerPenaltyWeight = slowPeerPenaltyWeight, + slowPeerPenaltyThreshold = slowPeerPenaltyThreshold, + slowPeerPenaltyDecay = slowPeerPenaltyDecay, + decayIntervalMs = decayIntervalMs, + decayToZero = decayToZero + #gossipThreshold = gossipThreshold, + #publishThreshold = publishThreshold, + #graylistThreshold = graylistThreshold + +proc subscribGossipsubTopic(gossipSub: GossipSub, topic: string) = + gossipSub.topicParams[topic] = TopicParams( + topicWeight: 1, + firstMessageDeliveriesWeight: 1, + firstMessageDeliveriesCap: 30, + firstMessageDeliveriesDecay: 0.9 + ) + + gossipSub.subscribe(topic, createMessageHandler()) + gossipSub.addValidator([topic], messageValidator) + + +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() + ) + else: + resolveTAddress(tAddress).mapIt(MultiAddress.init(it).tryGet()) + info "Address resolved", tAddress = tAddress, resolvedAddrs = resolvedAddrs + return ok(resolvedAddrs) + except CatchableError as exc: + if inShadow: + return err(exc.msg) + #keep trying for service mode + warn "Failed to resolve address", address = tAddress, error = exc.msg + await sleepAsync(15.seconds) + +proc connectGossipsubPeers( + switch: Switch, muxer: string, networkSize: int, myId: int, connectTo: int +): Future[Result[int, string]] {.async.} = + let rng = libp2p.newRng() + var + addrs: seq[MultiAddress] = @[] + tAddresses: seq[string] + connected = 0 + + if inShadow: + var peers = toSeq(0..= connectTo: break + try: + discard await switch.connect(peer, allowUnknownPeerId=true).wait(5.seconds) + connected.inc() + 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) + + if connected == 0: + return err("Failed to connect any peer") + elif connected < connectTo: + warn "Connected to fewer peers than target", connected = connected, target = connectTo + return ok(connected) + + +proc main {.async.} = + randomize() + let + rng = libp2p.newRng() + (myId, networkSize, connectTo, muxer, filePath, address) = getPeerDetails().valueOr: + error "Error reading peer settings ", err = error + return + + # Set global metric labels + gMuxer = muxer + + var + gossipSub: GossipSub + builder = SwitchBuilder + .new() + .withNoise() + .withAddress(MultiAddress.init(address).tryGet()) + .withMaxConnections(parseInt(getEnv("MAXCONNECTIONS", "250"))) + + builder = builder.withRng(rng) + + case muxer.toLowerAscii() + of "quic": + builder = builder.withQuicTransport() + of "yamux": + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) + .withYamux() + of "mplex": + builder = builder.withTcpTransport(flags = {ServerFlags.TcpNoDelay}) + .withMplex() + + let switch = builder.build() + + # Set peerId for metric labels + gPeerId = $switch.peerInfo.peerId + gossipSub = initializeGossipsub(switch, true) + + configureGossipsubParams(gossipSub) + subscribGossipsubTopic(gossipSub, "test") + switch.mount(gossipSub) + await switch.start() + + # Metrics + info "Starting metrics server" + let metricsServer = startMetricsServer(parseIpAddress("0.0.0.0"), prometheusPort) + if metricsServer.isErr: + error "Failed to initialize metrics server", err = metricsServer.error + elif inShadow: + asyncSpawn storeMetrics(myId) + + info "Listening on ", address = switch.peerInfo.addrs + info "Peer details ", peer = myId, peerId = switch.peerInfo.peerId + #Wait for node building + info "GossipSub codecs registered", codecs = gossipSub.codecs + await sleepAsync(60.seconds) + + #connect with peers + discard (await connectGossipsubPeers(switch, muxer, networkSize, myId, connectTo)).valueOr: + error "Failed to establish any connections", error = error + return + + 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 "Starting listening endpoint for publish controller" + discard gossipSub.startHttpServer(myId) + + await sleepAsync(2.days) + +waitFor(main()) \ No newline at end of file diff --git a/nim-test-node/gossipsub-queues/test_node.nimble b/nim-test-node/gossipsub-queues/test_node.nimble new file mode 100644 index 0000000..acbce1d --- /dev/null +++ b/nim-test-node/gossipsub-queues/test_node.nimble @@ -0,0 +1,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 = @[] + +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