{-# LANGUAGE OverloadedStrings #-}

-- | High-level client implementation for NATS.
module Client
  ( newClient
  , ConfigOption
  , withConnectName
  , withEcho
  , withAuthToken
  , withUserPass
  , withNKey
  , withJWT
  , withTLSCert
  , withMinimumLogLevel
  , withLogAction
  , withConnectionAttempts
  , withCallbackConcurrency
  , withBufferLimit
  , withExitAction
  , LogLevel (..)
  , LogEntry (..)
  , renderLogEntry
  , AuthTokenData
  , UserPassData
  , NKeyData
  , JWTTokenData
  , TLSPublicKey
  , TLSPrivateKey
  , TLSCertData
  , ClientExitReason (..)
  ) where

import           API                      (Client (..), MsgView (..))
import qualified Auth.Jwt                 as AuthJwt
import qualified Auth.NKey                as AuthNKey
import qualified Auth.None                as AuthNone
import qualified Auth.Token               as AuthToken
import           Auth.Types
    ( Auth
    , AuthTokenData
    , JWTTokenData
    , NKeyData
    , UserPassData
    )
import qualified Auth.UserPass            as AuthUserPass
import           Control.Concurrent       (forkIO)
import           Control.Concurrent.STM
import           Control.Exception        (SomeException, displayException)
import           Control.Monad            (void, when)
import qualified Data.ByteString          as BS
import           Engine                   (closeClient, resetClient, runEngine)
import           Lib.CallOption           (CallOption, applyCallOptions)
import           Lib.Logger
    ( LogEntry (..)
    , LogLevel (..)
    , LoggerConfig (..)
    , MonadLogger (..)
    , defaultLogger
    , newLogContext
    , renderLogEntry
    )
import           Network.Connection       (connectionApi)
import           Network.ConnectionAPI    (newConn)
import           Parser.Attoparsec        (parserApi)
import           Pipeline.Broadcasting    (broadcastingApi)
import           Pipeline.Streaming       (streamingApi)
import           Publish                  (defaultPublishConfig)
import           Publish.Config           (PublishConfig)
import           Queue.API                (QueueItem (QueueItem))
import           Queue.TransactionalQueue (newQueue)
import           State.Store
    ( ClientState
    , enqueue
    , newClientState
    , nextInbox
    , nextSid
    , pushPingAction
    , readStatus
    , runClient
    , setConnectName
    , waitForClosed
    , waitForNotRunning
    , waitForServerInfo
    )
import           State.Types
    ( ClientConfig (..)
    , ClientExitReason (..)
    , ClientStatus (..)
    , TLSCertData
    , TLSPrivateKey
    , TLSPublicKey
    )
import           Subscription.Store
    ( SubscriptionStore
    , awaitNoTrackedExpiries
    , hasTrackedExpiries
    , newSubscriptionStore
    , register
    , startExpiryWorker
    , startWorkers
    , unregister
    )
import           Subscription.Types
    ( SubscribeConfig (..)
    , SubscriptionMeta (SubscriptionMeta)
    )
import qualified Types.Connect            as Connect
import qualified Types.Msg                as Msg
import           Types.Ping               (Ping (..))
import qualified Types.Pub                as Pub
import qualified Types.Sub                as Sub
import qualified Types.Unsub              as Unsub

data ClientAuth = ClientAuthNone
                | ClientAuthToken AuthTokenData
                | ClientAuthUserPass UserPassData
                | ClientAuthNKey NKeyData
                | ClientAuthJWT JWTTokenData

data ClientOptions = ClientOptions
                       { ClientOptions -> Connect
optionConnectConfig       :: Connect.Connect
                       , ClientOptions -> ClientAuth
optionAuth                :: ClientAuth
                       , ClientOptions -> Maybe TLSCertData
optionTlsCert             :: Maybe TLSCertData
                       , ClientOptions -> LoggerConfig
optionLoggerConfig        :: LoggerConfig
                       , ClientOptions -> Int
optionConnectionAttempts  :: Int
                       , ClientOptions -> Int
optionCallbackConcurrency :: Int
                       , ClientOptions -> Int
optionBufferLimit         :: Int
                       , ClientOptions -> ClientExitReason -> IO ()
optionExitAction          :: ClientExitReason -> IO ()
                       , ClientOptions -> [([Char], Int)]
optionConnectOptions      :: [(String, Int)]
                       }

newClient :: [(String, Int)] -> [ConfigOption] -> IO Client
newClient :: [([Char], Int)] -> [ConfigOption] -> IO Client
newClient [([Char], Int)]
servers [ConfigOption]
configOptions = do
  LoggerConfig
loggerConfig' <- IO LoggerConfig
defaultLogger
  TVar LogContext
ctx <- IO (TVar LogContext)
newLogContext
  let defaultOptions :: ClientOptions
defaultOptions = [ConfigOption] -> ConfigOption
forall a. [CallOption a] -> CallOption a
applyCallOptions [ConfigOption]
configOptions ClientOptions
        { optionConnectConfig :: Connect
optionConnectConfig = Connect
defaultConnect
        , optionAuth :: ClientAuth
optionAuth = ClientAuth
ClientAuthNone
        , optionTlsCert :: Maybe TLSCertData
optionTlsCert = Maybe TLSCertData
forall a. Maybe a
Nothing
        , optionLoggerConfig :: LoggerConfig
optionLoggerConfig = LoggerConfig
loggerConfig'
        , optionConnectionAttempts :: Int
optionConnectionAttempts = Int
5
        , optionCallbackConcurrency :: Int
optionCallbackConcurrency = Int
1
        , optionBufferLimit :: Int
optionBufferLimit = Int
4096
        , optionExitAction :: ClientExitReason -> IO ()
optionExitAction = IO () -> ClientExitReason -> IO ()
forall a b. a -> b -> a
const (() -> IO ()
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ())
        , optionConnectOptions :: [([Char], Int)]
optionConnectOptions = [([Char], Int)]
servers
        }
      clientConfig :: ClientConfig
clientConfig =
        ClientConfig
          { connectionAttempts :: Int
connectionAttempts = ClientOptions -> Int
optionConnectionAttempts ClientOptions
defaultOptions
          , callbackConcurrency :: Int
callbackConcurrency = ClientOptions -> Int
optionCallbackConcurrency ClientOptions
defaultOptions
          , bufferLimit :: Int
bufferLimit = ClientOptions -> Int
optionBufferLimit ClientOptions
defaultOptions
          , connectConfig :: Connect
connectConfig = ClientOptions -> Connect
optionConnectConfig ClientOptions
defaultOptions
          , loggerConfig :: LoggerConfig
loggerConfig = ClientOptions -> LoggerConfig
optionLoggerConfig ClientOptions
defaultOptions
          , tlsCert :: Maybe TLSCertData
tlsCert = ClientOptions -> Maybe TLSCertData
optionTlsCert ClientOptions
defaultOptions
          , exitAction :: ClientExitReason -> IO ()
exitAction = ClientOptions -> ClientExitReason -> IO ()
optionExitAction ClientOptions
defaultOptions
          , connectOptions :: [([Char], Int)]
connectOptions = ClientOptions -> [([Char], Int)]
optionConnectOptions ClientOptions
defaultOptions
          }
      configuredAuth :: Auth
configuredAuth = ClientAuth -> Auth
selectAuth (ClientOptions -> ClientAuth
optionAuth ClientOptions
defaultOptions)

  Queue
queue <- IO Queue
newQueue
  Conn
conn <- ConnectionAPI -> IO Conn
newConn ConnectionAPI
connectionApi
  ClientState
clientState <- ClientConfig -> Queue -> Conn -> TVar LogContext -> IO ClientState
newClientState ClientConfig
clientConfig Queue
queue Conn
conn TVar LogContext
ctx
  SubscriptionStore
store <- IO SubscriptionStore
newSubscriptionStore

  ClientState -> Maybe SID -> IO ()
setConnectName ClientState
clientState (Connect -> Maybe SID
Connect.name (ClientOptions -> Connect
optionConnectConfig ClientOptions
defaultOptions))
  ClientState -> ClientOptions -> IO ()
logStaticConfiguration ClientState
clientState ClientOptions
defaultOptions

  Int
-> SubscriptionStore -> STM () -> (SomeException -> IO ()) -> IO ()
startWorkers
    (ClientConfig -> Int
callbackConcurrency ClientConfig
clientConfig)
    SubscriptionStore
store
    (do
        ClientState -> STM ()
waitForClosed ClientState
clientState
        SubscriptionStore -> STM ()
awaitNoTrackedExpiries SubscriptionStore
store)
    (ClientState -> SomeException -> IO ()
handleCallbackError ClientState
clientState)

  SubscriptionStore -> IO Bool -> IO ()
startExpiryWorker SubscriptionStore
store (IO Bool -> IO ()) -> IO Bool -> IO ()
forall a b. (a -> b) -> a -> b
$
    ClientState -> SubscriptionStore -> IO Bool
shouldStopExpiryWorker ClientState
clientState SubscriptionStore
store

  IO ThreadId -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO ThreadId -> IO ()) -> (IO () -> IO ThreadId) -> IO () -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. IO () -> IO ThreadId
