-- |
-- Module      : Relay.Reliable.Outbox
-- Description : Sender half of the reliable-delivery protocol.
module Relay.Reliable.Outbox
  ( Outbox (..),
    emptyOutbox,
    send,
    ack,
    unackedFrames,
  )
where

import Data.Map.Strict (Map)
import Data.Map.Strict qualified as Map
import Relay.Session (Frame (..), Nat, SessionId)

-- | Sender-local reliable state for one session. Collapses Monarch's
-- @Outbox@ (sequence assignment) and @Unacked@ (in-flight buffer): in
-- the ack-driven model a message is simply either acknowledged or
-- still awaiting acknowledgement.
data Outbox a = Outbox
  { -- | The session every frame is stamped with. One 'Outbox' is one
    -- sending session.
    forall a. Outbox a -> SessionId
outboxSession :: SessionId,
    -- | Next sequence number to assign. Starts at 1.
    forall a. Outbox a -> Nat
nextSeq :: Nat,
    -- | Frames sent but not yet acknowledged, keyed by sequence
    -- number. A 'Map' (rather than a deque) is a deliberate modelling
    -- choice: it gives out-of-order insertion and a simple cumulative
    -- prune, and carries unchanged into the multi-stream development.
    forall a. Outbox a -> Map Nat a
unacked :: Map Nat a,
    -- | The sender's view of delivery progress: the highest sequence
    -- number the receiver has confirmed, as learned from acks arriving
    -- on the reverse channel. 'Nothing' until the first ack arrives.
    -- This is the sender's (possibly lagging) copy of the receiver's
    -- frontier, not the receiver's own state.
    forall a. Outbox a -> Maybe Nat
largestAcked :: Maybe Nat
  }
  deriving (Int -> Outbox a -> ShowS
[Outbox a] -> ShowS
Outbox a -> String
(Int -> Outbox a -> ShowS)
-> (Outbox a -> String) -> ([Outbox a] -> ShowS) -> Show (Outbox a)
forall a. Show a => Int -> Outbox a -> ShowS
forall a. Show a => [Outbox a] -> ShowS
forall a. Show a => Outbox a -> String
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: forall a. Show a => Int -> Outbox a -> ShowS
showsPrec :: Int -> Outbox a -> ShowS
$cshow :: forall a. Show a => Outbox a -> String
show :: Outbox a -> String
$cshowList :: forall a. Show a => [Outbox a] -> ShowS
showList :: [Outbox a] -> ShowS
Show, Outbox a -> Outbox a -> Bool
(Outbox a -> Outbox a -> Bool)
-> (Outbox a -> Outbox a -> Bool) -> Eq (Outbox a)
forall a. Eq a => Outbox a -> Outbox a -> Bool
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: forall a. Eq a => Outbox a -> Outbox a -> Bool
== :: Outbox a -> Outbox a -> Bool
$c/= :: forall a. Eq a => Outbox a -> Outbox a -> Bool
/= :: Outbox a -> Outbox a -> Bool
Eq)

-- | Initial sender state for a session. The first assigned sequence
-- number is 1; nothing is in flight or acknowledged.
emptyOutbox :: SessionId -> Outbox a
emptyOutbox :: forall a. SessionId -> Outbox a
emptyOutbox SessionId
sid =
  Outbox
    { outboxSession :: SessionId
outboxSession = SessionId
sid,
      nextSeq :: Nat
nextSeq = Nat
1,
      unacked :: Map Nat a
unacked = Map Nat a
forall k a. Map k a
Map.empty,
      largestAcked :: Maybe Nat
largestAcked = Maybe Nat
forall a. Maybe a
Nothing
    }

-- | Assign the next sequence number to a payload, record it as
-- in-flight, and return the frame to put on the wire.
send :: a -> Outbox a -> (Outbox a, Frame a)
send :: forall a. a -> Outbox a -> (Outbox a, Frame a)
send a
v Outbox a
ob =
  ( Outbox a
ob
      { nextSeq = nextSeq ob + 1,
        unacked = Map.insert (nextSeq ob) v (unacked ob)
      },
    Frame
      { frameSession :: SessionId
frameSession = Outbox a -> SessionId
forall a. Outbox a -> SessionId
outboxSession Outbox a
ob,
        frameSeq :: Nat
frameSeq = Outbox a -> Nat
forall a. Outbox a -> Nat
nextSeq Outbox a
ob,
        framePayload :: a
framePayload = a
v
      }
  )

-- | Process a cumulative acknowledgement: drop every in-flight frame
-- with sequence number @<= n@ and return their payloads in sequence
-- order.
--
-- Total under impossible acks: the recorded watermark is clamped to
-- @min n (nextSeq - 1)@, so an ack for frames never sent (a dishonest
-- @n >= nextSeq@) cannot push 'largestAcked' past what exists. The
-- watermark never decreases, so an ack at or below it is a no-op.
ack :: Nat -> Outbox a -> (Outbox a, [a])
ack :: forall a. Nat -> Outbox a -> (Outbox a, [a])
ack Nat
n Outbox a
ob =
  let (Map Nat a
acked, Map Nat a
kept) = (Nat -> a -> Bool) -> Map Nat a -> (Map Nat a, Map Nat a)
forall k a. (k -> a -> Bool) -> Map k a -> (Map k a, Map k a)
Map.partitionWithKey (\Nat
k a
_ -> Nat
k Nat -> Nat -> Bool
forall a. Ord a => a -> a -> Bool
<= Nat
n) (Outbox a -> Map Nat a
forall a. Outbox a -> Map Nat a
unacked Outbox a
ob)
      clamped :: Nat
clamped = Nat -> Nat -> Nat
forall a. Ord a => a -> a -> a
min Nat
n (Outbox a -> Nat
forall a. Outbox a -> Nat
nextSeq Outbox a
ob Nat -> Nat -> Nat
forall a. Num a => a -> a -> a
- Nat
1)
      watermark :: Maybe Nat
watermark = Nat -> Maybe Nat
forall a. a -> Maybe a
Just (Nat -> (Nat -> Nat) -> Maybe Nat -> Nat
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Nat
clamped (Nat -> Nat -> Nat
forall a. Ord a => a -> a -> a
max Nat
clamped) (Outbox a -> Maybe Nat
forall a. Outbox a -> Maybe Nat
largestAcked Outbox a
ob))
   in (Outbox a
ob {unacked = kept, largestAcked = watermark}, Map Nat a -> [a]
forall k a. Map k a -> [a]
Map.elems Map Nat a
acked)

-- | The in-flight frames to retransmit, each carrying its original
-- sequence number, in ascending sequence order. Retransmission never
-- reassigns sequence numbers.
unackedFrames :: Outbox a -> [Frame a]
unackedFrames :: forall a. Outbox a -> [Frame a]
unackedFrames Outbox a
ob =
  [ Frame {frameSession :: SessionId
frameSession = Outbox a -> SessionId
forall a. Outbox a -> SessionId
outboxSession Outbox a
ob, frameSeq :: Nat
frameSeq = Nat
s, framePayload :: a
framePayload = a
v}
  | (Nat
s, a
v) <- Map Nat a -> [(Nat, a)]
forall k a. Map k a -> [(k, a)]
Map.toAscList (Outbox a -> Map Nat a
forall a. Outbox a -> Map Nat a
unacked Outbox a
ob)
  ]