{-# LANGUAGE ConstraintKinds #-}
{-# LANGUAGE DataKinds #-}
{-# LANGUAGE GADTs #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE StandaloneKindSignatures #-}
{-# LANGUAGE TypeFamilies #-}
module Circuit.Agent.Ends
(
Queue (..),
ChannelPolicy (..),
openChannel,
openChannelSTM,
openLinearChannel,
openLinearChannelSTM,
HaltChannel (..),
IsLinear,
openHaltChannel,
writeHaltChannel,
readHaltChannel,
openSTM,
openIO,
pipeEnds,
)
where
import Circuit.Category (K (..))
import Circuit.Poles (HasDual (..), Poles (..), commit, companion, conjoint, emit, open, polesK, splay0)
import Control.Applicative
import Control.Concurrent.Async (async, cancel)
import Control.Concurrent.STM
import Control.Monad (forever, void)
import Data.Kind (Constraint, Type)
import GHC.TypeLits (ErrorMessage (..), TypeError)
import Prelude
data Queue a
=
Unbounded
|
Bounded Int
|
Single
|
SwapQ
|
Latest a
|
Newest Int
deriving (Int -> Queue a -> ShowS
[Queue a] -> ShowS
Queue a -> String
(Int -> Queue a -> ShowS)
-> (Queue a -> String) -> ([Queue a] -> ShowS) -> Show (Queue a)
forall a. Show a => Int -> Queue a -> ShowS
forall a. Show a => [Queue a] -> ShowS
forall a. Show a => Queue a -> String
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: forall a. Show a => Int -> Queue a -> ShowS
showsPrec :: Int -> Queue a -> ShowS
$cshow :: forall a. Show a => Queue a -> String
show :: Queue a -> String
$cshowList :: forall a. Show a => [Queue a] -> ShowS
showList :: [Queue a] -> ShowS
Show, Queue a -> Queue a -> Bool
(Queue a -> Queue a -> Bool)
-> (Queue a -> Queue a -> Bool) -> Eq (Queue a)
forall a. Eq a => Queue a -> Queue a -> Bool
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: forall a. Eq a => Queue a -> Queue a -> Bool
== :: Queue a -> Queue a -> Bool
$c/= :: forall a. Eq a => Queue a -> Queue a -> Bool
/= :: Queue a -> Queue a -> Bool
Eq)
data ChannelPolicy a
=
Linear
|
SingleSlot
|
SwapOne
|
LatestValue a
|
BoundedN Int
|
NewestN Int
deriving (Int -> ChannelPolicy a -> ShowS
[ChannelPolicy a] -> ShowS
ChannelPolicy a -> String
(Int -> ChannelPolicy a -> ShowS)
-> (ChannelPolicy a -> String)
-> ([ChannelPolicy a] -> ShowS)
-> Show (ChannelPolicy a)
forall a. Show a => Int -> ChannelPolicy a -> ShowS
forall a. Show a => [ChannelPolicy a] -> ShowS
forall a. Show a => ChannelPolicy a -> String
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: forall a. Show a => Int -> ChannelPolicy a -> ShowS
showsPrec :: Int -> ChannelPolicy a -> ShowS
$cshow :: forall a. Show a => ChannelPolicy a -> String
show :: ChannelPolicy a -> String
$cshowList :: forall a. Show a => [ChannelPolicy a] -> ShowS
showList :: [ChannelPolicy a] -> ShowS
Show, ChannelPolicy a -> ChannelPolicy a -> Bool
(ChannelPolicy a -> ChannelPolicy a -> Bool)
-> (ChannelPolicy a -> ChannelPolicy a -> Bool)
-> Eq (ChannelPolicy a)
forall a. Eq a => ChannelPolicy a -> ChannelPolicy a -> Bool
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: forall a. Eq a => ChannelPolicy a -> ChannelPolicy a -> Bool
== :: ChannelPolicy a -> ChannelPolicy a -> Bool
$c/= :: forall a. Eq a => ChannelPolicy a -> ChannelPolicy a -> Bool
/= :: ChannelPolicy a -> ChannelPolicy a -> Bool
Eq)
policyToQueue :: ChannelPolicy a -> Queue a
policyToQueue :: forall a. ChannelPolicy a -> Queue a
policyToQueue = \case
ChannelPolicy a
Linear -> Queue a
forall a. Queue a
Unbounded
ChannelPolicy a
SingleSlot -> Queue a
forall a. Queue a
Single
ChannelPolicy a
SwapOne -> Queue a
forall a. Queue a
SwapQ
LatestValue a
a -> a -> Queue a
forall a. a -> Queue a
Latest a
a
BoundedN Int
n -> Int -> Queue a
forall a. Int -> Queue a
Bounded Int
n
NewestN Int
n -> Int -> Queue a
forall a. Int -> Queue a
Newest Int
n
openChannel :: ChannelPolicy a -> IO (Poles (K IO) a a)
openChannel :: forall a. ChannelPolicy a -> IO (Poles (K IO) a a)
openChannel = Queue a -> IO (Poles (K IO) a a)
forall a. Queue a -> IO (Poles (K IO) a a)
openIO (Queue a -> IO (Poles (K IO) a a))
-> (ChannelPolicy a -> Queue a)
-> ChannelPolicy a
-> IO (Poles (K IO) a a)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ChannelPolicy a -> Queue a
forall a. ChannelPolicy a -> Queue a
policyToQueue
openChannelSTM :: ChannelPolicy a -> STM (Poles (K STM) a a)
openChannelSTM :: forall a. ChannelPolicy a -> STM (Poles (K STM) a a)
openChannelSTM = Queue a -> STM (Poles (K STM) a a)
forall a. Queue a -> STM (Poles (K STM) a a)
openSTM (Queue a -> STM (Poles (K STM) a a))
-> (ChannelPolicy a -> Queue a)
-> ChannelPolicy a
-> STM (Poles (K STM) a a)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ChannelPolicy a -> Queue a
forall a. ChannelPolicy a -> Queue a
policyToQueue
openLinearChannel :: IO (Poles (K IO) a a)
openLinearChannel :: forall a. IO (Poles (K IO) a a)
openLinearChannel = ChannelPolicy a -> IO (Poles (K IO) a a)
forall a. ChannelPolicy a -> IO (Poles (K IO) a a)
openChannel ChannelPolicy a
forall a. ChannelPolicy a
Linear
openLinearChannelSTM :: STM (Poles (K STM) a a)
openLinearChannelSTM :: forall a. STM (Poles (K STM) a a)
openLinearChannelSTM = ChannelPolicy a -> STM (Poles (K STM) a a)
forall a. ChannelPolicy a -> STM (Poles (K STM) a a)
openChannelSTM ChannelPolicy a
forall a. ChannelPolicy a
Linear
type family IsLinear (p :: ChannelPolicy a) :: Constraint where
IsLinear 'Linear = ()
IsLinear p = TypeError ('Text "only 'Linear' channels can carry halt marks")
type HaltChannel :: ChannelPolicy a -> Type
data HaltChannel p where
HaltChannel :: (IsLinear p) => Poles (K STM) a a -> HaltChannel (p :: ChannelPolicy a)
openHaltChannel :: STM (HaltChannel 'Linear)
openHaltChannel :: forall {a}. STM (HaltChannel 'Linear)
openHaltChannel = Poles (K STM) a a -> HaltChannel 'Linear
forall a (p :: ChannelPolicy a).
IsLinear p =>
Poles (K STM) a a -> HaltChannel p
HaltChannel (Poles (K STM) a a -> HaltChannel 'Linear)
-> STM (Poles (K STM) a a) -> STM (HaltChannel 'Linear)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> STM (Poles (K STM) a a)
forall a. STM (Poles (K STM) a a)
openLinearChannelSTM
writeHaltChannel :: forall a (p :: ChannelPolicy a). HaltChannel p -> a -> STM ()
writeHaltChannel :: forall a (p :: ChannelPolicy a). HaltChannel p -> a -> STM ()
writeHaltChannel (HaltChannel Poles (K STM) a a
ends) = K STM a () -> a -> STM ()
forall {k} (m :: k -> *) a (b :: k). K m a b -> a -> m b
runK (In (K STM) a -> forall x. Out (K STM) x -> K STM a x
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1).
In arr a -> forall (x :: k2). Out arr x -> arr a x
commit (Poles (K STM) a a -> In (K STM) a
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> In arr a
conjoint Poles (K STM) a a
ends) Out (K STM) ()
haltOut)
where
haltOut :: Out (K STM) ()
haltOut = Poles (K STM) () () -> Out (K STM) ()
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> Out arr b
companion (Poles (K STM) () ()
unitEndsSTM :: Poles (K STM) () ())
readHaltChannel :: forall a (p :: ChannelPolicy a). HaltChannel p -> STM a
readHaltChannel :: forall a (p :: ChannelPolicy a). HaltChannel p -> STM a
readHaltChannel (HaltChannel Poles (K STM) a a
ends) = K STM () a -> () -> STM a
forall {k} (m :: k -> *) a (b :: k). K m a b -> a -> m b
runK (Out (K STM) a -> forall x. In (K STM) x -> K STM x a
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k2).
Out arr a -> forall (x :: k1). In arr x -> arr x a
emit (Poles (K STM) a a -> Out (K STM) a
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> Out arr b
companion Poles (K STM) a a
ends) In (K STM) ()
haltIn) ()
where
haltIn :: In (K STM) ()
haltIn = Poles (K STM) () () -> In (K STM) ()
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> In arr a
conjoint (Poles (K STM) () ()
unitEndsSTM :: Poles (K STM) () ())
unitEndsSTM :: Poles (K STM) () ()
unitEndsSTM :: Poles (K STM) () ()
unitEndsSTM = (() -> STM ()) -> STM () -> Poles (K STM) () ()
forall (m :: * -> *) a b.
Monad m =>
(a -> m ()) -> m b -> Poles (K m) a b
polesK (STM () -> () -> STM ()
forall a b. a -> b -> a
const (() -> STM ()
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ())) (() -> STM ()
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ())
endsSTM :: Queue a -> STM (a -> STM (), STM a)
endsSTM :: forall a. Queue a -> STM (a -> STM (), STM a)
endsSTM = \case
Bounded Int
n -> do
q <- Natural -> STM (TBQueue a)
forall a. Natural -> STM (TBQueue a)
newTBQueue (Int -> Natural
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
n)
pure (writeTBQueue q, readTBQueue q)
Queue a
Unbounded -> do
q <- STM (TQueue a)
forall a. STM (TQueue a)
newTQueue
pure (writeTQueue q, readTQueue q)
Queue a
Single -> do
m <- STM (TMVar a)
forall a. STM (TMVar a)
newEmptyTMVar
pure (putTMVar m, takeTMVar m)
Queue a
SwapQ -> do
v <- STM (TMVar a)
forall a. STM (TMVar a)
newEmptyTMVar
let write a
x = TMVar a -> a -> STM Bool
forall a. TMVar a -> a -> STM Bool
tryPutTMVar TMVar a
v a
x STM Bool -> (Bool -> STM ()) -> STM ()
forall a b. STM a -> (a -> STM b) -> STM b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case Bool
True -> () -> STM ()
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (); Bool
False -> STM a -> STM ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (TMVar a -> a -> STM a
forall a. TMVar a -> a -> STM a
swapTMVar TMVar a
v a
x)
pure (write, takeTMVar v)
Latest a
a -> do
t <- a -> STM (TVar a)
forall a. a -> STM (TVar a)
newTVar a
a
pure (writeTVar t, readTVar t)
Newest Int
n -> do
q <- Natural -> STM (TBQueue a)
forall a. Natural -> STM (TBQueue a)
newTBQueue (Int -> Natural
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
n)
let write a
x = TBQueue a -> a -> STM ()
forall a. TBQueue a -> a -> STM ()
writeTBQueue TBQueue a
q a
x STM () -> STM () -> STM ()
forall a. STM a -> STM a -> STM a
forall (f :: * -> *) a. Alternative f => f a -> f a -> f a
<|> (TBQueue a -> STM (Maybe a)
forall a. TBQueue a -> STM (Maybe a)
tryReadTBQueue TBQueue a
q STM (Maybe a) -> STM () -> STM ()
forall a b. STM a -> STM b -> STM b
forall (f :: * -> *) a b. Applicative f => f a -> f b -> f b
*> a -> STM ()
write a
x)
pure (write, readTBQueue q)
openSTM :: Queue a -> STM (Poles (K STM) a a)
openSTM :: forall a. Queue a -> STM (Poles (K STM) a a)
openSTM Queue a
q = do
(write, read') <- Queue a -> STM (a -> STM (), STM a)
forall a. Queue a -> STM (a -> STM (), STM a)
endsSTM Queue a
q
pure (polesK write read')
openIO :: Queue a -> IO (Poles (K IO) a a)
openIO :: forall a. Queue a -> IO (Poles (K IO) a a)
openIO Queue a
q = do
e <- STM (Poles (K STM) a a) -> IO (Poles (K STM) a a)
forall a. STM a -> IO a
atomically (Queue a -> STM (Poles (K STM) a a)
forall a. Queue a -> STM (Poles (K STM) a a)
openSTM Queue a
q)
let (K write, K receive) = splay0 e
pure (polesK (atomically . write) (atomically (receive ())))
pipeEnds ::
forall a b c.
Poles (K IO) a b ->
(TQueue b -> IO (Poles (K IO) b c)) ->
IO (Poles (K IO) a c, IO ())
pipeEnds :: forall a b c.
Poles (K IO) a b
-> (TQueue b -> IO (Poles (K IO) b c))
-> IO (Poles (K IO) a c, IO ())
pipeEnds Poles (K IO) a b
e1 TQueue b -> IO (Poles (K IO) b c)
makeE2 = do
q <- IO (TQueue b)
forall a. IO (TQueue a)
newTQueueIO
e2 <- makeE2 q
let unitEnds :: Poles (K IO) () ()
unitEnds = Poles (K IO) () ()
forall {k} (bot :: k) (arr :: k -> k -> *).
HasDual bot arr =>
Poles arr bot bot
open
readFromE1 :: IO b
readFromE1 = K IO () b -> () -> IO b
forall {k} (m :: k -> *) a (b :: k). K m a b -> a -> m b
runK (Out (K IO) b -> forall x. In (K IO) x -> K IO x b
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k2).
Out arr a -> forall (x :: k1). In arr x -> arr x a
emit (Poles (K IO) a b -> Out (K IO) b
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> Out arr b
companion Poles (K IO) a b
e1) (Poles (K IO) () () -> In (K IO) ()
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> In arr a
conjoint Poles (K IO) () ()
unitEnds)) ()
writeToE2 :: b -> IO ()
writeToE2 = K IO b () -> b -> IO ()
forall {k} (m :: k -> *) a (b :: k). K m a b -> a -> m b
runK (In (K IO) b -> forall x. Out (K IO) x -> K IO b x
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1).
In arr a -> forall (x :: k2). Out arr x -> arr a x
commit (Poles (K IO) b c -> In (K IO) b
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> In arr a
conjoint Poles (K IO) b c
e2) (Poles (K IO) () () -> Out (K IO) ()
forall {k1} {k2} (arr :: k1 -> k2 -> *) (a :: k1) (b :: k2).
Poles arr a b -> Out arr b
companion Poles (K IO) () ()
unitEnds))
pump <- async . forever $ do
x <- readFromE1
writeToE2 x
pure (Poles (conjoint e1) (companion e2), cancel pump)