forkIO (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$
    ConnectionAPI
-> StreamingAPI
-> BroadcastingAPI
-> ParserAPI ParsedMessage
-> ClientState
-> SubscriptionStore
-> Auth
-> IO ()
runEngine
      ConnectionAPI
connectionApi
      StreamingAPI
streamingApi
      BroadcastingAPI
broadcastingApi
      ParserAPI ParsedMessage
parserApi
      ClientState
clientState
      SubscriptionStore
store
      Auth
configuredAuth

  STM () -> IO ()
forall a. STM a -> IO a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$
    ClientState -> STM ()
waitForServerInfo ClientState
clientState
      STM () -> STM () -> STM ()
forall a. STM a -> STM a -> STM a
`orElse` ClientState -> STM ()
waitForClosed ClientState
clientState

  Client -> IO Client
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Client
    { publish :: SID -> [PublishOption] -> IO ()
publish = \SID
subject [PublishOption]
publishOptions -> do
        let cfg :: PublishConfig
cfg = [PublishOption] -> PublishOption
forall a. [CallOption a] -> CallOption a
applyCallOptions [PublishOption]
publishOptions PublishConfig
defaultPublishConfig
        ClientState -> SubscriptionStore -> SID -> PublishConfig -> IO ()
publishClient ClientState
clientState SubscriptionStore
store SID
subject PublishConfig
cfg
    , subscribe :: SID -> [SubscribeOption] -> (Maybe MsgView -> IO ()) -> IO SID
subscribe = \SID
subject [SubscribeOption]
subscribeOptions Maybe MsgView -> IO ()
callback -> do
        let cfg :: SubscribeConfig
cfg = [SubscribeOption] -> SubscribeOption
forall a. [CallOption a] -> CallOption a
applyCallOptions [SubscribeOption]
subscribeOptions SubscribeConfig
defaultSubscribeConfig
        ClientState
-> SubscriptionStore
-> Bool
-> SID
-> SubscribeConfig
-> (Maybe Msg -> IO ())
-> IO SID
subscribeClient ClientState
clientState SubscriptionStore
store Bool
False SID
subject SubscribeConfig
cfg ((Maybe MsgView -> IO ()) -> Maybe Msg -> IO ()
toInternalCallback Maybe MsgView -> IO ()
callback)
    , request :: SID -> [SubscribeOption] -> (Maybe MsgView -> IO ()) -> IO SID
request = \SID
subject [SubscribeOption]
subscribeOptions Maybe MsgView -> IO ()
callback -> do
        let cfg :: SubscribeConfig
cfg = [SubscribeOption] -> SubscribeOption
forall a. [CallOption a] -> CallOption a
applyCallOptions [SubscribeOption]
subscribeOptions SubscribeConfig
defaultSubscribeConfig
        ClientState
-> SubscriptionStore
-> Bool
-> SID
-> SubscribeConfig
-> (Maybe Msg -> IO ())
-> IO SID
subscribeClient ClientState
clientState SubscriptionStore
store Bool
True SID
subject SubscribeConfig
cfg ((Maybe MsgView -> IO ()) -> Maybe Msg -> IO ()
toInternalCallback Maybe MsgView -> IO ()
callback)
    , unsubscribe :: SID -> IO ()
unsubscribe = ClientState -> SubscriptionStore -> SID -> IO ()
unsubscribeClient ClientState
clientState SubscriptionStore
store
    , newInbox :: IO SID
newInbox = ClientState -> IO SID
nextInbox ClientState
clientState
    , ping :: IO () -> IO ()
ping = ClientState -> IO () -> IO ()
pingClient ClientState
clientState
    , flush :: IO ()
flush = ClientState -> IO ()
flushClient ClientState
clientState
    , reset :: IO ()
reset = ConnectionAPI -> ClientState -> SubscriptionStore -> IO ()
resetClient ConnectionAPI
connectionApi ClientState
clientState SubscriptionStore
store
    , close :: IO ()
close = ConnectionAPI -> ClientState -> SubscriptionStore -> IO ()
closeClient ConnectionAPI
connectionApi ClientState
clientState SubscriptionStore
store
    }

type ConfigOption = CallOption ClientOptions

withConnectName :: BS.ByteString -> ConfigOption
withConnectName :: SID -> ConfigOption
withConnectName SID
name ClientOptions
config =
  ClientOptions
config
    { optionConnectConfig =
        (optionConnectConfig config) { Connect.name = Just name }
    }

withEcho :: Bool -> ConfigOption
withEcho :: Bool -> ConfigOption
withEcho Bool
enabled ClientOptions
config =
  ClientOptions
config
    { optionConnectConfig =
        (optionConnectConfig config) { Connect.echo = Just enabled }
    }

withAuthToken :: AuthTokenData -> ConfigOption
withAuthToken :: SID -> ConfigOption
withAuthToken SID
token ClientOptions
config = ClientOptions
config { optionAuth = ClientAuthToken token }

withUserPass :: UserPassData -> ConfigOption
withUserPass :: TLSCertData -> ConfigOption
withUserPass TLSCertData
userPass ClientOptions
config = ClientOptions
config { optionAuth = ClientAuthUserPass userPass }

withNKey :: NKeyData -> ConfigOption
withNKey :: SID -> ConfigOption
withNKey SID
nkey ClientOptions
config = ClientOptions
config { optionAuth = ClientAuthNKey nkey }

withJWT :: JWTTokenData -> ConfigOption
withJWT :: SID -> ConfigOption
withJWT SID
jwt ClientOptions
config = ClientOptions
config { optionAuth = ClientAuthJWT jwt }

withTLSCert :: TLSCertData -> ConfigOption
withTLSCert :: TLSCertData -> ConfigOption
withTLSCert TLSCertData
cert ClientOptions
config = ClientOptions
config { optionTlsCert = Just cert }

withMinimumLogLevel :: LogLevel -> ConfigOption
withMinimumLogLevel :: LogLevel -> ConfigOption
withMinimumLogLevel LogLevel
minimumLogLevel ClientOptions
config =
  ClientOptions
config
    { optionLoggerConfig =
        (optionLoggerConfig config) { minLogLevel = minimumLogLevel }
    }

withLogAction :: (LogEntry -> IO ()) -> ConfigOption
withLogAction :: (LogEntry -> IO ()) -> ConfigOption
withLogAction LogEntry -> IO ()
logAction ClientOptions
config =
  ClientOptions
config
    { optionLoggerConfig =
        (optionLoggerConfig config) { logFn = logAction }
    }

withConnectionAttempts :: Int -> ConfigOption
withConnectionAttempts :: Int -> ConfigOption
withConnectionAttempts Int
attempts ClientOptions
config =
  ClientOptions
config { optionConnectionAttempts = attempts }

withCallbackConcurrency :: Int -> ConfigOption
withCallbackConcurrency :: Int -> ConfigOption
withCallbackConcurrency Int
concurrency ClientOptions
config =
  ClientOptions
config { optionCallbackConcurrency = concurrency }

withBufferLimit :: Int -> ConfigOption
withBufferLimit :: Int -> ConfigOption
withBufferLimit Int
limit ClientOptions
config =
  ClientOptions
config { optionBufferLimit = max 1 limit }

withExitAction :: (ClientExitReason -> IO ()) -> ConfigOption
withExitAction :: (ClientExitReason -> IO ()) -> ConfigOption
withExitAction ClientExitReason -> IO ()
action ClientOptions
config = ClientOptions
config { optionExitAction = action }

defaultConnect :: Connect.Connect
defaultConnect :: Connect
defaultConnect =
  Connect.Connect
    { verbose :: Bool
Connect.verbose = Bool
False
    , pedantic :: Bool
Connect.pedantic = Bool
True
    , tls_required :: Bool
Connect.tls_required = Bool
False
    , auth_token :: Maybe SID
Connect.auth_token = Maybe SID
forall a. Maybe a
Nothing
    , user :: Maybe SID
Connect.user = Maybe SID
forall a. Maybe a
Nothing
    , pass :: Maybe SID
Connect.pass = Maybe SID
forall a. Maybe a
Nothing
    , name :: Maybe SID
Connect.name = Maybe SID
forall a. Maybe a
Nothing
    , lang :: SID
Connect.lang = SID
"haskell"
    , version :: SID
Connect.version = SID
"0.1.0"
    , protocol :: Maybe Int
Connect.protocol = Maybe Int
forall a. Maybe a
Nothing
    , echo :: Maybe Bool
Connect.echo = Bool -> Maybe Bool
forall a. a -> Maybe a
Just Bool
True
    , sig :: Maybe SID
Connect.sig = Maybe SID
forall a. Maybe a
Nothing
    , jwt :: Maybe SID
Connect.jwt = Maybe SID
forall a. Maybe a
Nothing
    , nkey :: Maybe SID
Connect.nkey = Maybe SID
forall a. Maybe a
Nothing
    , no_responders :: Maybe Bool
Connect.no_responders = Bool -> Maybe Bool
forall a. a -> Maybe a
Just Bool
True
    , headers :: Maybe Bool
Connect.headers = Bool -> Maybe Bool
forall a. a -> Maybe a
Just Bool
True
    }

defaultSubscribeConfig :: SubscribeConfig
defaultSubscribeConfig :: SubscribeConfig
defaultSubscribeConfig = Maybe NominalDiffTime -> Maybe SID -> SubscribeConfig
SubscribeConfig Maybe NominalDiffTime
forall a. Maybe a
Nothing Maybe SID
forall a. Maybe a
Nothing

selectAuth :: ClientAuth -> Auth
selectAuth :: ClientAuth -> Auth
selectAuth ClientAuth
authSelection =
  case ClientAuth
authSelection of
    ClientAuth
ClientAuthNone ->
      Auth
AuthNone.auth
    ClientAuthToken SID
token ->
      SID -> Auth
AuthToken.auth SID
token
    ClientAuthUserPass TLSCertData
userPass ->
      TLSCertData -> Auth
AuthUserPass.auth TLSCertData
userPass
    ClientAuthNKey SID
seed ->
      SID -> Auth
AuthNKey.auth SID
seed
    ClientAuthJWT SID
creds ->
      SID -> Auth
AuthJwt.auth SID
creds

logStaticConfiguration :: ClientState -> ClientOptions -> IO ()
logStaticConfiguration :: ClientState -> ClientOptions -> IO ()
logStaticConfiguration ClientState
client ClientOptions
options =
  ClientState -> AppM () -> IO ()
forall a. ClientState -> AppM a -> IO a
runClient ClientState
client (AppM () -> IO ()) -> AppM () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
    case ClientOptions -> ClientAuth
optionAuth ClientOptions
options of
      ClientAuth
ClientAuthNone ->
        LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Info [Char]
"no authentication method provided"
      ClientAuthToken SID
_ ->
        LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Info [Char]
"using auth token"
      ClientAuthUserPass (SID
user, SID
_) ->
        LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Info ([Char]
"using user/pass: " [Char] -> [Char] -> [Char]
forall a. [a] -> [a] -> [a]
++ SID -> [Char]
forall a. Show a => a -> [Char]
show SID
user)
      ClientAuthNKey SID
_ ->
        LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Info [Char]
"using nkey"
      ClientAuthJWT SID
_ ->
        LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Info [Char]
"using jwt"
    case ClientOptions -> Maybe TLSCertData
optionTlsCert ClientOptions
options of
      Maybe TLSCertData
Nothing ->
        () -> AppM ()
forall a. a -> AppM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
      Just TLSCertData
_ ->
        LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Info [Char]
"using tls certificate"

handleCallbackError :: ClientState -> SomeException -> IO ()
handleCallbackError :: ClientState -> SomeException -> IO ()
handleCallbackError ClientState
client SomeException
err =
  ClientState -> AppM () -> IO ()
forall a. ClientState -> AppM a -> IO a
runClient ClientState
client (AppM () -> IO ()) -> AppM () -> IO ()
forall a b. (a -> b) -> a -> b
$
    LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Error ([Char]
"callback failed: " [Char] -> [Char] -> [Char]
forall a. [a] -> [a] -> [a]
++ SomeException -> [Char]
forall e. Exception e => e -> [Char]
displayException SomeException
err)

shouldStopExpiryWorker :: ClientState -> SubscriptionStore -> IO Bool
shouldStopExpiryWorker :: ClientState -> SubscriptionStore -> IO Bool
shouldStopExpiryWorker ClientState
client SubscriptionStore
store = do
  ClientStatus
status <- ClientState -> IO ClientStatus
readStatus ClientState
client
  Bool
tracked <- SubscriptionStore -> IO Bool
hasTrackedExpiries SubscriptionStore
store
  Bool -> IO Bool
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Bool -> IO Bool) -> Bool -> IO Bool
forall a b. (a -> b) -> a -> b
$
    case ClientStatus
