-- | Named runner circuits that tie free dual ends into a turn.
--
-- A turn is a /runner/ observation: it commits one token, then blocks on one
-- frame.  Framing lives port-side ('Circuit.Agent.StdPorts' stream marks);
-- the queue retry is the blocking boundary.  No polling, no backoff — the
-- halt is decided by content, not inferred from quiet.
--
-- The correlation between a command and its response is carried by the
-- 'TurnToken' envelope, not by pipe adjacency.  The mediator's residual
-- state holds the pending command; a response matches when its 'thread'
-- cites the command's id.  This makes zero-frame and two-frame failures
-- detectable in-band rather than by positional guess.
--
-- Mark-carrying turns must use a 'Linear' channel policy; weakening policies
-- can drop a command or response and break correlation.
--
-- @
--   turn e        :: IO (Loop (,) (K IO) Text Text)
--   turnTimeout u e :: IO (Loop (,) (K IO) Text (Maybe Text))
-- @
module Circuit.Agent.Turn
  ( -- * Turn envelope
    TurnToken (..),

    -- * Turn process
    TurnState (..),
    turnProcess,

    -- * Runner circuits
    turn,
    turnTimeout,
  )
where

import Circuit (Trace, base)
import Circuit.Agent (PostId)
import Circuit.Category (K (..))
import Circuit.Poles (HasDual (..), Poles (..), commit, emit, open)
import Circuit.Process (Process, mealy)
import Data.IORef
import Data.Text (Text)
import System.Timeout (timeout)
import Prelude

-- | A turn envelope.  Commands are sent with an empty 'turnThread'; the
-- turn assigns them a fresh id.  Responses must carry that id in their
-- 'turnThread' to be matched with the pending command.
data TurnToken a = TurnToken
  { forall a. TurnToken a -> a
turnBody :: a,
    forall a. TurnToken a -> [PostId]
turnThread :: [PostId]
  }
  deriving (TurnToken a -> TurnToken a -> Bool
(TurnToken a -> TurnToken a -> Bool)
-> (TurnToken a -> TurnToken a -> Bool) -> Eq (TurnToken a)
forall a. Eq a => TurnToken a -> TurnToken a -> Bool
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: forall a. Eq a => TurnToken a -> TurnToken a -> Bool
== :: TurnToken a -> TurnToken a -> Bool
$c/= :: forall a. Eq a => TurnToken a -> TurnToken a -> Bool
/= :: TurnToken a -> TurnToken a -> Bool
Eq, Int -> TurnToken a -> ShowS
[TurnToken a] -> ShowS
TurnToken a -> String
(Int -> TurnToken a -> ShowS)
-> (TurnToken a -> String)
-> ([TurnToken a] -> ShowS)
-> Show (TurnToken a)
forall a. Show a => Int -> TurnToken a -> ShowS
forall a. Show a => [TurnToken a] -> ShowS
forall a. Show a => TurnToken a -> String
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: forall a. Show a => Int -> TurnToken a -> ShowS
showsPrec :: Int -> TurnToken a -> ShowS
$cshow :: forall a. Show a => TurnToken a -> String
show :: TurnToken a -> String
$cshowList :: forall a. Show a => [TurnToken a] -> ShowS
showList :: [TurnToken a] -> ShowS
Show)

