diff --git a/app/src/main/java/org/thoughtcrime/securesms/jobs/StorageSyncJobV2.java b/app/src/main/java/org/thoughtcrime/securesms/jobs/StorageSyncJobV2.java
index 3c8f8e56f6dd70dd0a4a319aa9ac0f7885101080..3c3e2ce3ed13fd709e7b5ed3f76e7c020f0ad0cf 100644
--- a/app/src/main/java/org/thoughtcrime/securesms/jobs/StorageSyncJobV2.java
+++ b/app/src/main/java/org/thoughtcrime/securesms/jobs/StorageSyncJobV2.java
@@ -53,6 +53,7 @@ import org.whispersystems.signalservice.api.storage.StorageKey;
import org.whispersystems.signalservice.internal.storage.protos.ManifestRecord;
import java.io.IOException;
+import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
@@ -294,25 +295,37 @@ public class StorageSyncJobV2 extends BaseJob {
//noinspection unchecked Stop yelling at my beautiful method signatures
mergeWriteOperation = StorageSyncHelper.createWriteOperation(remoteManifest.get().getVersion(), localStorageIdsAfterMerge, contactResult, gv1Result, gv2Result, accountResult);
- List<StorageId> postMergeLocalOnlyIds = StorageSyncHelper.findKeyDifference(remoteManifest.get().getStorageIds(), mergeWriteOperation.getManifest().getStorageIds()).getLocalOnlyKeys();
- List<StorageId> remoteInsertIds = Stream.of(mergeWriteOperation.getInserts()).map(SignalStorageRecord::getId).toList();
- Set<StorageId> unhandledLocalOnlyIds = SetUtil.difference(postMergeLocalOnlyIds, remoteInsertIds);
+ KeyDifferenceResult postMergeKeyDifference = StorageSyncHelper.findKeyDifference(remoteManifest.get().getStorageIds(), mergeWriteOperation.getManifest().getStorageIds());
+ List<StorageId> postMergeLocalOnlyIds = postMergeKeyDifference.getLocalOnlyKeys();
+ List<ByteBuffer> postMergeRemoteOnlyIds = Stream.of(postMergeKeyDifference.getRemoteOnlyKeys()).map(StorageId::getRaw).map(ByteBuffer::wrap).toList();
+ List<StorageId> remoteInsertIds = Stream.of(mergeWriteOperation.getInserts()).map(SignalStorageRecord::getId).toList();
+ List<ByteBuffer> remoteDeleteIds = Stream.of(mergeWriteOperation.getDeletes()).map(ByteBuffer::wrap).toList();
+ Set<StorageId> unhandledLocalOnlyIds = SetUtil.difference(postMergeLocalOnlyIds, remoteInsertIds);
+ Set<ByteBuffer> unhandledRemoteOnlyIds = SetUtil.difference(postMergeRemoteOnlyIds, remoteDeleteIds);
if (unhandledLocalOnlyIds.size() > 0) {
Log.i(TAG, "[Remote Sync] After the conflict resolution, there are " + unhandledLocalOnlyIds.size() + " local-only records remaining that weren't otherwise inserted. Adding them as inserts.");
- List<SignalStorageRecord> localOnly = buildLocalStorageRecords(context, self, unhandledLocalOnlyIds);
+ List<SignalStorageRecord> unhandledInserts = buildLocalStorageRecords(context, self, unhandledLocalOnlyIds);
mergeWriteOperation = new WriteOperationResult(mergeWriteOperation.getManifest(),
- Util.concatenatedList(mergeWriteOperation.getInserts(), localOnly),
+ Util.concatenatedList(mergeWriteOperation.getInserts(), unhandledInserts),
mergeWriteOperation.getDeletes());
recipientDatabase.clearDirtyStateForStorageIds(unhandledLocalOnlyIds);
- } else {
- Log.i(TAG, "[Remote Sync] After the conflict resolution, there are no local-only records remaining.");
}
- StorageSyncValidations.validate(mergeWriteOperation, remoteManifest, needsForcePush);
+ if (unhandledRemoteOnlyIds.size() > 0) {
+ Log.i(TAG, "[Remote Sync] After the conflict resolution, there are " + unhandledRemoteOnlyIds.size() + " remote-only records remaining that weren't otherwise deleted. Adding them as deletes.");
+
+ List<byte[]> unhandledDeletes = Stream.of(unhandledRemoteOnlyIds).map(ByteBuffer::array).toList();
+
+ mergeWriteOperation = new WriteOperationResult(mergeWriteOperation.getManifest(),
+ mergeWriteOperation.getInserts(),
+ Util.concatenatedList(mergeWriteOperation.getDeletes(), unhandledDeletes));
+
+ }
+
db.setTransactionSuccessful();
} finally {
db.endTransaction();
@@ -323,10 +336,7 @@ public class StorageSyncJobV2 extends BaseJob {
Log.i(TAG, "[Remote Sync] WriteOperationResult :: " + mergeWriteOperation);
Log.i(TAG, "[Remote Sync] We have something to write remotely.");
- if (mergeWriteOperation.getManifest().getStorageIds().size() != remoteManifest.get().getStorageIds().size() + mergeWriteOperation.getInserts().size() - mergeWriteOperation.getDeletes().size()) {
- Log.w(TAG, String.format(Locale.US, "[Remote Sync] Bad storage key management! originalRemoteKeys: %d, newRemoteKeys: %d, insertedKeys: %d, deletedKeys: %d",
- remoteManifest.get().getStorageIds().size(), mergeWriteOperation.getManifest().getStorageIds().size(), mergeWriteOperation.getInserts().size(), mergeWriteOperation.getDeletes().size()));
- }
+ StorageSyncValidations.validate(mergeWriteOperation, remoteManifest, needsForcePush);
Optional<SignalStorageManifest> conflict = accountManager.writeStorageRecords(storageServiceKey, mergeWriteOperation.getManifest(), mergeWriteOperation.getInserts(), mergeWriteOperation.getDeletes());
@@ -372,7 +382,6 @@ public class StorageSyncJobV2 extends BaseJob {
Log.i(TAG, String.format(Locale.ENGLISH, "[Local Changes] Local changes present. %d updates, %d inserts, %d deletes, account update: %b, account insert: %b.", pendingUpdates.size(), pendingInsertions.size(), pendingDeletions.size(), pendingAccountUpdate.isPresent(), pendingAccountInsert.isPresent()));
WriteOperationResult localWrite = localWriteResult.get().getWriteResult();
- StorageSyncValidations.validate(localWrite, remoteManifest, needsForcePush);
Log.i(TAG, "[Local Changes] WriteOperationResult :: " + localWrite);
@@ -380,6 +389,8 @@ public class StorageSyncJobV2 extends BaseJob {
throw new AssertionError("Decided there were local writes, but our write result was empty!");
}
+ StorageSyncValidations.validate(localWrite, remoteManifest, needsForcePush);
+
Optional<SignalStorageManifest> conflict = accountManager.writeStorageRecords(storageServiceKey, localWrite.getManifest(), localWrite.getInserts(), localWrite.getDeletes());
if (conflict.isPresent()) {