status of
      Closed ClientExitReason
_ -> Bool -> Bool
not Bool
tracked
      ClientStatus
_        -> Bool
False

toInternalCallback :: (Maybe MsgView -> IO ()) -> Maybe Msg.Msg -> IO ()
toInternalCallback :: (Maybe MsgView -> IO ()) -> Maybe Msg -> IO ()
toInternalCallback Maybe MsgView -> IO ()
callback =
  Maybe MsgView -> IO ()
callback (Maybe MsgView -> IO ())
-> (Maybe Msg -> Maybe MsgView) -> Maybe Msg -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Msg -> MsgView) -> Maybe Msg -> Maybe MsgView
forall a b. (a -> b) -> Maybe a -> Maybe b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap Msg -> MsgView
toMsgView

toMsgView :: Msg.Msg -> MsgView
toMsgView :: Msg -> MsgView
toMsgView Msg
msg =
  MsgView
    { subject :: SID
subject = Msg -> SID
Msg.subject Msg
msg
    , sid :: SID
sid = Msg -> SID
Msg.sid Msg
msg
    , replyTo :: Maybe SID
replyTo = Msg -> Maybe SID
Msg.replyTo Msg
msg
    , payload :: Maybe SID
payload = Msg -> Maybe SID
Msg.payload Msg
msg
    , headers :: Maybe [TLSCertData]
headers = Msg -> Maybe [TLSCertData]
Msg.headers Msg
msg
    }