-- | Residual state of the turn mediator.  'nextId' supplies fresh command
-- ids; 'pending' holds the command awaiting its response.
data TurnState a = TurnState
  { forall a. TurnState a -> PostId
nextId :: PostId,
    forall a. TurnState a -> Maybe (PostId, TurnToken a)
pending :: Maybe (PostId, TurnToken a)
  }
  deriving (TurnState a -> TurnState a -> Bool
(TurnState a -> TurnState a -> Bool)
-> (TurnState a -> TurnState a -> Bool) -> Eq (TurnState a)
forall a. Eq a => TurnState a -> TurnState a -> Bool
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: forall a. Eq a => TurnState a -> TurnState a -> Bool
== :: TurnState a -> TurnState a -> Bool
$c/= :: forall a. Eq a => TurnState a -> TurnState a -> Bool
/= :: TurnState a -> TurnState a -> Bool
Eq, Int -> TurnState a -> ShowS
[TurnState a] -> ShowS
TurnState a -> String
(Int -> TurnState a -> ShowS)
-> (TurnState a -> String)
-> ([TurnState a] -> ShowS)
-> Show (TurnState a)
forall a. Show a => Int -> TurnState a -> ShowS
forall a. Show a => [TurnState a] -> ShowS
forall a. Show a => TurnState a -> String
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: forall a. Show a => Int -> TurnState a -> ShowS
showsPrec :: Int -> TurnState a -> ShowS
$cshow :: forall a. Show a => TurnState a -> String
show :: TurnState a -> String
$cshowList :: forall a. Show a => [TurnState a] -> ShowS
showList :: [TurnState a] -> ShowS
Show)

-- | Initial turn state.
emptyTurnState :: TurnState a
emptyTurnState :: forall a. TurnState a
emptyTurnState = PostId -> Maybe (PostId, TurnToken a) -> TurnState a
forall a. PostId -> Maybe (PostId, TurnToken a) -> TurnState a
TurnState PostId
0 Maybe (PostId, TurnToken a)
forall a. Maybe a
Nothing

-- | Inject a command into the turn state, assigning it the next id.
injectCommand :: TurnState a -> a -> (TurnState a, TurnToken a)
injectCommand :: forall a. TurnState a -> a -> (TurnState a, TurnToken a)
injectCommand TurnState a
st a
body =
  let cmdId :: PostId
cmdId = TurnState a -> PostId
forall a. TurnState a -> PostId
nextId TurnState a
st
      cmd :: TurnToken a
cmd = a -> [PostId] -> TurnToken a
forall a. a -> [PostId] -> TurnToken a
TurnToken a
body [PostId
cmdId]
   in (TurnState a
st {nextId = succ cmdId, pending = Just (cmdId, cmd)}, TurnToken a
cmd)

-- | Match a response against the pending command.
matchResponse :: TurnState a -> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a))
matchResponse :: forall a.
TurnState a
-> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a))
matchResponse st :: TurnState a
st@TurnState {pending :: forall a. TurnState a -> Maybe (PostId, TurnToken a)
pending = Maybe (PostId, TurnToken a)
Nothing} TurnToken a
_ = (TurnState a
st, Maybe (TurnToken a, TurnToken a)
forall a. Maybe a
Nothing)
matchResponse st :: TurnState a
st@TurnState {pending :: forall a. TurnState a -> Maybe (PostId, TurnToken a)
pending = Just (PostId
cmdId, TurnToken a
cmd)} TurnToken a
resp
  | TurnToken a -> [PostId]
forall a. TurnToken a -> [PostId]
turnThread TurnToken a
resp [PostId] -> [PostId] -> Bool
forall a. Eq a => a -> a -> Bool
== [PostId
cmdId] = (TurnState a
st {pending = Nothing}, (TurnToken a, TurnToken a) -> Maybe (TurnToken a, TurnToken a)
forall a. a -> Maybe a
Just (TurnToken a
cmd, TurnToken a
resp))
  | Bool
otherwise = (TurnState a
st, Maybe (TurnToken a, TurnToken a)
forall a. Maybe a
Nothing)

-- | Process that correlates responses with the pending command by thread.
--
-- * A token with empty thread is treated as a command: it receives the next
--   id and is stored as the pending command.
-- * A token with non-empty thread is treated as a response: it matches when
--   its thread equals the pending command's id, emitting the pair.
turnProcess :: Process (TurnToken a) (Maybe (TurnToken a, TurnToken a))
turnProcess :: forall a. Process (TurnToken a) (Maybe (TurnToken a, TurnToken a))
turnProcess =
  TurnState a
-> (TurnState a
    -> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a)))
