Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions changelog.d/5-internal/WPB-22954
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
Move mls commit locks from cassandra to postgresql advisory locks

Locks are now pg advisory locks held on a dedicated pooled connection for the duration of a commit; contention answers stale-message immediately. No backfill worker or storage flag is needed; the cassandra mls_commit_locks table becomes unread and can be dropped in a follow-up. During a rolling restart of galley there is a brief mixed-arbiter window (old pods lock via cassandra, new pods via postgres); avoid running commits against one group across old and new pods simultaneously, and avoid rolling back after cutover.
1 change: 1 addition & 0 deletions integration/integration.cabal
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,7 @@ library
Test.Migration.Util
Test.MLS
Test.MLS.Clients
Test.MLS.CommitLock
Test.MLS.History
Test.MLS.KeyPackage
Test.MLS.Keys
Expand Down
38 changes: 38 additions & 0 deletions integration/test/Test/MLS/CommitLock.hs
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
-- This file is part of the Wire Server implementation.
--
-- Copyright (C) 2026 Wire Swiss GmbH <opensource@wire.com>
--
-- This program is free software: you can redistribute it and/or modify it under
-- the terms of the GNU Affero General Public License as published by the Free
-- Software Foundation, either version 3 of the License, or (at your option) any
-- later version.
--
-- This program is distributed in the hope that it will be useful, but WITHOUT
-- ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS
-- FOR A PARTICULAR PURPOSE. See the GNU Affero General Public License for more
-- details.
--
-- You should have received a copy of the GNU Affero General Public License along
-- with this program. If not, see <https://www.gnu.org/licenses/>.

module Test.MLS.CommitLock where

import MLS.Util
import SetupHelpers
import Testlib.Prelude

-- | Every MLS commit acquires and releases the commit lock, so two successive
-- commits prove acquire -> release -> re-acquire through the pg advisory-lock
-- interpreter. A leaked lock would fail the second commit.
testMLSCommitLock :: (HasCallStack) => App ()
testMLSCommitLock = do
[alice, bob, charlie] <- createAndConnectUsers [OwnDomain, OwnDomain, OwnDomain]
alice1 <- createMLSClient def alice
bob1 <- createMLSClient def bob
void $ uploadNewKeyPackage def bob1
convId <- createNewGroup def alice1
void $ createAddCommit alice1 convId [bob] >>= sendAndConsumeCommitBundle

charlie1 <- createMLSClient def charlie
void $ uploadNewKeyPackage def charlie1
void $ createAddCommit alice1 convId [charlie] >>= sendAndConsumeCommitBundle
12 changes: 4 additions & 8 deletions libs/wire-subsystems/src/Wire/ConversationStore.hs
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@ import Data.Id
import Data.Misc
import Data.Qualified
import Data.Range
import Data.Time.Clock
import Imports
import Polysemy
import Wire.API.Conversation hiding (Conversation, Member)
Expand All @@ -44,14 +43,11 @@ import Wire.Sem.Paging.Cassandra
import Wire.StoredConversation
import Wire.UserList

data LockAcquired
= Acquired
| NotAcquired
deriving (Show, Eq)

data MLSCommitLockStore m a where
AcquireCommitLock :: GroupId -> Epoch -> NominalDiffTime -> MLSCommitLockStore m LockAcquired
ReleaseCommitLock :: GroupId -> Epoch -> MLSCommitLockStore m ()
-- | Runs the action while holding an exclusive lock for @(groupId, epoch)@.
-- Returns 'Nothing' without running the action when another holder is active
-- (callers respond 'MLSStaleMessage').
HoldCommitLock :: GroupId -> Epoch -> m a -> MLSCommitLockStore m (Maybe a)

