Skip to content
Draft
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
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
<!--
A new scriv changelog fragment.

Uncomment the section that is right (remove the HTML comment wrapper).
For top level release notes, leave all the headers commented out.
-->

<!--
### Breaking

- A bullet item for the Breaking category.

-->
### Non-Breaking

- Integrate weighted fair queue + burst mux

<!--
### Patch

- A bullet item for the Patch category.

-->
12 changes: 8 additions & 4 deletions cardano-diffusion/demo/chain-sync.hs
Original file line number Diff line number Diff line change
Expand Up @@ -209,7 +209,8 @@ rmIfExists path = do
maximumMiniProtocolLimits :: MiniProtocolLimits
maximumMiniProtocolLimits =
MiniProtocolLimits {
maximumIngressQueue = maxBound
maximumIngressQueue = maxBound,
burst = Nothing
}


Expand All @@ -226,7 +227,8 @@ demoProtocol2 chainSync =
miniProtocolNum = MiniProtocolNum 2,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = maximumMiniProtocolLimits,
miniProtocolRun = chainSync
miniProtocolRun = chainSync,
miniProtocolWeight = 1
}
]

Expand Down Expand Up @@ -336,13 +338,15 @@ demoProtocol3 chainSync blockFetch =
miniProtocolNum = MiniProtocolNum 2,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = maximumMiniProtocolLimits,
miniProtocolRun = chainSync
miniProtocolRun = chainSync,
miniProtocolWeight = 1
}
, MiniProtocol {
miniProtocolNum = MiniProtocolNum 3,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = maximumMiniProtocolLimits,
miniProtocolRun = blockFetch
miniProtocolRun = blockFetch,
miniProtocolWeight = 1
}
]

Expand Down
19 changes: 12 additions & 7 deletions cardano-diffusion/lib/Cardano/Network/NodeToClient.hs
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,8 @@ module Cardano.Network.NodeToClient
, Handshake
) where

import Control.Exception (SomeException)
import Control.DeepSeq (NFData)
import Control.Exception (SomeException)
import Control.Monad (forever)
import Control.Monad.Class.MonadAsync
import Control.Monad.Class.MonadTimer.SI
Expand Down Expand Up @@ -149,35 +149,40 @@ nodeToClientProtocols protocols _version _versionData =
miniProtocolNum = MiniProtocolNum 5,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = maximumMiniProtocolLimits,
miniProtocolRun = localChainSyncProtocol
miniProtocolRun = localChainSyncProtocol,
miniProtocolWeight = 1
}
localTxSubmissionMiniProtocol localTxSubmissionProtocol = MiniProtocol {
miniProtocolNum = MiniProtocolNum 6,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = maximumMiniProtocolLimits,
miniProtocolRun = localTxSubmissionProtocol
miniProtocolRun = localTxSubmissionProtocol,
miniProtocolWeight = 1
}
localStateQueryMiniProtocol localStateQueryProtocol = MiniProtocol {
miniProtocolNum = MiniProtocolNum 7,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = maximumMiniProtocolLimits,
miniProtocolRun = localStateQueryProtocol
miniProtocolRun = localStateQueryProtocol,
miniProtocolWeight = 1
}
localTxMonitorMiniProtocol localTxMonitorProtocol = MiniProtocol {
miniProtocolNum = MiniProtocolNum 9,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = maximumMiniProtocolLimits,
miniProtocolRun = localTxMonitorProtocol
miniProtocolRun = localTxMonitorProtocol,
miniProtocolWeight = 1
}

maximumMiniProtocolLimits :: MiniProtocolLimits
maximumMiniProtocolLimits =
MiniProtocolLimits {
#if !defined(wasm32_HOST_ARCH)
maximumIngressQueue = 0xffffffff
maximumIngressQueue = 0xffffffff,
#else
maximumIngressQueue = 0x7fffffff
maximumIngressQueue = 0x7fffffff,
#endif
burst = Nothing
}