publishClient :: ClientState -> SubscriptionStore -> Msg.Subject -> PublishConfig -> IO ()
publishClient :: ClientState -> SubscriptionStore -> SID -> PublishConfig -> IO ()
publishClient ClientState
client SubscriptionStore
store SID
subject (Maybe SID
payload, Maybe (Maybe Msg -> IO ())
callback, Maybe [TLSCertData]
headers, Maybe SID
configuredReplyTo) = do
  ClientState -> AppM () -> IO ()
forall a. ClientState -> AppM a -> IO a
runClient ClientState
client (AppM () -> IO ()) -> AppM () -> IO ()
forall a b. (a -> b) -> a -> b
$
    LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Debug ([Char]
"publishing to subject: " [Char] -> [Char] -> [Char]
forall a. [a] -> [a] -> [a]
++ SID -> [Char]
forall a. Show a => a -> [Char]
show SID
subject)
  Maybe SID
replyTo <- case Maybe (Maybe Msg -> IO ())
callback of
    Maybe (Maybe Msg -> IO ())
Nothing ->
      Maybe SID -> IO (Maybe SID)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe SID
configuredReplyTo
    Just Maybe Msg -> IO ()
replyCallback -> do
      SID
inbox <- IO SID -> (SID -> IO SID) -> Maybe SID -> IO SID
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (ClientState -> IO SID
nextInbox ClientState
client) SID -> IO SID
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe SID
configuredReplyTo
      SID
