Add a command to map alternate forms of existing Gambian phone numbers to the same PNI

This commit is contained in:
Jon Chambers authored and Jon Chambers committed 2026-09-02 16:55:49 -04:00
1 parent 2f712b3f93
commit 77e3675fc9
5 files changed
+262 -5

No files matched your search

@@ -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<WhisperServerConfiguration
bootstrap.addCommand(new PopulateAccountRecoveryPasswordsCommand());
bootstrap.addCommand(new UpdateGambiaPniMappingsCommand());
ServiceLoader.load(SpamFilter.class)
.stream()
.map(ServiceLoader.Provider::get)
@@ -6,7 +6,6 @@ package org.whispersystems.textsecuregcm.storage;
import static java.util.Objects.requireNonNull;
import static org.whispersystems.textsecuregcm.metrics.MetricsUtil.name;
import static org.whispersystems.textsecuregcm.util.Util.getAlternateForms;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectWriter;
@@ -468,7 +467,7 @@ public class Accounts {
&& accountToCreate.getNumber().isPresent()) {
final String existingAccountNumber = existingAccount.getNumber().get();
final String accountToCreateNumber = accountToCreate.getNumber().get();
if (getAlternateForms(existingAccountNumber).contains(accountToCreateNumber)) {
if (Util.getAlternateForms(existingAccountNumber).contains(accountToCreateNumber)) {
final AttributeValue uuidAttr = AttributeValues.fromUUID(existingAccount.getAccountIdentifier());
final AttributeValue numberAttr = AttributeValues.fromString(accountToCreateNumber);
final TransactWriteItem phoneNumberConstraintPut = buildConstraintTablePutIfAbsent(
@@ -1094,7 +1093,7 @@ public class Accounts {
setClauses.add("#number = :number");
final MembershipExpression membershipExpression = maybeExpectedExistingE164
.map(e164 -> 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());
@@ -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<UUID> setPni(final String originalPhoneNumber, final List<String> allPhoneNumberForms,
public CompletableFuture<UUID> setPni(final String originalPhoneNumber,
final List<String> allPhoneNumberForms,
final UUID pni) {
if (!originalPhoneNumber.equals(allPhoneNumberForms.getFirst())) {
throw new IllegalArgumentException("allPhoneNumberForms must start with the target phoneNumber");
@@ -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<Account> 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();
}
}
@@ -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<String> 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());
}
}