diff --git a/src/Simplex/Messaging/Agent.hs b/src/Simplex/Messaging/Agent.hs index dd163b255..50652d9f0 100644 --- a/src/Simplex/Messaging/Agent.hs +++ b/src/Simplex/Messaging/Agent.hs @@ -253,13 +253,15 @@ processSMPTransmission c@AgentClient {sndQ} st (srv, rId, cmd) = do logServer "<--" c srv rId "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 diff --git a/src/Simplex/Messaging/Agent/Store.hs b/src/Simplex/Messaging/Agent/Store.hs index 19c8c3f01..1aa90b5e9 100644 --- a/src/Simplex/Messaging/Agent/Store.hs +++ b/src/Simplex/Messaging/Agent/Store.hs @@ -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 diff --git a/src/Simplex/Messaging/Agent/Store/SQLite.hs b/src/Simplex/Messaging/Agent/Store/SQLite.hs index a10e19838..54586352d 100644 --- a/src/Simplex/Messaging/Agent/Store/SQLite.hs +++ b/src/Simplex/Messaging/Agent/Store/SQLite.hs @@ -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| diff --git a/tests/AgentTests/SQLiteTests.hs b/tests/AgentTests/SQLiteTests.hs index 0840da3fd..f46d27156 100644 --- a/tests/AgentTests/SQLiteTests.hs +++ b/tests/AgentTests/SQLiteTests.hs @@ -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