sid <- ClientState -> IO SID
nextSid ClientState
client
      let meta :: SubscriptionMeta
meta =
            SID -> Maybe SID -> Bool -> SubscriptionMeta
SubscriptionMeta SID
inbox Maybe SID
forall a. Maybe a
Nothing Bool
True
      SubscriptionStore
-> SID
-> SubscriptionMeta
-> SubscribeConfig
-> (Maybe Msg -> IO ())
-> IO ()
register SubscriptionStore
store SID
sid SubscriptionMeta
meta SubscribeConfig
defaultSubscribeConfig Maybe Msg -> IO ()
replyCallback
      ClientState -> QueueItem -> IO ()
enqueue ClientState
client (QueueItem -> IO ()) -> QueueItem -> IO ()
forall a b. (a -> b) -> a -> b
$
        Sub -> QueueItem
forall m. Transformer m => m -> QueueItem
QueueItem
          Sub.Sub
            { subject :: SID
Sub.subject = SID
inbox
            , queueGroup :: Maybe SID
Sub.queueGroup = Maybe SID
forall a. Maybe a
Nothing
            , sid :: SID
Sub.sid = SID
sid
            }
      ClientState -> QueueItem -> IO ()
enqueue ClientState
client (QueueItem -> IO ()) -> QueueItem -> IO ()
forall a b. (a -> b) -> a -> b
$
        Unsub -> QueueItem
