diff --git a/logos_delivery/api/kernel_conf.nim b/logos_delivery/api/kernel_conf.nim new file mode 100644 index 000000000..f2604c936 --- /dev/null +++ b/logos_delivery/api/kernel_conf.nim @@ -0,0 +1,6 @@ +## Kernel-layer configuration type. + +import tools/confutils/cli_args +export cli_args + +type KernelConf* = WakuNodeConf diff --git a/logos_delivery/api/logos_delivery_interface.nim b/logos_delivery/api/logos_delivery_interface.nim index b0d2bd975..11532384b 100644 --- a/logos_delivery/api/logos_delivery_interface.nim +++ b/logos_delivery/api/logos_delivery_interface.nim @@ -15,6 +15,12 @@ import ./reliable_channel_manager_interface as ireliablechannelmanager_iface export ikernel_iface, imessagingclient_iface, ireliablechannelmanager_iface +import ./messaging_conf # MessagingConf +export messaging_conf + +import logos_delivery/channels/channels_conf # ChannelsConf +export channels_conf + BrokerInterface(LogosDeliveryInterface): EventBroker: type ConnectionStatusChangeEvent* = object @@ -25,7 +31,10 @@ BrokerInterface(LogosDeliveryInterface): RequestBroker: proc startAsClient( - mode: WakuMode, preset: string + mode: WakuMode, + preset: string, + overrides: MessagingConf, + channelsConf: ChannelsConf, ): Future[Result[MessagingClientInterface, string]] {.async.} RequestBroker: diff --git a/logos_delivery/api/messaging_conf.nim b/logos_delivery/api/messaging_conf.nim new file mode 100644 index 000000000..f81d91ab5 --- /dev/null +++ b/logos_delivery/api/messaging_conf.nim @@ -0,0 +1,90 @@ +## Messaging-layer configuration. + +import std/options +import std/sequtils # mapIt +import results +import std/nativesockets # Port +import std/net # IpAddress + +import logos_delivery/api/types # WakuMode +import logos_delivery/api/kernel_conf # KernelConf, ConfResult, defaultWakuNodeConf, EthRpcUrl +export kernel_conf + +type MessagingConf* = object + ## Messaging-layer configuration as an all-`Option` partial. `toKernelConf` + ## applies each set field to the corresponding `KernelConf` field; unset fields + ## leave the kernel default unchanged. + clusterId*: Option[uint16] + ## Network cluster id. + numShardsInCluster*: Option[uint16] + ## Number of shards in the cluster. + p2pTcpPort*: Option[Port] + ## TCP listening port. + discv5UdpPort*: Option[Port] + ## discv5 UDP port. + listenIpv4*: Option[IpAddress] + ## Inbound bind address. + maxMessageSize*: Option[string] + ## Maximum accepted message size (e.g. "150 KiB"). + entryNodes*: Option[seq[string]] + ## Bootstrap / connectivity nodes (enrtree or multiaddr). + ethRpcEndpoints*: Option[seq[string]] + ## Ethereum RPC endpoints (required for RLN validation); multiple for fail-over. + rlnContractAddress*: Option[string] + ## RLN contract address; when set, RLN validation is enabled. + rlnChainId*: Option[uint] + ## Chain id the RLN contract is deployed on. + rlnEpochSizeSec*: Option[uint] + ## RLN epoch size, in seconds. + reliabilityEnabled*: Option[bool] + ## Enable store-based send reliability. + +proc toKernelConf*(m: MessagingConf, mode: WakuMode): ConfResult[KernelConf] = + ## Build a `KernelConf` from the operation mode and the messaging configuration. + ## The mode sets the kernel protocol flags; each set field is then written to its + ## kernel counterpart, and unset fields keep the kernel default. + var conf = ?defaultWakuNodeConf() + + case mode + of WakuMode.Core: + conf.relay = true + conf.filter = true + conf.lightpush = true + conf.discv5Discovery = some(true) + conf.peerExchange = true + conf.rendezvous = true + if conf.rateLimits.len == 0: + conf.rateLimits = @["filter:100/1s", "lightpush:5/1s", "px:5/1s"] + of WakuMode.Edge: + conf.peerExchange = true + conf.relay = false + conf.filter = false + conf.lightpush = false + conf.store = false + + if m.clusterId.isSome(): + conf.clusterId = m.clusterId + if m.numShardsInCluster.isSome(): + conf.numShardsInNetwork = m.numShardsInCluster.get() + if m.p2pTcpPort.isSome(): + conf.tcpPort = m.p2pTcpPort.get() + if m.discv5UdpPort.isSome(): + conf.discv5UdpPort = m.discv5UdpPort.get() + if m.listenIpv4.isSome(): + conf.listenAddress = m.listenIpv4.get() + if m.maxMessageSize.isSome(): + conf.maxMessageSize = m.maxMessageSize.get() + if m.entryNodes.isSome(): + conf.entryNodes = m.entryNodes.get() + if m.ethRpcEndpoints.isSome(): + conf.ethClientUrls = m.ethRpcEndpoints.get().mapIt(EthRpcUrl(it)) + if m.rlnContractAddress.isSome(): + conf.rlnRelayEthContractAddress = m.rlnContractAddress.get() + conf.rlnRelay = some(true) + if m.rlnChainId.isSome(): + conf.rlnRelayChainId = m.rlnChainId.get() + if m.rlnEpochSizeSec.isSome(): + conf.rlnEpochSizeSec = some(m.rlnEpochSizeSec.get().uint64) + if m.reliabilityEnabled.isSome(): + conf.reliabilityEnabled = m.reliabilityEnabled + ok(conf) diff --git a/logos_delivery/api/messaging_conf_preset.nim b/logos_delivery/api/messaging_conf_preset.nim new file mode 100644 index 000000000..96833147a --- /dev/null +++ b/logos_delivery/api/messaging_conf_preset.nim @@ -0,0 +1,58 @@ +## Preset resolution and override merge for the Messaging layer. +## +## `resolvePreset` turns a network preset name into the Messaging-layer fields it +## implies; `merge` applies a caller's overrides on top. `startAsClient` uses both +## before inferring the kernel config. The kernel still resolves the preset itself +## for the node-local fields a preset does not carry. + +import std/options +import results +import stint # UInt256.truncate + +import logos_delivery/api/messaging_conf + # MessagingConf, ConfResult, toNetworkPresetConf +import logos_delivery/waku/factory/networks_config # NetworkPresetConf, AutoSharding + +proc merge*(base, overrides: MessagingConf): MessagingConf = + ## Combine two messaging configs field by field: a set `overrides` field wins; + ## otherwise the `base` value is kept. + result = base + if overrides.clusterId.isSome(): result.clusterId = overrides.clusterId + if overrides.numShardsInCluster.isSome(): + result.numShardsInCluster = overrides.numShardsInCluster + if overrides.p2pTcpPort.isSome(): result.p2pTcpPort = overrides.p2pTcpPort + if overrides.discv5UdpPort.isSome(): result.discv5UdpPort = overrides.discv5UdpPort + if overrides.listenIpv4.isSome(): result.listenIpv4 = overrides.listenIpv4 + if overrides.maxMessageSize.isSome(): result.maxMessageSize = overrides.maxMessageSize + if overrides.entryNodes.isSome(): result.entryNodes = overrides.entryNodes + if overrides.ethRpcEndpoints.isSome(): + result.ethRpcEndpoints = overrides.ethRpcEndpoints + if overrides.rlnContractAddress.isSome(): + result.rlnContractAddress = overrides.rlnContractAddress + if overrides.rlnChainId.isSome(): result.rlnChainId = overrides.rlnChainId + if overrides.rlnEpochSizeSec.isSome(): + result.rlnEpochSizeSec = overrides.rlnEpochSizeSec + if overrides.reliabilityEnabled.isSome(): + result.reliabilityEnabled = overrides.reliabilityEnabled + +proc resolvePreset*(preset: string): ConfResult[MessagingConf] = + ## Resolve a network preset name into the Messaging-layer fields it implies. + ## Node-local fields (ports, bind address, RPC endpoints) are not preset-defined + ## and stay unset. An empty preset resolves to an empty config. + let npcOpt = ?toNetworkPresetConf(preset, none(uint16)) + if npcOpt.isNone(): + return ok(MessagingConf()) + let npc = npcOpt.get() + var m = MessagingConf() + m.clusterId = some(npc.clusterId) + if npc.shardingConf.kind == AutoSharding: + m.numShardsInCluster = some(npc.shardingConf.numShardsInCluster) + m.maxMessageSize = some(npc.maxMessageSize) + if npc.entryNodes.len > 0: + m.entryNodes = some(npc.entryNodes) + if npc.rlnRelay and npc.rlnRelayEthContractAddress.len > 0: + m.rlnContractAddress = some(npc.rlnRelayEthContractAddress) + m.rlnChainId = some(npc.rlnRelayChainId.truncate(uint)) + m.rlnEpochSizeSec = some(npc.rlnEpochSizeSec.uint) + m.reliabilityEnabled = some(npc.p2pReliability) + ok(m) diff --git a/logos_delivery/channels/channels_conf.nim b/logos_delivery/channels/channels_conf.nim new file mode 100644 index 000000000..5d817ec5f --- /dev/null +++ b/logos_delivery/channels/channels_conf.nim @@ -0,0 +1,33 @@ +## Reliable-channels configuration. + +import std/options + +import logos_delivery/channels/segmentation/segmentation_persistence + # SegmentationPersistence +import types/persistence # Persistence (nim-sds) + +type ChannelsConf* = object + ## Reliable-channels configuration as an all-`Option` partial. Unset fields fall + ## back to the defaults used by `createReliableChannel`. + # Segmentation + segmentationEnableReedSolomon*: Option[bool] + ## Add Reed-Solomon parity segments for recovery of lost segments. + segmentationSegmentSizeBytes*: Option[int] + ## Maximum segment size in bytes. + # SDS + sdsAcknowledgementTimeoutMs*: Option[int] + ## Time to wait before retransmitting an unacknowledged message. + sdsMaxRetransmissions*: Option[int] + ## Maximum retransmission attempts before delivery fails. + sdsCausalHistorySize*: Option[int] + ## Number of message ids kept in causal history. + # Rate limiting + rateLimitEnabled*: Option[bool] + ## Enable rate limiting. + rateLimitEpochPeriodSec*: Option[int] + ## Rate-limit epoch length in seconds. + # Pluggable backends (dependency injection) + segmentationPersistence*: Option[SegmentationPersistence] + ## Persists partial reassembly state across restarts. + sdsPersistence*: Option[Persistence] + ## Persists SDS local history. diff --git a/logos_delivery/channels/reliable_channel_manager.nim b/logos_delivery/channels/reliable_channel_manager.nim index ceb8902e1..ad8f3e12d 100644 --- a/logos_delivery/channels/reliable_channel_manager.nim +++ b/logos_delivery/channels/reliable_channel_manager.nim @@ -21,8 +21,10 @@ import logos_delivery/waku/persistency/sds_persistency import ./reliable_channel import ./encryption/noop_encryption +import ./channels_conf export reliable_channel +export channels_conf const SdsJobId = "sds" ## One persistency job shared by every channel's SDS state; rows are @@ -35,6 +37,9 @@ type ReliableChannelManager* = ref object of ReliableChannelManagerInterface ## Default egress dispatch for channels created through this manager. ## Constructed at mount time as a closure over `MessagingClient.send` ## so the channel layer itself stays callable-only. + channelsConf: ChannelsConf + ## Configuration applied to channels created by this manager; unset fields + ## fall back to the component defaults. proc start*(self: ReliableChannelManager): Result[void, string] = ## Placeholder: per-channel listeners are installed in `ReliableChannel.new`, @@ -63,11 +68,14 @@ proc sdsPersistence(): Option[Persistence] = BrokerImplement ReliableChannelManager of ReliableChannelManagerInterface: proc new( - T: typedesc[ReliableChannelManager], messagingClient: MessagingClientInterface + T: typedesc[ReliableChannelManager], + messagingClient: MessagingClientInterface, + channelsConf: ChannelsConf, ): T = T( channels: initTable[ChannelId, ReliableChannel](), messagingClient: messagingClient, + channelsConf: channelsConf, ) method createReliableChannel*( @@ -87,19 +95,26 @@ BrokerImplement ReliableChannelManager of ReliableChannelManagerInterface: if self.channels.hasKey(channelId): return err("channel already exists: " & channelId) + # Apply the channel configuration, falling back to component defaults for + # unset fields. + let cc = self.channelsConf let segConfig = SegmentationConfig( - segmentSizeBytes: DefaultSegmentSizeBytes, - enableReedSolomon: false, - persistence: nil, + segmentSizeBytes: cc.segmentationSegmentSizeBytes.get(DefaultSegmentSizeBytes), + enableReedSolomon: cc.segmentationEnableReedSolomon.get(false), + persistence: cc.segmentationPersistence.get(nil), ) let sdsConfig = SdsConfig( - acknowledgementTimeoutMs: DefaultAcknowledgementTimeoutMs, - maxRetransmissions: DefaultMaxRetransmissions, - causalHistorySize: DefaultCausalHistorySize, - persistence: sdsPersistence(), + acknowledgementTimeoutMs: + cc.sdsAcknowledgementTimeoutMs.get(DefaultAcknowledgementTimeoutMs), + maxRetransmissions: cc.sdsMaxRetransmissions.get(DefaultMaxRetransmissions), + causalHistorySize: cc.sdsCausalHistorySize.get(DefaultCausalHistorySize), + persistence: + if cc.sdsPersistence.isSome(): cc.sdsPersistence else: sdsPersistence(), ) let rateConfig = RateLimitConfig( - epochPeriodSec: DefaultEpochPeriodSec, messagesPerEpoch: DefaultMessagesPerEpoch + enabled: cc.rateLimitEnabled.get(false), + epochPeriodSec: cc.rateLimitEpochPeriodSec.get(DefaultEpochPeriodSec), + messagesPerEpoch: DefaultMessagesPerEpoch, ) let chn = ReliableChannel.new( diff --git a/logos_delivery/logos_delivery.nim b/logos_delivery/logos_delivery.nim index 5ddb2dec5..c8f938a63 100644 --- a/logos_delivery/logos_delivery.nim +++ b/logos_delivery/logos_delivery.nim @@ -10,13 +10,14 @@ ## does NOT start it. `kernel()` is a stub until a KernelImpl exists. import results, chronos -import std/options # some()/Option for WakuNodeConf.mode import std/json # newJArray/newJObject/%*/pretty in getAvailableConfigs import brokers/broker_context, brokers/broker_interface, brokers/broker_implement import logos_delivery/api/logos_delivery_interface as logosdelivery_iface import logos_delivery/api/types as api_types import logos_delivery/api/kernel_interface as kernel_iface +import logos_delivery/api/messaging_conf # MessagingConf +import logos_delivery/api/messaging_conf_preset # resolvePreset, merge import logos_delivery/waku/factory/waku, @@ -30,6 +31,7 @@ type LogosDelivery* = ref object of LogosDeliveryInterface waku: Waku ## the owned node facade (built in initializeRequest) messagingClient: MessagingClient reliableChannelManager: ReliableChannelManager + channelsConf: ChannelsConf ## set by startAsClient; applied when channels mount proc loadConf(configPath: string): Result[WakuNodeConf, string] = ## Delegates to cli_args' concrete loader so confutils' `load` macro expands in @@ -59,7 +61,7 @@ proc initReliableChannelManager( ): Result[ReliableChannelManager, string] = return ok( ReliableChannelManager.createUnderContext( - globalBrokerContext(), self.messagingClient + globalBrokerContext(), self.messagingClient, self.channelsConf ) ) @@ -100,6 +102,8 @@ BrokerImplement LogosDelivery of LogosDeliveryInterface: method startAsNode( self: LogosDelivery, config: string ): Future[Result[void, string]] {.async.} = + ## Load a complete node configuration from `config` and start a bare kernel + ## node without a messaging client. if not self.waku.isNil(): return err("already initialized") let conf = loadConf(config).valueOr: @@ -117,7 +121,11 @@ BrokerImplement LogosDelivery of LogosDeliveryInterface: return err("initialize failed: " & e.msg) method startAsClient( - self: LogosDelivery, mode: api_types.WakuMode, preset: string + self: LogosDelivery, + mode: api_types.WakuMode, + preset: string, + overrides: MessagingConf, + channelsConf: ChannelsConf, ): Future[Result[MessagingClientInterface, string]] {.async.} = if not self.messagingClient.isNil(): return err("already initialized") @@ -126,9 +134,14 @@ BrokerImplement LogosDelivery of LogosDeliveryInterface: "already started as node; cannot start as client, but you can use as client" ) + self.channelsConf = channelsConf try: - var conf: WakuNodeConf = ?defaultWakuNodeConf() - conf.mode = some(mode) + # Resolve the preset into the messaging fields it implies, apply the caller's + # overrides on top, then infer the kernel configuration. The kernel still + # resolves the preset itself for the node-local fields a preset does not carry. + let base = ?resolvePreset(preset) + let msgConf = base.merge(overrides) + var conf: KernelConf = ?msgConf.toKernelConf(mode) conf.preset = preset self.waku = (await createNode(conf)).valueOr: diff --git a/logos_delivery/waku/api/api.nim b/logos_delivery/waku/api/api.nim index 025a65179..a991e6106 100644 --- a/logos_delivery/waku/api/api.nim +++ b/logos_delivery/waku/api/api.nim @@ -62,7 +62,7 @@ proc mountReliableChannelManager*(w: Waku): Result[void, string] = return ok() w.reliableChannelManager = ReliableChannelManager.createUnderContext( - w.brokerCtx, MessagingClientInterface(w.messagingClient) + w.brokerCtx, MessagingClientInterface(w.messagingClient), ChannelsConf() ) return ok() diff --git a/tests/api/test_all.nim b/tests/api/test_all.nim index 56be19c27..ab4b18f77 100644 --- a/tests/api/test_all.nim +++ b/tests/api/test_all.nim @@ -3,6 +3,7 @@ import ./test_entry_nodes, ./test_node_conf, + ./test_messaging_conf, ./test_api_send, ./test_api_subscription, ./test_api_receive, diff --git a/tests/api/test_api_health.nim b/tests/api/test_api_health.nim index e448164af..01e81625b 100644 --- a/tests/api/test_api_health.nim +++ b/tests/api/test_api_health.nim @@ -16,6 +16,7 @@ import logos_delivery/waku/common/waku_protocol, logos_delivery/waku/factory/waku_conf import tools/confutils/cli_args +import logos_delivery/api/messaging_conf const TestTimeout = chronos.seconds(10) const DefaultShard = PubsubTopic("/waku/2/rs/3/0") @@ -92,9 +93,8 @@ suite "LM API health checking": serviceNode.wakuRelay.subscribe(DefaultShard, dummyHandler) lockNewGlobalBrokerContext: - var conf = defaultWakuNodeConf().valueOr: + var conf = MessagingConf().toKernelConf(WakuMode.Core).valueOr: raiseAssert error - conf.mode = some(WakuMode.Core) conf.listenAddress = parseIpAddress("0.0.0.0") conf.tcpPort = Port(0) conf.discv5UdpPort = Port(0) @@ -271,9 +271,8 @@ suite "LM API health checking": var edgeWaku: Waku lockNewGlobalBrokerContext: - var edgeConf = defaultWakuNodeConf().valueOr: + var edgeConf = MessagingConf().toKernelConf(WakuMode.Edge).valueOr: raiseAssert error - edgeConf.mode = some(WakuMode.Edge) edgeConf.listenAddress = parseIpAddress("0.0.0.0") edgeConf.tcpPort = Port(0) edgeConf.discv5UdpPort = Port(0) diff --git a/tests/api/test_api_receive.nim b/tests/api/test_api_receive.nim index 35061ec81..dc97c6b54 100644 --- a/tests/api/test_api_receive.nim +++ b/tests/api/test_api_receive.nim @@ -22,6 +22,7 @@ import ] import logos_delivery/waku/factory/waku_conf import tools/confutils/cli_args +import logos_delivery/api/messaging_conf const TestTimeout = chronos.seconds(60) @@ -83,9 +84,8 @@ proc waitForConnectionStatus( await health_events.ConnectionStatusChangeEvent.dropListener(brokerCtx, handle) proc createApiNodeConf(numShards: uint16 = 1): WakuNodeConf = - var conf = defaultWakuNodeConf().valueOr: + var conf = MessagingConf().toKernelConf(cli_args.WakuMode.Core).valueOr: raiseAssert error - conf.mode = some(cli_args.WakuMode.Core) conf.listenAddress = parseIpAddress("0.0.0.0") conf.tcpPort = Port(0) conf.discv5UdpPort = Port(0) diff --git a/tests/api/test_api_send.nim b/tests/api/test_api_send.nim index c406e16a9..6fe028343 100644 --- a/tests/api/test_api_send.nim +++ b/tests/api/test_api_send.nim @@ -8,6 +8,7 @@ import ../waku_archive/archive_utils import logos_delivery, logos_delivery/waku/[waku_node, waku_core, waku_relay/protocol] import logos_delivery/waku/factory/waku_conf import tools/confutils/cli_args +import logos_delivery/api/messaging_conf type SendEventOutcome {.pure.} = enum Sent @@ -120,9 +121,8 @@ proc validate( check requestId == expectedRequestId proc createApiNodeConf(mode: cli_args.WakuMode = cli_args.WakuMode.Core): WakuNodeConf = - var conf = defaultWakuNodeConf().valueOr: + var conf = MessagingConf().toKernelConf(mode).valueOr: raiseAssert error - conf.mode = some(mode) conf.listenAddress = parseIpAddress("0.0.0.0") conf.tcpPort = Port(0) conf.discv5UdpPort = Port(0) diff --git a/tests/api/test_api_subscription.nim b/tests/api/test_api_subscription.nim index 2a8108b39..d34484f9d 100644 --- a/tests/api/test_api_subscription.nim +++ b/tests/api/test_api_subscription.nim @@ -19,6 +19,7 @@ import ] import logos_delivery/waku/factory/waku_conf import tools/confutils/cli_args +import logos_delivery/api/messaging_conf const TestTimeout = chronos.seconds(10) const NegativeTestTimeout = chronos.seconds(2) @@ -71,9 +72,8 @@ type TestNetwork = ref object proc createApiNodeConf( mode: cli_args.WakuMode = cli_args.WakuMode.Core, numShards: uint16 = 1 ): WakuNodeConf = - var conf = defaultWakuNodeConf().valueOr: + var conf = MessagingConf().toKernelConf(mode).valueOr: raiseAssert error - conf.mode = some(mode) conf.listenAddress = parseIpAddress("0.0.0.0") conf.tcpPort = Port(0) conf.discv5UdpPort = Port(0) diff --git a/tests/api/test_messaging_conf.nim b/tests/api/test_messaging_conf.nim new file mode 100644 index 000000000..f6efb89b8 --- /dev/null +++ b/tests/api/test_messaging_conf.nim @@ -0,0 +1,108 @@ +{.used.} + +import std/options, results, testutils/unittests +import std/nativesockets # Port + +import logos_delivery/api/types # WakuMode +import logos_delivery/api/messaging_conf + # MessagingConf, toKernelConf, defaultWakuNodeConf +import logos_delivery/api/messaging_conf_preset # resolvePreset, merge + +suite "MessagingConf - toKernelConf inference": + test "Core mode enables relay and service protocols": + let conf = MessagingConf().toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check: + conf.relay == true + conf.filter == true + conf.lightpush == true + conf.discv5Discovery == some(true) + conf.peerExchange == true + conf.rendezvous == true + + test "Edge mode disables relay and service protocols": + let conf = MessagingConf().toKernelConf(WakuMode.Edge).valueOr: + raiseAssert error + check: + conf.relay == false + conf.filter == false + conf.lightpush == false + conf.store == false + conf.peerExchange == true + + test "clusterId is applied only when set": + let setConf = MessagingConf(clusterId: some(7'u16)).toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check setConf.clusterId == some(7'u16) + + let kernelDefault = defaultWakuNodeConf().valueOr: + raiseAssert error + let unsetConf = MessagingConf().toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check unsetConf.clusterId == kernelDefault.clusterId + + test "p2pTcpPort sets the kernel tcpPort": + let conf = MessagingConf(p2pTcpPort: some(Port(60123))).toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check conf.tcpPort.uint16 == 60123'u16 + + test "reliabilityEnabled is applied to the kernel": + let conf = MessagingConf(reliabilityEnabled: some(true)).toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check conf.reliabilityEnabled == some(true) + + test "rlnContractAddress sets the contract and enables rlnRelay": + let conf = MessagingConf(rlnContractAddress: some("0xabc")).toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check: + conf.rlnRelayEthContractAddress == "0xabc" + conf.rlnRelay == some(true) + + test "ethRpcEndpoints maps to the kernel ethClientUrls": + let mc = MessagingConf(ethRpcEndpoints: some(@["http://node:8545"])) + let conf = mc.toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check: + conf.ethClientUrls.len == 1 + string(conf.ethClientUrls[0]) == "http://node:8545" + + test "rlnChainId and rlnEpochSizeSec map to the kernel": + let mc = MessagingConf(rlnChainId: some(5'u), rlnEpochSizeSec: some(600'u)) + let conf = mc.toKernelConf(WakuMode.Core).valueOr: + raiseAssert error + check: + conf.rlnRelayChainId == 5'u + conf.rlnEpochSizeSec == some(600'u64) + +suite "MessagingConf - preset resolution and merge": + test "resolvePreset twn populates the network fields": + let m = resolvePreset("twn").valueOr: + raiseAssert error + check: + m.clusterId == some(1'u16) + m.numShardsInCluster == some(8'u16) + m.reliabilityEnabled == some(false) + m.rlnContractAddress.isSome() + + test "resolvePreset logosdev sets cluster 2 and reliability": + let m = resolvePreset("logosdev").valueOr: + raiseAssert error + check: + m.clusterId == some(2'u16) + m.reliabilityEnabled == some(true) + m.rlnContractAddress.isNone() + + test "resolvePreset empty is a no-op": + let m = resolvePreset("").valueOr: + raiseAssert error + check: + m.clusterId.isNone() + m.reliabilityEnabled.isNone() + + test "merge: overrides win over the preset base, base kept otherwise": + let base = resolvePreset("twn").valueOr: + raiseAssert error + let merged = base.merge(MessagingConf(clusterId: some(99'u16))) + check: + merged.clusterId == some(99'u16) + merged.numShardsInCluster == some(8'u16) diff --git a/tests/api/test_node_conf.nim b/tests/api/test_node_conf.nim index b89efb0e4..98244828a 100644 --- a/tests/api/test_node_conf.nim +++ b/tests/api/test_node_conf.nim @@ -4,6 +4,7 @@ import std/[options, strutils], results, stint, testutils/unittests import json_serialization, confutils, confutils/std/net import tools/confutils/cli_args, + logos_delivery/api/messaging_conf, tools/confutils/conf_from_json, logos_delivery/waku/api/api_conf, logos_delivery/waku/factory/waku_conf, @@ -11,18 +12,13 @@ import logos_delivery/waku/factory/conf_builder/conf_builder, logos_delivery/waku/common/logging -suite "WakuNodeConf - mode-driven toWakuConf": +suite "MessagingConf mode - toKernelConf - toWakuConf": test "Core mode enables service protocols": - ## Given - var conf = defaultWakuNodeConf().valueOr: + let conf = MessagingConf(clusterId: some(1'u16)).toKernelConf(WakuMode.Core).valueOr: raiseAssert error - conf.mode = some(WakuMode.Core) - conf.clusterId = some(1'u16) - ## When let wakuConfRes = conf.toWakuConf() - ## Then require wakuConfRes.isOk() let wakuConf = wakuConfRes.get() require wakuConf.validate().isOk() @@ -34,16 +30,11 @@ suite "WakuNodeConf - mode-driven toWakuConf": wakuConf.clusterId == 1 test "Edge mode disables service protocols": - ## Given - var conf = defaultWakuNodeConf().valueOr: + let conf = MessagingConf(clusterId: some(1'u16)).toKernelConf(WakuMode.Edge).valueOr: raiseAssert error - conf.mode = some(WakuMode.Edge) - conf.clusterId = some(1'u16) - ## When let wakuConfRes = conf.toWakuConf() - ## Then require wakuConfRes.isOk() let wakuConf = wakuConfRes.get() require wakuConf.validate().isOk() @@ -54,44 +45,6 @@ suite "WakuNodeConf - mode-driven toWakuConf": wakuConf.storeServiceConf.isSome() == false wakuConf.peerExchangeService == true - test "WakuMode.none uses explicit CLI flags as-is": - ## Given - var conf = defaultWakuNodeConf().valueOr: - raiseAssert error - conf.mode = none[WakuMode]() - conf.relay = true - conf.lightpush = false - conf.clusterId = some(5'u16) - - ## When - let wakuConfRes = conf.toWakuConf() - - ## Then - require wakuConfRes.isOk() - let wakuConf = wakuConfRes.get() - require wakuConf.validate().isOk() - check: - wakuConf.relay == true - wakuConf.lightPush == false - wakuConf.clusterId == 5 - - test "Core mode overrides individual protocol flags": - ## Given - user sets relay=false but mode=Core should override - var conf = defaultWakuNodeConf().valueOr: - raiseAssert error - conf.mode = some(WakuMode.Core) - conf.relay = false # will be overridden by Core mode - - ## When - let wakuConfRes = conf.toWakuConf() - - ## Then - require wakuConfRes.isOk() - let wakuConf = wakuConfRes.get() - require wakuConf.validate().isOk() - check: - wakuConf.relay == true # mode overrides - suite "WakuNodeConf - JSON parsing with fieldPairs": test "Empty JSON produces valid default conf": ## Given / When @@ -101,31 +54,9 @@ suite "WakuNodeConf - JSON parsing with fieldPairs": require confRes.isOk() let conf = confRes.get() check: - conf.mode == none[WakuMode]() conf.clusterId.isNone() conf.logLevel == logging.LogLevel.INFO - test "JSON with mode and clusterId": - ## Given / When - let confRes = parseNodeConfFromJson("""{"mode": "Core", "clusterId": 42}""") - - ## Then - require confRes.isOk() - let conf = confRes.get() - check: - conf.mode == some(WakuMode.Core) - conf.clusterId == some(42'u16) - - test "JSON with Edge mode": - ## Given / When - let confRes = parseNodeConfFromJson("""{"mode": "Edge"}""") - - ## Then - require confRes.isOk() - let conf = confRes.get() - check: - conf.mode == some(WakuMode.Edge) - test "JSON with logLevel": ## Given / When let confRes = parseNodeConfFromJson("""{"logLevel": "DEBUG"}""") @@ -254,49 +185,9 @@ suite "WakuNodeConf - preset integration": check wakuConfRes.isErr() suite "WakuNodeConf JSON -> WakuConf integration": - test "Core mode JSON config produces valid WakuConf": - ## Given - let confRes = parseNodeConfFromJson( - """{"mode": "Core", "clusterId": 55, "numShardsInNetwork": 6}""" - ) - require confRes.isOk() - let conf = confRes.get() - - ## When - let wakuConfRes = conf.toWakuConf() - - ## Then - require wakuConfRes.isOk() - let wakuConf = wakuConfRes.get() - require wakuConf.validate().isOk() - check: - wakuConf.relay == true - wakuConf.lightPush == true - wakuConf.peerExchangeService == true - wakuConf.clusterId == 55 - wakuConf.shardingConf.numShardsInCluster == 6 - - test "Edge mode JSON config produces valid WakuConf": - ## Given - let confRes = parseNodeConfFromJson("""{"mode": "Edge", "clusterId": 1}""") - require confRes.isOk() - let conf = confRes.get() - - ## When - let wakuConfRes = conf.toWakuConf() - - ## Then - require wakuConfRes.isOk() - let wakuConf = wakuConfRes.get() - require wakuConf.validate().isOk() - check: - wakuConf.relay == false - wakuConf.lightPush == false - wakuConf.peerExchangeService == true - test "JSON with preset produces valid WakuConf": ## Given - let confRes = parseNodeConfFromJson("""{"mode": "Core", "preset": "logosdev"}""") + let confRes = parseNodeConfFromJson("""{"preset": "logosdev"}""") require confRes.isOk() let conf = confRes.get() @@ -314,7 +205,7 @@ suite "WakuNodeConf JSON -> WakuConf integration": test "JSON with static nodes": ## Given let confRes = parseNodeConfFromJson( - """{"mode": "Core", "clusterId": 42, "staticnodes": ["/ip4/127.0.0.1/tcp/60000/p2p/16Uuu2HBmAcHvhLqQKwSSbX6BG5JLWUDRcaLVrehUVqpw7fz1hbYc"]}""" + """{"clusterId": 42, "staticnodes": ["/ip4/127.0.0.1/tcp/60000/p2p/16Uuu2HBmAcHvhLqQKwSSbX6BG5JLWUDRcaLVrehUVqpw7fz1hbYc"]}""" ) require confRes.isOk() let conf = confRes.get() diff --git a/tests/channels/test_reliable_channel_send_receive.nim b/tests/channels/test_reliable_channel_send_receive.nim index 1f90541b1..89f88fd95 100644 --- a/tests/channels/test_reliable_channel_send_receive.nim +++ b/tests/channels/test_reliable_channel_send_receive.nim @@ -12,6 +12,7 @@ import logos_delivery/waku/[waku_node, waku_core] import logos_delivery/waku/factory/waku_conf import logos_delivery/waku/events/message_events as waku_message_events import tools/confutils/cli_args +import logos_delivery/api/messaging_conf import logos_delivery/channels/reliable_channel_manager import logos_delivery/channels/encryption/noop_encryption @@ -28,9 +29,8 @@ import logos_delivery/api/messaging_client_interface const TestTimeout = chronos.seconds(15) proc createApiNodeConf(): WakuNodeConf = - var conf = defaultWakuNodeConf().valueOr: + var conf = MessagingConf().toKernelConf(cli_args.WakuMode.Core).valueOr: raiseAssert error - conf.mode = some(cli_args.WakuMode.Core) conf.listenAddress = parseIpAddress("0.0.0.0") conf.tcpPort = Port(0) conf.discv5UdpPort = Port(0) diff --git a/tests/test_waku.nim b/tests/test_waku.nim index d69b54e4b..9a085027a 100644 --- a/tests/test_waku.nim +++ b/tests/test_waku.nim @@ -6,15 +6,15 @@ import chronos, testutils/unittests import logos_delivery import tools/confutils/cli_args +import logos_delivery/api/messaging_conf import logos_delivery/waku/factory/networks_config import logos_delivery/waku/factory/conf_builder/conf_builder suite "Waku API - Create node": asyncTest "Create node with minimal configuration": ## Given - var nodeConf = defaultWakuNodeConf().valueOr: - raiseAssert "defaultWakuNodeConf failed: " & error - nodeConf.mode = some(WakuMode.Core) + var nodeConf = MessagingConf().toKernelConf(WakuMode.Core).valueOr: + raiseAssert "toKernelConf failed: " & error nodeConf.clusterId = some(3'u16) nodeConf.rest = false @@ -32,9 +32,8 @@ suite "Waku API - Create node": asyncTest "Create node with full configuration": ## Given - var nodeConf = defaultWakuNodeConf().valueOr: - raiseAssert "defaultWakuNodeConf failed: " & error - nodeConf.mode = some(WakuMode.Core) + var nodeConf = MessagingConf().toKernelConf(WakuMode.Core).valueOr: + raiseAssert "toKernelConf failed: " & error nodeConf.clusterId = some(99'u16) nodeConf.rest = false nodeConf.numShardsInNetwork = 16 @@ -64,9 +63,8 @@ suite "Waku API - Create node": asyncTest "Create node with mixed entry nodes (enrtree, multiaddr)": ## Given - var nodeConf = defaultWakuNodeConf().valueOr: - raiseAssert "defaultWakuNodeConf failed: " & error - nodeConf.mode = some(WakuMode.Core) + var nodeConf = MessagingConf().toKernelConf(WakuMode.Core).valueOr: + raiseAssert "toKernelConf failed: " & error nodeConf.clusterId = some(42'u16) nodeConf.rest = false nodeConf.entryNodes = @[ diff --git a/tools/confutils/cli_args.nim b/tools/confutils/cli_args.nim index 5a79a4945..5309e356f 100644 --- a/tools/confutils/cli_args.nim +++ b/tools/confutils/cli_args.nim @@ -163,13 +163,6 @@ type WakuNodeConf* = object .}: seq[ProtectedShard] ## General node config - mode* {. - desc: - "Node operation mode. 'Core' enables relay+service protocols. 'Edge' enables client-only protocols. Default (unset): explicit CLI flags used.", - defaultValue: none(WakuMode), - name: "mode" - .}: Option[WakuMode] - preset* {. desc: "Network preset to use. 'twn' is The RLN-protected Waku Network (cluster 1). 'logos.dev' is the Logos Dev Network (cluster 2). 'logos.test' is the Logos Test Network (cluster 2). Overrides other values.", @@ -993,7 +986,7 @@ proc toKeystoreGeneratorConf*(n: WakuNodeConf): RlnKeystoreGeneratorConf = credPassword: n.rlnRelayCredPassword, ) -proc toNetworkPresetConf( +proc toNetworkPresetConf*( preset: string, clusterId: Option[uint16] ): ConfResult[Option[NetworkPresetConf]] = var lcPreset = toLowerAscii(preset) @@ -1220,25 +1213,4 @@ proc toWakuConf*(n: WakuNodeConf): ConfResult[WakuConf] = chronos.seconds(n.kadServiceLookupIntervalSec.int64) ) - # Mode-driven configuration overrides. `none` (formerly `noMode`) means: - # use explicit CLI flags as-is. - if n.mode.isSome(): - case n.mode.get() - of WakuMode.Core: - b.withRelay(true) - b.filterServiceConf.withEnabled(true) - b.withLightPush(true) - b.discv5Conf.withEnabled(true) - b.withPeerExchange(true) - b.withRendezvous(true) - b.rateLimitConf.withRateLimitsIfNotAssigned( - @["filter:100/1s", "lightpush:5/1s", "px:5/1s"] - ) - of WakuMode.Edge: - b.withPeerExchange(true) - b.withRelay(false) - b.filterServiceConf.withEnabled(false) - b.withLightPush(false) - b.storeServiceConf.withEnabled(false) - return b.build()