data ConversationSearch = ConversationSearch
{ team :: TeamId,
Expand Down
46 changes: 2 additions & 44 deletions libs/wire-subsystems/src/Wire/ConversationStore/Cassandra.hs
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,7 @@
-- with this program. If not, see <https://www.gnu.org/licenses/>.

module Wire.ConversationStore.Cassandra
( interpretMLSCommitLockStoreToCassandra,
interpretConversationStoreToCassandra,
( interpretConversationStoreToCassandra,
interpretConversationStoreToCassandraAndPostgres,
interpretConversationStoreByMigration,
MigrationError (..),
Expand All @@ -41,7 +40,6 @@ import Data.Monoid
import Data.Qualified
import Data.Range
import Data.Set qualified as Set
import Data.Time
import Data.UUID.Util qualified as UUID
import Imports
import Network.HTTP.Types.Status (status500)
Expand Down Expand Up @@ -71,7 +69,7 @@ import Wire.API.MLS.GroupInfo
import Wire.API.MLS.LeafNode (LeafIndex)
import Wire.API.MLS.SubConversation
import Wire.API.Provider.Service
import Wire.ConversationStore (ConversationStore (..), LockAcquired (..), MLSCommitLockStore (..))
import Wire.ConversationStore (ConversationStore (..))
import Wire.ConversationStore qualified as ConvStore
import Wire.ConversationStore.Cassandra.Instances ()
import Wire.ConversationStore.Cassandra.Queries qualified as Cql
Expand Down Expand Up @@ -351,37 +349,6 @@ updateToMLSProtocol client cnv =
updateChannelAddPermissions :: ConvId -> AddPermission -> Client ()
updateChannelAddPermissions cid cap = retry x5 $ write Cql.updateChannelAddPermission (params LocalQuorum (cap, cid))

acquireCommitLock :: GroupId -> Epoch -> NominalDiffTime -> Client LockAcquired
acquireCommitLock groupId epoch ttl = do
rows <-
retry x5 $
trans
Cql.acquireCommitLock
( params
LocalQuorum
(groupId, epoch, round ttl)
)
{ serialConsistency = Just LocalSerialConsistency
}
pure $
if checkTransSuccess rows
then Acquired
else NotAcquired

releaseCommitLock :: GroupId -> Epoch -> Client ()
releaseCommitLock groupId epoch =
retry x5 $
write
Cql.releaseCommitLock
( params
LocalQuorum
(groupId, epoch)
)

checkTransSuccess :: [Row] -> Bool
checkTransSuccess [] = False
checkTransSuccess (row : _) = either (const False) (fromMaybe False) $ fromRow 0 row

removeTeamConv :: TeamId -> ConvId -> Client ()
removeTeamConv tid cid = liftClient $ do
retry x5 . batch $ do
Expand Down Expand Up @@ -880,15 +847,6 @@ isConversationOutOfSync cid =
maybe False (fromMaybe False . runIdentity)
<$> retry x1 (query1 Cql.lookupConvOutOfSync (params LocalQuorum (Identity cid)))

interpretMLSCommitLockStoreToCassandra :: (Member (Embed IO) r, Member TinyLog r) => ClientState -> InterpreterFor MLSCommitLockStore r
interpretMLSCommitLockStoreToCassandra client = interpret $ \case
AcquireCommitLock gId epoch ttl -> do
logEffect "MLSCommitLockStore.AcquireCommitLock"
embedClient client $ acquireCommitLock gId epoch ttl
ReleaseCommitLock gId epoch -> do
logEffect "MLSCommitLockStore.ReleaseCommitLock"
embedClient client $ releaseCommitLock gId epoch

interpretConversationStoreToCassandra ::
forall r a.
( PGConstraints r,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@
-- - conversation
-- - member
-- - member_remote_user
-- - mls_commit_locks
-- - mls_group_member_client
-- - subconversation
-- - team_conv
Expand Down Expand Up @@ -355,12 +354,6 @@ removeAllMLSClients = "DELETE FROM mls_group_member_client WHERE group_id = ?"
lookupMLSClients :: PrepQuery R (Identity GroupId) (Domain, UserId, ClientId, Int32, Bool)
lookupMLSClients = "select user_domain, user, client, leaf_node_index, removal_pending from mls_group_member_client where group_id = ?"

acquireCommitLock :: PrepQuery W (GroupId, Epoch, Int32) Row
acquireCommitLock = "insert into mls_commit_locks (group_id, epoch) values (?, ?) if not exists using ttl ?"

releaseCommitLock :: PrepQuery W (GroupId, Epoch) ()
releaseCommitLock = "delete from mls_commit_locks where group_id = ? and epoch = ?"

-- Bots ---------------------------------------------------------------------

insertBot :: PrepQuery W (ConvId, BotId, ServiceId, ProviderId) ()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ import Imports
import Polysemy
import Polysemy.Error
import Polysemy.Input
import Polysemy.Resource
import Polysemy.TinyLog qualified as P
import System.Logger.Class qualified as Log
import Wire.API.Conversation hiding (Member)
Expand Down Expand Up @@ -67,7 +66,6 @@ resetLocalMLSMainConversation ::
Member NotificationSubsystem r,
Member ProposalStore r,
Member Random r,
Member Resource r,
Member ConversationStore r,
Member P.TinyLog r,
Member MLSCommitLockStore r,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -593,7 +593,6 @@ sendMLSCommitBundle ::
Member (Input (Maybe (MLSKeysByPurpose MLSPrivateKeys))) r,
Member Now r,
Member LegalHoldStore r,
Member Resource r,
Member TeamStore r,
Member FederationSubsystem r,
Member TeamSubsystem r,
Expand Down Expand Up @@ -705,7 +704,6 @@ leaveSubConversation ::
( HasLeaveSubConversationEffects r,
Member (Error FederationError) r,
Member (Input (Local ())) r,
Member Resource r,
Member TeamSubsystem r,
Member E.MLSCommitLockStore r,
Member (Input ConversationSubsystemConfig) r
Expand All @@ -728,7 +726,6 @@ leaveSubConversation domain lscr = do
deleteSubConversationForRemoteUser ::
( Member E.ConversationStore r,
Member (Input (Local ())) r,
Member Resource r,
Member TeamSubsystem r,
Member E.MLSCommitLockStore r,
Member TinyLog r
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@ import Data.Set qualified as Set
import Imports
import Polysemy
import Polysemy.Error
import Polysemy.Resource (Resource)
import Polysemy.State
import Wire.API.Conversation.Protocol
import Wire.API.Error
Expand Down Expand Up @@ -136,7 +135,6 @@ processExternalCommit ::
Member (ErrorS MLSStaleMessage) r,
Member (ErrorS MLSIdentityMismatch) r,
Member (ErrorS MLSSubConvClientNotInParent) r,
Member Resource r,
HasProposalActionEffects r,
Member (ErrorS MLSInvalidLeafNodeSignature) r,
Member MLSCommitLockStore r
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,6 @@ import Imports
import Polysemy
import Polysemy.Error
import Polysemy.Input (Input)
import Polysemy.Resource (Resource)
import Wire.API.Conversation hiding (Member)
import Wire.API.Conversation.Action
import Wire.API.Conversation.Config (ConversationSubsystemConfig)
Expand Down Expand Up @@ -77,7 +76,6 @@ processInternalCommit ::
Member (ErrorS 'MLSIdentityMismatch) r,
Member (ErrorS 'MissingLegalholdConsent) r,
Member (ErrorS 'GroupIdVersionNotSupported) r,
Member Resource r,
Member Random r,
Member (ErrorS MLSInvalidLeafNodeSignature) r,
Member MLSCommitLockStore r,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,6 @@ import Polysemy
import Polysemy.Error
import Polysemy.Input
import Polysemy.Output
import Polysemy.Resource (Resource)
import Polysemy.TinyLog
import System.Logger qualified as Log
import Wire.API.Conversation hiding (Member)
Expand Down Expand Up @@ -175,7 +174,6 @@ postMLSCommitBundle ::
Member (ErrorS GroupIdVersionNotSupported) r,
Member (Input (Maybe GroupInfoCheckEnabled)) r,
Member Random r,
Member Resource r,
Members MLSMessageStaticErrors r,
Member (ErrorS 'MLSInvalidLeafNodeSignature) r,
HasProposalEffects r,
Expand Down Expand Up @@ -212,7 +210,6 @@ postMLSCommitBundleFromLocalUser ::
Member (Input (Maybe GroupInfoCheckEnabled)) r,
Member (Input (Maybe (MLSKeysByPurpose MLSPrivateKeys))) r,
Member Random r,
Member Resource r,
Members MLSMessageStaticErrors r,
Member (ErrorS 'MLSInvalidLeafNodeSignature) r,
HasProposalEffects r,
Expand Down Expand Up @@ -249,7 +246,6 @@ postMLSCommitBundleToLocalConv ::
Member (Input EnableOutOfSyncCheck) r,
Member (Input (Maybe GroupInfoCheckEnabled)) r,
Member Random r,
Member Resource r,
Members MLSMessageStaticErrors r,
Member (ErrorS 'MLSInvalidLeafNodeSignature) r,
HasProposalEffects r,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,6 @@ import Imports
import Polysemy
import Polysemy.Error
import Polysemy.Input
import Polysemy.Resource
import Polysemy.TinyLog
import Wire.API.Conversation hiding (Member)
import Wire.API.Conversation.Config (ConversationSubsystemConfig)
Expand Down Expand Up @@ -218,7 +217,6 @@ deleteSubConversation ::
Member (Error FederationError) r,
Member (FederationAPIAccess FederatorClient) r,
Member (Input (Maybe (MLSKeysByPurpose MLSPrivateKeys))) r,
Member Resource r,
Member Conversation.MLSCommitLockStore r,
Member TeamSubsystem r,
Member TinyLog r
Expand Down Expand Up @@ -292,7 +290,6 @@ leaveSubConversation ::
Member (Error FederationError) r,
Member (ErrorS 'MLSStaleMessage) r,
Member (ErrorS 'MLSNotEnabled) r,
Member Resource r,
Members LeaveSubConversationStaticErrors r,
Member Conversation.MLSCommitLockStore r,
Member TeamSubsystem r,
Expand All @@ -319,7 +316,6 @@ leaveLocalSubConversation ::
Member (ErrorS 'MLSStaleMessage) r,
Member (ErrorS 'MLSNotEnabled) r,
Member (Error FederationError) r,
Member Resource r,
Members LeaveSubConversationStaticErrors r,
Member Conversation.MLSCommitLockStore r,
Member TeamSubsystem r,
Expand Down Expand Up @@ -394,7 +390,6 @@ resetLocalSubConversation ::
Member (ErrorS 'ConvAccessDenied) r,
Member (ErrorS 'ConvNotFound) r,
Member (ErrorS 'MLSStaleMessage) r,
Member Resource r,
Member Conversation.MLSCommitLockStore r,
Member TeamSubsystem r,
Member TinyLog r
Expand Down
Loading