forall m. Transformer m => m -> QueueItem
QueueItem
          Unsub.Unsub
            { sid :: SID
Unsub.sid = SID
sid
            , maxMsg :: Maybe Int
Unsub.maxMsg = Int -> Maybe Int
forall a. a -> Maybe a
Just Int
1
            }
      Maybe SID -> IO (Maybe SID)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (SID -> Maybe SID
forall a. a -> Maybe a
Just SID
inbox)
  ClientState -> QueueItem -> IO ()
enqueue ClientState
client (QueueItem -> IO ()) -> QueueItem -> IO ()
forall a b. (a -> b) -> a -> b
$
    Pub -> QueueItem
forall m. Transformer m => m -> QueueItem
QueueItem
      Pub.Pub
        { subject :: SID
Pub.subject = SID
subject
        , payload :: Maybe SID
Pub.payload = Maybe SID
payload
        , replyTo :: Maybe SID
Pub.replyTo = Maybe SID
replyTo
        , headers :: Maybe [TLSCertData]
Pub.headers = Maybe [TLSCertData]
headers
        }

subscribeClient :: ClientState -> SubscriptionStore -> Bool -> Msg.Subject -> SubscribeConfig -> (Maybe Msg.Msg -> IO ()) -> IO Msg.SID
subscribeClient :: ClientState
-> SubscriptionStore
-> Bool
-> SID
-> SubscribeConfig
-> (Maybe Msg -> IO ())
-> IO SID
subscribeClient ClientState
client SubscriptionStore
store Bool
isReply SID
subject SubscribeConfig
cfg Maybe Msg -> IO ()
callback = do
  ClientState -> AppM () -> IO ()
