mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-08-09 21:33:25 +00:00
chore: move conf types to api/conf
* move MessagingClientConf to api/conf/messaging_conf * move ReliableChannelManagerConf to new api/conf/channels_conf * impl modules now import their conf module, not the reverse * logos_delivery_conf uses channels_conf, not reliable_channel_manager * cleanup import lists
This commit is contained in:
parent
9f2a1c89ff
commit
a26da34557
15
logos_delivery/api/conf/channels_conf.nim
Normal file
15
logos_delivery/api/conf/channels_conf.nim
Normal file
@ -0,0 +1,15 @@
|
|||||||
|
import std/options
|
||||||
|
|
||||||
|
type ReliableChannelManagerConf* = object
|
||||||
|
## All-`Option` partial; unset fields fall back to `createReliableChannel` defaults.
|
||||||
|
segmentationEnableReedSolomon*: Option[bool]
|
||||||
|
## Add Reed-Solomon parity segments for recovery of lost segments.
|
||||||
|
segmentationSegmentSizeBytes*: Option[int] ## Maximum segment size in bytes.
|
||||||
|
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.
|
||||||
|
rateLimitEnabled*: Option[bool] ## Enable rate limiting.
|
||||||
|
rateLimitEpochPeriodSec*: Option[int] ## Rate-limit epoch length in seconds.
|
||||||
|
rateLimitMessagesPerEpoch*: Option[int] ## Messages allowed per rate-limit epoch.
|
||||||
@ -4,9 +4,9 @@ import std/options
|
|||||||
import results
|
import results
|
||||||
|
|
||||||
import logos_delivery/api/conf/messaging_conf
|
import logos_delivery/api/conf/messaging_conf
|
||||||
import logos_delivery/channels/reliable_channel_manager
|
import logos_delivery/api/conf/channels_conf
|
||||||
|
|
||||||
export options, messaging_conf, reliable_channel_manager
|
export options, messaging_conf, channels_conf
|
||||||
|
|
||||||
type LogosDeliveryConf* = object
|
type LogosDeliveryConf* = object
|
||||||
## Aggregates the per-layer config objects. A layer is mounted iff its config
|
## Aggregates the per-layer config objects. A layer is mounted iff its config
|
||||||
|
|||||||
@ -1,18 +1,60 @@
|
|||||||
import std/options
|
import std/options
|
||||||
import std/net
|
import std/net
|
||||||
import results
|
import results
|
||||||
|
import libp2p/crypto/crypto
|
||||||
|
|
||||||
import logos_delivery/api/conf/kernel_conf
|
import logos_delivery/api/conf/kernel_conf
|
||||||
import logos_delivery/messaging/messaging_client
|
import logos_delivery/waku/common/logging
|
||||||
import logos_delivery/waku/factory/networks_config
|
import logos_delivery/waku/factory/networks_config
|
||||||
|
|
||||||
export kernel_conf, messaging_client
|
export kernel_conf
|
||||||
|
|
||||||
type LogosDeliveryMode* {.pure.} = enum
|
type LogosDeliveryMode* {.pure.} = enum
|
||||||
Edge # client-only node
|
Edge # client-only node
|
||||||
Core # full service node
|
Core # full service node
|
||||||
Fleet # kernel-only node from a raw kernel config
|
Fleet # kernel-only node from a raw kernel config
|
||||||
|
|
||||||
|
type MessagingClientConf* = object
|
||||||
|
clusterId* {.name: "cluster-id".}: Option[uint16] ## Network cluster id.
|
||||||
|
numShardsInCluster* {.name: "num-shards-in-network".}: Option[uint16]
|
||||||
|
## Number of shards in the cluster.
|
||||||
|
p2pTcpPort* {.name: "tcp-port".}: Option[Port] ## TCP listening port.
|
||||||
|
discv5UdpPort* {.name: "discv5-udp-port".}: Option[Port] ## discv5 UDP port.
|
||||||
|
websocketSupport* {.name: "websocket-support".}: Option[bool]
|
||||||
|
## Enable the websocket transport.
|
||||||
|
websocketPort* {.name: "websocket-port".}: Option[Port] ## Websocket listening port.
|
||||||
|
quicSupport* {.name: "quic-support".}: Option[bool] ## Enable the QUIC transport.
|
||||||
|
quicPort* {.name: "quic-port".}: Option[Port] ## QUIC (UDP) listening port.
|
||||||
|
listenIpv4* {.name: "listen-address".}: Option[IpAddress] ## Inbound bind address.
|
||||||
|
maxMessageSize* {.name: "max-msg-size".}: Option[string]
|
||||||
|
## Maximum accepted message size (e.g. "150 KiB").
|
||||||
|
entryNodes* {.name: "entry-node".}: Option[seq[string]]
|
||||||
|
## Bootstrap / connectivity nodes (enrtree or multiaddr).
|
||||||
|
ethRpcEndpoints* {.name: "rln-relay-eth-client-address".}: Option[seq[EthRpcUrl]]
|
||||||
|
## Ethereum RPC endpoints (required for RLN validation); multiple for fail-over.
|
||||||
|
rlnContractAddress* {.name: "rln-relay-eth-contract-address".}: Option[string]
|
||||||
|
## RLN contract address; when set, RLN validation is enabled.
|
||||||
|
rlnChainId* {.name: "rln-relay-chain-id".}: Option[uint]
|
||||||
|
## Chain id the RLN contract is deployed on.
|
||||||
|
rlnEpochSizeSec* {.name: "rln-relay-epoch-sec".}: Option[uint]
|
||||||
|
## RLN epoch size, in seconds.
|
||||||
|
reliabilityEnabled* {.name: "reliability".}: Option[bool]
|
||||||
|
## Enable store-based send reliability.
|
||||||
|
store*: Option[bool] ## Enable the store protocol.
|
||||||
|
storenode* {.name: "storenode".}: Option[string]
|
||||||
|
storeMessageDbUrl* {.name: "store-message-db-url".}: Option[string]
|
||||||
|
## Database connection URL for the store service's persistent storage.
|
||||||
|
storeMessageRetentionPolicy* {.name: "store-message-retention-policy".}:
|
||||||
|
Option[string] ## Store retention policy (e.g. "time:3600;size:1GB").
|
||||||
|
storeMaxNumDbConnections* {.name: "store-max-num-db-connections".}: Option[int]
|
||||||
|
## Maximum number of simultaneous store database connections.
|
||||||
|
logLevel* {.name: "log-level".}: Option[logging.LogLevel]
|
||||||
|
## Process log level (TRACE..FATAL); applied by the kernel on node creation.
|
||||||
|
logFormat* {.name: "log-format".}: Option[logging.LogFormat]
|
||||||
|
## Process log format (TEXT or JSON); applied by the kernel on node creation.
|
||||||
|
nodeKey* {.name: "nodekey".}: Option[crypto.PrivateKey]
|
||||||
|
## P2P node private key (64-char hex): stable identity / peerId across restarts.
|
||||||
|
|
||||||
proc applyMode*(conf: var WakuNodeConf, mode: LogosDeliveryMode): ConfResult[void] =
|
proc applyMode*(conf: var WakuNodeConf, mode: LogosDeliveryMode): ConfResult[void] =
|
||||||
## Sets the protocol flags implied by the mode.
|
## Sets the protocol flags implied by the mode.
|
||||||
case mode
|
case mode
|
||||||
|
|||||||
@ -5,7 +5,7 @@
|
|||||||
##
|
##
|
||||||
## See: https://lip.logos.co/messaging/raw/reliable-channel-api.html
|
## See: https://lip.logos.co/messaging/raw/reliable-channel-api.html
|
||||||
|
|
||||||
import std/[options, tables]
|
import std/tables
|
||||||
import results
|
import results
|
||||||
import chronos
|
import chronos
|
||||||
import chronicles
|
import chronicles
|
||||||
@ -15,30 +15,16 @@ import brokers/broker_context
|
|||||||
|
|
||||||
import logos_delivery/api/types
|
import logos_delivery/api/types
|
||||||
import logos_delivery/api/reliable_channel_manager_api
|
import logos_delivery/api/reliable_channel_manager_api
|
||||||
|
import logos_delivery/api/conf/channels_conf
|
||||||
|
|
||||||
import ./reliable_channel
|
import ./reliable_channel
|
||||||
|
|
||||||
export reliable_channel
|
export reliable_channel, channels_conf
|
||||||
|
|
||||||
type
|
type ReliableChannelManager* = ref object ## Implements `ReliableChannelApi`.
|
||||||
ReliableChannelManagerConf* = object
|
channels*: Table[ChannelId, ReliableChannel] ## read by `channels/api.nim`
|
||||||
## All-`Option` partial; unset fields fall back to `createReliableChannel` defaults.
|
conf*: ReliableChannelManagerConf
|
||||||
segmentationEnableReedSolomon*: Option[bool]
|
brokerCtx*: BrokerContext
|
||||||
## Add Reed-Solomon parity segments for recovery of lost segments.
|
|
||||||
segmentationSegmentSizeBytes*: Option[int] ## Maximum segment size in bytes.
|
|
||||||
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.
|
|
||||||
rateLimitEnabled*: Option[bool] ## Enable rate limiting.
|
|
||||||
rateLimitEpochPeriodSec*: Option[int] ## Rate-limit epoch length in seconds.
|
|
||||||
rateLimitMessagesPerEpoch*: Option[int] ## Messages allowed per rate-limit epoch.
|
|
||||||
|
|
||||||
ReliableChannelManager* = ref object ## Implements `ReliableChannelApi`.
|
|
||||||
channels*: Table[ChannelId, ReliableChannel] ## read by `channels/api.nim`
|
|
||||||
conf*: ReliableChannelManagerConf
|
|
||||||
brokerCtx*: BrokerContext
|
|
||||||
|
|
||||||
proc new*(
|
proc new*(
|
||||||
T: type ReliableChannelManager,
|
T: type ReliableChannelManager,
|
||||||
|
|||||||
@ -1,70 +1,23 @@
|
|||||||
## Messaging layer core: the `MessagingClient` type plus its construction and
|
## Messaging layer core: the `MessagingClient` type plus its construction and
|
||||||
## lifecycle. The public operations (subscribe / unsubscribe / send) live in
|
## lifecycle. The public operations (subscribe / unsubscribe / send) live in
|
||||||
## `messaging/api.nim`.
|
## `messaging/api.nim`.
|
||||||
import std/[options, net]
|
import std/options
|
||||||
import results, chronos
|
import results, chronos, chronicles
|
||||||
import confutils/defs
|
|
||||||
import libp2p/crypto/crypto
|
|
||||||
import logos_delivery/waku/common/logging
|
|
||||||
import logos_delivery/api/conf/kernel_conf
|
|
||||||
import chronicles
|
|
||||||
import
|
import
|
||||||
|
logos_delivery/api/conf/messaging_conf,
|
||||||
logos_delivery/api/messaging_client_api,
|
logos_delivery/api/messaging_client_api,
|
||||||
logos_delivery/waku/waku,
|
logos_delivery/waku/waku,
|
||||||
logos_delivery/waku/factory/conf_builder/waku_conf_builder,
|
logos_delivery/waku/factory/conf_builder/waku_conf_builder,
|
||||||
logos_delivery/messaging/delivery_service/[recv_service, send_service]
|
logos_delivery/messaging/delivery_service/[recv_service, send_service]
|
||||||
|
|
||||||
# Surfaces the messaging API interface (and its Message* events) to consumers.
|
export messaging_client_api, messaging_conf
|
||||||
export messaging_client_api
|
|
||||||
|
|
||||||
type
|
type MessagingClient* = ref object
|
||||||
MessagingClientConf* = object
|
brokerCtx*: BrokerContext
|
||||||
clusterId* {.name: "cluster-id".}: Option[uint16] ## Network cluster id.
|
waku*: Waku ## The Waku kernel this layer drives; read by `messaging/api/*`.
|
||||||
numShardsInCluster* {.name: "num-shards-in-network".}: Option[uint16]
|
sendService*: SendService
|
||||||
## Number of shards in the cluster.
|
recvService*: RecvService
|
||||||
p2pTcpPort* {.name: "tcp-port".}: Option[Port] ## TCP listening port.
|
started*: bool
|
||||||
discv5UdpPort* {.name: "discv5-udp-port".}: Option[Port] ## discv5 UDP port.
|
|
||||||
websocketSupport* {.name: "websocket-support".}: Option[bool]
|
|
||||||
## Enable the websocket transport.
|
|
||||||
websocketPort* {.name: "websocket-port".}: Option[Port] ## Websocket listening port.
|
|
||||||
quicSupport* {.name: "quic-support".}: Option[bool] ## Enable the QUIC transport.
|
|
||||||
quicPort* {.name: "quic-port".}: Option[Port] ## QUIC (UDP) listening port.
|
|
||||||
listenIpv4* {.name: "listen-address".}: Option[IpAddress] ## Inbound bind address.
|
|
||||||
maxMessageSize* {.name: "max-msg-size".}: Option[string]
|
|
||||||
## Maximum accepted message size (e.g. "150 KiB").
|
|
||||||
entryNodes* {.name: "entry-node".}: Option[seq[string]]
|
|
||||||
## Bootstrap / connectivity nodes (enrtree or multiaddr).
|
|
||||||
ethRpcEndpoints* {.name: "rln-relay-eth-client-address".}: Option[seq[EthRpcUrl]]
|
|
||||||
## Ethereum RPC endpoints (required for RLN validation); multiple for fail-over.
|
|
||||||
rlnContractAddress* {.name: "rln-relay-eth-contract-address".}: Option[string]
|
|
||||||
## RLN contract address; when set, RLN validation is enabled.
|
|
||||||
rlnChainId* {.name: "rln-relay-chain-id".}: Option[uint]
|
|
||||||
## Chain id the RLN contract is deployed on.
|
|
||||||
rlnEpochSizeSec* {.name: "rln-relay-epoch-sec".}: Option[uint]
|
|
||||||
## RLN epoch size, in seconds.
|
|
||||||
reliabilityEnabled* {.name: "reliability".}: Option[bool]
|
|
||||||
## Enable store-based send reliability.
|
|
||||||
store*: Option[bool] ## Enable the store protocol.
|
|
||||||
storenode* {.name: "storenode".}: Option[string]
|
|
||||||
storeMessageDbUrl* {.name: "store-message-db-url".}: Option[string]
|
|
||||||
## Database connection URL for the store service's persistent storage.
|
|
||||||
storeMessageRetentionPolicy* {.name: "store-message-retention-policy".}:
|
|
||||||
Option[string] ## Store retention policy (e.g. "time:3600;size:1GB").
|
|
||||||
storeMaxNumDbConnections* {.name: "store-max-num-db-connections".}: Option[int]
|
|
||||||
## Maximum number of simultaneous store database connections.
|
|
||||||
logLevel* {.name: "log-level".}: Option[logging.LogLevel]
|
|
||||||
## Process log level (TRACE..FATAL); applied by the kernel on node creation.
|
|
||||||
logFormat* {.name: "log-format".}: Option[logging.LogFormat]
|
|
||||||
## Process log format (TEXT or JSON); applied by the kernel on node creation.
|
|
||||||
nodeKey* {.name: "nodekey".}: Option[crypto.PrivateKey]
|
|
||||||
## P2P node private key (64-char hex): stable identity / peerId across restarts.
|
|
||||||
|
|
||||||
MessagingClient* = ref object
|
|
||||||
brokerCtx*: BrokerContext
|
|
||||||
waku*: Waku ## The Waku kernel this layer drives; read by `messaging/api/*`.
|
|
||||||
sendService*: SendService
|
|
||||||
recvService*: RecvService
|
|
||||||
started*: bool
|
|
||||||
|
|
||||||
proc new*(
|
proc new*(
|
||||||
T: type MessagingClient, conf: MessagingClientConf, waku: Waku
|
T: type MessagingClient, conf: MessagingClientConf, waku: Waku
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user