mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-08-28 20:08:16 +00:00
agent store: organize sender and broker fields into tuples (#67)
This commit is contained in:
@@ -253,13 +253,15 @@ processSMPTransmission c@AgentClient {sndQ} st (srv, rId, cmd) = do
|
||||
logServer "<--" c srv rId "MSG <MSG>"
|
||||
-- TODO check message status
|
||||
recipientTs <- liftIO getCurrentTime
|
||||
recipientId <- withStore $ createRcvMsg st connAlias body recipientTs senderMsgId senderTimestamp srvMsgId srvTs
|
||||
let m_sender = (senderMsgId, senderTimestamp)
|
||||
let m_broker = (srvMsgId, srvTs)
|
||||
recipientId <- withStore $ createRcvMsg st connAlias body recipientTs m_sender m_broker
|
||||
notify connAlias $
|
||||
MSG
|
||||
{ m_status = MsgOk,
|
||||
m_recipient = (recipientId, recipientTs),
|
||||
m_sender = (senderMsgId, senderTimestamp),
|
||||
m_broker = (srvMsgId, srvTs),
|
||||
m_sender,
|
||||
m_broker,
|
||||
m_body = body
|
||||
}
|
||||
sendAck c rq
|
||||
|
||||
@@ -40,7 +40,7 @@ class Monad m => MonadAgentStore s m where
|
||||
setSndQueueStatus :: s -> SndQueue -> QueueStatus -> m ()
|
||||
|
||||
-- Msg management
|
||||
createRcvMsg :: s -> ConnAlias -> MsgBody -> InternalTs -> ExternalSndId -> ExternalSndTs -> BrokerId -> BrokerTs -> m InternalId
|
||||
createRcvMsg :: s -> ConnAlias -> MsgBody -> InternalTs -> (ExternalSndId, ExternalSndTs) -> (BrokerId, BrokerTs) -> m InternalId
|
||||
createSndMsg :: s -> ConnAlias -> MsgBody -> InternalTs -> m InternalId
|
||||
getMsg :: s -> ConnAlias -> InternalId -> m Msg
|
||||
|
||||
|
||||
@@ -126,10 +126,10 @@ instance (MonadUnliftIO m, MonadError StoreError m) => MonadAgentStore SQLiteSto
|
||||
liftIO $
|
||||
updateSndQueueStatus dbConn sndQueue status
|
||||
|
||||
createRcvMsg :: SQLiteStore -> ConnAlias -> MsgBody -> InternalTs -> ExternalSndId -> ExternalSndTs -> BrokerId -> BrokerTs -> m InternalId
|
||||
createRcvMsg SQLiteStore {dbConn} connAlias msgBody internalTs externalSndId externalSndTs brokerId brokerTs =
|
||||
createRcvMsg :: SQLiteStore -> ConnAlias -> MsgBody -> InternalTs -> (ExternalSndId, ExternalSndTs) -> (BrokerId, BrokerTs) -> m InternalId
|
||||
createRcvMsg SQLiteStore {dbConn} connAlias msgBody internalTs (externalSndId, externalSndTs) (brokerId, brokerTs) =
|
||||
liftIOEither $
|
||||
insertRcvMsg dbConn connAlias msgBody internalTs externalSndId externalSndTs brokerId brokerTs
|
||||
insertRcvMsg dbConn connAlias msgBody internalTs (externalSndId, externalSndTs) (brokerId, brokerTs)
|
||||
|
||||
createSndMsg :: SQLiteStore -> ConnAlias -> MsgBody -> InternalTs -> m InternalId
|
||||
createSndMsg SQLiteStore {dbConn} connAlias msgBody internalTs =
|
||||
@@ -463,12 +463,10 @@ insertRcvMsg ::
|
||||
ConnAlias ->
|
||||
MsgBody ->
|
||||
InternalTs ->
|
||||
ExternalSndId ->
|
||||
ExternalSndTs ->
|
||||
BrokerId ->
|
||||
BrokerTs ->
|
||||
(ExternalSndId, ExternalSndTs) ->
|
||||
(BrokerId, BrokerTs) ->
|
||||
IO (Either StoreError InternalId)
|
||||
insertRcvMsg dbConn connAlias msgBody internalTs externalSndId externalSndTs brokerId brokerTs =
|
||||
insertRcvMsg dbConn connAlias msgBody internalTs (externalSndId, externalSndTs) (brokerId, brokerTs) =
|
||||
DB.withTransaction dbConn $ do
|
||||
queues <- retrieveConnQueues_ dbConn connAlias
|
||||
case queues of
|
||||
@@ -477,7 +475,7 @@ insertRcvMsg dbConn connAlias msgBody internalTs externalSndId externalSndTs bro
|
||||
let internalId = lastInternalId + 1
|
||||
let internalRcvId = lastInternalRcvId + 1
|
||||
insertRcvMsgBase_ dbConn connAlias internalId internalTs internalRcvId msgBody
|
||||
insertRcvMsgDetails_ dbConn connAlias internalRcvId internalId externalSndId externalSndTs brokerId brokerTs
|
||||
insertRcvMsgDetails_ dbConn connAlias internalRcvId internalId (externalSndId, externalSndTs) (brokerId, brokerTs)
|
||||
updateLastInternalIdsRcv_ dbConn connAlias internalId internalRcvId
|
||||
return $ Right internalId
|
||||
(Nothing, Just _sndQ) -> return $ Left (SEBadConnType CSnd)
|
||||
@@ -518,12 +516,10 @@ insertRcvMsgDetails_ ::
|
||||
ConnAlias ->
|
||||
InternalRcvId ->
|
||||
InternalId ->
|
||||
ExternalSndId ->
|
||||
ExternalSndTs ->
|
||||
BrokerId ->
|
||||
BrokerTs ->
|
||||
(ExternalSndId, ExternalSndTs) ->
|
||||
(BrokerId, BrokerTs) ->
|
||||
IO ()
|
||||
insertRcvMsgDetails_ dbConn connAlias internalRcvId internalId externalSndId externalSndTs brokerId brokerTs =
|
||||
insertRcvMsgDetails_ dbConn connAlias internalRcvId internalId (externalSndId, externalSndTs) (brokerId, brokerTs) =
|
||||
DB.executeNamed
|
||||
dbConn
|
||||
[sql|
|
||||
|
||||
@@ -303,18 +303,18 @@ testCreateRcvMsg = do
|
||||
`returnsResult` ()
|
||||
-- TODO getMsg to check message
|
||||
let ts = UTCTime (fromGregorian 2021 02 24) (secondsToDiffTime 0)
|
||||
createRcvMsg store "conn1" (encodeUtf8 "Hello world!") ts 1 ts "1" ts
|
||||
createRcvMsg store "conn1" (encodeUtf8 "Hello world!") ts (1, ts) ("1", ts)
|
||||
`returnsResult` (1 :: InternalId)
|
||||
|
||||
testCreateRcvMsgNoQueue :: SpecWith SQLiteStore
|
||||
testCreateRcvMsgNoQueue = do
|
||||
it "should throw error on attempt to create a RcvMsg w/t a RcvQueue" $ \store -> do
|
||||
let ts = UTCTime (fromGregorian 2021 02 24) (secondsToDiffTime 0)
|
||||
createRcvMsg store "conn1" (encodeUtf8 "Hello world!") ts 1 ts "1" ts
|
||||
createRcvMsg store "conn1" (encodeUtf8 "Hello world!") ts (1, ts) ("1", ts)
|
||||
`throwsError` SEBadConn
|
||||
createSndConn store sndQueue1
|
||||
`returnsResult` ()
|
||||
createRcvMsg store "conn1" (encodeUtf8 "Hello world!") ts 1 ts "1" ts
|
||||
createRcvMsg store "conn1" (encodeUtf8 "Hello world!") ts (1, ts) ("1", ts)
|
||||
`throwsError` SEBadConnType CSnd
|
||||
|
||||
testCreateSndMsg :: SpecWith SQLiteStore
|
||||
|
||||
Reference in New Issue
Block a user