mirror of
https://github.com/simplex-chat/simplexmq.git
synced 2026-08-28 02:54:51 +00:00
Co-authored-by: Evgeny Poberezkin <2769109+epoberezkin@users.noreply.github.com>
This commit is contained in:
co-authored by
Evgeny Poberezkin
parent
633b3a4bda
commit
b7902ee4c8
@@ -190,7 +190,7 @@ processCommand c@AgentClient {sndQ} st (corrId, connAlias, cmd) =
|
||||
msgHash = C.sha256Hash msgStr
|
||||
withStore $
|
||||
createSndMsg st sq $
|
||||
SndMsgData {internalId, internalSndId, internalTs, msgBody, msgHash}
|
||||
SndMsgData {internalId, internalSndId, internalTs, msgBody, internalHash = msgHash}
|
||||
sendAgentMessage c sq msgStr
|
||||
respond $ SENT (unId internalId)
|
||||
|
||||
@@ -309,7 +309,8 @@ processSMPTransmission c@AgentClient {sndQ} st (srv, rId, cmd) = do
|
||||
senderMeta,
|
||||
brokerMeta,
|
||||
msgBody,
|
||||
msgHash,
|
||||
internalHash = msgHash,
|
||||
externalPrevSndHash = receivedPrevMsgHash,
|
||||
msgIntegrity
|
||||
}
|
||||
notify connAlias $
|
||||
|
||||
@@ -147,6 +147,7 @@ type PrevRcvMsgHash = MsgHash
|
||||
-- | Corresponds to `last_snd_msg_hash` in `connections` table
|
||||
type PrevSndMsgHash = MsgHash
|
||||
|
||||
-- ? merge/replace these with RcvMsg and SndMsg
|
||||
-- * Message data containers - used on Msg creation to reduce number of parameters
|
||||
|
||||
data RcvMsgData = RcvMsgData
|
||||
@@ -156,7 +157,8 @@ data RcvMsgData = RcvMsgData
|
||||
senderMeta :: (ExternalSndId, ExternalSndTs),
|
||||
brokerMeta :: (BrokerId, BrokerTs),
|
||||
msgBody :: MsgBody,
|
||||
msgHash :: MsgHash,
|
||||
internalHash :: MsgHash,
|
||||
externalPrevSndHash :: MsgHash,
|
||||
msgIntegrity :: MsgIntegrity
|
||||
}
|
||||
|
||||
@@ -165,7 +167,7 @@ data SndMsgData = SndMsgData
|
||||
internalSndId :: InternalSndId,
|
||||
internalTs :: InternalTs,
|
||||
msgBody :: MsgBody,
|
||||
msgHash :: MsgHash
|
||||
internalHash :: MsgHash
|
||||
}
|
||||
|
||||
-- * Message types
|
||||
@@ -194,7 +196,10 @@ data RcvMsg = RcvMsg
|
||||
-- | Timestamp of acknowledgement to sender, corresponds to `AcknowledgedToSender` status.
|
||||
-- Do not mix up with `externalSndTs` - timestamp created at sender before sending,
|
||||
-- which in its turn corresponds to `internalTs` in sending agent.
|
||||
ackSenderTs :: AckSenderTs
|
||||
ackSenderTs :: AckSenderTs,
|
||||
-- | Hash of previous message as received from sender - stored for integrity forensics.
|
||||
externalPrevSndHash :: MsgHash,
|
||||
msgIntegrity :: MsgIntegrity
|
||||
}
|
||||
deriving (Eq, Show)
|
||||
|
||||
@@ -254,7 +259,9 @@ data MsgBase = MsgBase
|
||||
-- due to a possibility of implementation errors in different agents.
|
||||
internalId :: InternalId,
|
||||
internalTs :: InternalTs,
|
||||
msgBody :: MsgBody
|
||||
msgBody :: MsgBody,
|
||||
-- | Hash of the message as computed by agent.
|
||||
internalHash :: MsgHash
|
||||
}
|
||||
deriving (Eq, Show)
|
||||
|
||||
|
||||
@@ -40,6 +40,7 @@ import Network.Socket (ServiceName)
|
||||
import Simplex.Messaging.Agent.Store
|
||||
import Simplex.Messaging.Agent.Store.SQLite.Schema (createSchema)
|
||||
import Simplex.Messaging.Agent.Transmission
|
||||
import Simplex.Messaging.Parsers (parseAll)
|
||||
import qualified Simplex.Messaging.Protocol as SMP
|
||||
import Simplex.Messaging.Util (bshow, liftIOEither)
|
||||
import System.Exit (ExitCode (ExitFailure), exitWith)
|
||||
@@ -276,6 +277,16 @@ instance ToField RcvMsgStatus where toField = toField . show
|
||||
|
||||
instance ToField SndMsgStatus where toField = toField . show
|
||||
|
||||
instance ToField MsgIntegrity where toField = toField . serializeMsgIntegrity
|
||||
|
||||
instance FromField MsgIntegrity where
|
||||
fromField = \case
|
||||
f@(Field (SQLBlob b) _) ->
|
||||
case parseAll msgIntegrityP b of
|
||||
Right k -> Ok k
|
||||
Left e -> returnError ConversionFailed f ("can't parse msg integrity field: " ++ e)
|
||||
f -> returnError ConversionFailed f "expecting SQLBlob column type"
|
||||
|
||||
fromFieldToReadable_ :: forall a. (Read a, E.Typeable a) => Field -> Ok a
|
||||
fromFieldToReadable_ = \case
|
||||
f@(Field (SQLText t) _) ->
|
||||
@@ -528,10 +539,12 @@ insertRcvMsgDetails_ dbConn connAlias RcvMsgData {..} =
|
||||
[sql|
|
||||
INSERT INTO rcv_messages
|
||||
( conn_alias, internal_rcv_id, internal_id, external_snd_id, external_snd_ts,
|
||||
broker_id, broker_ts, rcv_status, ack_brocker_ts, ack_sender_ts)
|
||||
broker_id, broker_ts, rcv_status, ack_brocker_ts, ack_sender_ts,
|
||||
internal_hash, external_prev_snd_hash, integrity)
|
||||
VALUES
|
||||
(:conn_alias,:internal_rcv_id,:internal_id,:external_snd_id,:external_snd_ts,
|
||||
:broker_id,:broker_ts,:rcv_status, NULL, NULL);
|
||||
:broker_id,:broker_ts,:rcv_status, NULL, NULL,
|
||||
:internal_hash,:external_prev_snd_hash,:integrity);
|
||||
|]
|
||||
[ ":conn_alias" := connAlias,
|
||||
":internal_rcv_id" := internalRcvId,
|
||||
@@ -540,7 +553,10 @@ insertRcvMsgDetails_ dbConn connAlias RcvMsgData {..} =
|
||||
":external_snd_ts" := snd senderMeta,
|
||||
":broker_id" := fst brokerMeta,
|
||||
":broker_ts" := snd brokerMeta,
|
||||
":rcv_status" := Received
|
||||
":rcv_status" := Received,
|
||||
":internal_hash" := internalHash,
|
||||
":external_prev_snd_hash" := externalPrevSndHash,
|
||||
":integrity" := msgIntegrity
|
||||
]
|
||||
|
||||
updateHashRcv_ :: DB.Connection -> ConnAlias -> RcvMsgData -> IO ()
|
||||
@@ -556,7 +572,7 @@ updateHashRcv_ dbConn connAlias RcvMsgData {..} =
|
||||
AND last_internal_rcv_msg_id = :last_internal_rcv_msg_id;
|
||||
|]
|
||||
[ ":last_external_snd_msg_id" := fst senderMeta,
|
||||
":last_rcv_msg_hash" := msgHash,
|
||||
":last_rcv_msg_hash" := internalHash,
|
||||
":conn_alias" := connAlias,
|
||||
":last_internal_rcv_msg_id" := internalRcvId
|
||||
]
|
||||
@@ -616,14 +632,15 @@ insertSndMsgDetails_ dbConn connAlias SndMsgData {..} =
|
||||
dbConn
|
||||
[sql|
|
||||
INSERT INTO snd_messages
|
||||
( conn_alias, internal_snd_id, internal_id, snd_status, sent_ts, delivered_ts)
|
||||
( conn_alias, internal_snd_id, internal_id, snd_status, sent_ts, delivered_ts, internal_hash)
|
||||
VALUES
|
||||
(:conn_alias,:internal_snd_id,:internal_id,:snd_status, NULL, NULL);
|
||||
(:conn_alias,:internal_snd_id,:internal_id,:snd_status, NULL, NULL,:internal_hash);
|
||||
|]
|
||||
[ ":conn_alias" := connAlias,
|
||||
":internal_snd_id" := internalSndId,
|
||||
":internal_id" := internalId,
|
||||
":snd_status" := Created
|
||||
":snd_status" := Created,
|
||||
":internal_hash" := internalHash
|
||||
]
|
||||
|
||||
updateHashSnd_ :: DB.Connection -> ConnAlias -> SndMsgData -> IO ()
|
||||
@@ -637,7 +654,7 @@ updateHashSnd_ dbConn connAlias SndMsgData {..} =
|
||||
WHERE conn_alias = :conn_alias
|
||||
AND last_internal_snd_msg_id = :last_internal_snd_msg_id;
|
||||
|]
|
||||
[ ":last_snd_msg_hash" := msgHash,
|
||||
[ ":last_snd_msg_hash" := internalHash,
|
||||
":conn_alias" := connAlias,
|
||||
":last_internal_snd_msg_id" := internalSndId
|
||||
]
|
||||
|
||||
@@ -138,6 +138,9 @@ rcvMessages =
|
||||
rcv_status TEXT NOT NULL,
|
||||
ack_brocker_ts TEXT,
|
||||
ack_sender_ts TEXT,
|
||||
internal_hash BLOB NOT NULL,
|
||||
external_prev_snd_hash BLOB NOT NULL,
|
||||
integrity BLOB NOT NULL,
|
||||
PRIMARY KEY (conn_alias, internal_rcv_id),
|
||||
FOREIGN KEY (conn_alias, internal_id)
|
||||
REFERENCES messages (conn_alias, internal_id)
|
||||
@@ -155,6 +158,7 @@ sndMessages =
|
||||
snd_status TEXT NOT NULL,
|
||||
sent_ts TEXT,
|
||||
delivered_ts TEXT,
|
||||
internal_hash BLOB NOT NULL,
|
||||
PRIMARY KEY (conn_alias, internal_snd_id),
|
||||
FOREIGN KEY (conn_alias, internal_id)
|
||||
REFERENCES messages (conn_alias, internal_id)
|
||||
|
||||
@@ -315,7 +315,7 @@ commandP =
|
||||
sendCmd = ACmd SClient . SEND <$> A.takeByteString
|
||||
sentResp = ACmd SAgent . SENT <$> A.decimal
|
||||
message = do
|
||||
msgIntegrity <- integrity <* A.space
|
||||
msgIntegrity <- msgIntegrityP <* A.space
|
||||
recipientMeta <- "R=" *> partyMeta A.decimal
|
||||
brokerMeta <- "B=" *> partyMeta base64P
|
||||
senderMeta <- "S=" *> partyMeta A.decimal
|
||||
@@ -326,13 +326,16 @@ commandP =
|
||||
<|> A.space *> (ReplyVia <$> smpServerP)
|
||||
<|> pure ReplyOn
|
||||
partyMeta idParser = (,) <$> idParser <* "," <*> tsISO8601P <* A.space
|
||||
integrity = "OK" $> MsgOk <|> "ERR " *> (MsgError <$> msgErrorType)
|
||||
agentError = ACmd SAgent . ERR <$> agentErrorTypeP
|
||||
|
||||
msgIntegrityP :: Parser MsgIntegrity
|
||||
msgIntegrityP = "OK" $> MsgOk <|> "ERR " *> (MsgError <$> msgErrorType)
|
||||
where
|
||||
msgErrorType =
|
||||
"ID " *> (MsgBadId <$> A.decimal)
|
||||
<|> "IDS " *> (MsgSkipped <$> A.decimal <* A.space <*> A.decimal)
|
||||
<|> "HASH" $> MsgBadHash
|
||||
<|> "DUPLICATE" $> MsgDuplicate
|
||||
agentError = ACmd SAgent . ERR <$> agentErrorTypeP
|
||||
|
||||
parseCommand :: ByteString -> Either AgentErrorType ACmd
|
||||
parseCommand = parse commandP $ CMD SYNTAX
|
||||
@@ -369,16 +372,17 @@ serializeCommand = \case
|
||||
ReplyOn -> ""
|
||||
showTs :: UTCTime -> ByteString
|
||||
showTs = B.pack . formatISO8601Millis
|
||||
serializeMsgIntegrity :: MsgIntegrity -> ByteString
|
||||
serializeMsgIntegrity = \case
|
||||
MsgOk -> "OK"
|
||||
MsgError e ->
|
||||
"ERR " <> case e of
|
||||
MsgSkipped fromMsgId toMsgId ->
|
||||
B.unwords ["NO_ID", bshow fromMsgId, bshow toMsgId]
|
||||
MsgBadId aMsgId -> "ID " <> bshow aMsgId
|
||||
MsgBadHash -> "HASH"
|
||||
MsgDuplicate -> "DUPLICATE"
|
||||
|
||||
serializeMsgIntegrity :: MsgIntegrity -> ByteString
|
||||
serializeMsgIntegrity = \case
|
||||
MsgOk -> "OK"
|
||||
MsgError e ->
|
||||
"ERR " <> case e of
|
||||
MsgSkipped fromMsgId toMsgId ->
|
||||
B.unwords ["NO_ID", bshow fromMsgId, bshow toMsgId]
|
||||
MsgBadId aMsgId -> "ID " <> bshow aMsgId
|
||||
MsgBadHash -> "HASH"
|
||||
MsgDuplicate -> "DUPLICATE"
|
||||
|
||||
agentErrorTypeP :: Parser AgentErrorType
|
||||
agentErrorTypeP =
|
||||
|
||||
@@ -79,8 +79,8 @@ import Data.ByteString.Lazy (fromStrict, toStrict)
|
||||
import Data.String
|
||||
import Data.Typeable (Typeable)
|
||||
import Data.X509
|
||||
import Database.SQLite.Simple as DB
|
||||
import Database.SQLite.Simple.FromField
|
||||
import Database.SQLite.Simple (ResultError (..), SQLData (..))
|
||||
import Database.SQLite.Simple.FromField (FieldParser, FromField (..), returnError)
|
||||
import Database.SQLite.Simple.Internal (Field (..))
|
||||
import Database.SQLite.Simple.Ok (Ok (Ok))
|
||||
import Database.SQLite.Simple.ToField (ToField (..))
|
||||
|
||||
@@ -323,7 +323,7 @@ ts :: UTCTime
|
||||
ts = UTCTime (fromGregorian 2021 02 24) (secondsToDiffTime 0)
|
||||
|
||||
mkRcvMsgData :: InternalId -> InternalRcvId -> ExternalSndId -> BrokerId -> MsgHash -> RcvMsgData
|
||||
mkRcvMsgData internalId internalRcvId externalSndId brokerId msgHash =
|
||||
mkRcvMsgData internalId internalRcvId externalSndId brokerId internalHash =
|
||||
RcvMsgData
|
||||
{ internalId,
|
||||
internalRcvId,
|
||||
@@ -331,7 +331,8 @@ mkRcvMsgData internalId internalRcvId externalSndId brokerId msgHash =
|
||||
senderMeta = (externalSndId, ts),
|
||||
brokerMeta = (brokerId, ts),
|
||||
msgBody = hw,
|
||||
msgHash = msgHash,
|
||||
internalHash,
|
||||
externalPrevSndHash = "hash_from_sender",
|
||||
msgIntegrity = MsgOk
|
||||
}
|
||||
|
||||
@@ -351,13 +352,13 @@ testCreateRcvMsg = do
|
||||
testCreateRcvMsg' store 1 "hash_dummy" rcvQueue1 $ mkRcvMsgData (InternalId 2) (InternalRcvId 2) 2 "2" "new_hash_dummy"
|
||||
|
||||
mkSndMsgData :: InternalId -> InternalSndId -> MsgHash -> SndMsgData
|
||||
mkSndMsgData internalId internalSndId msgHash =
|
||||
mkSndMsgData internalId internalSndId internalHash =
|
||||
SndMsgData
|
||||
{ internalId,
|
||||
internalSndId,
|
||||
internalTs = ts,
|
||||
msgBody = hw,
|
||||
msgHash = msgHash
|
||||
internalHash
|
||||
}
|
||||
|
||||
testCreateSndMsg' :: SQLiteStore -> PrevSndMsgHash -> SndQueue -> SndMsgData -> Expectation
|
||||
|
||||
Reference in New Issue
Block a user