Expand Down
42 changes: 28 additions & 14 deletions cardano-diffusion/lib/Cardano/Network/NodeToNode.hs
Original file line number Diff line number Diff line change
Expand Up @@ -263,19 +263,22 @@ nodeToNodeProtocols _featureFlags miniProtocolParameters protocols
miniProtocolNum = chainSyncMiniProtocolNum,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = chainSyncProtocolLimits miniProtocolParameters,
miniProtocolRun = chainSyncProtocol
miniProtocolRun = chainSyncProtocol,
miniProtocolWeight = 1
}
, MiniProtocol {
miniProtocolNum = blockFetchMiniProtocolNum,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = blockFetchProtocolLimits miniProtocolParameters,
miniProtocolRun = blockFetchProtocol
miniProtocolRun = blockFetchProtocol,
miniProtocolWeight = 1
}
, MiniProtocol {
miniProtocolNum = txSubmissionMiniProtocolNum,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = txSubmissionProtocolLimits miniProtocolParameters,
miniProtocolRun = txSubmissionProtocol
miniProtocolRun = txSubmissionProtocol,
miniProtocolWeight = 1
}
]
<> case perasSupport of
Expand All @@ -287,13 +290,15 @@ nodeToNodeProtocols _featureFlags miniProtocolParameters protocols
miniProtocolNum = perasCertDiffusionMiniProtocolNum,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = perasCertDiffusionProtocolLimits miniProtocolParameters,
miniProtocolRun = perasCertDiffusionProtocol
miniProtocolRun = perasCertDiffusionProtocol,
miniProtocolWeight = 1
}
, MiniProtocol {
miniProtocolNum = perasVoteDiffusionMiniProtocolNum,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = perasVoteDiffusionProtocolLimits miniProtocolParameters,
miniProtocolRun = perasVoteDiffusionProtocol
miniProtocolRun = perasVoteDiffusionProtocol,
miniProtocolWeight = 1
}
])

Expand All @@ -309,15 +314,17 @@ nodeToNodeProtocols _featureFlags miniProtocolParameters protocols
miniProtocolNum = keepAliveMiniProtocolNum,
miniProtocolStart = StartOnDemandAny,
miniProtocolLimits = keepAliveProtocolLimits miniProtocolParameters,
miniProtocolRun = keepAliveProtocol
miniProtocolRun = keepAliveProtocol,
miniProtocolWeight = 1
}
: case peerSharing of
PeerSharingEnabled ->
[ MiniProtocol {
miniProtocolNum = peerSharingMiniProtocolNum,
miniProtocolStart = StartOnDemand,
miniProtocolLimits = peerSharingProtocolLimits miniProtocolParameters,
miniProtocolRun = peerSharingProtocol
miniProtocolRun = peerSharingProtocol,
miniProtocolWeight = 1
}
]
PeerSharingDisabled ->
Expand All @@ -342,7 +349,8 @@ chainSyncProtocolLimits MiniProtocolParameters { chainSyncPipeliningHighMark } =
-- TODO: 1400 comes from maxBlockHeaderSize in genesis, but should come
-- from consensus rather than being hard coded.
maximumIngressQueue = addSafetyMargin $
fromIntegral chainSyncPipeliningHighMark * 1400
fromIntegral chainSyncPipeliningHighMark * 1400,
burst = Nothing
}

blockFetchProtocolLimits MiniProtocolParameters { blockFetchPipeliningMax } = MiniProtocolLimits {
Expand All @@ -364,7 +372,8 @@ blockFetchProtocolLimits MiniProtocolParameters { blockFetchPipeliningMax } = Mi
-- relaxed limit here.
--
maximumIngressQueue = addSafetyMargin $
max (10 * 2_097_154 :: Int) (fromIntegral blockFetchPipeliningMax * 90_112)
max (10 * 2_097_154 :: Int) (fromIntegral blockFetchPipeliningMax * 90_112),
burst = Just $ Mx.ProtocolBurst 90_112 10_000
}

