Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@ import java.time.format.DateTimeFormatter
import javax.inject.Inject
import javax.inject.Singleton
import kotlin.time.Duration
import kotlin.time.Duration.Companion.milliseconds

@Singleton
class FileServerApi @Inject constructor(
Expand Down Expand Up @@ -63,7 +62,7 @@ class FileServerApi @Inject constructor(
val useOnionRouting: Boolean = true,

// Computed fresh for each attempt (after clock resync etc.)
val dynamicHeaders: (suspend () -> Map<String, String>)? = null
val dynamicHeaders: (() -> Map<String, String>)? = null
)

private fun createBody(body: ByteArray?, parameters: Any?): RequestBody? {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -269,7 +269,7 @@ object OpenGroupApi {
* this when running over Lokinet.
*/
val useOnionRouting: Boolean = true,
val dynamicHeaders: (suspend () -> Map<String, String>)? = null
val dynamicHeaders: (() -> Map<String, String>)? = null
)

private fun createBody(body: ByteArray?, parameters: Any?): RequestBody? {
Expand Down

This file was deleted.

50 changes: 36 additions & 14 deletions app/src/main/java/org/session/libsession/network/ServerClient.kt
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@ package org.session.libsession.network

import okhttp3.Request
import org.session.libsession.network.model.OnionDestination
import org.session.libsession.network.model.OnionError
import org.session.libsession.network.model.OnionResponse
import org.session.libsession.network.onion.Version
import org.session.libsession.network.utilities.getBodyForOnionRequest
Expand All @@ -23,33 +22,32 @@ class ServerClient @Inject constructor(
) {

/**
* The request is sent as a lambda in order to be recalculated as part of the retry strategy.
* This is useful for things like timestamps that might have been updated
* as part of a clock resync
* Send a request to a server destination. The request is sent as a lambda in order to be
* recalculated as part of the retry strategy.
*
* This version of the method returns both the data generated by the requestFactory
* and the OnionResponse from the network send.
*/
suspend fun send(
requestFactory: suspend () -> Request,
suspend fun <T> sendWithData(
requestFactory: suspend () -> Pair<T, Request>,
serverBaseUrl: String,
x25519PublicKey: String,
version: Version = Version.V4,
operationName: String = "ServerClient.send",
): OnionResponse {
val initialRequest = requestFactory()
val url = initialRequest.url

): Pair<T, OnionResponse> {
return retryWithBackOff(
operationName = operationName,
classifier = { error, previous ->
errorManager.onFailure(
error = error,
ctx = ServerClientFailureContext(
url = url,
url = serverBaseUrl,
previousError = previous
)
)
}
) { attempt ->
val request = if (attempt == 1) initialRequest else requestFactory()
) { _ ->
val (data, request) = requestFactory()
val url = request.url

val destination = OnionDestination.ServerDestination(
Expand All @@ -62,7 +60,7 @@ class ServerClient @Inject constructor(

val payload = generatePayload(request, serverBaseUrl, version)

sessionNetwork.sendWithRetry(
data to sessionNetwork.sendWithRetry(
destination = destination,
payload = payload,
version = version,
Expand All @@ -73,6 +71,30 @@ class ServerClient @Inject constructor(
}
}

/**
* The request is sent as a lambda in order to be recalculated as part of the retry strategy.
* This is useful for things like timestamps that might have been updated
* as part of a clock resync
*/
suspend fun send(
requestFactory: suspend () -> Request,
serverBaseUrl: String,
x25519PublicKey: String,
version: Version = Version.V4,
operationName: String = "ServerClient.send",
): OnionResponse {
return sendWithData(
requestFactory = {
val request = requestFactory()
Unit to request
},
serverBaseUrl = serverBaseUrl,
x25519PublicKey = x25519PublicKey,
version = version,
operationName = operationName
).second
}

private fun generatePayload(request: Request, server: String, version: Version): ByteArray {
val headers = request.getHeadersForOnionRequest().toMutableMap()
val url = request.url
Expand Down
Original file line number Diff line number Diff line change
@@ -1,14 +1,8 @@
package org.session.libsession.network

import okhttp3.HttpUrl
import org.session.libsession.network.model.FailureDecision
import org.session.libsession.network.model.OnionDestination
import org.session.libsession.network.model.OnionError
import org.session.libsession.network.onion.PathManager
import org.session.libsession.network.snode.SnodeDirectory
import org.session.libsession.network.snode.SwarmDirectory
import org.session.libsignal.utilities.Log
import org.session.libsignal.utilities.Snode
import javax.inject.Inject
import javax.inject.Singleton

Expand Down Expand Up @@ -53,6 +47,6 @@ class ServerClientErrorManager @Inject constructor(
}

data class ServerClientFailureContext(
val url: HttpUrl,
val url: String,
val previousError: OnionError? = null // in some situations we could be coming from a retry to a previous error
)
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ import network.loki.messenger.libsession_util.util.UserPic
import org.session.libsession.avatars.AvatarCacheCleaner
import org.session.libsession.database.StorageProtocol
import org.session.libsession.messaging.sending_receiving.notifications.MessageNotifier
import org.session.libsession.messaging.sending_receiving.notifications.PushRegistryV1
import org.session.libsession.network.SnodeClient
import org.session.libsession.network.SnodeClock
import org.session.libsession.snode.OwnedSwarmAuth
Expand Down Expand Up @@ -205,8 +204,6 @@ class ConfigToDatabaseSync @Inject constructor(
// Store the encryption key pair
val keyPair = ECKeyPair(DjbECPublicKey(group.encPubKey.data), DjbECPrivateKey(group.encSecKey.data))
storage.addClosedGroupEncryptionKeyPair(keyPair, group.accountId, clock.currentTimeMillis())
// Notify the PN server
PushRegistryV1.subscribeGroup(group.accountId, publicKey = myAccountId.hexString)
threadDatabase.setCreationDate(threadId, formationTimestamp)
}

Expand Down Expand Up @@ -256,8 +253,6 @@ class ConfigToDatabaseSync @Inject constructor(
// Remove the key pairs
storage.removeAllClosedGroupEncryptionKeyPairs(address.groupPublicKeyHex)
storage.removeMember(address.address, myAccountId.toAddress())
// Notify the PN server
PushRegistryV1.unsubscribeGroup(closedGroupPublicKey = address.groupPublicKeyHex, publicKey = myAccountId.hexString)
messageNotifier.updateNotification(context)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,11 +70,11 @@ class GroupLeavingWorker @AssistedInject constructor(

if (groupAuth != null) {
val resp = pushRegistryV2.unregister {
listOf(pushRegistryV2.buildUnregisterRequest(currentToken, groupAuth))
listOf(runCatching { pushRegistryV2.buildUnregisterRequest(currentToken, groupAuth) })
}.firstOrNull()

check(resp?.success == true) {
"Unsubscription failed: code = ${resp?.error}, message = ${resp?.message}"
check(resp?.getOrNull()?.success == true) {
"Unsubscription failed: $resp"
}
Log.d(TAG, "Unsubscribed from group $groupId successfully")
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -170,59 +170,39 @@ class PushRegistrationWorker @AssistedInject constructor(
private suspend inline fun <T, Req, Res : Response> batchRequest(
items: List<T>,
crossinline buildRequest: (T) -> Req,
crossinline sendBatchRequest: suspend (suspend () -> Collection<Req>) -> List<Res>,
sendBatchRequest: suspend (requestsBuilder: () -> Collection<kotlin.Result<Req>>) -> List<kotlin.Result<Res>>,
): List<Pair<T, kotlin.Result<Unit>>> {
if (items.isEmpty()) return emptyList()

val results = ArrayList<Pair<T, kotlin.Result<Unit>>>(items.size)

// Items that are valid to send, and their per-attempt builders
val batchItems = mutableListOf<T>()
val requestBuilders = mutableListOf<() -> Req>()

for (item in items) {
try {
//todo ONION I have to double the buildRequest here, once for validation and again to recompute... Is this ok? FANCHAO
buildRequest(item)

batchItems += item
requestBuilders += { buildRequest(item) } // <- rebuilt each retry attempt
} catch (ec: Exception) {
results += item to kotlin.Result.failure(
NonRetryableException("Failed to build a request", ec)
)
}
}

if (batchItems.isEmpty()) return results

try {
return try {
val responses = sendBatchRequest {
requestBuilders.map { it() }
items.map { item ->
try {
kotlin.Result.success(buildRequest(item))
} catch (e: Exception) {
kotlin.Result.failure(NonRetryableException("Error building request", e))
}
}
}

responses.forEachIndexed { idx, response ->
val item = batchItems[idx]
results += item to when {
response.isSuccess() -> kotlin.Result.success(Unit)
response.error == 403 -> kotlin.Result.failure(
NonRetryableException("Request failed: code = ${response.error}, message = ${response.message}")
)
else -> kotlin.Result.failure(
RuntimeException("Request failed: code = ${response.error}, message = ${response.message}")
)
responses.mapIndexed { idx, result ->
val item = items[idx]
item to result.map { response ->
when {
response.isSuccess() -> Unit
response.error == 403 -> throw NonRetryableException("Request failed: code = ${response.error}, message = ${response.message}")
else -> throw RuntimeException("Request failed: code = ${response.error}, message = ${response.message}")
}
}
}
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
// Batch call failed -> mark all *sent* items as failed
batchItems.forEach { item ->
results += item to kotlin.Result.failure(e)
items.map { item ->
item to kotlin.Result.failure(e)
}
}

return results
}


Expand Down
Loading