| Safe Haskell | Safe-Inferred |
|---|---|
| Language | GHC2021 |
Hydra.Node
Description
Top-level module to run a single Hydra node.
Checkout Hydra
Documentation
for some details about the overall architecture of the Node.
Synopsis
- initEnvironment :: RunOptions -> IO Environment
- checkHeadState :: MonadThrow m => Tracer m (HydraNodeLog tx) -> Environment -> HeadState tx -> m ()
- data DraftHydraNode tx m = DraftHydraNode {
- tracer :: Tracer m (HydraNodeLog tx)
- env :: Environment
- ledger :: Ledger tx
- nodeStateHandler :: NodeStateHandler tx m
- inputQueue :: InputQueue m (Input tx)
- eventSource :: EventSource (StateEvent tx) m
- eventSinks :: [EventSink (StateEvent tx) m]
- networkOutbox :: Outbox m
- chainStateHistory :: ChainStateHistory tx
- hydrate :: (IsChainState tx, MonadDelay m, MonadLabelledSTM m, MonadAsync m, MonadThrow m, MonadUnliftIO m) => Tracer m (HydraNodeLog tx) -> Environment -> Ledger tx -> ChainStateType tx -> EventStore (StateEvent tx) m -> [EventSink (StateEvent tx) m] -> m (DraftHydraNode tx m)
- wireChainInput :: DraftHydraNode tx m -> ChainEvent tx -> m ()
- wireClientInput :: DraftHydraNode tx m -> ClientInput tx -> m ()
- wireNetworkInput :: DraftHydraNode tx m -> NetworkCallback (Authenticated (Message tx)) m
- mkNetworkInput :: Party -> Message tx -> Input tx
- connect :: Monad m => Chain tx m -> Network m (Message tx) -> Server tx m -> DraftHydraNode tx m -> m (HydraNode tx m)
- data HydraNode tx m = HydraNode {
- tracer :: Tracer m (HydraNodeLog tx)
- env :: Environment
- ledger :: Ledger tx
- nodeStateHandler :: NodeStateHandler tx m
- inputQueue :: InputQueue m (Input tx)
- eventSource :: EventSource (StateEvent tx) m
- eventSinks :: [EventSink (StateEvent tx) m]
- oc :: Chain tx m
- hn :: Network m (Message tx)
- server :: Server tx m
- networkOutbox :: Outbox m
- withNetworkOutbox :: (MonadAsync m, MonadDelay m) => HydraNode tx m -> m () -> m ()
- monitorBroadcast :: (MonadDelay m, MonadSTM m) => HydraNode tx m -> m ()
- broadcastStallBounds :: StallBounds
- runHydraNode :: (MonadCatch m, MonadAsync m, MonadDelay m, MonadTime m, IsChainState tx) => HydraNode tx m -> m ()
- stepHydraNode :: (MonadCatch m, MonadAsync m, MonadTime m, IsChainState tx) => UTCTime -> HydraNode tx m -> m ()
- growsBroadcastBacklog :: ClientInput tx -> Bool
- inOpenHeadAndSynced :: NodeState tx -> Bool
- defaultTTL :: TTL
- defaultTxTTL :: TTL
- waitDelay :: DiffTime
- processNextInput :: IsChainState tx => HydraNode tx m -> Input tx -> UTCTime -> STM m (Outcome tx)
- processStateChanges :: (MonadSTM m, MonadTime m) => HydraNode tx m -> [StateChanged tx] -> m ()
- processEffects :: (MonadAsync m, MonadCatch m, IsChainState tx) => HydraNode tx m -> Tracer m (HydraNodeLog tx) -> Word64 -> [Effect tx] -> m ()
- data NodeStateHandler tx m = NodeStateHandler {
- modifyNodeState :: forall a. (NodeState tx -> (a, NodeState tx)) -> STM m a
- queryNodeState :: STM m (NodeState tx)
- getNextEventId :: STM m EventId
- createNodeStateHandler :: MonadLabelledSTM m => Maybe EventId -> NodeState tx -> m (NodeStateHandler tx m)
- data HydraNodeLog tx
- = BeginInput { }
- | EndInput { }
- | BeginEffect { }
- | EndEffect { }
- | LogicOutcome { }
- | DroppedFromQueue { }
- | LoadingState
- | LoadedState {
- lastEventId :: Last EventId
- nodeState :: NodeState tx
- | LoadedChainState {
- lastKnownChainPoint :: ChainPointType tx
- | ReplayingState
- | Misconfiguration { }
- | DiscardedBroadcasts { }
- | BroadcastBacklog { }
Environment Handling
initEnvironment :: RunOptions -> IO Environment Source #
Initialize the Environment from command line options.
checkHeadState :: MonadThrow m => Tracer m (HydraNodeLog tx) -> Environment -> HeadState tx -> m () Source #
Checks that command line options match a given HeadState. This function
takes Environment because it is derived from RunOptions via
initEnvironment.
Throws: ParameterMismatch when state not matching the environment.
Create and run a hydra node
data DraftHydraNode tx m Source #
A draft version of the HydraNode that holds state, but is not yet
connected (see connect). This is commonly created by the hydrate smart
constructor.
Constructors
| DraftHydraNode | |
Fields
| |
hydrate :: (IsChainState tx, MonadDelay m, MonadLabelledSTM m, MonadAsync m, MonadThrow m, MonadUnliftIO m) => Tracer m (HydraNodeLog tx) -> Environment -> Ledger tx -> ChainStateType tx -> EventStore (StateEvent tx) m -> [EventSink (StateEvent tx) m] -> m (DraftHydraNode tx m) Source #
Hydrate a DraftHydraNode by loading events from source, re-aggregate node
state and sending events to sinks while doing so.
wireChainInput :: DraftHydraNode tx m -> ChainEvent tx -> m () Source #
wireClientInput :: DraftHydraNode tx m -> ClientInput tx -> m () Source #
wireNetworkInput :: DraftHydraNode tx m -> NetworkCallback (Authenticated (Message tx)) m Source #
mkNetworkInput :: Party -> Message tx -> Input tx Source #
Create a network input with corresponding default ttl from given sender.
connect :: Monad m => Chain tx m -> Network m (Message tx) -> Server tx m -> DraftHydraNode tx m -> m (HydraNode tx m) Source #
Connect chain, network and API to a hydrated DraftHydraNode to get a fully
connected HydraNode.
Fully connected hydra node with everything wired in.
Constructors
| HydraNode | |
Fields
| |
withNetworkOutbox :: (MonadAsync m, MonadDelay m) => HydraNode tx m -> m () -> m () Source #
Run the network hand-off and its stall monitor concurrently with the given action, stopping when any of them does. An effect throwing therefore still takes the node down, as it did when effects ran on the main loop.
NOTE: it takes it down asynchronously, though, where running inline threw
synchronously from between two effects. The main loop is now cancelled
wherever it happens to be, which can be mid-processStateChanges. That is
the same class of interruption the surrounding withChain/withNetwork
brackets could already deliver, but this is a new source of it.
monitorBroadcast :: (MonadDelay m, MonadSTM m) => HydraNode tx m -> m () Source #
Report the network hand-off stalling and recovering as Connectivity
events, which reach the event log and from there clients already listening.
A client connecting mid-stall is told instead by
NetworkInfo, which reads the outbox live rather
than replaying these - past outputs are only replayed on request, and a
replayed stall may long since have ended.
These reports travel as network inputs, so while the node is catching up
updateCatchingUpHead parks them and a stalled/resumed pair only reaches
clients once it is in sync, by which time the stall it describes may be
over. NetworkInfo is unaffected, being read live.
Only a status seen on two consecutive polls is reported, so clients are not flooded (as with the sync status, see #2749). Without that, a network completing something every so often but less often than the stall period flaps between the two reports forever: with completions 15s apart and a 10s poll, the observed gaps cycle 9s, 4s, 14s, giving a report every ~30s indefinitely.
broadcastStallBounds :: StallBounds Source #
When the node starts refusing the client transactions that grow the outbound backlog, and when it reports the backlog to clients. Ten seconds of no progress, or of backlog at the recent drain rate, is comfortably above the etcd broadcast loop's one second retry and far below any contestation period. The cap on queued messages is the memory backstop: each holds at most a maximum size transaction, so ~160MB at the cap, and it sits well above the bursts a client can fire at a node whose network is keeping up.
Deliberately not operator-configurable: the useful range is narrow, nothing observable would tell an operator which value to pick, and the natural guess for "off" (zero) is the most aggressive setting rather than the least.
NOTE: this measures the hand-off, and the shipped $sel:broadcast:Network completes as
soon as the message is in the network component's own 100-slot
pending-broadcast queue. So during an outage the first ~100 messages still
complete promptly and nothing is reported; the stall only becomes visible
once both queues are saturated, which also puts the real in-flight bound
around $sel:maxPending:StallBounds plus that queue rather than at $sel:maxPending:StallBounds.
runHydraNode :: (MonadCatch m, MonadAsync m, MonadDelay m, MonadTime m, IsChainState tx) => HydraNode tx m -> m () Source #
stepHydraNode :: (MonadCatch m, MonadAsync m, MonadTime m, IsChainState tx) => UTCTime -> HydraNode tx m -> m () Source #
growsBroadcastBacklog :: ClientInput tx -> Bool Source #
Client inputs that turn into a broadcast directly, and so are the ones
worth refusing: the protocol's own messages cannot pile up while the
network is down, because $sel:broadcast:Network is self-delivering, so our own AckSn
never comes back, the confirmed snapshot number freezes and
snapshotInFlight caps both ReqSn and AckSn at one apiece.
NOTE: that cap is not airtight, and this is not a complete bound.
onOpenChainTick emits a ReqSn on a chain tick with no network input at
all, and SideLoadSnapshot - which broadcasts nothing itself, so is not
listed here - clears exactly the state that cap reads. A client looping
SideLoadSnapshot can therefore draw one further ReqSn per tick. Left
ungated on purpose: side-loading is the documented recovery for a ReqSn
or AckSn lost from the hand-off, so refusing it while stalled would block
the way out. The leak is one small message per chain tick.
inOpenHeadAndSynced :: NodeState tx -> Bool Source #
Whether growsBroadcastBacklog inputs would reach the code that
broadcasts, rather than being answered by the sync or head-state checks.
defaultTTL :: TTL Source #
The maximum number of times to re-enqueue a network messages upon Wait.
outcome.
defaultTxTTL :: TTL Source #
processNextInput :: IsChainState tx => HydraNode tx m -> Input tx -> UTCTime -> STM m (Outcome tx) Source #
Monadic interface around update.
processStateChanges :: (MonadSTM m, MonadTime m) => HydraNode tx m -> [StateChanged tx] -> m () Source #
processEffects :: (MonadAsync m, MonadCatch m, IsChainState tx) => HydraNode tx m -> Tracer m (HydraNodeLog tx) -> Word64 -> [Effect tx] -> m () Source #
Manage state
data NodeStateHandler tx m Source #
Handle to access and modify the state in the Hydra Node.
Constructors
| NodeStateHandler | |
Fields
| |
createNodeStateHandler Source #
Arguments
| :: MonadLabelledSTM m | |
| => Maybe EventId | Last seen |
| -> NodeState tx | |
| -> m (NodeStateHandler tx m) |
Initialize a new NodeStateHandler.
Logging
data HydraNodeLog tx Source #
Constructors
| BeginInput | |
| EndInput | |
| BeginEffect | |
| EndEffect | |
| LogicOutcome | |
| DroppedFromQueue | |
| LoadingState | |
| LoadedState | |
Fields
| |
| LoadedChainState | |
Fields
| |
| ReplayingState | |
| Misconfiguration | |
Fields | |
| DiscardedBroadcasts | Outbound messages accepted from the head logic but never handed to the network, dropped because the node is stopping. |
| BroadcastBacklog | How much the outbound hand-off is holding, and for how long it has
completed nothing. Emitted by |
Fields | |