agent: commit ratchet key hash with resync state

This commit is contained in:
shum
2026-10-02 11:02:05 +00:00
parent 053e83b704
commit 90bec26ada
11 changed files with 163 additions and 29 deletions
+2
View File
@@ -195,6 +195,7 @@ library
Simplex.Messaging.Agent.Store.Postgres.Migrations.M20260712_address_dr_rpc
Simplex.Messaging.Agent.Store.Postgres.Migrations.M20260823_snd_files_entitlement
Simplex.Messaging.Agent.Store.Postgres.Migrations.M20260919_ratchet_verify_codes
Simplex.Messaging.Agent.Store.Postgres.Migrations.M20261001_ratchet_key_hashes_unique
else
exposed-modules:
Simplex.Messaging.Agent.Store.SQLite
@@ -250,6 +251,7 @@ library
Simplex.Messaging.Agent.Store.SQLite.Migrations.M20260712_address_dr_rpc
Simplex.Messaging.Agent.Store.SQLite.Migrations.M20260823_snd_files_entitlement
Simplex.Messaging.Agent.Store.SQLite.Migrations.M20260919_ratchet_verify_codes
Simplex.Messaging.Agent.Store.SQLite.Migrations.M20261001_ratchet_key_hashes_unique
Simplex.Messaging.Agent.Store.SQLite.Util
if flag(client_postgres) || flag(server_postgres)
exposed-modules:
+17 -17
View File
@@ -4150,10 +4150,7 @@ processSMPTransmissions c@AgentClient {subQ} (tSess@(userId, srv, _), THandlePar
duplicateRequest = case e2eSndParams of
CR.AE2ERatchetParams _ (CR.E2ERatchetParams _ k1 k2 _) -> do
let rkHash = C.sha256Hash $ C.pubKeyBytes k1 <> C.pubKeyBytes k2
withStore' c $ \db -> do
exists <- checkRatchetKeyHashExists db connId rkHash
unless exists $ addProcessedRatchetKeyHash db connId rkHash
pure exists
withStore' c $ \db -> not <$> addProcessedRatchetKeyHash db connId rkHash
qDuplex :: Connection c -> String -> (Connection 'CDuplex -> AM a) -> AM a
qDuplex conn' name action = case conn' of
@@ -4167,16 +4164,18 @@ processSMPTransmissions c@AgentClient {subQ} (tSess@(userId, srv, _), THandlePar
unless (e2eVersion `isCompatible` e2eEncryptVRange) (throwE $ AGENT A_VERSION)
keys <- getSendRatchetKeys
let rcVs = CR.RatchetVersions {current = e2eVersion, maxSupported = maxVersion e2eEncryptVRange}
initRatchet rcVs keys
notifyAgreed
whenM (initRatchet rcVs keys) notifyAgreed
where
rkHashRcv = rkHash k1Rcv k2Rcv
rkHash k1 k2 = C.sha256Hash $ C.pubKeyBytes k1 <> C.pubKeyBytes k2
ratchetExists :: AM Bool
ratchetExists = withStore' c $ \db -> do
exists <- checkRatchetKeyHashExists db connId rkHashRcv
unless exists $ addProcessedRatchetKeyHash db connId rkHashRcv
pure exists
ratchetExists = withStore' c $ \db -> checkRatchetKeyHashExists db connId rkHashRcv
-- the hash is committed with the sync state change, so a key whose processing failed is not treated as a replay
markKeyProcessed :: DB.Connection -> IO () -> IO Bool
markKeyProcessed db changeState = do
added <- addProcessedRatchetKeyHash db connId rkHashRcv
when added changeState
pure added
getSendRatchetKeys :: AM (CR.RcvE2EPrivRatchetParams 'C.X448)
getSendRatchetKeys = case rss of
RSOk -> sendReplyKey -- receiving client
@@ -4184,8 +4183,8 @@ processSMPTransmissions c@AgentClient {subQ} (tSess@(userId, srv, _), THandlePar
RSRequired -> sendReplyKey
RSStarted -> withStore c (`getRatchetX3dhKeys` connId) -- initiating client
RSAgreed -> do
withStore' c $ \db -> setConnRatchetSync db connId RSRequired
notifyRatchetSyncError
marked <- withStore' c $ \db -> markKeyProcessed db $ setConnRatchetSync db connId RSRequired
when marked notifyRatchetSyncError
-- can communicate for other client to reset to RSRequired
-- - need to add new AgentMsgEnvelope, AgentMessage, AgentMessageType
-- - need to deduplicate on receiving side
@@ -4207,14 +4206,14 @@ processSMPTransmissions c@AgentClient {subQ} (tSess@(userId, srv, _), THandlePar
conn'' = updateConnection cData'' conn'
cStats <- connectionStats c conn''
notify $ RSYNC RSAgreed Nothing cStats
recreateRatchet :: CR.Ratchet 'C.X448 -> AM ()
recreateRatchet rc = withStore' c $ \db -> do
recreateRatchet :: CR.Ratchet 'C.X448 -> AM Bool
recreateRatchet rc = withStore' c $ \db -> markKeyProcessed db $ do
setConnRatchetSync db connId RSAgreed
deleteRatchet db connId
createRatchet db connId rc
-- compare public keys `k1` in AgentRatchetKey messages sent by self and other party
-- to determine ratchet initilization ordering
initRatchet :: CR.RatchetVersions -> CR.RcvE2EPrivRatchetParams 'C.X448 -> AM ()
initRatchet :: CR.RatchetVersions -> CR.RcvE2EPrivRatchetParams 'C.X448 -> AM Bool
initRatchet rcVs (pk1, pk2, pKem)
| rkHash (C.publicKey pk1) (C.publicKey pk2) <= rkHashRcv = do
rcParams <- liftError cryptoError $ CR.pqX3dhRcv (pk1, pk2, pKem) e2eOtherPartyParams
@@ -4222,8 +4221,9 @@ processSMPTransmissions c@AgentClient {subQ} (tSess@(userId, srv, _), THandlePar
| otherwise = do
(_, rcDHRs) <- atomically . C.generateKeyPair =<< asks random
rcParams <- liftEitherWith cryptoError $ CR.pqX3dhSnd (pk1, pk2, CR.APRKP CR.SRKSProposed <$> pKem) e2eOtherPartyParams
recreateRatchet $ CR.initSndRatchet rcVs k2Rcv rcDHRs rcParams
void . enqueueMessages' c cData' sqs SMP.MsgFlags {notification = True} $ EREADY lastExternalSndId
recreated <- recreateRatchet $ CR.initSndRatchet rcVs k2Rcv rcDHRs rcParams
when recreated $ void . enqueueMessages' c cData' sqs SMP.MsgFlags {notification = True} $ EREADY lastExternalSndId
pure recreated
checkMsgIntegrity :: PrevExternalSndId -> ExternalSndId -> PrevRcvMsgHash -> ByteString -> MsgIntegrity
checkMsgIntegrity prevExtSndId extSndId internalPrevMsgHash receivedPrevMsgHash
@@ -2797,9 +2797,15 @@ setConnRatchetSync :: DB.Connection -> ConnId -> RatchetSyncState -> IO ()
setConnRatchetSync db connId ratchetSyncState =
DB.execute db "UPDATE connections SET ratchet_sync_state = ? WHERE conn_id = ?" (ratchetSyncState, connId)
addProcessedRatchetKeyHash :: DB.Connection -> ConnId -> ByteString -> IO ()
addProcessedRatchetKeyHash db connId hash =
DB.execute db "INSERT INTO processed_ratchet_key_hashes (conn_id, hash) VALUES (?,?)" (connId, Binary hash)
-- | Returns False if the hash was already processed for this connection.
addProcessedRatchetKeyHash :: DB.Connection -> ConnId -> ByteString -> IO Bool
addProcessedRatchetKeyHash db connId hash = do
(rs :: [Only Int]) <-
DB.query
db
"INSERT INTO processed_ratchet_key_hashes (conn_id, hash) VALUES (?,?) ON CONFLICT (conn_id, hash) DO NOTHING RETURNING 1"
(connId, Binary hash)
pure $ not $ null rs
checkRatchetKeyHashExists :: DB.Connection -> ConnId -> ByteString -> IO Bool
checkRatchetKeyHashExists db connId hash =
@@ -16,6 +16,7 @@ import Simplex.Messaging.Agent.Store.Postgres.Migrations.M20260411_service_certs
import Simplex.Messaging.Agent.Store.Postgres.Migrations.M20260712_address_dr_rpc
import Simplex.Messaging.Agent.Store.Postgres.Migrations.M20260823_snd_files_entitlement
import Simplex.Messaging.Agent.Store.Postgres.Migrations.M20260919_ratchet_verify_codes
import Simplex.Messaging.Agent.Store.Postgres.Migrations.M20261001_ratchet_key_hashes_unique
import Simplex.Messaging.Agent.Store.Shared (Migration (..))
schemaMigrations :: [(String, Text, Maybe Text)]
@@ -31,7 +32,8 @@ schemaMigrations =
("20260411_service_certs", m20260411_service_certs, Just down_m20260411_service_certs),
("20260712_address_dr_rpc", m20260712_address_dr_rpc, Just down_m20260712_address_dr_rpc),
("20260823_snd_files_entitlement", m20260823_snd_files_entitlement, Just down_m20260823_snd_files_entitlement),
("20260919_ratchet_verify_codes", m20260919_ratchet_verify_codes, Just down_m20260919_ratchet_verify_codes)
("20260919_ratchet_verify_codes", m20260919_ratchet_verify_codes, Just down_m20260919_ratchet_verify_codes),
("20261001_ratchet_key_hashes_unique", m20261001_ratchet_key_hashes_unique, Just down_m20261001_ratchet_key_hashes_unique)
]
-- | The list of migrations in ascending order by date
@@ -0,0 +1,28 @@
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Simplex.Messaging.Agent.Store.Postgres.Migrations.M20261001_ratchet_key_hashes_unique where
import Data.Text (Text)
import Text.RawString.QQ (r)
m20261001_ratchet_key_hashes_unique :: Text
m20261001_ratchet_key_hashes_unique =
[r|
DELETE FROM processed_ratchet_key_hashes
WHERE processed_ratchet_key_hash_id NOT IN (
SELECT MIN(processed_ratchet_key_hash_id)
FROM processed_ratchet_key_hashes
GROUP BY conn_id, hash
);
DROP INDEX idx_processed_ratchet_key_hashes_hash;
CREATE UNIQUE INDEX idx_processed_ratchet_key_hashes_hash ON processed_ratchet_key_hashes(conn_id, hash);
|]
down_m20261001_ratchet_key_hashes_unique :: Text
down_m20261001_ratchet_key_hashes_unique =
[r|
DROP INDEX idx_processed_ratchet_key_hashes_hash;
CREATE INDEX idx_processed_ratchet_key_hashes_hash ON processed_ratchet_key_hashes(conn_id, hash);
|]
@@ -1173,7 +1173,7 @@ CREATE INDEX idx_processed_ratchet_key_hashes_created_at ON smp_agent_test_proto
CREATE INDEX idx_processed_ratchet_key_hashes_hash ON smp_agent_test_protocol_schema.processed_ratchet_key_hashes USING btree (conn_id, hash);
CREATE UNIQUE INDEX idx_processed_ratchet_key_hashes_hash ON smp_agent_test_protocol_schema.processed_ratchet_key_hashes USING btree (conn_id, hash);
@@ -52,6 +52,7 @@ import Simplex.Messaging.Agent.Store.SQLite.Migrations.M20260411_service_certs
import Simplex.Messaging.Agent.Store.SQLite.Migrations.M20260712_address_dr_rpc
import Simplex.Messaging.Agent.Store.SQLite.Migrations.M20260823_snd_files_entitlement
import Simplex.Messaging.Agent.Store.SQLite.Migrations.M20260919_ratchet_verify_codes
import Simplex.Messaging.Agent.Store.SQLite.Migrations.M20261001_ratchet_key_hashes_unique
import Simplex.Messaging.Agent.Store.Shared (Migration (..))
schemaMigrations :: [(String, Query, Maybe Query)]
@@ -103,7 +104,8 @@ schemaMigrations =
("m20260411_service_certs", m20260411_service_certs, Just down_m20260411_service_certs),
("m20260712_address_dr_rpc", m20260712_address_dr_rpc, Just down_m20260712_address_dr_rpc),
("m20260823_snd_files_entitlement", m20260823_snd_files_entitlement, Just down_m20260823_snd_files_entitlement),
("m20260919_ratchet_verify_codes", m20260919_ratchet_verify_codes, Just down_m20260919_ratchet_verify_codes)
("m20260919_ratchet_verify_codes", m20260919_ratchet_verify_codes, Just down_m20260919_ratchet_verify_codes),
("m20261001_ratchet_key_hashes_unique", m20261001_ratchet_key_hashes_unique, Just down_m20261001_ratchet_key_hashes_unique)
]
-- | The list of migrations in ascending order by date
@@ -0,0 +1,27 @@
{-# LANGUAGE QuasiQuotes #-}
module Simplex.Messaging.Agent.Store.SQLite.Migrations.M20261001_ratchet_key_hashes_unique where
import Database.SQLite.Simple (Query)
import Database.SQLite.Simple.QQ (sql)
m20261001_ratchet_key_hashes_unique :: Query
m20261001_ratchet_key_hashes_unique =
[sql|
DELETE FROM processed_ratchet_key_hashes
WHERE processed_ratchet_key_hash_id NOT IN (
SELECT MIN(processed_ratchet_key_hash_id)
FROM processed_ratchet_key_hashes
GROUP BY conn_id, hash
);
DROP INDEX idx_processed_ratchet_key_hashes_hash;
CREATE UNIQUE INDEX idx_processed_ratchet_key_hashes_hash ON processed_ratchet_key_hashes(conn_id, hash);
|]
down_m20261001_ratchet_key_hashes_unique :: Query
down_m20261001_ratchet_key_hashes_unique =
[sql|
DROP INDEX idx_processed_ratchet_key_hashes_hash;
CREATE INDEX idx_processed_ratchet_key_hashes_hash ON processed_ratchet_key_hashes(conn_id, hash);
|]
@@ -573,10 +573,6 @@ CREATE INDEX idx_encrypted_rcv_message_hashes_hash ON encrypted_rcv_message_hash
conn_id,
hash
);
CREATE INDEX idx_processed_ratchet_key_hashes_hash ON processed_ratchet_key_hashes(
conn_id,
hash
);
CREATE INDEX idx_snd_messages_rcpt_internal_id ON snd_messages(
conn_id,
rcpt_internal_id
@@ -641,6 +637,10 @@ CREATE INDEX idx_connections_deleted ON connections(deleted);
CREATE INDEX idx_connections_service_request_expires_at ON connections(
service_request_expires_at
);
CREATE UNIQUE INDEX idx_processed_ratchet_key_hashes_hash ON processed_ratchet_key_hashes(
conn_id,
hash
);
CREATE TRIGGER tr_rcv_queue_insert
AFTER INSERT ON rcv_queues
FOR EACH ROW
+51 -1
View File
@@ -92,7 +92,7 @@ import Simplex.Messaging.Agent.Env.SQLite (AgentConfig (..), Env (..), InitialAg
import Simplex.Messaging.Agent.Protocol hiding (CON, CONF, INFO, REQ, SENT)
import qualified Simplex.Messaging.Agent.Protocol as A
import Simplex.Messaging.Agent.Store (Connection' (..), SomeConn' (..), StoredRcvQueue (..))
import Simplex.Messaging.Agent.Store.AgentStore (getConn)
import Simplex.Messaging.Agent.Store.AgentStore (addProcessedRatchetKeyHash, checkRatchetKeyHashExists, getConn)
import Simplex.Messaging.Agent.Store.Common (DBStore (..), withTransaction)
import Simplex.Messaging.Agent.Store.Interface
import qualified Simplex.Messaging.Agent.Store.DB as DB
@@ -459,6 +459,8 @@ functionalAPITests ps = do
testRatchetSyncSuspendForeground ps
it "should synchronize ratchets when clients start synchronization simultaneously" $
testRatchetSyncSimultaneous ps
it "should not mark ratchet key as processed when ratchet recreation fails" $
testRatchetSyncFailedKeyNotProcessed ps
#endif
describe "Subscription mode OnlyCreate" $ do
it "messages delivered only when polled" $
@@ -587,6 +589,9 @@ functionalAPITests ps = do
describe "getConnectionVerifyCodes" $
it "should return the same codes for both peers" $
withSmpServer ps testConnectionVerifyCodes
describe "processed ratchet key hashes" $
it "should add ratchet key hash once per connection" $
withSmpServer ps testAddProcessedRatchetKeyHash
describe "Delivery receipts" $ do
it "should send and receive delivery receipt" $ withSmpServer ps testDeliveryReceipts
it "send delivery receipts concurrently with messages" $ testDeliveryReceiptsConcurrent ps
@@ -2744,6 +2749,34 @@ testRatchetSyncSimultaneous ps = do
disposeAgentClient bob
disposeAgentClient bob2
testRatchetSyncFailedKeyNotProcessed :: HasCallStack => (ASrvTransport, AStoreType) -> IO ()
testRatchetSyncFailedKeyNotProcessed ps = withAgentClients2 $ \alice bob -> do
(aliceId, bobId, bob2) <- withSmpServerStoreMsgLogOn ps testPort $ \_ ->
setupDesynchronizedRatchet alice bob
("", "", DOWN _ _) <- nGet alice
("", "", DOWN _ _) <- nGet bob2
ConnectionStats {ratchetSyncState} <- runRight $ synchronizeRatchet bob2 aliceId PQSupportOn False
ratchetSyncState `shouldBe` RSStarted
withTransaction (store $ agentEnv bob2) $ \db ->
DB.execute_ db "UPDATE ratchets SET x3dh_priv_key_1 = NULL"
withSmpServerStoreMsgLogOn ps testPort $ \_ ->
concurrently_
(getInAnyOrder alice [ratchetSyncP' bobId RSAgreed, serverUpP])
(getInAnyOrder bob2 [x3dhKeysNotFoundP aliceId, serverUpP])
map fst <$> processedRatchetKeyHashes alice `shouldReturn` [bobId]
processedRatchetKeyHashes bob2 `shouldReturn` []
disposeAgentClient bob2
where
x3dhKeysNotFoundP :: ConnId -> ATransmission -> Bool
x3dhKeysNotFoundP cId = \case
(_, cId', AEvt SAEConn (ERR (A.INTERNAL e))) -> cId' == cId && "SEX3dhKeysNotFound" `isPrefixOf` e
_ -> False
processedRatchetKeyHashes :: AgentClient -> IO [(ConnId, ByteString)]
processedRatchetKeyHashes c = withTransaction (store $ agentEnv c) (`DB.query_` "SELECT conn_id, hash FROM processed_ratchet_key_hashes")
getMsg :: AgentClient -> ConnId -> ExceptT AgentErrorType IO a -> ExceptT AgentErrorType IO a
getMsg c cId action = do
liftIO $ noMessages c "nothing should be delivered before GET"
@@ -4087,6 +4120,23 @@ testConnectionVerifyCodes =
codes1'' <- getConnectionVerifyCodes a bId
liftIO $ codes1'' `shouldBe` codes1
testAddProcessedRatchetKeyHash :: HasCallStack => IO ()
testAddProcessedRatchetKeyHash =
withAgentClients2 $ \a b -> do
(bId1, bId2) <- runRight $ do
(_, bId1) <- makeConnection a b
(_, bId2) <- makeConnection a b
pure (bId1, bId2)
withTransaction (store $ agentEnv a) $ \db -> do
checkRatchetKeyHashExists db bId1 "hash1" `shouldReturn` False
addProcessedRatchetKeyHash db bId1 "hash1" `shouldReturn` True
addProcessedRatchetKeyHash db bId1 "hash1" `shouldReturn` False
addProcessedRatchetKeyHash db bId1 "hash2" `shouldReturn` True
addProcessedRatchetKeyHash db bId2 "hash1" `shouldReturn` True
checkRatchetKeyHashExists db bId1 "hash1" `shouldReturn` True
DB.query_ db "SELECT conn_id, hash FROM processed_ratchet_key_hashes ORDER BY processed_ratchet_key_hash_id"
`shouldReturn` [(bId1, "hash1" :: ByteString), (bId1, "hash2"), (bId2, "hash1")]
testDeliveryReceipts :: HasCallStack => IO ()
testDeliveryReceipts =
withAgentClients2 $ \a b -> runRight_ $ do
+18 -1
View File
@@ -48,6 +48,7 @@ schemaDumpTest = do
it "verify strict tables" testVerifyStrict
it "should NOT create user record for new database" testUsersMigrationNew
it "should create user record for old database" testUsersMigrationOld
it "should remove duplicate ratchet key hashes before adding unique index" testRatchetKeyHashesUniqueMigration
testVerifySchemaDump :: IO ()
testVerifySchemaDump = do
@@ -114,12 +115,28 @@ testUsersMigrationOld = do
`shouldReturn` ([Only (1 :: Int)])
closeDBStore st'
testRatchetKeyHashesUniqueMigration :: IO ()
testRatchetKeyHashesUniqueMigration = do
let beforeUnique = takeWhile (("m20261001_ratchet_key_hashes_unique" /=) . name) appMigrations
Right st <- createDBStore (DBOpts testDB [] "" False True TQOff) beforeUnique (MigrationConfig MCError Nothing)
withTransaction' st $ \db -> do
SQL.execute_ db "INSERT INTO users (user_id) VALUES (1)"
SQL.execute_ db "INSERT INTO connections (conn_id, conn_mode, user_id) VALUES (x'01', 'INV', 1), (x'02', 'INV', 1)"
SQL.execute_ db "INSERT INTO processed_ratchet_key_hashes (conn_id, hash) VALUES (x'01', x'aa'), (x'01', x'aa'), (x'01', x'bb'), (x'02', x'aa'), (x'01', x'aa')"
closeDBStore st
Right st' <- createDBStore (DBOpts testDB [] "" False True TQOff) appMigrations (MigrationConfig MCYesUp Nothing)
withTransaction' st' (`SQL.query_` "SELECT processed_ratchet_key_hash_id FROM processed_ratchet_key_hashes ORDER BY processed_ratchet_key_hash_id")
`shouldReturn` [Only (1 :: Int), Only 3, Only 4]
closeDBStore st'
skipComparisonForDownMigrations :: [String]
skipComparisonForDownMigrations =
[ -- on down migration idx_messages_internal_snd_id_ts index moves down to the end of the file
"m20230814_indexes",
-- snd_secure and last_broker_ts columns swap order on down migration
"m20250322_short_links"
"m20250322_short_links",
-- on down migration idx_processed_ratchet_key_hashes_hash index moves down to the end of the file
"m20261001_ratchet_key_hashes_unique"
]
getSchema :: FilePath -> FilePath -> IO String