txSubmissionProtocolLimits MiniProtocolParameters
Expand Down Expand Up @@ -432,13 +441,15 @@ txSubmissionProtocolLimits MiniProtocolParameters
-- 10% as a safety margin.
--
maximumIngressQueue = addSafetyMargin $
fromIntegral maxUnacknowledgedTxIds * (44 + fromIntegral @SizeInBytes @Int max_TX_SIZE)
fromIntegral maxUnacknowledgedTxIds * (44 + fromIntegral @SizeInBytes @Int max_TX_SIZE),
burst = Nothing
}

keepAliveProtocolLimits _ =
MiniProtocolLimits {
-- One small outstanding message.
maximumIngressQueue = addSafetyMargin 1280
maximumIngressQueue = addSafetyMargin 1280,
burst = Nothing
}

peerSharingProtocolLimits _ =
Expand All @@ -449,7 +460,8 @@ peerSharingProtocolLimits _ =
-- window size of 4 and a TCP segment is 1440, which gives us 4 * 1440 =
-- 5760 bytes to fit into a single RTT. So setting the maximum ingress
-- queue to be a single RTT should be enough to cover for CBOR overhead.
maximumIngressQueue = 4 * 1440
maximumIngressQueue = 4 * 1440,
burst = Nothing
}

perasCertDiffusionProtocolLimits MiniProtocolParameters { perasCertDiffusionMaxObjectsUnacknowledged } =
Expand All @@ -460,7 +472,8 @@ perasCertDiffusionProtocolLimits MiniProtocolParameters { perasCertDiffusionMaxO
-- even much smaller.
-- See https://github.com/tweag/cardano-peras/issues/97
maximumIngressQueue = addSafetyMargin $
fromIntegral perasCertDiffusionMaxObjectsUnacknowledged * 20_000
fromIntegral perasCertDiffusionMaxObjectsUnacknowledged * 20_000,
burst = Nothing
}

perasVoteDiffusionProtocolLimits MiniProtocolParameters { perasVoteDiffusionMaxObjectsUnacknowledged } =
Expand All @@ -469,7 +482,8 @@ perasVoteDiffusionProtocolLimits MiniProtocolParameters { perasVoteDiffusionMaxO
-- We assume an upper bound of 1 kB per vote.
-- See https://github.com/tweag/cardano-peras/issues/97
maximumIngressQueue = addSafetyMargin $
fromIntegral perasVoteDiffusionMaxObjectsUnacknowledged * 1_000
fromIntegral perasVoteDiffusionMaxObjectsUnacknowledged * 1_000,
burst = Nothing
}

chainSyncMiniProtocolNum :: MiniProtocolNum
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -340,14 +340,16 @@ applications debugTracer txSubmissionInboundTracer txSubmissionInboundDebug node
-> MiniProtocolWithExpandedCtx Mx.InitiatorMode NtNAddr PeerTrustable ByteString m () Void
f MiniProtocol { miniProtocolNum
, miniProtocolLimits
, miniProtocolRun } =
, miniProtocolRun
, miniProtocolWeight } =
MiniProtocol { miniProtocolNum
, miniProtocolStart = StartEagerly
, miniProtocolLimits
, miniProtocolRun =
case miniProtocolRun of
InitiatorAndResponderProtocol initiator _respnder ->
InitiatorProtocolOnly initiator
InitiatorProtocolOnly initiator,
miniProtocolWeight
}

