{-# LANGUAGE OverloadedStrings #-}
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