From e695a954223f44a7b985c8964f96fd0ce10d609a Mon Sep 17 00:00:00 2001 From: Gautier DI FOLCO Date: Mon, 3 Aug 2026 20:22:47 +0200 Subject: [PATCH] WPB-22965: migrate blacklist to PostGreSQL Migrate the brig blacklist (BlockListStore) from Cassandra to PostgreSQL using the dual-write + background-worker pattern, following the shipped domainRegistration migration as the reference. Adds the Postgres/DualWrite/Migration interpreters, the blockList StorageLocation config field, background-worker wiring, the SQL migration plus regenerated postgres-schema.sql, chart/config/gotmpl/docs updates, and an integration test. --- changelog.d/5-internal/WPB-22965 | 1 + .../background-worker/configmap.yaml | 1 + charts/wire-server/values.yaml | 5 + .../src/developer/reference/config-options.md | 10 ++ hack/helm_vars/common.yaml.gotmpl | 1 + hack/helm_vars/wire-server/values.yaml.gotmpl | 1 + integration/integration.cabal | 1 + integration/test/API/BrigInternal.hs | 15 ++ integration/test/Test/Migration/BlockList.hs | 68 ++++++++ .../20260803164250-blacklist.sql | 3 + .../src/Wire/BlockListStore/Cassandra.hs | 4 + .../src/Wire/BlockListStore/DualWrite.hs | 45 +++++ .../src/Wire/BlockListStore/Migration.hs | 155 ++++++++++++++++++ .../src/Wire/BlockListStore/Postgres.hs | 74 +++++++++ .../src/Wire/PostgresMigrationOpts.hs | 4 +- libs/wire-subsystems/wire-subsystems.cabal | 3 + postgres-schema.sql | 18 ++ .../background-worker.integration.yaml | 2 + .../src/Wire/BackgroundWorker.hs | 10 +- .../src/Wire/BackgroundWorker/Options.hs | 1 + .../src/Wire/PostgresMigrations.hs | 19 +++ .../Wire/BackendNotificationPusherSpec.hs | 6 +- .../background-worker/test/Test/Wire/Util.hs | 3 +- services/brig/brig.integration.yaml | 1 + .../brig/src/Brig/CanonicalInterpreter.hs | 11 +- services/galley/galley.integration.yaml | 1 + 26 files changed, 456 insertions(+), 7 deletions(-) create mode 100644 changelog.d/5-internal/WPB-22965 create mode 100644 integration/test/Test/Migration/BlockList.hs create mode 100644 libs/wire-subsystems/postgres-migrations/20260803164250-blacklist.sql create mode 100644 libs/wire-subsystems/src/Wire/BlockListStore/DualWrite.hs create mode 100644 libs/wire-subsystems/src/Wire/BlockListStore/Migration.hs create mode 100644 libs/wire-subsystems/src/Wire/BlockListStore/Postgres.hs diff --git a/changelog.d/5-internal/WPB-22965 b/changelog.d/5-internal/WPB-22965 new file mode 100644 index 00000000000..46a3360e74e --- /dev/null +++ b/changelog.d/5-internal/WPB-22965 @@ -0,0 +1 @@ +brig: migrate blacklist (BlockListStore) from Cassandra to PostgreSQL diff --git a/charts/wire-server/templates/background-worker/configmap.yaml b/charts/wire-server/templates/background-worker/configmap.yaml index 299d0703d3a..545f1c3f268 100644 --- a/charts/wire-server/templates/background-worker/configmap.yaml +++ b/charts/wire-server/templates/background-worker/configmap.yaml @@ -85,6 +85,7 @@ data: migrateTeamFeatures: {{ .migrateTeamFeatures }} migrateDomainRegistration: {{ .migrateDomainRegistration }} migrateUsers: {{ .migrateUsers }} + migrateBlockList: {{ .migrateBlockList }} migrationOptions: {{ toYaml .migrationOptions | indent 6 }} diff --git a/charts/wire-server/values.yaml b/charts/wire-server/values.yaml index d2a7deafdc5..ea989ea232a 100644 --- a/charts/wire-server/values.yaml +++ b/charts/wire-server/values.yaml @@ -90,6 +90,7 @@ galley: teamFeatures: cassandra domainRegistration: cassandra user: cassandra + blockList: cassandra settings: httpPoolSize: 128 maxTeamSize: 10000 @@ -1029,6 +1030,10 @@ background-worker: # It's important to set `settings.postgresMigration.users` to `migration-to-postgresql` # before starting the migration. migrateUsers: false + # This will start the migration of blacklist data. + # It's important to set `settings.postgresMigration.blockList` to `migration-to-postgresql` + # before starting the migration. + migrateBlockList: false backendNotificationPusher: pushBackoffMinWait: 10000 # in microseconds, so 10ms diff --git a/docs/src/developer/reference/config-options.md b/docs/src/developer/reference/config-options.md index 062f891f03c..4f99108b922 100644 --- a/docs/src/developer/reference/config-options.md +++ b/docs/src/developer/reference/config-options.md @@ -2182,12 +2182,14 @@ galley: teamFeatures: postgresql domainRegistration: postgresql user: postgresql + blockList: postgresql background-worker: config: migrateConversations: false migrateConversationCodes: false migrateTeamFeatures: false migrateDomainRegistration: false + migrateBlockList: false ``` #### Migration for existing installations @@ -2219,6 +2221,7 @@ The current settings and their background-worker flags are: - `teamFeatures` -> `migrateTeamFeatures` - `domainRegistration` -> `migrateDomainRegistration` - `user` -> `migrateUsers` +- `blockList` -> `migrateBlockList` **Migration pattern per migration setting** @@ -2239,6 +2242,7 @@ The current settings and their background-worker flags are: teamFeatures: migration-to-postgresql domainRegistration: migration-to-postgresql user: migration-to-postgresql + blockList: cassandra background-worker: config: migrateConversations: false @@ -2246,6 +2250,7 @@ The current settings and their background-worker flags are: migrateTeamFeatures: false migrateDomainRegistration: false migrateUsers: false + migrateBlockList: false ``` This change should restart the affected pods, and new writes will follow the @@ -2261,6 +2266,7 @@ The current settings and their background-worker flags are: migrateTeamFeatures: true migrateDomainRegistration: true migrateUsers: true + migrateBlockList: true ``` During migration, Cassandra rows are not deleted. Writes and migration share @@ -2286,6 +2292,7 @@ The current settings and their background-worker flags are: > to be saved, the operator must insert some value as `name` and/or > `activated` and then re-trigger the migration **after** the background > worker finishes migrating the valid users. + - `blockList`: `wire_block_list_migration_finished` 3. Cut over reads and writes to PostgreSQL for the selected migration setting(s). This configuration must be used from now on for every new @@ -2300,6 +2307,7 @@ The current settings and their background-worker flags are: teamFeatures: postgresql domainRegistration: postgresql user: postgresql + blockList: cassandra background-worker: config: migrateConversations: false @@ -2307,6 +2315,7 @@ The current settings and their background-worker flags are: migrateTeamFeatures: false migrateDomainRegistration: false migrateUsers: false + migrateBlockList: false ``` **How to run migrations independently or in batches** @@ -2395,6 +2404,7 @@ migrateConversations: false migrateConversationCodes: false migrateTeamFeatures: false migrateDomainRegistration: false +migrateBlockList: false # migration settings migrationOptions: diff --git a/hack/helm_vars/common.yaml.gotmpl b/hack/helm_vars/common.yaml.gotmpl index 2276355e2a9..4d36ffac030 100644 --- a/hack/helm_vars/common.yaml.gotmpl +++ b/hack/helm_vars/common.yaml.gotmpl @@ -19,6 +19,7 @@ conversationCodesStore: {{ $preferredStore }} teamFeaturesStore: {{ $preferredStore }} domainRegistration: {{ $preferredStore }} userStore: {{ $preferredStore }} +blockListStore: {{ $preferredStore }} {{- if (eq (env "UPLOAD_XML_S3_BASE_URL") "") }} uploadXml: {} diff --git a/hack/helm_vars/wire-server/values.yaml.gotmpl b/hack/helm_vars/wire-server/values.yaml.gotmpl index f05d42a98e4..cbd37ce4237 100644 --- a/hack/helm_vars/wire-server/values.yaml.gotmpl +++ b/hack/helm_vars/wire-server/values.yaml.gotmpl @@ -306,6 +306,7 @@ galley: teamFeatures: {{ .Values.teamFeaturesStore }} domainRegistration: {{ .Values.domainRegistration }} user: {{ .Values.userStore }} + blockList: {{ .Values.blockListStore }} settings: maxConvAndTeamSize: 16 maxTeamSize: 32 diff --git a/integration/integration.cabal b/integration/integration.cabal index 990fc21e84f..049fdf30d3b 100644 --- a/integration/integration.cabal +++ b/integration/integration.cabal @@ -177,6 +177,7 @@ library Test.Login Test.Meetings Test.MessageTimer + Test.Migration.BlockList Test.Migration.Conversation Test.Migration.ConversationCodes Test.Migration.DomainRegistration diff --git a/integration/test/API/BrigInternal.hs b/integration/test/API/BrigInternal.hs index 4407e1a043b..b720115bfa3 100644 --- a/integration/test/API/BrigInternal.hs +++ b/integration/test/API/BrigInternal.hs @@ -353,6 +353,21 @@ getPasswordResetCode domain email = do req <- baseRequest domain Brig Unversioned "i/users/password-reset-code" submit "GET" $ req & addQueryParams [("email", email)] +addBlacklist :: (HasCallStack, MakesValue domain) => domain -> String -> App Response +addBlacklist domain email = do + req <- baseRequest domain Brig Unversioned "i/users/blacklist" + submit "POST" $ req & addQueryParams [("email", email)] + +deleteBlacklist :: (HasCallStack, MakesValue domain) => domain -> String -> App Response +deleteBlacklist domain email = do + req <- baseRequest domain Brig Unversioned "i/users/blacklist" + submit "DELETE" $ req & addQueryParams [("email", email)] + +checkBlacklist :: (HasCallStack, MakesValue domain) => domain -> String -> App Response +checkBlacklist domain email = do + req <- baseRequest domain Brig Unversioned "i/users/blacklist" + submit "GET" $ req & addQueryParams [("email", email)] + data PutSSOId = PutSSOId { scimExternalId :: Maybe String, subject :: Maybe String, diff --git a/integration/test/Test/Migration/BlockList.hs b/integration/test/Test/Migration/BlockList.hs new file mode 100644 index 00000000000..824c9763853 --- /dev/null +++ b/integration/test/Test/Migration/BlockList.hs @@ -0,0 +1,68 @@ +-- This file is part of the Wire Server implementation. +-- +-- Copyright (C) 2026 Wire Swiss GmbH +-- +-- 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 . + +module Test.Migration.BlockList (testBlockListMigration) where + +import qualified API.BrigInternal as BrigInternal +import API.Common +import Control.Monad.Codensity +import Control.Monad.Reader +import Test.Migration.Util (waitForMigration) +import Testlib.Prelude +import Testlib.ResourcePool + +-- | Migrate the 'blacklist' store (brig) from Cassandra to PostgreSQL. +-- +-- The blacklist holds email keys with no read-back payload, so the migration is +-- a straight copy: a key blacklisted in Cassandra must survive the cutover and +-- remain deletable once PostgreSQL is the sole source of truth. +testBlockListMigration :: (HasCallStack) => App () +testBlockListMigration = do + resourcePool <- asks (.resourcePool) + email <- randomEmail + runCodensity (acquireResources 1 resourcePool) $ \[backend] -> do + let domain = backend.berDomain + + -- Cassandra: blacklist an email key and confirm it is reported as such. + runCodensity (startDynamicBackend backend (conf "cassandra" False)) $ \_ -> do + assertSuccess =<< BrigInternal.addBlacklist domain email + assertStatus 200 =<< BrigInternal.checkBlacklist domain email + + -- migration-to-postgresql with the worker running: backfill the existing key + -- and confirm it is still blacklisted once the migration is finished. + runCodensity (startDynamicBackend backend (conf "migration-to-postgresql" True)) $ \_ -> do + waitForMigration domain counterName + assertStatus 200 =<< BrigInternal.checkBlacklist domain email + + -- PostgreSQL only: the migrated key must persist, and deleting it must + -- remove it. + runCodensity (startDynamicBackend backend (conf "postgresql" False)) $ \_ -> do + assertStatus 200 =<< BrigInternal.checkBlacklist domain email + assertSuccess =<< BrigInternal.deleteBlacklist domain email + assertStatus 404 =<< BrigInternal.checkBlacklist domain email + where + conf :: String -> Bool -> ServiceOverrides + conf db runMigration = + def + { brigCfg = setField "postgresMigration.blockList" db, + backgroundWorkerCfg = + setField "postgresMigration.blockList" db + >=> setField "migrateBlockList" runMigration + } + + counterName :: String + counterName = "^wire_block_list_migration_finished" diff --git a/libs/wire-subsystems/postgres-migrations/20260803164250-blacklist.sql b/libs/wire-subsystems/postgres-migrations/20260803164250-blacklist.sql new file mode 100644 index 00000000000..91c8fc3a144 --- /dev/null +++ b/libs/wire-subsystems/postgres-migrations/20260803164250-blacklist.sql @@ -0,0 +1,3 @@ +CREATE TABLE IF NOT EXISTS blacklist ( + key text PRIMARY KEY +); diff --git a/libs/wire-subsystems/src/Wire/BlockListStore/Cassandra.hs b/libs/wire-subsystems/src/Wire/BlockListStore/Cassandra.hs index 0bf4f46d3ac..7cfbda00b31 100644 --- a/libs/wire-subsystems/src/Wire/BlockListStore/Cassandra.hs +++ b/libs/wire-subsystems/src/Wire/BlockListStore/Cassandra.hs @@ -17,6 +17,7 @@ module Wire.BlockListStore.Cassandra ( interpretBlockListStoreToCassandra, + selectAllBlacklist, ) where @@ -60,3 +61,6 @@ keySelect = "SELECT key FROM blacklist WHERE key = ?" keyDelete :: PrepQuery W (Identity Text) () keyDelete = "DELETE FROM blacklist WHERE key = ?" + +selectAllBlacklist :: PrepQuery R () (Identity Text) +selectAllBlacklist = "SELECT key FROM blacklist" diff --git a/libs/wire-subsystems/src/Wire/BlockListStore/DualWrite.hs b/libs/wire-subsystems/src/Wire/BlockListStore/DualWrite.hs new file mode 100644 index 00000000000..b95371b7c34 --- /dev/null +++ b/libs/wire-subsystems/src/Wire/BlockListStore/DualWrite.hs @@ -0,0 +1,45 @@ +-- This file is part of the Wire Server implementation. +-- +-- Copyright (C) 2026 Wire Swiss GmbH +-- +-- 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 . + +module Wire.BlockListStore.DualWrite + ( interpretBlockListStoreToCassandraAndPostgres, + ) +where + +import Cassandra (ClientState) +import Imports +import Polysemy +import Wire.BlockListStore +import Wire.BlockListStore qualified as BlockListStore +import Wire.BlockListStore.Cassandra qualified as Cassandra +import Wire.BlockListStore.Postgres qualified as Postgres +import Wire.Postgres (PGConstraints) + +-- | Cassandra is the source of truth during migration; writes are mirrored to Postgres. +interpretBlockListStoreToCassandraAndPostgres :: + (PGConstraints r) => + ClientState -> + InterpreterFor BlockListStore r +interpretBlockListStoreToCassandraAndPostgres cs = interpret $ \case + Insert key -> do + Cassandra.interpretBlockListStoreToCassandra cs $ BlockListStore.insert key + Postgres.interpretBlockListStoreToPostgres $ BlockListStore.insert key + Exists key -> + Cassandra.interpretBlockListStoreToCassandra cs $ BlockListStore.exists key + Delete key -> do + Cassandra.interpretBlockListStoreToCassandra cs $ BlockListStore.delete key + Postgres.interpretBlockListStoreToPostgres $ BlockListStore.delete key diff --git a/libs/wire-subsystems/src/Wire/BlockListStore/Migration.hs b/libs/wire-subsystems/src/Wire/BlockListStore/Migration.hs new file mode 100644 index 00000000000..5a8c3ae39a1 --- /dev/null +++ b/libs/wire-subsystems/src/Wire/BlockListStore/Migration.hs @@ -0,0 +1,155 @@ +-- This file is part of the Wire Server implementation. +-- +-- Copyright (C) 2026 Wire Swiss GmbH +-- +-- 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 . + +module Wire.BlockListStore.Migration (migrateBlacklistLoop) where + +import Cassandra hiding (Value) +import Data.ByteString.Conversion (toByteString') +import Data.Conduit +import Data.Conduit.List qualified as C +import Data.IORef qualified as IORef +import Data.Time +import Hasql.Pool.Extended qualified as Hasql +import Imports +import Polysemy +import Polysemy.Async +import Polysemy.AtomicState +import Polysemy.Conc (interpretRace) +import Polysemy.Conc qualified as Conc +import Polysemy.Conc.Effect.Race hiding (Timeout) +import Polysemy.Input +import Polysemy.Resource (Resource, bracket, resourceToIOFinal) +import Polysemy.TinyLog +import Prometheus qualified +import System.Logger qualified as Log +import UnliftIO qualified +import Wire.BlockListStore.Cassandra qualified as Cql +import Wire.BlockListStore.Postgres qualified as Postgres +import Wire.Migration +import Wire.Postgres +import Wire.Sem.Logger (mapLogger) +import Wire.Sem.Logger.TinyLog (loggerToTinyLog) + +type EffectStack = + [ AtomicState Int, + Input ClientState, + Input Hasql.Pool, + Resource, + Async, + Race, + TinyLog, + Embed IO, + Final IO + ] + +migrateBlacklistLoop :: + MigrationOptions -> + ClientState -> + Hasql.Pool -> + Log.Logger -> + Prometheus.Counter -> + Prometheus.Counter -> + Prometheus.Counter -> + Prometheus.Vector Text Prometheus.Histogram -> + IO () +migrateBlacklistLoop migOpts cassClient pgPool logger migCounter migFinished migFailed migDuration = + migrationLoop + logger + "blacklist" + migFinished + migFailed + (interpreter cassClient pgPool logger "blacklist") + (migrateAllBlacklist migOpts migCounter migDuration) + +interpreter :: ClientState -> Hasql.Pool -> Log.Logger -> ByteString -> Sem EffectStack a -> IO (Int, a) +interpreter cassClient pgPool logger name = + runFinal + . embedToFinal + . loggerToTinyLog logger + . mapLogger (Log.field "migration" (Log.val name) .) + . raiseUnder + . interpretRace + . asyncToIOFinal + . resourceToIOFinal + . runInputConst pgPool + . runInputConst cassClient + . atomicStateToIO 0 + +migrateAllBlacklist :: + ( Member (Input Hasql.Pool) r, + Member (Embed IO) r, + Member (Input ClientState) r, + Member TinyLog r, + Member (AtomicState Int) r, + Member Resource r, + Member Race r + ) => + MigrationOptions -> + Prometheus.Counter -> + Prometheus.Vector Text Prometheus.Histogram -> + ConduitM () Void (Sem r) () +migrateAllBlacklist migOpts migCounter migDuration = do + lift $ info $ Log.msg (Log.val "migrateAllBlacklist") + withCount + (paginateSem Cql.selectAllBlacklist (paramsP LocalQuorum () migOpts.pageSize) x5) + .| logRetrievedPage migOpts.pageSize id + .| C.mapM_ + ( traverse_ + ( \row@(Identity key) -> + handleErrors (toByteString' key) (migrateBlacklistRow migOpts migCounter migDuration row) + ) + ) + +migrateBlacklistRow :: + ( PGConstraints r, + Member TinyLog r, + Member Resource r, + Member Race r + ) => + MigrationOptions -> + Prometheus.Counter -> + Prometheus.Vector Text Prometheus.Histogram -> + Identity Text -> + Sem r () +migrateBlacklistRow migOpts migCounter migDuration (Identity key) = do + outcomeRef <- liftIO $ IORef.newIORef @Text "error" + bracket + (liftIO getCurrentTime) + (observeDuration migDuration outcomeRef) + ( const $ do + timeoutResult <- Conc.timeout (migOpts.timeout <$ handleTimeout) migOpts.timeout $ Postgres.insertKey key + case timeoutResult of + Left timedOutAfter -> do + markOutcome outcomeRef "timeout" + liftIO . UnliftIO.throwIO $ MigrationTimedOut key timedOutAfter + Right () -> do + markOutcome outcomeRef "success" + liftIO $ Prometheus.incCounter migCounter + ) + where + handleTimeout = + err $ + Log.msg (Log.val "blacklist migration timed out") + . Log.field "key" (show key) + . Log.field "timeout" (show migOpts.timeout) + + markOutcome ref outcome = liftIO $ IORef.writeIORef ref outcome + + observeDuration metric outcomeRef start = do + outcome <- liftIO $ IORef.readIORef outcomeRef + end <- liftIO getCurrentTime + liftIO $ Prometheus.withLabel metric outcome (`Prometheus.observe` realToFrac (diffUTCTime end start)) diff --git a/libs/wire-subsystems/src/Wire/BlockListStore/Postgres.hs b/libs/wire-subsystems/src/Wire/BlockListStore/Postgres.hs new file mode 100644 index 00000000000..5a7d1d8f283 --- /dev/null +++ b/libs/wire-subsystems/src/Wire/BlockListStore/Postgres.hs @@ -0,0 +1,74 @@ +-- This file is part of the Wire Server implementation. +-- +-- Copyright (C) 2026 Wire Swiss GmbH +-- +-- 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 . + +module Wire.BlockListStore.Postgres + ( interpretBlockListStoreToPostgres, + insertKey, + ) +where + +import Hasql.Statement qualified as Hasql +import Hasql.TH +import Imports +import Polysemy +import Wire.BlockListStore (BlockListStore (..)) +import Wire.Postgres +import Wire.UserKeyStore (EmailKey, emailKeyUniq) + +interpretBlockListStoreToPostgres :: + (PGConstraints r) => + InterpreterFor BlockListStore r +interpretBlockListStoreToPostgres = interpret $ \case + Insert key -> insertImpl key + Exists key -> existsImpl key + Delete key -> deleteImpl key + +insertImpl :: (PGConstraints r) => EmailKey -> Sem r () +insertImpl = insertKey . emailKeyUniq + +insertKey :: (PGConstraints r) => Text -> Sem r () +insertKey key = + runStatement key insertStatement + where + insertStatement :: Hasql.Statement Text () + insertStatement = + [resultlessStatement|INSERT INTO blacklist (key) + VALUES ($1 :: text) + ON CONFLICT DO NOTHING + |] + +existsImpl :: (PGConstraints r) => EmailKey -> Sem r Bool +existsImpl key = + runStatement (emailKeyUniq key) existsStatement + where + existsStatement :: Hasql.Statement Text Bool + existsStatement = + [singletonStatement|SELECT EXISTS ( + SELECT 1 + FROM blacklist + WHERE key = ($1 :: text) + ) :: bool|] + +deleteImpl :: (PGConstraints r) => EmailKey -> Sem r () +deleteImpl key = + runStatement (emailKeyUniq key) deleteStatement + where + deleteStatement :: Hasql.Statement Text () + deleteStatement = + [resultlessStatement|DELETE FROM blacklist + WHERE key = ($1 :: text) + |] diff --git a/libs/wire-subsystems/src/Wire/PostgresMigrationOpts.hs b/libs/wire-subsystems/src/Wire/PostgresMigrationOpts.hs index 327862f7cd5..d5e350e608f 100644 --- a/libs/wire-subsystems/src/Wire/PostgresMigrationOpts.hs +++ b/libs/wire-subsystems/src/Wire/PostgresMigrationOpts.hs @@ -56,7 +56,8 @@ data PostgresMigrationOpts = PostgresMigrationOpts conversationCodes :: StorageLocation, teamFeatures :: StorageLocation, domainRegistration :: StorageLocation, - user :: StorageLocation + user :: StorageLocation, + blockList :: StorageLocation } deriving (Show) @@ -68,3 +69,4 @@ instance FromJSON PostgresMigrationOpts where <*> o .: "teamFeatures" <*> o .: "domainRegistration" <*> o .: "user" + <*> o .: "blockList" diff --git a/libs/wire-subsystems/wire-subsystems.cabal b/libs/wire-subsystems/wire-subsystems.cabal index a8e33868abd..19ccab495c1 100644 --- a/libs/wire-subsystems/wire-subsystems.cabal +++ b/libs/wire-subsystems/wire-subsystems.cabal @@ -241,6 +241,9 @@ library Wire.BackgroundJobsRunner.Interpreter Wire.BlockListStore Wire.BlockListStore.Cassandra + Wire.BlockListStore.DualWrite + Wire.BlockListStore.Migration + Wire.BlockListStore.Postgres Wire.BoundedQueue Wire.BoundedQueue.STM Wire.BrigAPIAccess diff --git a/postgres-schema.sql b/postgres-schema.sql index de4a57a0f2e..7f77a0acc60 100644 --- a/postgres-schema.sql +++ b/postgres-schema.sql @@ -1272,6 +1272,17 @@ CREATE TABLE public.asset ( ALTER TABLE public.asset OWNER TO "wire-server"; +-- +-- Name: blacklist; Type: TABLE; Schema: public; Owner: wire-server +-- + +CREATE TABLE public.blacklist ( + key text NOT NULL +); + + +ALTER TABLE public.blacklist OWNER TO "wire-server"; + -- -- Name: bot_conv; Type: TABLE; Schema: public; Owner: wire-server -- @@ -1854,6 +1865,13 @@ ALTER TABLE ONLY arbiter.meetings_results ALTER TABLE ONLY public.apps ADD CONSTRAINT apps_pkey PRIMARY KEY (user_id); +-- +-- Name: blacklist blacklist_pkey; Type: CONSTRAINT; Schema: public; Owner: wire-server +-- + +ALTER TABLE ONLY public.blacklist + ADD CONSTRAINT blacklist_pkey PRIMARY KEY (key); + -- -- Name: bot_conv bot_conv_pkey; Type: CONSTRAINT; Schema: public; Owner: wire-server diff --git a/services/background-worker/background-worker.integration.yaml b/services/background-worker/background-worker.integration.yaml index b0bd0d172e4..dc202db6c00 100644 --- a/services/background-worker/background-worker.integration.yaml +++ b/services/background-worker/background-worker.integration.yaml @@ -59,6 +59,7 @@ migrateConversationCodes: false migrateTeamFeatures: false migrateDomainRegistration: false migrateUsers: false +migrateBlockList: false # Background jobs consumer configuration for integration backgroundJobs: @@ -93,3 +94,4 @@ postgresMigration: teamFeatures: postgresql domainRegistration: postgresql user: postgresql + blockList: postgresql diff --git a/services/background-worker/src/Wire/BackgroundWorker.hs b/services/background-worker/src/Wire/BackgroundWorker.hs index 6c12b02e816..c16c92db1dd 100644 --- a/services/background-worker/src/Wire/BackgroundWorker.hs +++ b/services/background-worker/src/Wire/BackgroundWorker.hs @@ -86,6 +86,13 @@ run opts galleyOpts = do Migrations.users opts.migrationOptions else pure $ pure () + cleanupBlockListMigration <- + if opts.migrateBlockList + then + runAppT env $ + withNamedLogger "migrate-block-list" $ + Migrations.blockList opts.migrationOptions + else pure $ pure () cleanupJobs <- runAppT env $ withNamedLogger "background-job-consumer" $ @@ -97,7 +104,7 @@ run opts galleyOpts = do let cleanup = void $ runConcurrently $ - (,,,,,,,,) + (,,,,,,,,,) <$> Concurrently cleanupDeadUserNotifWatcher <*> Concurrently cleanupBackendNotifPusher <*> Concurrently cleanupConvMigration @@ -105,6 +112,7 @@ run opts galleyOpts = do <*> Concurrently cleanupTeamFeaturesMigration <*> Concurrently cleanupDomainRegistrationMigration <*> Concurrently cleanupUsersMigration + <*> Concurrently cleanupBlockListMigration <*> Concurrently cleanupJobRunner <*> Concurrently cleanupJobs diff --git a/services/background-worker/src/Wire/BackgroundWorker/Options.hs b/services/background-worker/src/Wire/BackgroundWorker/Options.hs index 035460cfc32..9b6eadbd885 100644 --- a/services/background-worker/src/Wire/BackgroundWorker/Options.hs +++ b/services/background-worker/src/Wire/BackgroundWorker/Options.hs @@ -56,6 +56,7 @@ data Opts = Opts migrateTeamFeatures :: !Bool, migrateDomainRegistration :: !Bool, migrateUsers :: !Bool, + migrateBlockList :: !Bool, jobs :: JobConfig, meetingsCleanup :: MeetingsCleanupConfig, backgroundJobs :: BackgroundJobsConfig diff --git a/services/background-worker/src/Wire/PostgresMigrations.hs b/services/background-worker/src/Wire/PostgresMigrations.hs index 28c6a789a4a..b629fb8f0c8 100644 --- a/services/background-worker/src/Wire/PostgresMigrations.hs +++ b/services/background-worker/src/Wire/PostgresMigrations.hs @@ -23,6 +23,7 @@ import System.Logger qualified as Log import UnliftIO import Wire.BackgroundWorker.Env import Wire.BackgroundWorker.Util +import Wire.BlockListStore.Migration import Wire.CodeStore.Migration import Wire.ConversationStore.Migration qualified as ConversationStore import Wire.DomainRegistrationStore.Migration @@ -126,3 +127,21 @@ users migOpts = do pure $ do Log.info logger $ Log.msg (Log.val "cancelling user migration") cancel migrationLoop + +blockList :: MigrationOptions -> AppT IO CleanupAction +blockList migOpts = do + cassClient <- asks (.cassandraBrig) + pgPool <- asks (.hasqlPool) + logger <- asks (.logger) + Log.info logger $ Log.msg (Log.val "starting blacklist migration") + count <- register $ counter $ Prometheus.Info "wire_block_list_migrated_to_pg" "Number of blacklist keys migrated to Postgresql" + finished <- register $ counter $ Prometheus.Info "wire_block_list_migration_finished" "Whether the blacklist migration to Postgresql is finished successfully" + failed <- register $ counter $ Prometheus.Info "wire_block_list_migration_failed" "Whether the blacklist migration to Postgresql has failed" + duration <- register $ vector "outcome" $ histogram (Prometheus.Info "wire_block_list_migration_duration_seconds" "Duration of blacklist migration attempts") defaultBuckets + + migrationLoop <- async . lift $ migrateBlacklistLoop migOpts cassClient pgPool logger count finished failed duration + + Log.info logger $ Log.msg (Log.val "started blacklist migration") + pure $ do + Log.info logger $ Log.msg (Log.val "cancelling blacklist migration") + cancel migrationLoop diff --git a/services/background-worker/test/Test/Wire/BackendNotificationPusherSpec.hs b/services/background-worker/test/Test/Wire/BackendNotificationPusherSpec.hs index be46a03c648..2d6e7c18c57 100644 --- a/services/background-worker/test/Test/Wire/BackendNotificationPusherSpec.hs +++ b/services/background-worker/test/Test/Wire/BackendNotificationPusherSpec.hs @@ -496,7 +496,8 @@ spec = do conversationCodes = CassandraStorage, teamFeatures = CassandraStorage, domainRegistration = CassandraStorage, - user = CassandraStorage + user = CassandraStorage, + blockList = CassandraStorage } gundeckEndpoint = undefined brigEndpoint = undefined @@ -560,7 +561,8 @@ spec = do conversationCodes = CassandraStorage, teamFeatures = CassandraStorage, domainRegistration = CassandraStorage, - user = CassandraStorage + user = CassandraStorage, + blockList = CassandraStorage } gundeckEndpoint = undefined brigEndpoint = undefined diff --git a/services/background-worker/test/Test/Wire/Util.hs b/services/background-worker/test/Test/Wire/Util.hs index 5d89532bfec..74355c339fc 100644 --- a/services/background-worker/test/Test/Wire/Util.hs +++ b/services/background-worker/test/Test/Wire/Util.hs @@ -50,7 +50,8 @@ testEnv = do conversationCodes = CassandraStorage, teamFeatures = CassandraStorage, domainRegistration = CassandraStorage, - user = CassandraStorage + user = CassandraStorage, + blockList = CassandraStorage } statuses <- newIORef mempty backendNotificationMetrics <- mkBackendNotificationMetrics diff --git a/services/brig/brig.integration.yaml b/services/brig/brig.integration.yaml index 8be11f028bd..d8a8a9957e2 100644 --- a/services/brig/brig.integration.yaml +++ b/services/brig/brig.integration.yaml @@ -176,6 +176,7 @@ postgresMigration: teamFeatures: postgresql domainRegistration: postgresql user: postgresql + blockList: postgresql optSettings: setActivationTimeout: 4 diff --git a/services/brig/src/Brig/CanonicalInterpreter.hs b/services/brig/src/Brig/CanonicalInterpreter.hs index 05365ca28a8..de04b0620de 100644 --- a/services/brig/src/Brig/CanonicalInterpreter.hs +++ b/services/brig/src/Brig/CanonicalInterpreter.hs @@ -68,6 +68,8 @@ import Wire.BackgroundJobsPublisher (BackgroundJobPublisher) import Wire.BackgroundJobsPublisher.RabbitMQ (interpretBackgroundJobPublisherRabbitMQ) import Wire.BlockListStore import Wire.BlockListStore.Cassandra +import Wire.BlockListStore.DualWrite +import Wire.BlockListStore.Postgres import Wire.BudgetStore import Wire.BudgetStore.Cassandra import Wire.ClientStore (ClientStore) @@ -219,6 +221,7 @@ type BrigLowerLevelEffects = UserGroupStore, DomainRegistrationStore, DomainVerificationChallengeStore, + BlockListStore, Error AppSubsystemError, Error TeamCollaboratorsError, Error UsageError, @@ -261,7 +264,6 @@ type BrigLowerLevelEffects = FederationConfigStore, Jwk, JwtTools, - BlockListStore, BudgetStore, UserPendingActivationStore InternalPaging, Now, @@ -411,6 +413,11 @@ runBrigToIO e (AppT ma) = do PostgresqlStorage -> interpretDomainRegistrationStoreToPostgres MigrationToPostgresql -> interpretDomainRegistrationStoreToCassandraAndPostgres e.casClient + blockListStore = case e.postgresMigration.blockList of + CassandraStorage -> interpretBlockListStoreToCassandra e.casClient + PostgresqlStorage -> interpretBlockListStoreToPostgres + MigrationToPostgresql -> interpretBlockListStoreToCassandraAndPostgres e.casClient + domainVerificationChallengeStore = case e.postgresMigration.domainRegistration of CassandraStorage -> interpretDomainVerificationChallengeStoreToCassandra e.settings.challengeTTL PostgresqlStorage -> interpretDomainVerificationChallengeStoreToPostgres e.settings.challengeTTL @@ -449,7 +456,6 @@ runBrigToIO e (AppT ma) = do . nowToIOAction e.currentTime . userPendingActivationStoreToCassandra . budgetStoreToCassandra @Cas.Client - . interpretBlockListStoreToCassandra e.casClient . interpretJwtTools . interpretJwk . interpretFederationDomainConfig e.casClient e.settings.federationStrategy (foldMap (remotesMapFromCfgFile . fmap (.federationDomainConfig)) e.settings.federationDomainConfigs) @@ -492,6 +498,7 @@ runBrigToIO e (AppT ma) = do . mapError postgresUsageErrorToHttpError . mapError teamCollaboratorsSubsystemErrorToHttpError . mapError appSubsystemErrorToHttpError + . blockListStore . domainVerificationChallengeStore . domainRegistrationStore . interpretUserGroupStoreToPostgres diff --git a/services/galley/galley.integration.yaml b/services/galley/galley.integration.yaml index 34762762965..842c1b5ba18 100644 --- a/services/galley/galley.integration.yaml +++ b/services/galley/galley.integration.yaml @@ -268,3 +268,4 @@ postgresMigration: teamFeatures: postgresql domainRegistration: postgresql user: postgresql + blockList: postgresql