diff --git a/core/src/main/java/com/segment/analytics/kotlin/core/Settings.kt b/core/src/main/java/com/segment/analytics/kotlin/core/Settings.kt index a0efb34c..3a92235c 100644 --- a/core/src/main/java/com/segment/analytics/kotlin/core/Settings.kt +++ b/core/src/main/java/com/segment/analytics/kotlin/core/Settings.kt @@ -6,6 +6,7 @@ import com.segment.analytics.kotlin.core.platform.plugins.logger.log import com.segment.analytics.kotlin.core.utilities.LenientJson import com.segment.analytics.kotlin.core.utilities.safeJsonObject import kotlinx.coroutines.launch +import kotlinx.coroutines.runInterruptible import kotlinx.coroutines.withContext import kotlinx.serialization.DeserializationStrategy import kotlinx.serialization.Serializable @@ -89,7 +90,7 @@ suspend fun Analytics.checkSettings() { val settingsObj = withContext(networkIODispatcher) { log("Fetching settings on ${Thread.currentThread().name}") - return@withContext fetchSettings(writeKey, cdnHost) + return@withContext runInterruptible { fetchSettings(writeKey, cdnHost) } } settingsObj?.let { diff --git a/core/src/main/java/com/segment/analytics/kotlin/core/Telemetry.kt b/core/src/main/java/com/segment/analytics/kotlin/core/Telemetry.kt index 6e14b6fd..8fddc0ad 100644 --- a/core/src/main/java/com/segment/analytics/kotlin/core/Telemetry.kt +++ b/core/src/main/java/com/segment/analytics/kotlin/core/Telemetry.kt @@ -263,15 +263,19 @@ object Telemetry: Subscriber { // We're using this to leave off the 'log' parameter if unset. val payload = Json.encodeToString(mapOf("series" to sendQueue)) - val connection = httpClient.upload(host) - connection.outputStream?.use { outputStream -> - // Write the JSON string to the outputStream. - outputStream.write(payload.toByteArray(Charsets.UTF_8)) - outputStream.flush() // Ensure all data is written + runBlocking { + runInterruptible { + val connection = httpClient.upload(host) + connection.outputStream?.use { outputStream -> + // Write the JSON string to the outputStream. + outputStream.write(payload.toByteArray(Charsets.UTF_8)) + outputStream.flush() // Ensure all data is written + } + connection.inputStream?.close() + connection.outputStream?.close() + connection.close() + } } - connection.inputStream?.close() - connection.outputStream?.close() - connection.close() } catch (e: HTTPException) { errorHandler?.invoke(e) if (e.responseCode == 429) { diff --git a/core/src/main/java/com/segment/analytics/kotlin/core/platform/EventPipeline.kt b/core/src/main/java/com/segment/analytics/kotlin/core/platform/EventPipeline.kt index b966291c..9a12085f 100644 --- a/core/src/main/java/com/segment/analytics/kotlin/core/platform/EventPipeline.kt +++ b/core/src/main/java/com/segment/analytics/kotlin/core/platform/EventPipeline.kt @@ -10,6 +10,7 @@ import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED import kotlinx.coroutines.channels.consumeEach import kotlinx.coroutines.launch +import kotlinx.coroutines.runInterruptible import kotlinx.coroutines.withContext import kotlinx.serialization.encodeToString import kotlinx.serialization.json.Json @@ -135,14 +136,16 @@ open class EventPipeline( var shouldCleanup = true storage.readAsStream(url)?.use { data -> try { - val connection = httpClient.upload(apiHost) - connection.outputStream?.let { - // Write the payloads into the OutputStream - data.copyTo(connection.outputStream) - connection.outputStream.close() - - // Upload the payloads. - connection.close() + runInterruptible { + val connection = httpClient.upload(apiHost) + connection.outputStream?.let { + // Write the payloads into the OutputStream + data.copyTo(connection.outputStream) + connection.outputStream.close() + + // Upload the payloads. + connection.close() + } } // Cleanup uploaded payloads analytics.log("$logTag uploaded $url")