-> Process (TurnToken a) (Maybe (TurnToken a, TurnToken a))
forall ch a b.
ch -> (ch -> a -> (ch, Maybe b)) -> Process a (Maybe b)
mealy TurnState a
forall a. TurnState a
emptyTurnState ((TurnState a
  -> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a)))
 -> Process (TurnToken a) (Maybe (TurnToken a, TurnToken a)))
-> (TurnState a
    -> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a)))
-> Process (TurnToken a) (Maybe (TurnToken a, TurnToken a))
forall a b. (a -> b) -> a -> b
$ \TurnState a
st TurnToken a
tok ->
    if [PostId] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null (TurnToken a -> [PostId]
forall a. TurnToken a -> [PostId]
turnThread TurnToken a
tok)
      then ((TurnState a, TurnToken a) -> TurnState a
forall a b. (a, b) -> a
fst (TurnState a -> a -> (TurnState a, TurnToken a)
forall a. TurnState a -> a -> (TurnState a, TurnToken a)
injectCommand TurnState a
st (TurnToken a -> a
forall a. TurnToken a -> a
turnBody TurnToken a
tok)), Maybe (TurnToken a, TurnToken a)
forall a. Maybe a
Nothing)
      else TurnState a
-> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a))
forall a.
TurnState a
-> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a))
matchResponse TurnState a
st TurnToken a
tok

-- | Read the unit poles used to plug the unused side of a commit or emit.
unitEnds :: Poles (K IO) () ()
unitEnds :: Poles (K IO) () ()
unitEnds = Poles (K IO) () ()
forall {k} (bot :: k) (arr :: k -> k -> *).
HasDual bot arr =>
Poles arr bot bot
open

-- | Run one turn: commit a body, then block until a matching response
-- arrives.  The correlation is by 'thread', not by position.
turn ::
  Poles (K IO) (TurnToken Text) (TurnToken Text) ->
  IO (Trace (,) (K IO) Text Text)
turn :: Poles (K IO) (TurnToken Text) (TurnToken Text)
-> IO (Trace (,) (K IO) Text Text)
turn Poles (K IO) (TurnToken Text) (TurnToken Text)
e = do
  ref <- TurnState Text -> IO (IORef (TurnState Text))
forall a. a -> IO (IORef a)
newIORef TurnState Text
forall a. TurnState a
emptyTurnState
  pure . base . K $ \Text
body -> do
    cmd <- IORef (TurnState Text)
-> (TurnState Text -> (TurnState Text, TurnToken Text))
-> IO (TurnToken Text)
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef (TurnState Text)
ref ((TurnState Text -> (TurnState Text, TurnToken Text))
 -> IO (TurnToken Text))
-> (TurnState Text -> (TurnState Text, TurnToken Text))
-> IO (TurnToken Text)
forall a b. (a -> b) -> a -> b
$ \TurnState Text
st -> TurnState Text -> Text -> (TurnState Text, TurnToken Text)
forall a. TurnState a -> a -> (TurnState a, TurnToken a)
injectCommand TurnState Text
st Text
body
    runK (commit (conjoint e) outU) cmd
    let loop = do
          resp <- K IO () (TurnToken Text) -> () -> IO (TurnToken Text)
forall {k} (m :: k -> *) a (b :: k). K m a b -> a -> m b
runK (Out (K IO) (TurnToken Text)
-> forall x. In (K IO) x -> K IO x (TurnToken Text)
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k2).
Out arr a -> forall (x :: k1). In arr x -> arr x a
emit (Poles (K IO) (TurnToken Text) (TurnToken Text)
-> Out (K IO) (TurnToken Text)
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> Out arr b
companion Poles (K IO) (TurnToken Text) (TurnToken Text)
e) In (K IO) ()
inU) ()
          mResult <- atomicModifyIORef' ref $ \TurnState Text