forall a. ClientState -> AppM a -> IO a
runClient ClientState
client (AppM () -> IO ()) -> AppM () -> IO ()
forall a b. (a -> b) -> a -> b
$
    LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Debug ([Char]
"subscribing to subject: " [Char] -> [Char] -> [Char]
forall a. [a] -> [a] -> [a]
++ SID -> [Char]
forall a. Show a => a -> [Char]
show SID
subject)
  SID
sid <- ClientState -> IO SID
nextSid ClientState
client
  let queueGroup :: Maybe SID
queueGroup =
        SubscribeConfig -> Maybe SID
subscribeQueueGroup SubscribeConfig
cfg
  let meta :: SubscriptionMeta
meta =
        SID -> Maybe SID -> Bool -> SubscriptionMeta
SubscriptionMeta SID
subject Maybe SID
queueGroup Bool
isReply
  SubscriptionStore
-> SID
-> SubscriptionMeta
-> SubscribeConfig
-> (Maybe Msg -> IO ())
-> IO ()
register SubscriptionStore
store SID
sid SubscriptionMeta
meta SubscribeConfig
cfg Maybe Msg -> IO ()
callback
  ClientState -> QueueItem -> IO ()
enqueue ClientState
client (QueueItem -> IO ()) -> QueueItem -> IO ()
forall a b. (a -> b) -> a -> b
$
    Sub -> QueueItem
