Dedupe native transfer/notification code and split ServerProvider into focused controllers
Build APK / build (push) Successful in 5m20s
Build APK / build (push) Successful in 5m20s
Consolidates duplicated GET/PUT/notification-channel logic across DownloadService/ShareUploadService/SyncEngine into shared Kotlin helpers, gives upload/download real batch queueing instead of dropping a second concurrent batch, and dedupes repeated Dart channel-argument boilerplate. Replaces the 2300+ line ServerProvider god object with ten focused ChangeNotifiers (SessionController, SettingsController, FilesController, PhotosController, FavoritesController, TrashController, SharesController, RecentController, SyncStatusController, PickController) plus ItemOperations, a plain coordinator for cross-domain item mutations - fixing the coupling where device-sync status, per-tab data, and global UI prefs all lived in one object. Updates every view/widget call site accordingly and refreshes the architecture/server/standards/styling docs to match. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,140 @@
|
||||
package dev.ayushya.noo
|
||||
|
||||
import java.io.File
|
||||
import java.io.OutputStream
|
||||
import java.net.HttpURLConnection
|
||||
import java.net.URL
|
||||
|
||||
/**
|
||||
* Shared low-level WebDAV GET/PUT primitives (plain `HttpURLConnection`, no
|
||||
* new HTTP dependency) - used by [DownloadService], [ShareUploadService],
|
||||
* and [SyncEngine], which each independently re-implemented this same
|
||||
* ~20-line connect/stream/disconnect shape before this existed. Progress
|
||||
* reporting and cancellation are optional: the two foreground Services
|
||||
* pass both (user-visible, cancellable transfers); [SyncEngine]'s
|
||||
* background callers pass neither.
|
||||
*
|
||||
* Deliberately does *not* know about `OCS-APIRequest` - that header is for
|
||||
* OCS REST endpoints (shares/activity/quota), not plain WebDAV GET/PUT,
|
||||
* and was previously sent by mistake on some of these calls; dropped here
|
||||
* so every caller gets the same (correct) minimal header set.
|
||||
*/
|
||||
object DavTransfer {
|
||||
const val CONNECT_TIMEOUT_MS = 15000
|
||||
const val READ_TIMEOUT_MS = 30000
|
||||
const val UPLOAD_READ_TIMEOUT_MS = 60000
|
||||
private const val BUFFER_SIZE = 256 * 1024
|
||||
private const val PROGRESS_THROTTLE_MS = 200
|
||||
|
||||
/** GETs [url] into [destination] (a plain file), creating parent dirs as needed. */
|
||||
fun get(
|
||||
url: URL,
|
||||
authHeader: String,
|
||||
destination: File,
|
||||
onProgress: ((received: Long, total: Long?) -> Unit)? = null,
|
||||
isCancelled: (() -> Boolean)? = null,
|
||||
): Boolean {
|
||||
destination.parentFile?.mkdirs()
|
||||
return getInto(url, authHeader, { destination.outputStream() }, onProgress, isCancelled)
|
||||
}
|
||||
|
||||
/**
|
||||
* GETs [url] into whatever [openOutput] returns - a plain `File` output
|
||||
* stream or, for [DownloadService]'s MediaStore-staged entries, a
|
||||
* `ContentResolver`-provided one the caller owns.
|
||||
*/
|
||||
fun getInto(
|
||||
url: URL,
|
||||
authHeader: String,
|
||||
openOutput: () -> OutputStream?,
|
||||
onProgress: ((received: Long, total: Long?) -> Unit)? = null,
|
||||
isCancelled: (() -> Boolean)? = null,
|
||||
): Boolean {
|
||||
val connection = url.openConnection() as HttpURLConnection
|
||||
return try {
|
||||
connection.requestMethod = "GET"
|
||||
connection.setRequestProperty("Authorization", authHeader)
|
||||
connection.connectTimeout = CONNECT_TIMEOUT_MS
|
||||
connection.readTimeout = READ_TIMEOUT_MS
|
||||
connection.connect()
|
||||
if (connection.responseCode !in 200..299) return false
|
||||
|
||||
val total = connection.contentLengthLong.takeIf { it > 0 }
|
||||
val output = openOutput() ?: return false
|
||||
output.use { out ->
|
||||
connection.inputStream.use { input ->
|
||||
copyWithProgress(input, out, total, onProgress, isCancelled)
|
||||
}
|
||||
}
|
||||
isCancelled?.invoke() != true
|
||||
} catch (e: Exception) {
|
||||
false
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
}
|
||||
}
|
||||
|
||||
/** PUTs [source]'s bytes to [url]. */
|
||||
fun put(
|
||||
url: URL,
|
||||
authHeader: String,
|
||||
source: File,
|
||||
contentType: String? = null,
|
||||
onProgress: ((sent: Long, total: Long) -> Unit)? = null,
|
||||
isCancelled: (() -> Boolean)? = null,
|
||||
): Boolean {
|
||||
val length = source.length()
|
||||
val connection = url.openConnection() as HttpURLConnection
|
||||
return try {
|
||||
connection.requestMethod = "PUT"
|
||||
connection.doOutput = true
|
||||
connection.setFixedLengthStreamingMode(length)
|
||||
connection.setRequestProperty("Authorization", authHeader)
|
||||
if (contentType != null) connection.setRequestProperty("Content-Type", contentType)
|
||||
connection.connectTimeout = CONNECT_TIMEOUT_MS
|
||||
connection.readTimeout = UPLOAD_READ_TIMEOUT_MS
|
||||
|
||||
connection.outputStream.use { output ->
|
||||
source.inputStream().use { input ->
|
||||
copyWithProgress(input, output, length, { sent, total ->
|
||||
onProgress?.invoke(sent, total ?: length)
|
||||
}, isCancelled)
|
||||
}
|
||||
}
|
||||
if (isCancelled?.invoke() == true) {
|
||||
false
|
||||
} else {
|
||||
connection.responseCode in 200..299
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
false
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
}
|
||||
}
|
||||
|
||||
private fun copyWithProgress(
|
||||
input: java.io.InputStream,
|
||||
output: OutputStream,
|
||||
total: Long?,
|
||||
onProgress: ((Long, Long?) -> Unit)?,
|
||||
isCancelled: (() -> Boolean)?,
|
||||
) {
|
||||
val buffer = ByteArray(BUFFER_SIZE)
|
||||
var transferred = 0L
|
||||
var lastEmit = 0L
|
||||
while (isCancelled?.invoke() != true) {
|
||||
val read = input.read(buffer)
|
||||
if (read == -1) break
|
||||
output.write(buffer, 0, read)
|
||||
transferred += read
|
||||
if (onProgress != null) {
|
||||
val now = System.currentTimeMillis()
|
||||
if (now - lastEmit >= PROGRESS_THROTTLE_MS) {
|
||||
lastEmit = now
|
||||
onProgress(transferred, total)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,7 +1,6 @@
|
||||
package dev.ayushya.noo
|
||||
|
||||
import android.app.Notification
|
||||
import android.app.NotificationChannel
|
||||
import android.app.NotificationManager
|
||||
import android.app.PendingIntent
|
||||
import android.app.Service
|
||||
@@ -18,24 +17,22 @@ import androidx.core.app.NotificationCompat
|
||||
import androidx.core.app.ServiceCompat
|
||||
import org.json.JSONArray
|
||||
import java.io.File
|
||||
import java.net.HttpURLConnection
|
||||
import java.net.URL
|
||||
import java.util.concurrent.atomic.AtomicBoolean
|
||||
|
||||
/**
|
||||
* Foreground service that downloads one or more files over WebDAV straight
|
||||
* into the device's public Downloads folder, independent of MainActivity/
|
||||
* the Flutter engine being alive - the download counterpart of
|
||||
* ShareUploadService.kt (see its doc comment for the full rationale: only
|
||||
* a real Android Service survives the app being closed mid-transfer, the
|
||||
* same guarantee a real file-manager app's download notification gives
|
||||
* you). Re-implements a plain WebDAV GET here in Kotlin for the same
|
||||
* reason uploads do - keep NextcloudService.downloadToFile in sync if
|
||||
* download semantics change.
|
||||
* ShareUploadService.kt (see its doc comment: only a real Android Service
|
||||
* survives the app being closed mid-transfer, the same guarantee a real
|
||||
* file-manager app's download notification gives you). GET itself is
|
||||
* [DavTransfer.getInto]; this file owns only what's specific to
|
||||
* downloading - MediaStore staging and the notification/queue plumbing.
|
||||
*
|
||||
* Started via the `dev.ayushya.noo/download_service` MethodChannel
|
||||
* (MainActivity.kt). Shows one persistent, cancellable notification for
|
||||
* the whole batch, same UX as ShareUploadService's.
|
||||
* the whole batch; a second batch arriving mid-download queues behind the
|
||||
* first (see [TransferQueue]) instead of being dropped.
|
||||
*/
|
||||
class DownloadService : Service() {
|
||||
companion object {
|
||||
@@ -49,9 +46,6 @@ class DownloadService : Service() {
|
||||
private const val NOTIFICATION_ID = 4301
|
||||
}
|
||||
|
||||
private val cancelled = AtomicBoolean(false)
|
||||
private var downloadThread: Thread? = null
|
||||
|
||||
private data class DownloadFile(
|
||||
val path: String,
|
||||
val name: String,
|
||||
@@ -59,11 +53,23 @@ class DownloadService : Service() {
|
||||
val size: Long?,
|
||||
)
|
||||
|
||||
private data class DownloadBatch(
|
||||
val files: List<DownloadFile>,
|
||||
val serverUrl: String,
|
||||
val username: String,
|
||||
val authHeader: String,
|
||||
)
|
||||
|
||||
@Volatile
|
||||
private var cancelled = false
|
||||
private val queue = TransferQueue<DownloadBatch> { runDownloads(it) }
|
||||
|
||||
override fun onBind(intent: Intent?): IBinder? = null
|
||||
|
||||
override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int {
|
||||
if (intent?.action == ACTION_CANCEL) {
|
||||
cancelled.set(true)
|
||||
cancelled = true
|
||||
queue.clearPending()
|
||||
return START_NOT_STICKY
|
||||
}
|
||||
|
||||
@@ -76,7 +82,13 @@ class DownloadService : Service() {
|
||||
return START_NOT_STICKY
|
||||
}
|
||||
|
||||
createNotificationChannel()
|
||||
NooNotificationChannels.ensure(
|
||||
this,
|
||||
CHANNEL_ID,
|
||||
"File downloads",
|
||||
NotificationManager.IMPORTANCE_LOW,
|
||||
"Progress for files downloaded from Noo",
|
||||
)
|
||||
ServiceCompat.startForeground(
|
||||
this,
|
||||
NOTIFICATION_ID,
|
||||
@@ -88,16 +100,7 @@ class DownloadService : Service() {
|
||||
},
|
||||
)
|
||||
|
||||
// A second batch arriving mid-download is dropped rather than
|
||||
// queued - same tradeoff ShareUploadService makes, rare in
|
||||
// practice and not worth a real queue for.
|
||||
if (downloadThread == null) {
|
||||
val files = parseFiles(filesJson)
|
||||
downloadThread = Thread {
|
||||
runDownloads(files, serverUrl, username, authHeader)
|
||||
}.also { it.start() }
|
||||
}
|
||||
|
||||
queue.enqueue(DownloadBatch(parseFiles(filesJson), serverUrl, username, authHeader))
|
||||
return START_NOT_STICKY
|
||||
}
|
||||
|
||||
@@ -114,25 +117,25 @@ class DownloadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
private fun runDownloads(
|
||||
files: List<DownloadFile>,
|
||||
serverUrl: String,
|
||||
username: String,
|
||||
authHeader: String,
|
||||
) {
|
||||
private fun runDownloads(batch: DownloadBatch) {
|
||||
cancelled = false
|
||||
var succeeded = 0
|
||||
var failed = 0
|
||||
val cleanServer = serverUrl.trimEnd('/')
|
||||
val cleanServer = batch.serverUrl.trimEnd('/')
|
||||
|
||||
for ((index, file) in files.withIndex()) {
|
||||
if (cancelled.get()) break
|
||||
val label = if (files.size == 1) file.name else "${file.name} (${index + 1}/${files.size})"
|
||||
for ((index, file) in batch.files.withIndex()) {
|
||||
if (cancelled) break
|
||||
val label = if (batch.files.size == 1) {
|
||||
file.name
|
||||
} else {
|
||||
"${file.name} (${index + 1}/${batch.files.size})"
|
||||
}
|
||||
try {
|
||||
var cleanPath = file.path.trim()
|
||||
if (!cleanPath.startsWith("/")) cleanPath = "/$cleanPath"
|
||||
val encodedPath = cleanPath.split("/").joinToString("/") { Uri.encode(it) }
|
||||
val url = URL("$cleanServer/remote.php/dav/files/$username$encodedPath")
|
||||
val ok = downloadFile(file, url, authHeader) { received, total ->
|
||||
val url = URL("$cleanServer/remote.php/dav/files/${batch.username}$encodedPath")
|
||||
val ok = downloadFile(file, url, batch.authHeader) { received, total ->
|
||||
notify(
|
||||
buildProgressNotification(
|
||||
"Downloading $label…",
|
||||
@@ -147,17 +150,18 @@ class DownloadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
val manager = getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
val finalText = when {
|
||||
cancelled.get() -> "Download cancelled"
|
||||
failed == 0 && files.size == 1 -> "Downloaded ${files.first().name}"
|
||||
failed == 0 -> "Downloaded $succeeded of ${files.size} files"
|
||||
else -> "Downloaded $succeeded of ${files.size} files - $failed failed"
|
||||
cancelled -> "Download cancelled"
|
||||
failed == 0 && batch.files.size == 1 -> "Downloaded ${batch.files.first().name}"
|
||||
failed == 0 -> "Downloaded $succeeded of ${batch.files.size} files"
|
||||
else -> "Downloaded $succeeded of ${batch.files.size} files - $failed failed"
|
||||
}
|
||||
manager.notify(NOTIFICATION_ID, buildFinalNotification(finalText))
|
||||
manager().notify(NOTIFICATION_ID, buildFinalNotification(finalText))
|
||||
|
||||
stopForeground(STOP_FOREGROUND_DETACH)
|
||||
stopSelf()
|
||||
if (queue.pendingCount == 0) {
|
||||
stopForeground(STOP_FOREGROUND_DETACH)
|
||||
stopSelf()
|
||||
}
|
||||
}
|
||||
|
||||
private fun progressFraction(received: Long, total: Long?): Float? {
|
||||
@@ -177,57 +181,32 @@ class DownloadService : Service() {
|
||||
authHeader: String,
|
||||
onProgress: (Long, Long?) -> Unit,
|
||||
): Boolean {
|
||||
val connection = url.openConnection() as HttpURLConnection
|
||||
return try {
|
||||
connection.requestMethod = "GET"
|
||||
connection.setRequestProperty("Authorization", authHeader)
|
||||
connection.setRequestProperty("OCS-APIRequest", "true")
|
||||
connection.connectTimeout = 15000
|
||||
connection.readTimeout = 30000
|
||||
connection.connect()
|
||||
|
||||
if (connection.responseCode !in 200..299) return false
|
||||
|
||||
val total = file.size ?: connection.contentLengthLong.takeIf { it > 0 }
|
||||
val mimeType = file.mimeType ?: "application/octet-stream"
|
||||
val outputUri = createDownloadsEntry(file.name, mimeType) ?: return false
|
||||
|
||||
val stream = contentResolver.openOutputStream(outputUri) ?: return false
|
||||
stream.use { output ->
|
||||
connection.inputStream.use { input ->
|
||||
val buffer = ByteArray(256 * 1024)
|
||||
var received = 0L
|
||||
var lastEmit = 0L
|
||||
while (!cancelled.get()) {
|
||||
val read = input.read(buffer)
|
||||
if (read == -1) break
|
||||
output.write(buffer, 0, read)
|
||||
received += read
|
||||
val now = System.currentTimeMillis()
|
||||
if (now - lastEmit >= 200) {
|
||||
lastEmit = now
|
||||
onProgress(received, total)
|
||||
}
|
||||
}
|
||||
val mimeType = file.mimeType ?: "application/octet-stream"
|
||||
var outputUri: Uri? = null
|
||||
val ok = DavTransfer.getInto(
|
||||
url,
|
||||
authHeader,
|
||||
openOutput = {
|
||||
val uri = createDownloadsEntry(file.name, mimeType)
|
||||
if (uri != null) {
|
||||
outputUri = uri
|
||||
contentResolver.openOutputStream(uri)
|
||||
} else {
|
||||
null
|
||||
}
|
||||
}
|
||||
},
|
||||
onProgress = { received, total -> onProgress(received, file.size ?: total) },
|
||||
isCancelled = { cancelled },
|
||||
)
|
||||
|
||||
outputUri?.let { uri ->
|
||||
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.Q) {
|
||||
val values = ContentValues().apply { put(MediaStore.Downloads.IS_PENDING, 0) }
|
||||
contentResolver.update(outputUri, values, null, null)
|
||||
contentResolver.update(uri, values, null, null)
|
||||
}
|
||||
|
||||
if (cancelled.get()) {
|
||||
contentResolver.delete(outputUri, null, null)
|
||||
false
|
||||
} else {
|
||||
true
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
false
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
if (!ok) contentResolver.delete(uri, null, null)
|
||||
}
|
||||
return ok
|
||||
}
|
||||
|
||||
private fun createDownloadsEntry(name: String, mimeType: String): Uri? {
|
||||
@@ -247,21 +226,11 @@ class DownloadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
private fun notify(notification: Notification) {
|
||||
val manager = getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
manager.notify(NOTIFICATION_ID, notification)
|
||||
}
|
||||
private fun manager() =
|
||||
getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
|
||||
private fun createNotificationChannel() {
|
||||
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) {
|
||||
val channel = NotificationChannel(
|
||||
CHANNEL_ID,
|
||||
"File downloads",
|
||||
NotificationManager.IMPORTANCE_LOW,
|
||||
).apply { description = "Progress for files downloaded from Noo" }
|
||||
val manager = getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
manager.createNotificationChannel(channel)
|
||||
}
|
||||
private fun notify(notification: Notification) {
|
||||
manager().notify(NOTIFICATION_ID, notification)
|
||||
}
|
||||
|
||||
private fun buildProgressNotification(
|
||||
@@ -301,7 +270,7 @@ class DownloadService : Service() {
|
||||
}
|
||||
|
||||
override fun onDestroy() {
|
||||
cancelled.set(true)
|
||||
cancelled = true
|
||||
super.onDestroy()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
package dev.ayushya.noo
|
||||
|
||||
import android.app.NotificationChannel
|
||||
import android.app.NotificationManager
|
||||
import android.content.Context
|
||||
import android.os.Build
|
||||
|
||||
/**
|
||||
* Shared notification-channel creation - [DownloadService],
|
||||
* [ShareUploadService], and [SyncWorker] each independently re-implemented
|
||||
* this same `if (SDK_INT >= O) { NotificationChannel(...); create }` shape
|
||||
* before this existed.
|
||||
*/
|
||||
object NooNotificationChannels {
|
||||
fun ensure(
|
||||
context: Context,
|
||||
id: String,
|
||||
name: String,
|
||||
importance: Int,
|
||||
description: String,
|
||||
) {
|
||||
if (Build.VERSION.SDK_INT < Build.VERSION_CODES.O) return
|
||||
val manager = context.getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
val channel = NotificationChannel(id, name, importance).apply {
|
||||
this.description = description
|
||||
}
|
||||
manager.createNotificationChannel(channel)
|
||||
}
|
||||
}
|
||||
@@ -1,7 +1,6 @@
|
||||
package dev.ayushya.noo
|
||||
|
||||
import android.app.Notification
|
||||
import android.app.NotificationChannel
|
||||
import android.app.NotificationManager
|
||||
import android.app.PendingIntent
|
||||
import android.app.Service
|
||||
@@ -16,9 +15,7 @@ import androidx.core.app.ServiceCompat
|
||||
import org.json.JSONArray
|
||||
import java.io.File
|
||||
import java.io.FileOutputStream
|
||||
import java.net.HttpURLConnection
|
||||
import java.net.URL
|
||||
import java.util.concurrent.atomic.AtomicBoolean
|
||||
|
||||
/**
|
||||
* Foreground service that prepares (copies from a content:// Uri) and
|
||||
@@ -28,19 +25,18 @@ import java.util.concurrent.atomic.AtomicBoolean
|
||||
* doesn't interrupt the upload, the same way a file-manager app's own
|
||||
* upload notification survives the app being closed.
|
||||
*
|
||||
* Re-implements a plain WebDAV PUT here in Kotlin (HttpURLConnection, no
|
||||
* new HTTP dependency) rather than reusing NextcloudService/Dio from Dart,
|
||||
* since a Dart isolate doesn't keep running once the Flutter engine/
|
||||
* Activity are gone - only a real Android Service does. This does mean the
|
||||
* PUT request itself is duplicated logic (see NextcloudService.
|
||||
* uploadFileFromPath); keep both in sync if the upload semantics change.
|
||||
* PUT itself is [DavTransfer.put]; this file owns only what's specific to
|
||||
* a share-upload - materializing the source content:// Uri into a real
|
||||
* file first (WebDAV PUT needs a known Content-Length, which streaming
|
||||
* straight from a content:// Uri can't always provide) and the
|
||||
* notification/queue plumbing.
|
||||
*
|
||||
* Started via the `dev.ayushya.noo/upload_service` MethodChannel
|
||||
* (MainActivity.kt) with credentials/destination passed as Intent extras -
|
||||
* never has an Activity in the loop after that. Shows one persistent,
|
||||
* cancellable notification for the whole batch; the Cancel action re-enters
|
||||
* this same running service instance with [ACTION_CANCEL], which the
|
||||
* copy/upload loops poll.
|
||||
* cancellable notification for the whole batch; a second share arriving
|
||||
* mid-upload queues behind the first (see [TransferQueue]) instead of
|
||||
* being dropped.
|
||||
*/
|
||||
class ShareUploadService : Service() {
|
||||
companion object {
|
||||
@@ -55,16 +51,26 @@ class ShareUploadService : Service() {
|
||||
private const val NOTIFICATION_ID = 4201
|
||||
}
|
||||
|
||||
private val cancelled = AtomicBoolean(false)
|
||||
private var uploadThread: Thread? = null
|
||||
|
||||
private data class ShareFile(val uri: String, val name: String, val size: Long?)
|
||||
|
||||
private data class UploadBatch(
|
||||
val files: List<ShareFile>,
|
||||
val serverUrl: String,
|
||||
val username: String,
|
||||
val authHeader: String,
|
||||
val remoteFolder: String,
|
||||
)
|
||||
|
||||
@Volatile
|
||||
private var cancelled = false
|
||||
private val queue = TransferQueue<UploadBatch> { runUploads(it) }
|
||||
|
||||
override fun onBind(intent: Intent?): IBinder? = null
|
||||
|
||||
override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int {
|
||||
if (intent?.action == ACTION_CANCEL) {
|
||||
cancelled.set(true)
|
||||
cancelled = true
|
||||
queue.clearPending()
|
||||
return START_NOT_STICKY
|
||||
}
|
||||
|
||||
@@ -80,7 +86,13 @@ class ShareUploadService : Service() {
|
||||
return START_NOT_STICKY
|
||||
}
|
||||
|
||||
createNotificationChannel()
|
||||
NooNotificationChannels.ensure(
|
||||
this,
|
||||
CHANNEL_ID,
|
||||
"File uploads",
|
||||
NotificationManager.IMPORTANCE_LOW,
|
||||
"Progress for files shared to Noo",
|
||||
)
|
||||
ServiceCompat.startForeground(
|
||||
this,
|
||||
NOTIFICATION_ID,
|
||||
@@ -92,16 +104,7 @@ class ShareUploadService : Service() {
|
||||
},
|
||||
)
|
||||
|
||||
// A second share arriving mid-upload is dropped rather than queued -
|
||||
// rare in practice (sharing again before the first batch finishes),
|
||||
// not worth a real queue for.
|
||||
if (uploadThread == null) {
|
||||
val files = parseFiles(filesJson)
|
||||
uploadThread = Thread {
|
||||
runUploads(files, serverUrl, username, authHeader, remoteFolder)
|
||||
}.also { it.start() }
|
||||
}
|
||||
|
||||
queue.enqueue(UploadBatch(parseFiles(filesJson), serverUrl, username, authHeader, remoteFolder))
|
||||
return START_NOT_STICKY
|
||||
}
|
||||
|
||||
@@ -117,23 +120,22 @@ class ShareUploadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
private fun runUploads(
|
||||
files: List<ShareFile>,
|
||||
serverUrl: String,
|
||||
username: String,
|
||||
authHeader: String,
|
||||
remoteFolder: String,
|
||||
) {
|
||||
private fun runUploads(batch: UploadBatch) {
|
||||
cancelled = false
|
||||
var succeeded = 0
|
||||
var failed = 0
|
||||
val cleanServer = serverUrl.trimEnd('/')
|
||||
var cleanFolder = remoteFolder.trim()
|
||||
val cleanServer = batch.serverUrl.trimEnd('/')
|
||||
var cleanFolder = batch.remoteFolder.trim()
|
||||
if (!cleanFolder.startsWith("/")) cleanFolder = "/$cleanFolder"
|
||||
if (!cleanFolder.endsWith("/")) cleanFolder = "$cleanFolder/"
|
||||
|
||||
for ((index, file) in files.withIndex()) {
|
||||
if (cancelled.get()) break
|
||||
val label = if (files.size == 1) file.name else "${file.name} (${index + 1}/${files.size})"
|
||||
for ((index, file) in batch.files.withIndex()) {
|
||||
if (cancelled) break
|
||||
val label = if (batch.files.size == 1) {
|
||||
file.name
|
||||
} else {
|
||||
"${file.name} (${index + 1}/${batch.files.size})"
|
||||
}
|
||||
var tempFile: File? = null
|
||||
try {
|
||||
notify(buildProgressNotification("Preparing $label…", null, indeterminate = true))
|
||||
@@ -146,19 +148,26 @@ class ShareUploadService : Service() {
|
||||
),
|
||||
)
|
||||
}
|
||||
if (cancelled.get()) break
|
||||
if (cancelled) break
|
||||
|
||||
val encodedName = Uri.encode(file.name)
|
||||
val url = URL("$cleanServer/remote.php/dav/files/$username$cleanFolder$encodedName")
|
||||
val ok = uploadFile(tempFile, url, authHeader) { sent, total ->
|
||||
notify(
|
||||
buildProgressNotification(
|
||||
"Uploading $label…",
|
||||
progressFraction(sent, total),
|
||||
indeterminate = false,
|
||||
),
|
||||
)
|
||||
}
|
||||
val url = URL("$cleanServer/remote.php/dav/files/${batch.username}$cleanFolder$encodedName")
|
||||
val ok = DavTransfer.put(
|
||||
url,
|
||||
batch.authHeader,
|
||||
tempFile,
|
||||
contentType = "application/octet-stream",
|
||||
onProgress = { sent, total ->
|
||||
notify(
|
||||
buildProgressNotification(
|
||||
"Uploading $label…",
|
||||
progressFraction(sent, total),
|
||||
indeterminate = false,
|
||||
),
|
||||
)
|
||||
},
|
||||
isCancelled = { cancelled },
|
||||
)
|
||||
if (ok) succeeded++ else failed++
|
||||
} catch (e: Exception) {
|
||||
failed++
|
||||
@@ -167,17 +176,18 @@ class ShareUploadService : Service() {
|
||||
}
|
||||
}
|
||||
|
||||
val manager = getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
val finalText = when {
|
||||
cancelled.get() -> "Upload cancelled"
|
||||
failed == 0 && files.size == 1 -> "Uploaded ${files.first().name}"
|
||||
failed == 0 -> "Uploaded $succeeded of ${files.size} files"
|
||||
else -> "Uploaded $succeeded of ${files.size} files - $failed failed"
|
||||
cancelled -> "Upload cancelled"
|
||||
failed == 0 && batch.files.size == 1 -> "Uploaded ${batch.files.first().name}"
|
||||
failed == 0 -> "Uploaded $succeeded of ${batch.files.size} files"
|
||||
else -> "Uploaded $succeeded of ${batch.files.size} files - $failed failed"
|
||||
}
|
||||
manager.notify(NOTIFICATION_ID, buildFinalNotification(finalText))
|
||||
manager().notify(NOTIFICATION_ID, buildFinalNotification(finalText))
|
||||
|
||||
stopForeground(STOP_FOREGROUND_DETACH)
|
||||
stopSelf()
|
||||
if (queue.pendingCount == 0) {
|
||||
stopForeground(STOP_FOREGROUND_DETACH)
|
||||
stopSelf()
|
||||
}
|
||||
}
|
||||
|
||||
private fun progressFraction(sent: Long, total: Long?): Float? {
|
||||
@@ -193,7 +203,7 @@ class ShareUploadService : Service() {
|
||||
val buffer = ByteArray(256 * 1024)
|
||||
var sent = 0L
|
||||
var lastEmit = 0L
|
||||
while (!cancelled.get()) {
|
||||
while (!cancelled) {
|
||||
val read = input.read(buffer)
|
||||
if (read == -1) break
|
||||
output.write(buffer, 0, read)
|
||||
@@ -209,70 +219,11 @@ class ShareUploadService : Service() {
|
||||
return target
|
||||
}
|
||||
|
||||
private fun uploadFile(
|
||||
file: File,
|
||||
url: URL,
|
||||
authHeader: String,
|
||||
onProgress: (Long, Long) -> Unit,
|
||||
): Boolean {
|
||||
val length = file.length()
|
||||
val connection = url.openConnection() as HttpURLConnection
|
||||
return try {
|
||||
connection.requestMethod = "PUT"
|
||||
connection.doOutput = true
|
||||
connection.setFixedLengthStreamingMode(length)
|
||||
connection.setRequestProperty("Authorization", authHeader)
|
||||
connection.setRequestProperty("OCS-APIRequest", "true")
|
||||
connection.setRequestProperty("Content-Type", "application/octet-stream")
|
||||
connection.connectTimeout = 15000
|
||||
connection.readTimeout = 30000
|
||||
|
||||
connection.outputStream.use { output ->
|
||||
file.inputStream().use { input ->
|
||||
val buffer = ByteArray(256 * 1024)
|
||||
var sent = 0L
|
||||
var lastEmit = 0L
|
||||
while (!cancelled.get()) {
|
||||
val read = input.read(buffer)
|
||||
if (read == -1) break
|
||||
output.write(buffer, 0, read)
|
||||
sent += read
|
||||
val now = System.currentTimeMillis()
|
||||
if (now - lastEmit >= 200) {
|
||||
lastEmit = now
|
||||
onProgress(sent, length)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if (cancelled.get()) {
|
||||
false
|
||||
} else {
|
||||
val code = connection.responseCode
|
||||
code == 201 || code == 204 || code == 200
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
false
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
}
|
||||
}
|
||||
private fun manager() =
|
||||
getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
|
||||
private fun notify(notification: Notification) {
|
||||
val manager = getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
manager.notify(NOTIFICATION_ID, notification)
|
||||
}
|
||||
|
||||
private fun createNotificationChannel() {
|
||||
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) {
|
||||
val channel = NotificationChannel(
|
||||
CHANNEL_ID,
|
||||
"File uploads",
|
||||
NotificationManager.IMPORTANCE_LOW,
|
||||
).apply { description = "Progress for files shared to Noo" }
|
||||
val manager = getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
manager.createNotificationChannel(channel)
|
||||
}
|
||||
manager().notify(NOTIFICATION_ID, notification)
|
||||
}
|
||||
|
||||
private fun buildProgressNotification(text: String, progress: Float?, indeterminate: Boolean): Notification {
|
||||
@@ -308,7 +259,7 @@ class ShareUploadService : Service() {
|
||||
}
|
||||
|
||||
override fun onDestroy() {
|
||||
cancelled.set(true)
|
||||
cancelled = true
|
||||
super.onDestroy()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,7 +10,6 @@ import okhttp3.RequestBody.Companion.toRequestBody
|
||||
import org.json.JSONArray
|
||||
import org.json.JSONObject
|
||||
import java.io.File
|
||||
import java.net.HttpURLConnection
|
||||
import java.net.URL
|
||||
import java.text.SimpleDateFormat
|
||||
import java.util.Locale
|
||||
@@ -258,33 +257,17 @@ object SyncEngine {
|
||||
remotePath: String,
|
||||
destination: File,
|
||||
): Boolean {
|
||||
val encodedPath = remotePath.split("/").joinToString("/") { Uri.encode(it) }
|
||||
val url = URL("${serverUrl.trimEnd('/')}/remote.php/dav/files/$username$encodedPath")
|
||||
val connection = url.openConnection() as HttpURLConnection
|
||||
return try {
|
||||
connection.requestMethod = "GET"
|
||||
connection.setRequestProperty("Authorization", authHeader)
|
||||
connection.connectTimeout = 15000
|
||||
connection.readTimeout = 30000
|
||||
connection.connect()
|
||||
val code = connection.responseCode
|
||||
if (code !in 200..299) {
|
||||
Log.w(TAG, "GET $url failed: $code")
|
||||
return false
|
||||
}
|
||||
|
||||
destination.parentFile?.mkdirs()
|
||||
connection.inputStream.use { input ->
|
||||
destination.outputStream().use { output -> input.copyTo(output) }
|
||||
}
|
||||
Log.d(TAG, "GET $url -> saved to ${destination.absolutePath} (${destination.length()} bytes)")
|
||||
true
|
||||
} catch (e: Exception) {
|
||||
Log.e(TAG, "GET $url threw", e)
|
||||
false
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
}
|
||||
val url = davUrl(serverUrl, username, remotePath)
|
||||
val ok = DavTransfer.get(url, authHeader, destination)
|
||||
Log.d(
|
||||
TAG,
|
||||
if (ok) {
|
||||
"GET $url -> saved to ${destination.absolutePath} (${destination.length()} bytes)"
|
||||
} else {
|
||||
"GET $url failed"
|
||||
},
|
||||
)
|
||||
return ok
|
||||
}
|
||||
|
||||
fun uploadFile(
|
||||
@@ -294,26 +277,13 @@ object SyncEngine {
|
||||
remotePath: String,
|
||||
source: File,
|
||||
): Boolean {
|
||||
val url = davUrl(serverUrl, username, remotePath)
|
||||
return DavTransfer.put(url, authHeader, source)
|
||||
}
|
||||
|
||||
private fun davUrl(serverUrl: String, username: String, remotePath: String): URL {
|
||||
val encodedPath = remotePath.split("/").joinToString("/") { Uri.encode(it) }
|
||||
val url = URL("${serverUrl.trimEnd('/')}/remote.php/dav/files/$username$encodedPath")
|
||||
val connection = url.openConnection() as HttpURLConnection
|
||||
return try {
|
||||
connection.requestMethod = "PUT"
|
||||
connection.setRequestProperty("Authorization", authHeader)
|
||||
connection.setFixedLengthStreamingMode(source.length())
|
||||
connection.doOutput = true
|
||||
connection.connectTimeout = 15000
|
||||
connection.readTimeout = 60000
|
||||
connection.connect()
|
||||
source.inputStream().use { input ->
|
||||
connection.outputStream.use { output -> input.copyTo(output) }
|
||||
}
|
||||
connection.responseCode in 200..299
|
||||
} catch (e: Exception) {
|
||||
false
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
}
|
||||
return URL("${serverUrl.trimEnd('/')}/remote.php/dav/files/$username$encodedPath")
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -1,12 +1,10 @@
|
||||
package dev.ayushya.noo
|
||||
|
||||
import android.app.Notification
|
||||
import android.app.NotificationChannel
|
||||
import android.app.NotificationManager
|
||||
import android.app.PendingIntent
|
||||
import android.content.Context
|
||||
import android.content.Intent
|
||||
import android.os.Build
|
||||
import android.util.Log
|
||||
import androidx.core.app.NotificationCompat
|
||||
import androidx.work.CoroutineWorker
|
||||
@@ -260,13 +258,12 @@ class SyncWorker(appContext: Context, params: WorkerParameters) :
|
||||
applicationContext.getSystemService(Context.NOTIFICATION_SERVICE) as NotificationManager
|
||||
|
||||
private fun createChannel() {
|
||||
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) {
|
||||
val channel = NotificationChannel(
|
||||
CHANNEL_ID,
|
||||
"Device sync",
|
||||
NotificationManager.IMPORTANCE_DEFAULT,
|
||||
).apply { description = "Updates and conflicts for folders synced to this device" }
|
||||
manager().createNotificationChannel(channel)
|
||||
}
|
||||
NooNotificationChannels.ensure(
|
||||
applicationContext,
|
||||
CHANNEL_ID,
|
||||
"Device sync",
|
||||
NotificationManager.IMPORTANCE_DEFAULT,
|
||||
"Updates and conflicts for folders synced to this device",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
package dev.ayushya.noo
|
||||
|
||||
import java.util.concurrent.LinkedBlockingQueue
|
||||
|
||||
/**
|
||||
* A tiny sequential work queue for [DownloadService]/[ShareUploadService]:
|
||||
* each enqueued batch runs to completion before the next starts, instead
|
||||
* of a second batch arriving mid-transfer being silently dropped - the
|
||||
* previous behavior both services independently accepted as a known
|
||||
* limitation (see their old doc comments). One daemon worker thread,
|
||||
* started lazily on first use and parked on the queue for the life of the
|
||||
* process.
|
||||
*/
|
||||
class TransferQueue<T>(private val process: (T) -> Unit) {
|
||||
private val queue = LinkedBlockingQueue<T>()
|
||||
|
||||
@Volatile
|
||||
private var started = false
|
||||
|
||||
/** How many batches (including whatever's currently running) are queued. */
|
||||
val pendingCount: Int
|
||||
get() = queue.size
|
||||
|
||||
@Synchronized
|
||||
fun enqueue(batch: T) {
|
||||
queue.put(batch)
|
||||
if (!started) {
|
||||
started = true
|
||||
Thread {
|
||||
while (true) {
|
||||
process(queue.take())
|
||||
}
|
||||
}.apply { isDaemon = true }.start()
|
||||
}
|
||||
}
|
||||
|
||||
/** Drops every batch that hasn't started running yet (not the current one). */
|
||||
fun clearPending() {
|
||||
queue.clear()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user