initiatorAndResponderApp
Expand All @@ -362,7 +364,8 @@ applications debugTracer txSubmissionInboundTracer txSubmissionInboundDebug node
, miniProtocolRun =
InitiatorAndResponderProtocol
chainSyncInitiator
chainSyncResponder
chainSyncResponder,
miniProtocolWeight = 1
}
, MiniProtocol
{ miniProtocolNum = blockFetchMiniProtocolNum
Expand All @@ -371,7 +374,8 @@ applications debugTracer txSubmissionInboundTracer txSubmissionInboundDebug node
, miniProtocolRun =
InitiatorAndResponderProtocol
blockFetchInitiator
blockFetchResponder
blockFetchResponder,
miniProtocolWeight = 1
}

, MiniProtocol {
Expand All @@ -384,7 +388,8 @@ applications debugTracer txSubmissionInboundTracer txSubmissionInboundDebug node
(txSubmissionResponder (nkMempool nodeKernel)
(nkTxChannelsVar nodeKernel)
(nkTxMempoolSem nodeKernel)
(nkSharedTxStateVar nodeKernel))
(nkSharedTxStateVar nodeKernel)),
miniProtocolWeight = 1
}
]
, withWarm = WithWarm
Expand All @@ -395,7 +400,8 @@ applications debugTracer txSubmissionInboundTracer txSubmissionInboundDebug node
, miniProtocolRun =
InitiatorAndResponderProtocol
pingPongInitiator
pingPongResponder
pingPongResponder,
miniProtocolWeight = 1
}
]
, withEstablished = WithEstablished $
Expand All @@ -406,7 +412,8 @@ applications debugTracer txSubmissionInboundTracer txSubmissionInboundDebug node
, miniProtocolRun =
InitiatorAndResponderProtocol
keepAliveInitiator
keepAliveResponder
keepAliveResponder,
miniProtocolWeight = 1
}
: case peerSharing of
PSTypes.PeerSharingEnabled ->
Expand All @@ -417,7 +424,8 @@ applications debugTracer txSubmissionInboundTracer txSubmissionInboundDebug node
, miniProtocolRun =
InitiatorAndResponderProtocol
peerSharingInitiator
(peerSharingResponder (nkPeerSharingAPI nodeKernel))
(peerSharingResponder (nkPeerSharingAPI nodeKernel)),
miniProtocolWeight = 1
}
]
PSTypes.PeerSharingDisabled ->
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,10 @@ module Test.Cardano.Network.Diffusion.Testnet.Simulation
, module PeerSelection
) where

import Control.Applicative (Alternative)
import Control.Concurrent.Class.MonadMVar (MonadMVar)
import Control.Concurrent.Class.MonadSTM qualified as LazySTM
import Control.Concurrent.Class.MonadSTM.Strict
import Control.Monad (forM, when)
import Control.Monad (MonadPlus, forM, when)
import Control.Monad.Class.MonadAsync
import Control.Monad.Class.MonadFork
import Control.Monad.Class.MonadSay
Expand Down Expand Up @@ -1033,8 +1032,7 @@ data Churn = CardanoChurn | OuroborosChurn
-- | Run an arbitrary topology in a generic monad `m`.
--
diffusionSimulationM
:: forall m. ( Alternative (STM m)
, MonadAsync m
:: forall m. ( MonadAsync m
, MonadDelay m
, MonadFix m
, MonadEvaluate m
Expand All @@ -1045,6 +1043,7 @@ diffusionSimulationM
, MonadLabelledSTM m
, MonadTraceSTM m
, MonadMask m
, MonadPlus (STM m)
, MonadTime m
, MonadTimer m
, MonadThrow (STM m)
Expand Down Expand Up @@ -1210,7 +1209,7 @@ diffusionSimulationM
acceptVersion = acceptableVersion
defaultMiniProtocolsLimit :: MiniProtocolLimits
defaultMiniProtocolsLimit =
MiniProtocolLimits { maximumIngressQueue = 64000 }
MiniProtocolLimits { maximumIngressQueue = 64000, burst = Nothing }

blockGeneratorArgs :: Node.BlockGeneratorArgs Block StdGen
blockGeneratorArgs =
Expand Down
Loading
Loading