forall m. Transformer m => m -> QueueItem
QueueItem
      Sub.Sub
        { subject :: SID
Sub.subject = SID
subject
        , queueGroup :: Maybe SID
Sub.queueGroup = Maybe SID
queueGroup
        , sid :: SID
Sub.sid = SID
sid
        }
  Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when Bool
isReply (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$
    ClientState -> QueueItem -> IO ()
enqueue ClientState
client
      (Unsub -> QueueItem
forall m. Transformer m => m -> QueueItem
QueueItem
        Unsub.Unsub
          { sid :: SID
Unsub.sid = SID
sid
          , maxMsg :: Maybe Int
Unsub.maxMsg = Int -> Maybe Int
forall a. a -> Maybe a
Just Int
1
          })
  SID -> IO SID
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure SID
sid

unsubscribeClient :: ClientState -> SubscriptionStore -> Msg.SID -> IO ()
unsubscribeClient :: ClientState -> SubscriptionStore -> SID -> IO ()
unsubscribeClient ClientState
client SubscriptionStore
store SID
sid = do
  ClientState -> AppM () -> IO ()
forall a. ClientState -> AppM a -> IO a
runClient ClientState
client (AppM () -> IO ()) -> AppM () -> IO ()
forall a b. (a -> b) -> a -> b
$
    LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Debug ([Char]
"unsubscribing SID: " [Char] -> [Char] -> [Char]
forall a. [a] -> [a] -> [a]
++ SID -> [Char]
forall a. Show a => a -> [Char]
show SID
sid)
  SubscriptionStore -> SID -> IO ()
unregister SubscriptionStore
store SID
sid
  ClientState -> QueueItem -> IO ()
enqueue ClientState
client (QueueItem -> IO ()) -> QueueItem -> IO ()
forall a b. (a -> b) -> a -> b
$
    Unsub -> QueueItem
forall m. Transformer m => m -> QueueItem
QueueItem
      Unsub.Unsub
        { sid :: SID
Unsub.sid = SID
sid
        , maxMsg :: Maybe Int
Unsub.maxMsg = Maybe Int
forall a. Maybe a
Nothing
        }

pingClient :: ClientState -> IO () -> IO ()
pingClient :: ClientState -> IO () -> IO ()
pingClient ClientState
client IO ()
action = do
  ClientState -> AppM () -> IO ()
forall a. ClientState -> AppM a -> IO a
runClient ClientState
client (AppM () -> IO ()) -> AppM () -> IO ()
forall a b. (a -> b) -> a -> b
$
    LogLevel -> [Char] -> AppM ()
forall (m :: * -> *). MonadLogger m => LogLevel -> [Char] -> m ()
logMessage LogLevel
Debug [Char]
"sending ping to server"
  ClientState -> IO () -> IO ()
pushPingAction ClientState
client IO ()
action
  ClientState -> QueueItem -> IO ()
enqueue ClientState
client (Ping -> QueueItem
forall m. Transformer m => m -> QueueItem
QueueItem Ping
Ping)

flushClient :: ClientState -> IO ()
flushClient :: ClientState -> IO ()
flushClient ClientState
client = do
  TMVar ()
ponged <- IO (TMVar ())
forall a. IO (TMVar a)
newEmptyTMVarIO
  ClientState -> IO () -> IO ()
pingClient ClientState
client (STM () -> IO ()
forall a. STM a -> IO a
atomically (STM Bool -> STM ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (TMVar () -> () -> STM Bool
forall a. TMVar a -> a -> STM Bool
tryPutTMVar TMVar ()
ponged ())))
  STM () -> IO ()
forall a. STM a -> IO a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$
    TMVar () -> STM ()
forall a. TMVar a -> STM a
readTMVar TMVar ()
ponged STM () -> STM () -> STM ()
forall a. STM a -> STM a -> STM a
`orElse` ClientState -> STM ()
waitForNotRunning ClientState
client