From 3c17e01c377a487581d169e0fa8f9e8a90e0c882 Mon Sep 17 00:00:00 2001 From: AsamK Date: Fri, 21 Aug 2026 21:39:51 +0200 Subject: [PATCH] Fix storage sync issues --- .../signal/manager/helper/StorageHelper.java | 9 +++- .../storage/recipients/RecipientStore.java | 42 +++++++++++++++---- .../storage/stickers/StickerStore.java | 2 +- .../syncStorage/ContactRecordProcessor.java | 24 +++++++++++ .../StickerPackRecordProcessor.java | 9 ++-- .../StorageRecordProcessorTest.java | 33 +++++++++++++++ 6 files changed, 106 insertions(+), 13 deletions(-) create mode 100644 lib/src/test/java/org/asamk/signal/manager/syncStorage/StorageRecordProcessorTest.java diff --git a/lib/src/main/java/org/asamk/signal/manager/helper/StorageHelper.java b/lib/src/main/java/org/asamk/signal/manager/helper/StorageHelper.java index a7aa7892..b47c97c7 100644 --- a/lib/src/main/java/org/asamk/signal/manager/helper/StorageHelper.java +++ b/lib/src/main/java/org/asamk/signal/manager/helper/StorageHelper.java @@ -452,7 +452,7 @@ public class StorageHelper { final var stickerPacks = account.getStickerStore() .getStickerPacks(connection) .stream() - .filter(pack -> pack.isInstalled() || pack.deletedTimestamp() > 0) + .filter(pack -> pack.storageId() != null) .toList(); newStickerPackStorageIds = generateStickerPackStorageIds(stickerPacks); for (final var stickerPack : stickerPacks) { @@ -716,6 +716,13 @@ public class StorageHelper { final var contactRecordProcessor = new ContactRecordProcessor(account, connection, context.getJobExecutor()); final var stickerPackRecordProcessor = new StickerPackRecordProcessor(account, connection); + final var contactRecords = records.stream() + .filter(record -> record.getProto().contact != null) + .map(record -> StorageRecordConvertersKt.toSignalContactRecord(record.getProto().contact, + record.getId())) + .toList(); + contactRecordProcessor.prepare(contactRecords); + for (final var record : records) { if (record.getProto().account != null) { logger.debug("Reading record {} of type account", record.getId()); diff --git a/lib/src/main/java/org/asamk/signal/manager/storage/recipients/RecipientStore.java b/lib/src/main/java/org/asamk/signal/manager/storage/recipients/RecipientStore.java index 15685223..97410474 100644 --- a/lib/src/main/java/org/asamk/signal/manager/storage/recipients/RecipientStore.java +++ b/lib/src/main/java/org/asamk/signal/manager/storage/recipients/RecipientStore.java @@ -923,11 +923,34 @@ public class RecipientStore implements RecipientIdCreator, RecipientResolver, Re statement.executeUpdate(); } if (contact != null && contact.unregisteredTimestamp() != null) { - markUnregisteredAndSplitIfNecessary(connection, recipientId); + markUnregisteredAndSplitIfNecessary(connection, recipientId, contact.unregisteredTimestamp()); } rotateStorageId(connection, recipientId); } + public void splitForStorageSyncIfNecessary(final Connection connection, final ACI aci) throws SQLException { + final var recipient = findByServiceId(connection, aci); + if (recipient.isEmpty()) { + return; + } + + final var recipientId = recipient.get().id(); + final var address = recipient.get().address(); + if (address.pni().isEmpty() && address.number().isEmpty()) { + return; + } + + logger.debug("Splitting {} for storage sync", recipientId); + final var splitAddress = new RecipientAddress(Optional.empty(), + address.pni(), + address.number(), + Optional.empty()); + updateRecipientAddress(connection, + recipientId, + new RecipientAddress(address.aci(), Optional.empty(), Optional.empty(), address.username())); + resolveRecipientTrusted(connection, splitAddress); + } + public int removeStorageIdsFromLocalOnlyUnregisteredRecipients( final Connection connection, final Collection storageIds @@ -1063,7 +1086,7 @@ public class RecipientStore implements RecipientIdCreator, RecipientResolver, Re if (recipientAddress.get().address().aci().isEmpty() || ( contact != null && contact.unregisteredTimestamp() != null )) { - markUnregisteredAndSplitIfNecessary(connection, recipientId); + markUnregisteredAndSplitIfNecessary(connection, recipientId, System.currentTimeMillis()); } } } @@ -1097,7 +1120,7 @@ public class RecipientStore implements RecipientIdCreator, RecipientResolver, Re if (registered) { markRegistered(connection, recipientId); } else { - markUnregisteredAndSplitIfNecessary(connection, recipientId); + markUnregisteredAndSplitIfNecessary(connection, recipientId, System.currentTimeMillis()); } connection.commit(); } catch (SQLException e) { @@ -1107,9 +1130,10 @@ public class RecipientStore implements RecipientIdCreator, RecipientResolver, Re private void markUnregisteredAndSplitIfNecessary( final Connection connection, - final RecipientId recipientId + final RecipientId recipientId, + final long unregisteredTimestamp ) throws SQLException { - markUnregistered(connection, recipientId); + markUnregistered(connection, recipientId, unregisteredTimestamp); final var address = resolveRecipientAddress(connection, recipientId); final var needSplit = address.aci().isPresent() && address.pni().isPresent(); logger.trace("Marking unregistered recipient {} as unregistered (and split={}): {}", @@ -1156,7 +1180,11 @@ public class RecipientStore implements RecipientIdCreator, RecipientResolver, Re } } - private void markUnregistered(final Connection connection, final RecipientId recipientId) throws SQLException { + private void markUnregistered( + final Connection connection, + final RecipientId recipientId, + final long unregisteredTimestamp + ) throws SQLException { final var sql = ( """ UPDATE %s @@ -1165,7 +1193,7 @@ public class RecipientStore implements RecipientIdCreator, RecipientResolver, Re """ ).formatted(TABLE_RECIPIENT); try (final var statement = connection.prepareStatement(sql)) { - statement.setLong(1, System.currentTimeMillis()); + statement.setLong(1, unregisteredTimestamp); statement.setLong(2, recipientId.id()); statement.executeUpdate(); } diff --git a/lib/src/main/java/org/asamk/signal/manager/storage/stickers/StickerStore.java b/lib/src/main/java/org/asamk/signal/manager/storage/stickers/StickerStore.java index c5289cd1..14fda15e 100644 --- a/lib/src/main/java/org/asamk/signal/manager/storage/stickers/StickerStore.java +++ b/lib/src/main/java/org/asamk/signal/manager/storage/stickers/StickerStore.java @@ -298,7 +298,7 @@ public class StickerStore { """ SELECT s.pack_id FROM %s s - WHERE s.storage_id IS NULL AND (s.installed = TRUE OR s.deleted_timestamp > 0) + WHERE s.storage_id IS NULL AND s.installed = TRUE """ ).formatted(TABLE_STICKER); final var updateSql = ( diff --git a/lib/src/main/java/org/asamk/signal/manager/syncStorage/ContactRecordProcessor.java b/lib/src/main/java/org/asamk/signal/manager/syncStorage/ContactRecordProcessor.java index f0101c8b..3829f8d2 100644 --- a/lib/src/main/java/org/asamk/signal/manager/syncStorage/ContactRecordProcessor.java +++ b/lib/src/main/java/org/asamk/signal/manager/syncStorage/ContactRecordProcessor.java @@ -24,6 +24,7 @@ import org.whispersystems.signalservice.internal.storage.protos.ContactRecord.Id import java.sql.Connection; import java.sql.SQLException; import java.util.Arrays; +import java.util.Collection; import java.util.Objects; import java.util.Optional; import java.util.regex.Pattern; @@ -56,6 +57,29 @@ public class ContactRecordProcessor extends DefaultStorageRecordProcessor remoteRecords) throws SQLException { + for (final var remoteRecord : remoteRecords) { + if (isInvalid(remoteRecord)) { + continue; + } + final var remote = remoteRecord.getProto(); + final var aci = ACI.parseOrNull(remote.aci, remote.aciBinary); + final var pni = PNI.parseOrNull(remote.pni, remote.pniBinary); + if (shouldSplitForStorageSync(remote.unregisteredAtTimestamp, aci, pni, remote.e164)) { + account.getRecipientStore().splitForStorageSyncIfNecessary(connection, aci); + } + } + } + + static boolean shouldSplitForStorageSync( + final long unregisteredAtTimestamp, + final ACI aci, + final PNI pni, + final String e164 + ) { + return unregisteredAtTimestamp > 0 && aci != null && pni == null && e164.isEmpty(); + } + /** * Error cases: * - You can't have a contact record without an ACI or PNI. diff --git a/lib/src/main/java/org/asamk/signal/manager/syncStorage/StickerPackRecordProcessor.java b/lib/src/main/java/org/asamk/signal/manager/syncStorage/StickerPackRecordProcessor.java index add19af4..52856f36 100644 --- a/lib/src/main/java/org/asamk/signal/manager/syncStorage/StickerPackRecordProcessor.java +++ b/lib/src/main/java/org/asamk/signal/manager/syncStorage/StickerPackRecordProcessor.java @@ -63,16 +63,17 @@ public class StickerPackRecordProcessor extends DefaultStorageRecordProcessor 0; - final var isLocalDeleted = local.deletedAtTimestamp > 0; - - if (isRemoteDeleted && isLocalDeleted && local.deletedAtTimestamp > remote.deletedAtTimestamp) { + if (shouldKeepLocalDeletion(remote.deletedAtTimestamp, local.deletedAtTimestamp)) { return localRecord; } return remoteRecord; } + static boolean shouldKeepLocalDeletion(final long remoteDeletedAt, final long localDeletedAt) { + return remoteDeletedAt > 0 && localDeletedAt > 0 && localDeletedAt < remoteDeletedAt; + } + @Override protected void insertLocal(final SignalStickerPackRecord record) throws SQLException { account.getStickerStore().upsertFromStorageSync(connection, record); diff --git a/lib/src/test/java/org/asamk/signal/manager/syncStorage/StorageRecordProcessorTest.java b/lib/src/test/java/org/asamk/signal/manager/syncStorage/StorageRecordProcessorTest.java new file mode 100644 index 00000000..2a0c6b75 --- /dev/null +++ b/lib/src/test/java/org/asamk/signal/manager/syncStorage/StorageRecordProcessorTest.java @@ -0,0 +1,33 @@ +package org.asamk.signal.manager.syncStorage; + +import org.junit.jupiter.api.Test; +import org.signal.core.models.ServiceId.ACI; +import org.signal.core.models.ServiceId.PNI; + +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class StorageRecordProcessorTest { + + @Test + void splitsOnlyUnregisteredAciOnlyRecords() { + final var aci = ACI.from(UUID.randomUUID()); + final var pni = PNI.from(UUID.randomUUID()); + + assertTrue(ContactRecordProcessor.shouldSplitForStorageSync(1, aci, null, "")); + assertFalse(ContactRecordProcessor.shouldSplitForStorageSync(0, aci, null, "")); + assertFalse(ContactRecordProcessor.shouldSplitForStorageSync(1, null, null, "")); + assertFalse(ContactRecordProcessor.shouldSplitForStorageSync(1, aci, pni, "")); + assertFalse(ContactRecordProcessor.shouldSplitForStorageSync(1, aci, null, "+12025550123")); + } + + @Test + void keepsOlderLocalStickerDeletion() { + assertTrue(StickerPackRecordProcessor.shouldKeepLocalDeletion(200, 100)); + assertFalse(StickerPackRecordProcessor.shouldKeepLocalDeletion(100, 200)); + assertFalse(StickerPackRecordProcessor.shouldKeepLocalDeletion(200, 0)); + assertFalse(StickerPackRecordProcessor.shouldKeepLocalDeletion(0, 100)); + } +} \ No newline at end of file