From c82638862824cc7b6559f8aae5fb4a2320743401 Mon Sep 17 00:00:00 2001 From: Simon-Pierre Date: Thu, 11 Jun 2026 12:35:47 -0400 Subject: [PATCH] chat2disco init --- Makefile | 6 +- apps/chat2disco/chat2disco.nim | 439 ++++++++++++++++++++++++++ apps/chat2disco/config_chat2disco.nim | 103 ++++++ apps/chat2disco/nim.cfg | 4 + logos_delivery.nimble | 11 + 5 files changed, 562 insertions(+), 1 deletion(-) create mode 100644 apps/chat2disco/chat2disco.nim create mode 100644 apps/chat2disco/config_chat2disco.nim create mode 100644 apps/chat2disco/nim.cfg diff --git a/Makefile b/Makefile index daa8a8ad7..60c26751b 100644 --- a/Makefile +++ b/Makefile @@ -219,7 +219,7 @@ testcommon: | build-deps build ########## ## Waku ## ########## -.PHONY: testwaku wakunode2 testwakunode2 example2 chat2 chat2bridge liteprotocoltester +.PHONY: testwaku wakunode2 testwakunode2 example2 chat2 chat2mix chat2disco chat2bridge liteprotocoltester testwaku: | build-deps build rln-deps librln echo -e $(BUILD_MSG) "build/$@" && \ @@ -256,6 +256,10 @@ chat2mix: | build-deps build deps librln echo -e $(BUILD_MSG) "build/$@" && \ $(NIMBLE) chat2mix +chat2disco: | build-deps build deps librln + echo -e $(BUILD_MSG) "build/$@" && \ + $(NIMBLE) chat2disco + rln-db-inspector: | build-deps build deps librln echo -e $(BUILD_MSG) "build/$@" && \ $(NIMBLE) rln_db_inspector diff --git a/apps/chat2disco/chat2disco.nim b/apps/chat2disco/chat2disco.nim new file mode 100644 index 000000000..83151e857 --- /dev/null +++ b/apps/chat2disco/chat2disco.nim @@ -0,0 +1,439 @@ +# chat2disco is a minimal chat app for testing Waku Kademlia service discovery. +# /room subscribes to the pubsub topic of that name and starts + +when not (compileOption("threads")): + {.fatal: "Please, compile this program with the --threads:on option!".} + +{.push raises: [].} + +import std/[strformat, strutils, times, options, sequtils] +import + confutils, + chronicles, + chronos, + eth/keys, + bearssl, + results, + stew/[byteutils], + metrics, + metrics/chronos_httpserver +import + libp2p/[ + switch, + crypto/crypto, + stream/connection, + multiaddress, + peerinfo, + peerid, + protobuf/minprotobuf, + protocols/kademlia/types, + protocols/service_discovery, + protocols/service_discovery/types as sd_types, + ] +import + logos_delivery/waku/[ + waku_core, + waku_enr, + discovery/waku_kademlia, + waku_node, + node/waku_metrics, + node/peer_manager, + factory/builder, + factory/conf_builder/kademlia_discovery_conf_builder, + common/utils/nat, + common/logging, + discovery/autonat_service, + ], + ./config_chat2disco + +import libp2p/protocols/pubsub/rpc/messages, libp2p/protocols/pubsub/pubsub +import libp2p/extended_peer_record # for ServiceInfo + +logScope: + topics = "chat2disco" + +const Help = """ + Commands: /[?|help|connect|nick|room|exit] + help: Prints this help + connect: dials a remote peer + nick: change nickname for current chat session + room : subscribe to pubsub topic and advertise + discover the kademlia service of the same name (kademlia always active) + exit: exits chat session +""" + +const DefaultContentTopic = "/chat2disco/1/room/proto" + +type Chat = ref object + node: WakuNode + transp: StreamTransport + subscribed: bool + connected: bool + started: bool + nick: string + prompt: bool + currentPubsubTopic: string + +type + PrivateKey* = crypto.PrivateKey + Topic* = waku_core.PubsubTopic + +##################### +## chat2 protobufs ## +##################### + +type + SelectResult*[T] = Result[T, string] + + Chat2Message* = object + timestamp*: int64 + nick*: string + payload*: seq[byte] + +proc init*(T: type Chat2Message, buffer: seq[byte]): ProtoResult[T] = + var msg = Chat2Message() + let pb = initProtoBuffer(buffer) + + var timestamp: uint64 + discard ?pb.getField(1, timestamp) + msg.timestamp = int64(timestamp) + + discard ?pb.getField(2, msg.nick) + discard ?pb.getField(3, msg.payload) + + ok(msg) + +proc encode*(message: Chat2Message): ProtoBuffer = + var serialised = initProtoBuffer() + + serialised.write(1, uint64(message.timestamp)) + serialised.write(2, message.nick) + serialised.write(3, message.payload) + + return serialised + +proc `$`*(message: Chat2Message): string = + let time = message.timestamp.fromUnix().local().format("'<'MMM' 'dd,' 'HH:mm'>'") + return time & " " & message.nick & ": " & string.fromBytes(message.payload) + +##################### + +proc connectToNodes(c: Chat, nodes: seq[string]) {.async.} = + echo "Connecting to nodes" + await c.node.connectToNodes(nodes) + c.connected = true + +proc showChatPrompt(c: Chat) = + if not c.prompt: + try: + stdout.write(">> ") + stdout.flushFile() + c.prompt = true + except IOError: + discard + +proc getChatLine(payload: seq[byte]): string = + let pb = Chat2Message.init(payload).valueOr: + return string.fromBytes(payload) + return $pb + +proc printReceivedMessage(c: Chat, pubsubTopic: PubsubTopic, msg: WakuMessage) = + let chatLine = getChatLine(msg.payload) + try: + echo &"[{pubsubTopic}] {chatLine}" + except ValueError: + echo chatLine + c.prompt = false + showChatPrompt(c) + trace "Printing message", chatLine, pubsubTopic, contentTopic = msg.contentTopic + +proc readNick(transp: StreamTransport): Future[string] {.async.} = + stdout.write("Choose a nickname >> ") + stdout.flushFile() + return await transp.readLine() + +proc readRoom(transp: StreamTransport): Future[string] {.async.} = + stdout.write("Choose a room >> ") + stdout.flushFile() + return await transp.readLine() + +proc startMetricsServer( + serverIp: IpAddress, serverPort: Port +): Result[MetricsHttpServerRef, string] = + info "Starting metrics HTTP server", serverIp = $serverIp, serverPort = $serverPort + + let server = MetricsHttpServerRef.new($serverIp, serverPort).valueOr: + return err("metrics HTTP server start failed: " & $error) + + try: + waitFor server.start() + except CatchableError: + return err("metrics HTTP server start failed: " & getCurrentExceptionMsg()) + + info "Metrics HTTP server started", serverIp = $serverIp, serverPort = $serverPort + ok(server) + +proc publish(c: Chat, line: string) {.async.} = + if c.currentPubsubTopic.len == 0: + echo "No active room. Use /room to join one first." + return + + let time = getTime().toUnix() + let chat2pb = + Chat2Message(timestamp: time, nick: c.nick, payload: line.toBytes()).encode() + + var message = WakuMessage( + payload: chat2pb.buffer, + contentTopic: DefaultContentTopic, + version: 0, + timestamp: getNanosecondTime(time), + ) + + try: + (await c.node.publish(some(c.currentPubsubTopic), message)).isOkOr: + error "failed to publish message", error = error + except CatchableError: + error "caught error publishing message: ", error = getCurrentExceptionMsg() + +proc joinRoom(c: Chat, roomName: string) {.async.} = + if c.currentPubsubTopic.len > 0 and c.currentPubsubTopic != roomName: + discard c.node.unsubscribe((kind: PubsubSub, topic: c.currentPubsubTopic)) + + c.currentPubsubTopic = roomName + + proc roomHandler( + pubsubTopic: PubsubTopic, msg: WakuMessage + ): Future[void] {.async, gcsafe.} = + c.printReceivedMessage(pubsubTopic, msg) + + c.node.subscribe((kind: PubsubSub, topic: roomName), WakuRelayHandler(roomHandler)).isOkOr: + error "failed to subscribe to pubsub topic for room", + topic = roomName, error = error + + echo "subscribed to pubsub topic: ", roomName + + # Kademlia service discovery + advertisement for the room. + # Kademlia is always mounted (seed node if no --kad-bootstrap-node provided). + let svc = roomName + discard c.node.wakuKademlia.protocol.registerInterest(svc) + let svcInfo = ServiceInfo(id: svc, data: @[]) + c.node.wakuKademlia.protocol.startAdvertising(svcInfo) + echo "advertising and discovering kademlia service: ", svc + + asyncSpawn ( + proc() {.async.} = + let res = await c.node.wakuKademlia.lookupServicePeers(svc) + if res.isOk(): + let peers = res.get() + if peers.len > 0: + echo "discovered room peers via kademlia: ", peers.mapIt($it.peerId) + await c.node.connectToNodes(peers) + else: + error "service lookup failed for room", service = svc, error = res.error + )() + +proc readAndPrint(c: Chat) {.async.} = + while true: + await sleepAsync(100.millis) + +proc writeAndPrint(c: Chat) {.async.} = + while true: + showChatPrompt(c) + + let line = await c.transp.readLine() + if line.startsWith("/help") or line.startsWith("/?") or not c.started: + echo Help + continue + elif line.startsWith("/connect"): + if c.connected: + echo "already connected to at least one peer" + continue + echo "enter address of remote peer" + let address = await c.transp.readLine() + if address.len > 0: + await c.connectToNodes(@[address]) + elif line.startsWith("/nick"): + c.nick = await readNick(c.transp) + echo "You are now known as " & c.nick + elif line.startsWith("/room"): + let parts = line.split(maxsplit = 1) + if parts.len < 2 or parts[1].strip().len == 0: + echo "usage: /room " + continue + let roomName = parts[1].strip() + await c.joinRoom(roomName) + elif line.startsWith("/exit"): + echo "quitting..." + + try: + await c.node.stop() + except: + echo "exception happened when stopping: " & getCurrentExceptionMsg() + + quit(QuitSuccess) + else: + if c.started: + if c.currentPubsubTopic.len > 0: + echo "publishing message to ", c.currentPubsubTopic, ": ", line + await c.publish(line) + else: + echo "Join a room first with /room before sending messages" + else: + try: + if line.startsWith("/") and "p2p" in line: + await c.connectToNodes(@[line]) + except: + echo &"unable to dial remote peer {line}" + echo getCurrentExceptionMsg() + +proc readWriteLoop(c: Chat) {.async.} = + asyncSpawn c.writeAndPrint() + asyncSpawn c.readAndPrint() + +proc readInput(wfd: AsyncFD) {.thread, raises: [Defect, CatchableError].} = + let transp = fromPipe(wfd) + + while true: + let line = stdin.readLine() + discard waitFor transp.write(line & "\r\n") + +{.pop.} + # @TODO confutils.nim(775, 17) Error: can raise an unlisted exception: ref IOError +proc processInput(rfd: AsyncFD, rng: crypto.Rng) {.async.} = + let + transp = fromPipe(rfd) + conf = Chat2Conf.load() + nodekey = + if conf.nodekey.isSome(): + conf.nodekey.get() + else: + PrivateKey.random(Secp256k1, rng).tryGet() + + # set log level + if conf.logLevel != LogLevel.NONE: + setLogLevel(conf.logLevel) + + let (extIp, extTcpPort, extUdpPort) = setupNat( + conf.nat, + clientId, + Port(uint16(conf.tcpPort) + conf.portsShift), + Port(uint16(conf.udpPort) + conf.portsShift), + ).valueOr: + raise newException(ValueError, "setupNat error " & error) + + var enrBuilder = EnrBuilder.init(nodeKey) + + let record = enrBuilder.build().valueOr: + error "failed to create enr record", error = error + quit(QuitFailure) + + let node = block: + var builder = WakuNodeBuilder.init() + builder.withNodeKey(nodeKey) + builder.withRecord(record) + + builder + .withNetworkConfigurationDetails( + conf.listenAddress, + Port(uint16(conf.tcpPort) + conf.portsShift), + extIp, + extTcpPort, + ) + .tryGet() + builder.build().tryGet() + + (await node.mountRelay()).isOkOr: + error "failed to mount relay: ", error = error + quit(QuitFailure) + + await node.mountLibp2pPing() + + var kadBootstrapPeers: seq[(PeerId, seq[MultiAddress])] + for nodeStr in conf.kadBootstrapNodes: + let (peerId, ma) = block: + let r = parseFullAddress(nodeStr) + if r.isErr: + error "Failed to parse kademlia bootstrap node", node = nodeStr, error = r.error + continue + r.get() + kadBootstrapPeers.add((peerId, @[ma])) + + node.mountKademlia( + KademliaDiscoveryConf( + bootstrapNodes: kadBootstrapPeers, + servicesToAdvertise: @[], + servicesToDiscover: @[], + randomLookupInterval: chronos.seconds(60), + serviceLookupInterval: chronos.seconds(60), + kadDhtConfig: KadDHTConfig.new(), + discoConfig: + sd_types.ServiceDiscoveryConfig.new(advertExpiry = chronos.seconds(60)), + clientMode: false, + xprPublishing: false, + ) + ).isOkOr: + error "failed to setup kademlia service discovery", error = error + quit(QuitFailure) + + let autonatService = getAutonatService(rng) + node.switch.services = @[Service(autonatService)] + for service in node.switch.services: + try: + service.setup(node.switch) + except ServiceSetupError as e: + error "failed to set up libp2p switch service", error = e.msg + + await node.start() + + node.peerManager.start() + + let nick = await readNick(transp) + let room = await readRoom(transp) + echo "Welcome, " & nick & "! Joined room '" & room & "'." + + var chat = Chat( + node: node, + transp: transp, + subscribed: false, + connected: false, + started: true, + nick: nick, + prompt: false, + currentPubsubTopic: "", + ) + + let peerInfo = node.switch.peerInfo + let listenStr = $peerInfo.addrs[0] & "/p2p/" & $peerInfo.peerId + echo &"Listening on\n {listenStr}" + + if conf.metricsLogging: + startMetricsLog() + + if conf.metricsServer: + let metricsServer = startMetricsServer( + conf.metricsServerAddress, Port(conf.metricsServerPort + conf.portsShift) + ) + + await chat.joinRoom(room) + + await chat.readWriteLoop() + + runForever() + +proc main(rng: crypto.Rng) {.async.} = + let (rfd, wfd) = createAsyncPipe() + if rfd == asyncInvalidPipe or wfd == asyncInvalidPipe: + raise newException(ValueError, "Could not initialize pipe!") + + var thread: Thread[AsyncFD] + thread.createThread(readInput, wfd) + try: + await processInput(rfd, rng) + except ConfigurationError as e: + raise e + +when isMainModule: + let rng = crypto.newRng() + try: + waitFor(main(rng)) + except CatchableError as e: + raise e diff --git a/apps/chat2disco/config_chat2disco.nim b/apps/chat2disco/config_chat2disco.nim new file mode 100644 index 000000000..9f1a8290a --- /dev/null +++ b/apps/chat2disco/config_chat2disco.nim @@ -0,0 +1,103 @@ +import chronicles, chronos, std/strutils + +import + eth/keys, + libp2p/crypto/crypto, + libp2p/crypto/secp, + libp2p/multiaddress, + nimcrypto/utils, + confutils, + confutils/defs, + confutils/std/net + +type Chat2Conf* = object ## General node config + logLevel* {. + desc: "Sets the log level.", defaultValue: LogLevel.DEBUG, name: "log-level" + .}: LogLevel + + nodekey* {.desc: "P2P node private key as 64 char hex string.", name: "nodekey".}: + Option[crypto.PrivateKey] + + listenAddress* {. + defaultValue: defaultListenAddress(config), + desc: "Listening address for the LibP2P traffic.", + name: "listen-address" + .}: IpAddress + + tcpPort* {.desc: "TCP listening port.", defaultValue: 60000, name: "tcp-port".}: Port + + udpPort* {.desc: "UDP listening port.", defaultValue: 60000, name: "udp-port".}: Port + + portsShift* {. + desc: "Add a shift to all port numbers.", defaultValue: 0, name: "ports-shift" + .}: uint16 + + nat* {. + desc: + "Specify method to use for determining public address. " & + "Must be one of: any, none, upnp, pmp, extip:.", + defaultValue: "any" + .}: string + + ## Metrics config + metricsServer* {. + desc: "Enable the metrics server: true|false", + defaultValue: false, + name: "metrics-server" + .}: bool + + metricsServerAddress* {. + desc: "Listening address of the metrics server.", + defaultValue: IpAddress(family: IpAddressFamily.IPv4, address_v4: [127'u8, 0, 0, 1]), + name: "metrics-server-address" + .}: IpAddress + + metricsServerPort* {. + desc: "Listening HTTP port of the metrics server.", + defaultValue: 8008, + name: "metrics-server-port" + .}: uint16 + + metricsLogging* {. + desc: "Enable metrics logging: true|false", + defaultValue: true, + name: "metrics-logging" + .}: bool + + ## Kademlia Discovery config (core for chat2disco) + kadBootstrapNodes* {. + desc: + "Optional peer multiaddr for kademlia bootstrap (must include /p2p/). Argument may be repeated.", + name: "kad-bootstrap-node" + .}: seq[string] + +proc parseCmdArg*(T: type crypto.PrivateKey, p: string): T = + try: + let key = SkPrivateKey.init(utils.fromHex(p)).tryGet() + result = crypto.PrivateKey(scheme: Secp256k1, skkey: key) + except CatchableError as e: + raise newException(ValueError, "Invalid private key") + +proc completeCmdArg*(T: type crypto.PrivateKey, val: string): seq[string] = + return @[] + +proc parseCmdArg*(T: type IpAddress, p: string): T = + try: + result = parseIpAddress(p) + except CatchableError as e: + raise newException(ValueError, "Invalid IP address") + +proc completeCmdArg*(T: type IpAddress, val: string): seq[string] = + return @[] + +proc parseCmdArg*(T: type Port, p: string): T = + try: + result = Port(parseInt(p)) + except CatchableError as e: + raise newException(ValueError, "Invalid Port number") + +proc completeCmdArg*(T: type Port, val: string): seq[string] = + return @[] + +func defaultListenAddress*(conf: Chat2Conf): IpAddress = + (static parseIpAddress("0.0.0.0")) diff --git a/apps/chat2disco/nim.cfg b/apps/chat2disco/nim.cfg new file mode 100644 index 000000000..2231f2ebe --- /dev/null +++ b/apps/chat2disco/nim.cfg @@ -0,0 +1,4 @@ +-d:chronicles_line_numbers +-d:chronicles_runtime_filtering:on +-d:discv5_protocol_id:d5waku +path = "../.." diff --git a/logos_delivery.nimble b/logos_delivery.nimble index efe118595..9c9ef7b96 100644 --- a/logos_delivery.nimble +++ b/logos_delivery.nimble @@ -439,6 +439,17 @@ task chat2mix, "Build example Waku chat mix usage": "-d:chronicles_sinks=textlines[file] -d:chronicles_log_level=TRACE " # -d:ssl - cause unlisted exception error in libp2p/utility... +task chat2disco, "Build chat2disco (kademlia service discovery test chat)": + # This is a testing/discovery tool, so we want debug logs on stdout by default + # (the config sets the default to DEBUG, and runtime filtering is enabled). + # Unlike the pure chat examples, we do not redirect to a file sink. + + let name = "chat2disco" + buildBinary name, + "apps/chat2disco/", + "-d:chronicles_log_level=TRACE " + # -d:ssl - cause unlisted exception error in libp2p/utility... + task chat2bridge, "Build chat2bridge": let name = "chat2bridge" buildBinary name, "apps/chat2bridge/"