st -> TurnState Text
-> TurnToken Text
-> (TurnState Text, Maybe (TurnToken Text, TurnToken Text))
forall a.
TurnState a
-> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a))
matchResponse TurnState Text
st TurnToken Text
resp
          case mResult of
            Just (TurnToken Text
_, TurnToken Text
resp') -> Text -> IO Text
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (TurnToken Text -> Text
forall a. TurnToken a -> a
turnBody TurnToken Text
resp')
            Maybe (TurnToken Text, TurnToken Text)
Nothing -> IO Text
loop
    loop
  where
    Poles In (K IO) ()
_ Out (K IO) ()
outU = Poles (K IO) () ()
unitEnds
    Poles In (K IO) ()
inU Out (K IO) ()
_ = Poles (K IO) () ()
unitEnds

-- | 'turn' under a deadline (microseconds).  'Nothing' on expiry; the
-- unarrived response is not lost — the next emit still receives it.
turnTimeout ::
  Int ->
  Poles (K IO) (TurnToken Text) (TurnToken Text) ->
  IO (Trace (,) (K IO) Text (Maybe Text))
turnTimeout :: Int
-> Poles (K IO) (TurnToken Text) (TurnToken Text)
-> IO (Trace (,) (K IO) Text (Maybe Text))
turnTimeout Int
us Poles (K IO) (TurnToken Text) (TurnToken Text)
e = do
  ref <- TurnState Text -> IO (IORef (TurnState Text))
forall a. a -> IO (IORef a)
newIORef TurnState Text
forall a. TurnState a
emptyTurnState
  pure . base . K $ \Text
body -> do
    cmd <- IORef (TurnState Text)
-> (TurnState Text -> (TurnState Text, TurnToken Text))
-> IO (TurnToken Text)
forall a b. IORef a -> (a -> (a, b)) -> IO b
atomicModifyIORef' IORef (TurnState Text)
ref ((TurnState Text -> (TurnState Text, TurnToken Text))
 -> IO (TurnToken Text))
-> (TurnState Text -> (TurnState Text, TurnToken Text))
-> IO (TurnToken Text)
forall a b. (a -> b) -> a -> b
$ \TurnState Text
st -> TurnState Text -> Text -> (TurnState Text, TurnToken Text)
forall a. TurnState a -> a -> (TurnState a, TurnToken a)
injectCommand TurnState Text
st Text
body
    runK (commit (conjoint e) outU) cmd
    timeout us $ do
      let loop = do
            resp <- K IO () (TurnToken Text) -> () -> IO (TurnToken Text)
forall {k} (m :: k -> *) a (b :: k). K m a b -> a -> m b
runK (Out (K IO) (TurnToken Text)
-> forall x. In (K IO) x -> K IO x (TurnToken Text)
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k2).
Out arr a -> forall (x :: k1). In arr x -> arr x a
emit (Poles (K IO) (TurnToken Text) (TurnToken Text)
-> Out (K IO) (TurnToken Text)
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> Out arr b
companion Poles (K IO) (TurnToken Text) (TurnToken Text)
e) In (K IO) ()
inU) ()
            mResult <- atomicModifyIORef' ref $ \TurnState Text
st -> TurnState Text
-> TurnToken Text
-> (TurnState Text, Maybe (TurnToken Text, TurnToken Text))
forall a.
TurnState a
-> TurnToken a -> (TurnState a, Maybe (TurnToken a, TurnToken a))
matchResponse TurnState Text
st TurnToken Text
resp
            case mResult of
              Just (TurnToken Text
_, TurnToken Text
resp') -> Text -> IO Text
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (TurnToken Text -> Text
forall a. TurnToken a -> a
turnBody TurnToken Text
resp')
              Maybe (TurnToken Text, TurnToken Text)
Nothing -> IO Text
loop
      loop
  where
    Poles In (K IO) ()
_ Out (K IO) ()
outU = Poles (K IO) () ()
unitEnds
    Poles In (K IO) ()
inU Out (K IO) ()
_ = Poles (K IO) () ()
unitEnds