From 77e3675fc99c2875ee69cafe9be0aa2f659a2cde Mon Sep 17 00:00:00 2001 From: Jon Chambers Date: Wed, 2 Sep 2026 11:55:42 -0400 Subject: [PATCH] Add a command to map alternate forms of existing Gambian phone numbers to the same PNI --- .../textsecuregcm/WhisperServerService.java | 3 + .../textsecuregcm/storage/Accounts.java | 5 +- .../storage/PhoneNumberIdentifiers.java | 5 +- .../UpdateGambiaPniMappingsCommand.java | 43 ++++ .../UpdateGambiaPniMappingsCommandTest.java | 211 ++++++++++++++++++ 5 files changed, 262 insertions(+), 5 deletions(-) create mode 100644 service/src/main/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommand.java create mode 100644 service/src/test/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommandTest.java diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/WhisperServerService.java b/service/src/main/java/org/whispersystems/textsecuregcm/WhisperServerService.java index 50a0406bf..6d75dd06a 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/WhisperServerService.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/WhisperServerService.java @@ -348,6 +348,7 @@ import org.whispersystems.textsecuregcm.workers.SetUserDiscoverabilityCommand; import org.whispersystems.textsecuregcm.workers.TrimOversizedFoundationDbMessageQueuesCommand; import org.whispersystems.textsecuregcm.workers.UnlinkDeviceCommand; import org.whispersystems.textsecuregcm.workers.UnlinkDevicesWithIdlePrimaryCommand; +import org.whispersystems.textsecuregcm.workers.UpdateGambiaPniMappingsCommand; import org.whispersystems.textsecuregcm.workers.ZkParamsCommand; import org.whispersystems.websocket.WebSocketResourceProviderFactory; import org.whispersystems.websocket.setup.WebSocketEnvironment; @@ -416,6 +417,8 @@ public class WhisperServerService extends Application MembershipExpression.build(getAlternateForms(e164))) + .map(e164 -> MembershipExpression.build(Util.getAlternateForms(e164))) .orElseThrow(() -> new IllegalArgumentException("E164 must be present on existing account")); attrValues.putAll(membershipExpression.values()); diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/storage/PhoneNumberIdentifiers.java b/service/src/main/java/org/whispersystems/textsecuregcm/storage/PhoneNumberIdentifiers.java index 3604e31dd..fde7fe5a4 100644 --- a/service/src/main/java/org/whispersystems/textsecuregcm/storage/PhoneNumberIdentifiers.java +++ b/service/src/main/java/org/whispersystems/textsecuregcm/storage/PhoneNumberIdentifiers.java @@ -43,7 +43,7 @@ public class PhoneNumberIdentifiers { private final String tableName; @VisibleForTesting - static final String KEY_E164 = "P"; + public static final String KEY_E164 = "P"; @VisibleForTesting static final String INDEX_NAME = "pni_to_p"; @VisibleForTesting @@ -183,7 +183,8 @@ public class PhoneNumberIdentifiers { * @return The provided PNI if the update occurred, or the existing PNI associated with originalPhoneNumber */ @VisibleForTesting - CompletableFuture setPni(final String originalPhoneNumber, final List allPhoneNumberForms, + public CompletableFuture setPni(final String originalPhoneNumber, + final List allPhoneNumberForms, final UUID pni) { if (!originalPhoneNumber.equals(allPhoneNumberForms.getFirst())) { throw new IllegalArgumentException("allPhoneNumberForms must start with the target phoneNumber"); diff --git a/service/src/main/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommand.java b/service/src/main/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommand.java new file mode 100644 index 000000000..0530a97ec --- /dev/null +++ b/service/src/main/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommand.java @@ -0,0 +1,43 @@ +/* + * Copyright 2026 Signal Messenger, LLC + * SPDX-License-Identifier: AGPL-3.0-only + */ + +package org.whispersystems.textsecuregcm.workers; + +import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.Metrics; +import org.whispersystems.textsecuregcm.metrics.MetricsUtil; +import org.whispersystems.textsecuregcm.storage.Account; +import org.whispersystems.textsecuregcm.storage.AccountsManager; +import org.whispersystems.textsecuregcm.storage.PhoneNumberIdentifiers; +import org.whispersystems.textsecuregcm.util.Util; +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; + +public class UpdateGambiaPniMappingsCommand extends AbstractSinglePassCrawlAccountsCommand { + + public UpdateGambiaPniMappingsCommand() { + super("update-gambia-pni-mappings", "Maps legacy Gambian phone numbers and their new forms to the same PNI"); + } + + @Override + protected void crawlAccounts(final Flux accounts) { + final PhoneNumberIdentifiers phoneNumberIdentifiers = getCommandDependencies().phoneNumberIdentifiers(); + + final Counter updatedPniMappingCounter = Metrics.counter(MetricsUtil.name(getClass(), "updatedPniMapping")); + + accounts + .filter(account -> "GM".equalsIgnoreCase(Util.getRegion(account))) + .flatMap(accountWithGambianNumber -> { + final String e164 = accountWithGambianNumber.getNumber().orElseThrow(); + + return Mono.fromFuture(() -> phoneNumberIdentifiers.setPni(e164, + Util.getAlternateForms(e164), + accountWithGambianNumber.getPhoneNumberIdentifier().orElseThrow())); + }) + .doOnNext(_ -> updatedPniMappingCounter.increment()) + .then() + .block(); + } +} diff --git a/service/src/test/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommandTest.java b/service/src/test/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommandTest.java new file mode 100644 index 000000000..89742468f --- /dev/null +++ b/service/src/test/java/org/whispersystems/textsecuregcm/workers/UpdateGambiaPniMappingsCommandTest.java @@ -0,0 +1,211 @@ +/* + * Copyright 2026 Signal Messenger, LLC + * SPDX-License-Identifier: AGPL-3.0-only + */ + +package org.whispersystems.textsecuregcm.workers; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; + +import java.nio.charset.StandardCharsets; +import java.time.Clock; +import java.time.Duration; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.RegisterExtension; +import org.whispersystems.textsecuregcm.auth.DisconnectionRequestManager; +import org.whispersystems.textsecuregcm.redis.FaultTolerantRedisClient; +import org.whispersystems.textsecuregcm.redis.RedisClusterExtension; +import org.whispersystems.textsecuregcm.securestorage.SecureStorageClient; +import org.whispersystems.textsecuregcm.securevaluerecovery.SecureValueRecoveryClient; +import org.whispersystems.textsecuregcm.storage.AccountLockManager; +import org.whispersystems.textsecuregcm.storage.Accounts; +import org.whispersystems.textsecuregcm.storage.AccountsManager; +import org.whispersystems.textsecuregcm.storage.ChangeNumberWaitingPeriodManager; +import org.whispersystems.textsecuregcm.storage.DynamoDbExtension; +import org.whispersystems.textsecuregcm.storage.DynamoDbExtensionSchema; +import org.whispersystems.textsecuregcm.storage.KeysManager; +import org.whispersystems.textsecuregcm.storage.MessagesManager; +import org.whispersystems.textsecuregcm.storage.PhoneNumberIdentifiers; +import org.whispersystems.textsecuregcm.storage.PhoneNumberRecoveryPasswords; +import org.whispersystems.textsecuregcm.storage.PhoneNumberRecoveryPasswordsManager; +import org.whispersystems.textsecuregcm.storage.ProfilesManager; +import org.whispersystems.textsecuregcm.storage.RedeemedReceiptsManager; +import org.whispersystems.textsecuregcm.storage.S3LocalStackExtension; +import org.whispersystems.textsecuregcm.tests.util.AccountsHelper; +import org.whispersystems.textsecuregcm.util.AttributeValues; +import org.whispersystems.textsecuregcm.util.Util; +import reactor.core.scheduler.Schedulers; +import software.amazon.awssdk.services.dynamodb.model.DeleteItemRequest; +import software.amazon.awssdk.services.dynamodb.model.GetItemRequest; + +class UpdateGambiaPniMappingsCommandTest { + + @RegisterExtension + static final DynamoDbExtension DYNAMO_DB_EXTENSION = new DynamoDbExtension( + DynamoDbExtensionSchema.Tables.ACCOUNTS, + DynamoDbExtensionSchema.Tables.DELETED_ACCOUNTS, + DynamoDbExtensionSchema.Tables.DELETED_ACCOUNTS_LOCK, + DynamoDbExtensionSchema.Tables.NUMBERS, + DynamoDbExtensionSchema.Tables.PNI, + DynamoDbExtensionSchema.Tables.PNI_ASSIGNMENTS, + DynamoDbExtensionSchema.Tables.USERNAMES, + DynamoDbExtensionSchema.Tables.PHONE_NUMBER_RECOVERY_PASSWORDS, + DynamoDbExtensionSchema.Tables.REDEEMED_RECEIPTS); + + @RegisterExtension + static final RedisClusterExtension CACHE_CLUSTER_EXTENSION = RedisClusterExtension.builder().build(); + + @RegisterExtension + static final S3LocalStackExtension S3_EXTENSION = new S3LocalStackExtension("testbucket"); + + private ScheduledExecutorService executor; + private AccountsManager accountsManager; + + private TestUpdateGambiaPniMappingsCommand updateGambiaPniMappingsCommand; + + private static class TestUpdateGambiaPniMappingsCommand extends UpdateGambiaPniMappingsCommand { + + private final PhoneNumberIdentifiers phoneNumberIdentifiers; + + private TestUpdateGambiaPniMappingsCommand(final PhoneNumberIdentifiers phoneNumberIdentifiers) { + this.phoneNumberIdentifiers = phoneNumberIdentifiers; + } + + @Override + protected CommandDependencies getCommandDependencies() { + return new CommandDependencies(null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + null, + phoneNumberIdentifiers, + null, + null); + } + } + + @BeforeEach + void setup() { + final Accounts accounts = new Accounts( + Clock.systemUTC(), + DYNAMO_DB_EXTENSION.getDynamoDbClient(), + DYNAMO_DB_EXTENSION.getDynamoDbAsyncClient(), + new RedeemedReceiptsManager(Clock.systemUTC(), DynamoDbExtensionSchema.Tables.REDEEMED_RECEIPTS.tableName(), + DYNAMO_DB_EXTENSION.getDynamoDbClient()), + DynamoDbExtensionSchema.Tables.ACCOUNTS.tableName(), + DynamoDbExtensionSchema.Tables.NUMBERS.tableName(), + DynamoDbExtensionSchema.Tables.PNI_ASSIGNMENTS.tableName(), + DynamoDbExtensionSchema.Tables.USERNAMES.tableName(), + DynamoDbExtensionSchema.Tables.DELETED_ACCOUNTS.tableName(), + DynamoDbExtensionSchema.Tables.USED_LINK_DEVICE_TOKENS.tableName()); + + executor = Executors.newSingleThreadScheduledExecutor(); + + final AccountLockManager accountLockManager = new AccountLockManager(DYNAMO_DB_EXTENSION.getDynamoDbClient(), + DynamoDbExtensionSchema.Tables.DELETED_ACCOUNTS_LOCK.tableName()); + + final PhoneNumberIdentifiers phoneNumberIdentifiers = + new PhoneNumberIdentifiers(DYNAMO_DB_EXTENSION.getDynamoDbAsyncClient(), + DynamoDbExtensionSchema.Tables.PNI.tableName()); + + final PhoneNumberRecoveryPasswords phoneNumberRecoveryPasswords = + new PhoneNumberRecoveryPasswords(DynamoDbExtensionSchema.Tables.PHONE_NUMBER_RECOVERY_PASSWORDS.tableName(), + Duration.ofDays(1), + DYNAMO_DB_EXTENSION.getDynamoDbClient(), + Clock.systemUTC()); + + final PhoneNumberRecoveryPasswordsManager phoneNumberRecoveryPasswordsManager = + new PhoneNumberRecoveryPasswordsManager(phoneNumberRecoveryPasswords); + + accountsManager = new AccountsManager( + accounts, + phoneNumberIdentifiers, + CACHE_CLUSTER_EXTENSION.getRedisCluster(), + mock(FaultTolerantRedisClient.class), + accountLockManager, + mock(KeysManager.class), + mock(MessagesManager.class), + mock(ProfilesManager.class), + mock(ChangeNumberWaitingPeriodManager.class), + mock(SecureStorageClient.class), + mock(SecureValueRecoveryClient.class), + mock(DisconnectionRequestManager.class), + phoneNumberRecoveryPasswordsManager, + executor, + executor, + executor, + mock(Clock.class), + "link-device-secret".getBytes(StandardCharsets.UTF_8), + AccountsManager.TOTP.getTimeStep().dividedBy(2)); + + updateGambiaPniMappingsCommand = new TestUpdateGambiaPniMappingsCommand(phoneNumberIdentifiers); + } + + @AfterEach + void tearDown() throws InterruptedException { + executor.shutdown(); + + //noinspection ResultOfMethodCallIgnored + executor.awaitTermination(1, TimeUnit.SECONDS); + } + + @Test + void crawlAccounts() { + final String legacyGambianNumber = "+2203123456"; + final Set alternateForms = new HashSet<>(Util.getAlternateForms(legacyGambianNumber)); + alternateForms.remove(legacyGambianNumber); + + assertFalse(alternateForms.isEmpty()); + + AccountsHelper.createAccount(accountsManager, legacyGambianNumber); + + // To simulate existing numbers with only a single form mapped to a PNI, artificially remove alternate PNI mappings + alternateForms.forEach(e164 -> DYNAMO_DB_EXTENSION.getDynamoDbClient().deleteItem(DeleteItemRequest.builder() + .tableName(DynamoDbExtensionSchema.Tables.PNI.tableName()) + .key(Map.of(PhoneNumberIdentifiers.KEY_E164, AttributeValues.fromString(e164))) + .build())); + + alternateForms.forEach(e164 -> assertFalse(DYNAMO_DB_EXTENSION.getDynamoDbClient().getItem(GetItemRequest.builder() + .tableName(DynamoDbExtensionSchema.Tables.PNI.tableName()) + .key(Map.of(PhoneNumberIdentifiers.KEY_E164, AttributeValues.fromString(e164))) + .build()) + .hasItem())); + + updateGambiaPniMappingsCommand.crawlAccounts(accountsManager.streamAllFromDynamo(1, Schedulers.boundedElastic())); + + alternateForms.forEach(e164 -> assertTrue(DYNAMO_DB_EXTENSION.getDynamoDbClient().getItem(GetItemRequest.builder() + .tableName(DynamoDbExtensionSchema.Tables.PNI.tableName()) + .key(Map.of(PhoneNumberIdentifiers.KEY_E164, AttributeValues.fromString(e164))) + .build()) + .hasItem())); + + assertTrue(accountsManager.getByE164(legacyGambianNumber).isPresent()); + } +}