Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 35 additions & 0 deletions .github/workflows/linters.yml
Original file line number Diff line number Diff line change
@@ -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
35 changes: 23 additions & 12 deletions nim-test-node/connmanager/env.nim
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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", "")

Expand Down
45 changes: 28 additions & 17 deletions nim-test-node/connmanager/main.nim
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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

Expand All @@ -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:
Expand Down Expand Up @@ -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..<cfg.numHubs:
let myIdx =
try:
parseInt(parts[^1])
except CatchableError:
0
for i in 0 ..< cfg.numHubs:
if i != myIdx:
let hubAddr = "hub-" & $i & ".nimp2p-service." & cfg.hubNamespace & ".svc.cluster.local:5000"
let hubAddr =
"hub-" & $i & ".nimp2p-service." & cfg.hubNamespace & ".svc.cluster.local:5000"
info "Dialing peer hub", target = hubAddr
asyncSpawn resolveAndConnect(switch, hubAddr)

while true: await sleepAsync(1.hours)
while true:
await sleepAsync(1.hours)

proc runPeer(cfg: PeerConfig) {.async.} =
var builder = SwitchBuilder
Expand Down Expand Up @@ -129,19 +134,25 @@ proc runPeer(cfg: PeerConfig) {.async.} =
await resolveAndConnect(switch, addr)
await sleepAsync(cfg.reconnectIntervalS.seconds)
for peerId in switch.connectedPeers(Direction.Out):
try: await switch.disconnect(peerId)
except CatchableError: discard
try:
await switch.disconnect(peerId)
except CatchableError:
discard
info "Cycled connection (grace abuse)", intervalS = cfg.reconnectIntervalS
else:
for addr in cfg.hubAddrs:
await resolveAndConnect(switch, addr)
while true: await sleepAsync(1.hours)
while true:
await sleepAsync(1.hours)
else:
while true: await sleepAsync(1.hours)
while true:
await sleepAsync(1.hours)

proc main() {.async.} =
case getRole()
of RoleHub: await runHub(parseHubConfig())
of RolePeer: await runPeer(parsePeerConfig())
of RoleHub:
await runHub(parseHubConfig())
of RolePeer:
await runPeer(parsePeerConfig())

waitFor main()
13 changes: 6 additions & 7 deletions nim-test-node/connmanager/test_node.nimble
Original file line number Diff line number Diff line change
@@ -1,11 +1,10 @@
mode = ScriptMode.Verbose

packageName = "test_node"
version = "0.1.0"
author = "Status Research & Development GmbH"
description = "Connection manager test node"
license = "MIT"
packageName = "test_node"
version = "0.1.0"
author = "Status Research & Development GmbH"
description = "Connection manager test node"
license = "MIT"

requires "nim >= 2.2.0",
"nimcrypto 0.6.4",
"libp2p#7cc4280e2efd5e6c2ebd732ce33d309376a9627e"
"nimcrypto 0.6.4", "libp2p#7cc4280e2efd5e6c2ebd732ce33d309376a9627e"
46 changes: 29 additions & 17 deletions nim-test-node/gossipsub-queues/env.nim
Original file line number Diff line number Diff line change
Expand Up @@ -6,35 +6,46 @@ 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"))
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

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))

Expand All @@ -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
Expand Down
Loading
Loading