Add device sync: mirror folders/files locally, background WorkManager engine
Build APK / build (push) Successful in 7m33s

- New native sync engine (SyncEngine/SyncWorker/ConflictResolveWorker,
  WorkManager-based) that mirrors selected Nextcloud folders or files to
  app-private local storage, periodically and on-demand ("Sync now"),
  with never-auto-resolved conflict notifications.
- Single files (not just folders) can now be marked "Sync to device".
- Live sync status (syncing/synced/conflicts) pushes from native to Dart
  over a new EventChannel; the persistent header chip/panel now reflects
  device-sync status instead of the WebDAV-refresh loading state, with a
  cloud_off/cloud_sync/cloud_done/cloud_alert icon set and an expandable
  conflicts list with in-app "Keep local"/"Use server" resolution.
  Per-item cloud_done/sync badges show on Files tiles.
- Share/Download short-circuit to the local copy for already-synced
  files instead of a fresh network fetch.
- Settings gained a Device Sync section (per-folder list, "Sync
  everything", Wi-Fi-only toggle, manual "Sync now").
- Fixed two release-only bugs found via on-device testing: Android's
  HttpURLConnection silently rejects the PROPFIND method (switched to
  OkHttp), and R8 was stripping WorkManager's reflection-instantiated
  internals (broadened proguard-rules.pro to keep androidx.work.**
  wholesale rather than chasing individual classes).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
2026-09-19 00:52:56 -04:00
co-authored by Claude Sonnet 5
parent 722b66f7b1
commit faaec5365f
21 changed files with 2023 additions and 42 deletions
+12
View File
@@ -84,6 +84,18 @@ kotlin {
}
}
dependencies {
// Device sync's periodic/one-off background jobs (SyncWorker,
// ConflictResolveWorker) - see their doc comments for why this is
// plain WorkManager rather than a Dart-side background-task plugin.
implementation("androidx.work:work-runtime-ktx:2.9.1")
// PROPFIND (WebDAV directory listing) - Android's HttpURLConnection
// hard-rejects any method outside {OPTIONS,GET,HEAD,POST,PUT,DELETE,
// TRACE,PATCH} (ProtocolException), unlike plain OpenJDK. OkHttp has
// no such whitelist. See SyncEngine.kt.
implementation("com.squareup.okhttp3:okhttp:4.12.0")
}
flutter {
source = "../.."
}
+18
View File
@@ -0,0 +1,18 @@
# Device sync (SyncWorker/ConflictResolveWorker) - WorkManager instantiates
# a lot of its own internals via reflection (Workers by class name recorded
# at enqueue time, its bundled Room database's *_Impl class off the
# abstract database class's own name, its default InputMerger, etc). R8's
# member-level shrinking silently breaks any of these - e.g. stripping an
# "unused" no-arg constructor - even when the class itself survives a
# plain `-keep class` with no wildcard, and narrowly keeping only the
# classes hit by one test pass just means the next reflection path (a
# different WorkManager-internal class) breaks instead. Two separate
# instances of this already bit real testing (WorkDatabase, then
# OverwritingInputMerger) before landing on this broad keep - see
# android/app/build.gradle.kts (androidx.work dependency) and
# .claude/context/server.md's "Device sync" section.
-keep class androidx.work.** { *; }
-keep class * extends androidx.room.RoomDatabase { *; }
-keep class dev.ayushya.noo.SyncWorker { *; }
-keep class dev.ayushya.noo.ConflictResolveWorker { *; }
-keep class dev.ayushya.noo.SyncConflictReceiver { *; }
+7
View File
@@ -76,6 +76,13 @@
android:name=".DownloadService"
android:exported="false"
android:foregroundServiceType="dataSync"/>
<!-- The "Keep local"/"Use server" actions on a sync-conflict
notification (SyncWorker/SyncEngine) - must be a manifest-
registered receiver, not a runtime-registered one, since it
needs to work even when the app process isn't running. -->
<receiver
android:name=".SyncConflictReceiver"
android:exported="false"/>
<!-- Hands back the file(s) a caller picked via the GET_CONTENT filter
above - only the app's own cache/picker/ subfolder is exposed,
see file_paths.xml. -->
@@ -0,0 +1,101 @@
package dev.ayushya.noo
import android.content.Context
import androidx.work.CoroutineWorker
import androidx.work.OneTimeWorkRequestBuilder
import androidx.work.WorkManager
import androidx.work.WorkerParameters
import androidx.work.workDataOf
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import java.io.File
/**
* Resolves one sync conflict, triggered either by [SyncConflictReceiver]
* (a notification action) or directly from Dart (in-app resolution via the
* `dev.ayushya.noo/sync_service` MethodChannel's `resolveConflict`, see
* `MainActivity.kt`) - "local" pushes the on-device copy up (overwriting
* the server), "server" pulls the server copy down (overwriting the local
* mirror). Either way, the sync-state entry for this file is refreshed
* afterward so the next [SyncWorker] pass doesn't immediately re-flag it.
*/
class ConflictResolveWorker(appContext: Context, params: WorkerParameters) :
CoroutineWorker(appContext, params) {
companion object {
const val KEY_ACCOUNT_ID = "accountId"
const val KEY_SERVER_URL = "serverUrl"
const val KEY_USERNAME = "username"
const val KEY_AUTH_HEADER = "authHeader"
const val KEY_FILE_ID = "fileId"
const val KEY_REMOTE_PATH = "remotePath"
const val KEY_REL_PATH = "relPath"
const val KEY_RESOLUTION = "resolution"
/** The one place a resolution gets enqueued - both callers above share it. */
fun enqueue(
context: Context,
accountId: String?,
serverUrl: String?,
username: String?,
authHeader: String?,
fileId: String?,
remotePath: String?,
relPath: String?,
resolution: String?,
) {
val data = workDataOf(
KEY_ACCOUNT_ID to accountId,
KEY_SERVER_URL to serverUrl,
KEY_USERNAME to username,
KEY_AUTH_HEADER to authHeader,
KEY_FILE_ID to fileId,
KEY_REMOTE_PATH to remotePath,
KEY_REL_PATH to relPath,
KEY_RESOLUTION to resolution,
)
val request = OneTimeWorkRequestBuilder<ConflictResolveWorker>()
.setInputData(data)
.build()
WorkManager.getInstance(context).enqueue(request)
}
}
override suspend fun doWork(): Result = withContext(Dispatchers.IO) {
val accountId = inputData.getString(KEY_ACCOUNT_ID) ?: return@withContext Result.failure()
val serverUrl = inputData.getString(KEY_SERVER_URL) ?: return@withContext Result.failure()
val username = inputData.getString(KEY_USERNAME) ?: return@withContext Result.failure()
val authHeader = inputData.getString(KEY_AUTH_HEADER) ?: return@withContext Result.failure()
val fileId = inputData.getString(KEY_FILE_ID) ?: return@withContext Result.failure()
val remotePath = inputData.getString(KEY_REMOTE_PATH) ?: return@withContext Result.failure()
val relPath = inputData.getString(KEY_REL_PATH) ?: return@withContext Result.failure()
val resolution = inputData.getString(KEY_RESOLUTION) ?: return@withContext Result.failure()
val syncRoot = SyncEngine.syncRoot(applicationContext, accountId)
val localFile = File(syncRoot, relPath)
val ok = when (resolution) {
"local" -> localFile.exists() &&
SyncEngine.uploadFile(serverUrl, username, authHeader, remotePath, localFile)
"server" -> SyncEngine.downloadFile(serverUrl, username, authHeader, remotePath, localFile)
else -> false
}
if (!ok) return@withContext Result.retry()
// Refresh the recorded state from the server's post-resolution
// etag, so this file isn't immediately re-flagged as a conflict on
// the next sync pass.
val fresh = SyncEngine.propfindSelf(serverUrl, username, authHeader, remotePath)
val state = SyncEngine.loadState(applicationContext, accountId).toMutableMap()
state[fileId] = SyncEngine.FileState(
relPath = relPath,
etag = fresh?.etag ?: "",
lastModified = fresh?.lastModified ?: 0L,
size = localFile.length(),
localMTime = localFile.lastModified(),
)
SyncEngine.saveState(applicationContext, accountId, state)
SyncStatusBus.removeConflict(accountId, fileId)
Result.success()
}
}
@@ -12,12 +12,22 @@ import android.provider.OpenableColumns
import androidx.core.app.ActivityCompat
import androidx.core.content.ContextCompat
import androidx.core.content.FileProvider
import androidx.work.Constraints
import androidx.work.ExistingPeriodicWorkPolicy
import androidx.work.ExistingWorkPolicy
import androidx.work.NetworkType
import androidx.work.OneTimeWorkRequestBuilder
import androidx.work.PeriodicWorkRequestBuilder
import androidx.work.WorkManager
import androidx.work.workDataOf
import io.flutter.embedding.android.FlutterFragmentActivity
import io.flutter.embedding.engine.FlutterEngine
import io.flutter.plugin.common.EventChannel
import io.flutter.plugin.common.MethodCall
import io.flutter.plugin.common.MethodChannel
import org.json.JSONArray
import java.io.File
import java.util.concurrent.TimeUnit
// FlutterFragmentActivity (not the default FlutterActivity) is required by
// local_auth's Android implementation, which hosts its biometric/device
@@ -44,6 +54,8 @@ class MainActivity : FlutterFragmentActivity() {
private val downloadServiceChannelName = "dev.ayushya.noo/download_service"
private val pickIntentChannelName = "dev.ayushya.noo/pick_intent"
private val newPickChannelName = "dev.ayushya.noo/pick_intent/new"
private val syncServiceChannelName = "dev.ayushya.noo/sync_service"
private val syncStatusChannelName = "dev.ayushya.noo/sync_service/status"
private val notificationPermissionRequestCode = 4202
// Lazy, not a field initializer - `packageName` reads through the
// Activity's base Context, which isn't attached yet while this class's
@@ -53,6 +65,7 @@ class MainActivity : FlutterFragmentActivity() {
private val mainHandler = Handler(Looper.getMainLooper())
private var newShareSink: EventChannel.EventSink? = null
private var newPickSink: EventChannel.EventSink? = null
private var syncStatusListener: ((SyncStatusBus.Status) -> Unit)? = null
override fun configureFlutterEngine(flutterEngine: FlutterEngine) {
super.configureFlutterEngine(flutterEngine)
@@ -114,6 +127,74 @@ class MainActivity : FlutterFragmentActivity() {
newPickSink = null
}
})
MethodChannel(flutterEngine.dartExecutor.binaryMessenger, syncServiceChannelName)
.setMethodCallHandler { call, result ->
when (call.method) {
"reschedule" -> rescheduleSyncWork(call, result)
"cancel" -> cancelSyncWork(result)
"syncNow" -> syncNow(call, result)
"getSyncStatus" -> result.success(syncStatusMap(SyncStatusBus.snapshot()))
"resolveConflict" -> resolveConflict(call, result)
else -> result.notImplemented()
}
}
EventChannel(flutterEngine.dartExecutor.binaryMessenger, syncStatusChannelName)
.setStreamHandler(object : EventChannel.StreamHandler {
override fun onListen(arguments: Any?, events: EventChannel.EventSink) {
val listener: (SyncStatusBus.Status) -> Unit = { status ->
mainHandler.post { events.success(syncStatusMap(status)) }
}
syncStatusListener = listener
SyncStatusBus.subscribe(listener)
}
override fun onCancel(arguments: Any?) {
syncStatusListener?.let { SyncStatusBus.unsubscribe(it) }
syncStatusListener = null
}
})
}
private fun syncStatusMap(status: SyncStatusBus.Status): Map<String, Any?> {
val syncedFileIds = status.accountId?.let {
SyncEngine.loadState(applicationContext, it).keys.toList()
} ?: emptyList()
return mapOf(
"accountId" to status.accountId,
"syncing" to status.syncing,
"syncingFileIds" to status.syncingFileIds.toList(),
"syncedFileIds" to syncedFileIds,
"conflicts" to status.conflicts.map {
mapOf(
"accountId" to it.accountId,
"fileId" to it.fileId,
"remotePath" to it.remotePath,
"relPath" to it.relPath,
"name" to it.name,
)
},
)
}
/// In-app conflict resolution (the header's "Keep local"/"Use server"
/// buttons) - shares [ConflictResolveWorker.enqueue] with
/// [SyncConflictReceiver], the only other caller, so there's one
/// resolution code path regardless of whether it's triggered from a
/// notification or from inside the app.
private fun resolveConflict(call: MethodCall, result: MethodChannel.Result) {
ConflictResolveWorker.enqueue(
this,
accountId = call.argument<String>("accountId"),
serverUrl = call.argument<String>("serverUrl"),
username = call.argument<String>("username"),
authHeader = call.argument<String>("authHeader"),
fileId = call.argument<String>("fileId"),
remotePath = call.argument<String>("remotePath"),
relPath = call.argument<String>("relPath"),
resolution = call.argument<String>("resolution"),
)
result.success(null)
}
override fun onNewIntent(intent: Intent) {
@@ -262,6 +343,85 @@ class MainActivity : FlutterFragmentActivity() {
result.success(null)
}
/// Shared arg-parsing for `reschedule`/`syncNow` - both need the same
/// account/credentials/folder-list shape, just enqueue differently.
private fun syncWorkData(call: MethodCall): androidx.work.Data? {
val accountId = call.argument<String>("accountId")
val serverUrl = call.argument<String>("serverUrl")
val username = call.argument<String>("username")
val authHeader = call.argument<String>("authHeader")
val foldersJson = call.argument<String>("folders")
if (accountId == null || serverUrl == null || username == null ||
authHeader == null || foldersJson == null
) {
return null
}
return workDataOf(
SyncWorker.KEY_ACCOUNT_ID to accountId,
SyncWorker.KEY_SERVER_URL to serverUrl,
SyncWorker.KEY_USERNAME to username,
SyncWorker.KEY_AUTH_HEADER to authHeader,
SyncWorker.KEY_FOLDERS to foldersJson,
)
}
/// Re-enqueues (or cancels, if the folder list is now empty) the
/// periodic device-sync job - called from Dart whenever the synced-
/// folder list, active account, or Wi-Fi-only setting changes, since a
/// periodic WorkRequest's input Data/constraints are fixed at enqueue
/// time and can only be changed by cancelling and re-enqueueing.
private fun rescheduleSyncWork(call: MethodCall, result: MethodChannel.Result) {
val data = syncWorkData(call)
val foldersJson = call.argument<String>("folders")
val wifiOnly = call.argument<Boolean>("wifiOnly") ?: true
val workManager = WorkManager.getInstance(this)
if (data == null || foldersJson == null || JSONArray(foldersJson).length() == 0) {
workManager.cancelUniqueWork(SyncWorker.UNIQUE_PERIODIC_NAME)
result.success(null)
return
}
val constraints = Constraints.Builder()
.setRequiredNetworkType(if (wifiOnly) NetworkType.UNMETERED else NetworkType.CONNECTED)
.build()
// 1 hour is WorkManager's own practical floor for a "battery-
// friendly" cadence well above its hard 15-minute minimum; there's
// no per-user interval setting for this in v1.
val request = PeriodicWorkRequestBuilder<SyncWorker>(1, TimeUnit.HOURS)
.setInputData(data)
.setConstraints(constraints)
.build()
workManager.enqueueUniquePeriodicWork(
SyncWorker.UNIQUE_PERIODIC_NAME,
ExistingPeriodicWorkPolicy.UPDATE,
request,
)
result.success(null)
}
private fun cancelSyncWork(result: MethodChannel.Result) {
WorkManager.getInstance(this).cancelUniqueWork(SyncWorker.UNIQUE_PERIODIC_NAME)
result.success(null)
}
/// One-off immediate run (Settings' "Sync now"), independent of the
/// periodic schedule.
private fun syncNow(call: MethodCall, result: MethodChannel.Result) {
val data = syncWorkData(call)
if (data == null) {
result.error("bad_args", "Missing required sync arguments", null)
return
}
val request = OneTimeWorkRequestBuilder<SyncWorker>().setInputData(data).build()
WorkManager.getInstance(this).enqueueUniqueWork(
SyncWorker.UNIQUE_ONE_OFF_NAME,
ExistingWorkPolicy.REPLACE,
request,
)
result.success(null)
}
/// Non-null only when this Activity was launched (or re-delivered a new
/// Intent) as another app's GET_CONTENT picker - the mimeType filter and
/// multi-select flag the caller asked for, plus a best-effort display
@@ -0,0 +1,49 @@
package dev.ayushya.noo
import android.content.BroadcastReceiver
import android.content.Context
import android.content.Intent
import androidx.core.app.NotificationManagerCompat
/**
* Handles the "Keep local" / "Use server" actions on a sync-conflict
* notification (see `SyncWorker.notifyConflicts`). A `BroadcastReceiver`
* can't block on network itself (~10s execution budget), so this only
* dismisses the notification and hands off to a one-shot
* [ConflictResolveWorker] for the actual upload/download.
*/
class SyncConflictReceiver : BroadcastReceiver() {
companion object {
const val ACTION_RESOLVE = "dev.ayushya.noo.action.RESOLVE_SYNC_CONFLICT"
const val EXTRA_ACCOUNT_ID = "accountId"
const val EXTRA_SERVER_URL = "serverUrl"
const val EXTRA_USERNAME = "username"
const val EXTRA_AUTH_HEADER = "authHeader"
const val EXTRA_FILE_ID = "fileId"
const val EXTRA_REMOTE_PATH = "remotePath"
const val EXTRA_REL_PATH = "relPath"
const val EXTRA_RESOLUTION = "resolution" // "local" | "server"
const val EXTRA_NOTIFICATION_ID = "notificationId"
}
override fun onReceive(context: Context, intent: Intent) {
if (intent.action != ACTION_RESOLVE) return
val notificationId = intent.getIntExtra(EXTRA_NOTIFICATION_ID, -1)
if (notificationId != -1) {
NotificationManagerCompat.from(context).cancel(notificationId)
}
ConflictResolveWorker.enqueue(
context,
accountId = intent.getStringExtra(EXTRA_ACCOUNT_ID),
serverUrl = intent.getStringExtra(EXTRA_SERVER_URL),
username = intent.getStringExtra(EXTRA_USERNAME),
authHeader = intent.getStringExtra(EXTRA_AUTH_HEADER),
fileId = intent.getStringExtra(EXTRA_FILE_ID),
remotePath = intent.getStringExtra(EXTRA_REMOTE_PATH),
relPath = intent.getStringExtra(EXTRA_REL_PATH),
resolution = intent.getStringExtra(EXTRA_RESOLUTION),
)
}
}
@@ -0,0 +1,399 @@
package dev.ayushya.noo
import android.content.Context
import android.net.Uri
import android.util.Log
import okhttp3.MediaType.Companion.toMediaType
import okhttp3.OkHttpClient
import okhttp3.Request
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
import java.util.concurrent.TimeUnit
import javax.xml.parsers.DocumentBuilderFactory
import org.w3c.dom.Element
/**
* The actual WebDAV walk/diff/GET/PUT engine for device sync - shared by
* [SyncWorker] (periodic/manual full-folder sync) and
* [ConflictResolveWorker] (single-file resolution from a notification
* action), since both need the same PROPFIND/GET/PUT primitives and the
* same on-disk sync-state bookkeeping.
*
* Deliberately plain Kotlin over `HttpURLConnection`, not
* `NextcloudService`'s Dart/Dio code - this has to run from a
* [androidx.work.CoroutineWorker], independent of the Flutter engine even
* being loaded, the same reason `DownloadService.kt`/`ShareUploadService.kt`
* re-implement GET/PUT natively instead of calling back into Dart.
*/
object SyncEngine {
data class RemoteEntry(
val path: String,
val fileId: String,
val etag: String,
val lastModified: Long,
val size: Long,
val isFolder: Boolean,
)
data class FileState(
val relPath: String,
val etag: String,
val lastModified: Long,
val size: Long,
val localMTime: Long,
)
sealed class SyncAction {
data class Download(val entry: RemoteEntry) : SyncAction()
data class Upload(val relPath: String, val fileId: String) : SyncAction()
data class Delete(val relPath: String, val fileId: String) : SyncAction()
data class Conflict(val entry: RemoteEntry, val relPath: String) : SyncAction()
}
private const val TAG = "NooSync"
private const val STATE_PREFS = "noo_sync_state"
private val rfc1123 =
SimpleDateFormat("EEE, dd MMM yyyy HH:mm:ss zzz", Locale.US)
// Only used for PROPFIND - unlike HttpURLConnection, OkHttp doesn't
// reject non-standard HTTP methods.
private val httpClient = OkHttpClient.Builder()
.connectTimeout(15, TimeUnit.SECONDS)
.readTimeout(30, TimeUnit.SECONDS)
.build()
fun syncRoot(context: Context, accountId: String): File {
val base = context.getExternalFilesDir(null) ?: context.filesDir
return File(base, "sync/$accountId").apply { mkdirs() }
}
private fun stateKey(accountId: String) = "state_$accountId"
fun loadState(context: Context, accountId: String): MutableMap<String, FileState> {
val prefs = context.getSharedPreferences(STATE_PREFS, Context.MODE_PRIVATE)
val json = prefs.getString(stateKey(accountId), null) ?: return mutableMapOf()
val obj = JSONObject(json)
val result = mutableMapOf<String, FileState>()
for (fileId in obj.keys()) {
val entry = obj.getJSONObject(fileId)
result[fileId] = FileState(
relPath = entry.getString("relPath"),
etag = entry.getString("etag"),
lastModified = entry.getLong("lastModified"),
size = entry.getLong("size"),
localMTime = entry.getLong("localMTime"),
)
}
return result
}
fun saveState(context: Context, accountId: String, state: Map<String, FileState>) {
val obj = JSONObject()
for ((fileId, s) in state) {
obj.put(
fileId,
JSONObject().apply {
put("relPath", s.relPath)
put("etag", s.etag)
put("lastModified", s.lastModified)
put("size", s.size)
put("localMTime", s.localMTime)
},
)
}
val prefs = context.getSharedPreferences(STATE_PREFS, Context.MODE_PRIVATE)
prefs.edit().putString(stateKey(accountId), obj.toString()).apply()
}
/** Depth-1 PROPFIND of [remotePath], returning its direct children only. */
fun propfindChildren(
serverUrl: String,
username: String,
authHeader: String,
remotePath: String,
): List<RemoteEntry> = propfind(serverUrl, username, authHeader, remotePath, depth = "1", skipSelf = true)
/** Depth-0 PROPFIND of [remotePath] itself (e.g. to refresh its etag after a PUT/GET). */
fun propfindSelf(
serverUrl: String,
username: String,
authHeader: String,
remotePath: String,
): RemoteEntry? = propfind(serverUrl, username, authHeader, remotePath, depth = "0", skipSelf = false).firstOrNull()
private fun propfind(
serverUrl: String,
username: String,
authHeader: String,
remotePath: String,
depth: String,
skipSelf: Boolean,
): List<RemoteEntry> {
var cleanPath = remotePath.trim()
if (!cleanPath.startsWith("/")) cleanPath = "/$cleanPath"
if (depth == "1" && !cleanPath.endsWith("/")) cleanPath = "$cleanPath/"
val encodedPath = cleanPath.split("/").joinToString("/") { Uri.encode(it) }
val cleanServer = serverUrl.trimEnd('/')
val url = "$cleanServer/remote.php/dav/files/$username$encodedPath"
val body = """<?xml version="1.0" encoding="utf-8" ?>
<d:propfind xmlns:d="DAV:" xmlns:oc="http://owncloud.org/ns">
<d:prop>
<d:getlastmodified/>
<d:getcontentlength/>
<d:resourcetype/>
<d:getetag/>
<oc:fileid/>
</d:prop>
</d:propfind>"""
val request = Request.Builder()
.url(url)
.method("PROPFIND", body.toRequestBody("application/xml".toMediaType()))
.header("Authorization", authHeader)
.header("Depth", depth)
.build()
return try {
httpClient.newCall(request).execute().use { response ->
Log.d(TAG, "PROPFIND $url -> ${response.code}")
if (!response.isSuccessful) {
Log.w(TAG, "PROPFIND $url failed: ${response.code} ${response.body?.string()}")
return emptyList()
}
val doc = DocumentBuilderFactory.newInstance()
.apply { isNamespaceAware = true }
.newDocumentBuilder()
.parse(response.body!!.byteStream())
parsePropfindResponse(doc, username, cleanPath, skipSelf)
}
} catch (e: Exception) {
Log.e(TAG, "PROPFIND $url threw", e)
emptyList()
}
}
private fun parsePropfindResponse(
doc: org.w3c.dom.Document,
username: String,
cleanPath: String,
skipSelf: Boolean,
): List<RemoteEntry> {
val marker = "/remote.php/dav/files/$username"
val selfPath = cleanPath.trimEnd('/')
val responses = doc.getElementsByTagNameNS("DAV:", "response")
val entries = mutableListOf<RemoteEntry>()
for (i in 0 until responses.length) {
val responseEl = responses.item(i) as? Element ?: continue
val hrefRaw = responseEl.getElementsByTagNameNS("DAV:", "href")
.item(0)?.textContent ?: continue
val decodedHref = Uri.decode(hrefRaw)
val idx = decodedHref.indexOf(marker)
val hrefPath = if (idx >= 0) decodedHref.substring(idx + marker.length) else decodedHref
val hrefPathNorm = (if (hrefPath.isEmpty()) "/" else hrefPath).trimEnd('/')
if (skipSelf && hrefPathNorm == selfPath) continue // self entry, not a child
val propEl = responseEl.getElementsByTagNameNS("DAV:", "prop").item(0) as? Element
?: continue
val resourceTypeEl =
propEl.getElementsByTagNameNS("DAV:", "resourcetype").item(0) as? Element
val isFolder =
(resourceTypeEl?.getElementsByTagNameNS("DAV:", "collection")?.length ?: 0) > 0
val etag = propEl.getElementsByTagNameNS("DAV:", "getetag")
.item(0)?.textContent?.trim('"') ?: ""
val lastModStr =
propEl.getElementsByTagNameNS("DAV:", "getlastmodified").item(0)?.textContent
val sizeStr =
propEl.getElementsByTagNameNS("DAV:", "getcontentlength").item(0)?.textContent
val fileId = propEl.getElementsByTagNameNS("http://owncloud.org/ns", "fileid")
.item(0)?.textContent ?: hrefPathNorm
entries.add(
RemoteEntry(
path = hrefPathNorm,
fileId = fileId,
etag = etag,
lastModified = lastModStr?.let { runCatching { rfc1123.parse(it)?.time }.getOrNull() } ?: 0L,
size = sizeStr?.toLongOrNull() ?: 0L,
isFolder = isFolder,
),
)
}
Log.d(TAG, "PROPFIND $cleanPath -> ${entries.size} entries")
return entries
}
/** Recursively walks [rootPath] (Depth-1 PROPFINDs, breadth-first) into a flat manifest. */
fun walkRemoteTree(
serverUrl: String,
username: String,
authHeader: String,
rootPath: String,
): List<RemoteEntry> {
val result = mutableListOf<RemoteEntry>()
val queue = ArrayDeque<String>()
queue.add(rootPath)
while (queue.isNotEmpty()) {
val current = queue.removeFirst()
val children = propfindChildren(serverUrl, username, authHeader, current)
for (child in children) {
result.add(child)
if (child.isFolder) queue.add(child.path)
}
}
return result
}
fun downloadFile(
serverUrl: String,
username: String,
authHeader: String,
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()
}
}
fun uploadFile(
serverUrl: String,
username: String,
authHeader: String,
remotePath: String,
source: 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 = "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()
}
}
/**
* Diffs one synced folder's remote manifest against the persisted
* state, returning what to do for each item - callers own actually
* doing the GET/PUT/delete and updating state afterward, so this stays
* pure/testable-in-principle. Also returns entries no longer present
* remotely (deleted server-side, relative to the [priorFolderRelPaths]
* this folder previously produced) as [SyncAction.Delete].
*/
fun diffFolder(
entries: List<RemoteEntry>,
state: Map<String, FileState>,
priorFolderFileIds: Set<String>,
syncRoot: File,
): List<SyncAction> {
val actions = mutableListOf<SyncAction>()
val seenFileIds = mutableSetOf<String>()
for (entry in entries) {
if (entry.isFolder) continue
seenFileIds.add(entry.fileId)
val relPath = entry.path.removePrefix("/")
val localFile = File(syncRoot, relPath)
val prior = state[entry.fileId]
if (prior == null) {
// First time this folder's been walked with sync enabled -
// nothing locally recorded yet to protect, so just pull it.
actions.add(SyncAction.Download(entry))
continue
}
val serverChanged = entry.etag != prior.etag
val localExists = localFile.exists()
val localChanged = localExists &&
(localFile.lastModified() != prior.localMTime || localFile.length() != prior.size)
when {
!localExists && !serverChanged -> {
// User deleted the local mirror copy themselves and the
// server hasn't changed - respect that deletion rather
// than silently re-creating it.
actions.add(SyncAction.Delete(relPath, entry.fileId))
}
serverChanged && localChanged -> actions.add(SyncAction.Conflict(entry, relPath))
serverChanged -> actions.add(SyncAction.Download(entry))
localChanged -> actions.add(SyncAction.Upload(relPath, entry.fileId))
else -> {} // unchanged, nothing to do
}
}
// Anything this folder had state for last time but that didn't show
// up in this walk at all was deleted server-side.
for (fileId in priorFolderFileIds) {
if (fileId !in seenFileIds) {
val prior = state[fileId] ?: continue
actions.add(SyncAction.Delete(prior.relPath, fileId))
}
}
Log.d(
TAG,
"diffFolder: ${entries.size} entries -> " +
"${actions.count { it is SyncAction.Download }} downloads, " +
"${actions.count { it is SyncAction.Upload }} uploads, " +
"${actions.count { it is SyncAction.Delete }} deletes, " +
"${actions.count { it is SyncAction.Conflict }} conflicts",
)
return actions
}
fun folderIsSyncedUnder(itemPath: String, syncedFolders: List<String>): String? {
val normalizedItem = itemPath.trimEnd('/')
for (folder in syncedFolders) {
val normalizedFolder = folder.trimEnd('/')
if (normalizedItem == normalizedFolder || normalizedItem.startsWith("$normalizedFolder/")) {
return folder
}
}
return null
}
}
@@ -0,0 +1,80 @@
package dev.ayushya.noo
/**
* In-memory, in-process pub/sub for live device-sync status -
* [SyncWorker]/[ConflictResolveWorker] publish into this, `MainActivity`'s
* `dev.ayushya.noo/sync_service/status` EventChannel forwards it to Dart.
* Plain in-memory state is enough (no IPC/persistence needed) since the
* workers and the Activity always run in the same process; contrast with
* [SyncEngine]'s on-disk per-account state map, which *is* durable and is
* what actually answers "is this file synced" (this bus only ever tracks
* the transient "syncing right now" / "unresolved conflict" parts of that
* picture).
*/
object SyncStatusBus {
data class Conflict(
val accountId: String,
val fileId: String,
val remotePath: String,
val relPath: String,
val name: String,
)
data class Status(
val accountId: String?,
val syncing: Boolean,
val syncingFileIds: Set<String>,
val conflicts: List<Conflict>,
)
@Volatile
private var current = Status(accountId = null, syncing = false, syncingFileIds = emptySet(), conflicts = emptyList())
private val listeners = mutableListOf<(Status) -> Unit>()
@Synchronized
fun subscribe(listener: (Status) -> Unit) {
listeners.add(listener)
listener(current)
}
@Synchronized
fun unsubscribe(listener: (Status) -> Unit) {
listeners.remove(listener)
}
fun snapshot(): Status = current
@Synchronized
fun setSyncing(accountId: String, syncing: Boolean) {
current = current.copy(accountId = accountId, syncing = syncing)
if (!syncing) current = current.copy(syncingFileIds = emptySet())
publish()
}
@Synchronized
fun markFileSyncing(accountId: String, fileId: String, syncing: Boolean) {
val ids = current.syncingFileIds.toMutableSet()
if (syncing) ids.add(fileId) else ids.remove(fileId)
current = current.copy(accountId = accountId, syncingFileIds = ids)
publish()
}
@Synchronized
fun addConflicts(accountId: String, newConflicts: List<Conflict>) {
if (newConflicts.isEmpty()) return
val existingIds = current.conflicts.map { it.fileId }.toSet()
val merged = current.conflicts + newConflicts.filter { it.fileId !in existingIds }
current = current.copy(accountId = accountId, conflicts = merged)
publish()
}
@Synchronized
fun removeConflict(accountId: String, fileId: String) {
current = current.copy(accountId = accountId, conflicts = current.conflicts.filter { it.fileId != fileId })
publish()
}
private fun publish() {
listeners.toList().forEach { it(current) }
}
}
@@ -0,0 +1,272 @@
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
import androidx.work.WorkerParameters
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import org.json.JSONArray
import java.io.File
/**
* The periodic (and manual "Sync now") device-sync job. Runs entirely
* independent of the Flutter engine (see [SyncEngine]'s doc comment) -
* enqueued/re-enqueued from `MainActivity.kt`'s
* `dev.ayushya.noo/sync_service` MethodChannel whenever the synced-folder
* list, account, or network setting changes, since a periodic
* `WorkRequest`'s input `Data` is fixed at enqueue time.
*/
class SyncWorker(appContext: Context, params: WorkerParameters) :
CoroutineWorker(appContext, params) {
companion object {
const val KEY_ACCOUNT_ID = "accountId"
const val KEY_SERVER_URL = "serverUrl"
const val KEY_USERNAME = "username"
const val KEY_AUTH_HEADER = "authHeader"
const val KEY_FOLDERS = "folders" // JSON array of remote file/folder paths
const val UNIQUE_PERIODIC_NAME = "noo_sync_periodic"
const val UNIQUE_ONE_OFF_NAME = "noo_sync_now"
private const val CHANNEL_ID = "device_sync"
private const val SUMMARY_NOTIFICATION_ID = 4401
private const val CONFLICT_NOTIFICATION_ID_BASE = 4500
private const val TAG = "NooSync"
}
override suspend fun doWork(): Result = withContext(Dispatchers.IO) {
val accountId = inputData.getString(KEY_ACCOUNT_ID) ?: return@withContext Result.failure()
val serverUrl = inputData.getString(KEY_SERVER_URL) ?: return@withContext Result.failure()
val username = inputData.getString(KEY_USERNAME) ?: return@withContext Result.failure()
val authHeader = inputData.getString(KEY_AUTH_HEADER) ?: return@withContext Result.failure()
val foldersJson = inputData.getString(KEY_FOLDERS) ?: "[]"
val folders = (0 until JSONArray(foldersJson).length()).map { JSONArray(foldersJson).getString(it) }
Log.d(TAG, "doWork: account=$accountId server=$serverUrl user=$username folders=$folders")
if (folders.isEmpty()) return@withContext Result.success()
val syncRoot = SyncEngine.syncRoot(applicationContext, accountId)
val state = SyncEngine.loadState(applicationContext, accountId).toMutableMap()
var downloaded = 0
var uploaded = 0
var deleted = 0
val conflicts = mutableListOf<Pair<SyncEngine.RemoteEntry, String>>()
SyncStatusBus.setSyncing(accountId, true)
try {
for (path in folders) {
// A configured path can be a file or a folder now - check
// which before deciding whether to walk it recursively or
// just diff the single item.
val self = SyncEngine.propfindSelf(serverUrl, username, authHeader, path)
val entries = when {
self == null -> emptyList()
self.isFolder -> SyncEngine.walkRemoteTree(serverUrl, username, authHeader, path)
else -> listOf(self)
}
val pathPrefix = path.trimEnd('/') + "/"
val priorFileIdsForPath = state.filterValues {
it.relPath == path.trimStart('/') || it.relPath.startsWith(pathPrefix.trimStart('/'))
}.keys
val actions = SyncEngine.diffFolder(entries, state, priorFileIdsForPath, syncRoot)
for (action in actions) {
when (action) {
is SyncEngine.SyncAction.Download -> {
SyncStatusBus.markFileSyncing(accountId, action.entry.fileId, true)
val relPath = action.entry.path.removePrefix("/")
val dest = File(syncRoot, relPath)
if (SyncEngine.downloadFile(serverUrl, username, authHeader, action.entry.path, dest)) {
state[action.entry.fileId] = SyncEngine.FileState(
relPath = relPath,
etag = action.entry.etag,
lastModified = action.entry.lastModified,
size = action.entry.size,
localMTime = dest.lastModified(),
)
downloaded++
}
SyncStatusBus.markFileSyncing(accountId, action.entry.fileId, false)
}
is SyncEngine.SyncAction.Upload -> {
SyncStatusBus.markFileSyncing(accountId, action.fileId, true)
val localFile = File(syncRoot, action.relPath)
val remotePath = "/${action.relPath}"
if (localFile.exists() &&
SyncEngine.uploadFile(serverUrl, username, authHeader, remotePath, localFile)
) {
val prior = state[action.fileId]
state[action.fileId] = SyncEngine.FileState(
relPath = action.relPath,
etag = prior?.etag ?: "",
lastModified = localFile.lastModified(),
size = localFile.length(),
localMTime = localFile.lastModified(),
)
uploaded++
}
SyncStatusBus.markFileSyncing(accountId, action.fileId, false)
}
is SyncEngine.SyncAction.Delete -> {
File(syncRoot, action.relPath).delete()
state.remove(action.fileId)
deleted++
}
is SyncEngine.SyncAction.Conflict -> conflicts.add(action.entry to action.relPath)
}
}
}
} finally {
SyncStatusBus.setSyncing(accountId, false)
}
SyncEngine.saveState(applicationContext, accountId, state)
Log.d(
TAG,
"doWork done: downloaded=$downloaded uploaded=$uploaded deleted=$deleted conflicts=${conflicts.size}",
)
if (downloaded > 0 || uploaded > 0 || deleted > 0) {
notifySummary(downloaded, uploaded, deleted)
}
if (conflicts.isNotEmpty()) {
SyncStatusBus.addConflicts(
accountId,
conflicts.map { (entry, relPath) ->
SyncStatusBus.Conflict(
accountId = accountId,
fileId = entry.fileId,
remotePath = entry.path,
relPath = relPath,
name = relPath.substringAfterLast('/'),
)
},
)
notifyConflicts(accountId, serverUrl, username, authHeader, conflicts)
}
Result.success()
}
private fun notifySummary(downloaded: Int, uploaded: Int, deleted: Int) {
createChannel()
val parts = mutableListOf<String>()
if (downloaded > 0) parts.add("$downloaded updated")
if (uploaded > 0) parts.add("$uploaded uploaded")
if (deleted > 0) parts.add("$deleted removed")
val text = parts.joinToString(", ")
val notification = NotificationCompat.Builder(applicationContext, CHANNEL_ID)
.setSmallIcon(android.R.drawable.stat_notify_sync)
.setContentTitle("Noo sync")
.setContentText(text)
.setAutoCancel(true)
.build()
manager().notify(SUMMARY_NOTIFICATION_ID, notification)
}
private fun notifyConflicts(
accountId: String,
serverUrl: String,
username: String,
authHeader: String,
conflicts: List<Pair<SyncEngine.RemoteEntry, String>>,
) {
createChannel()
for ((index, conflict) in conflicts.withIndex()) {
val (entry, relPath) = conflict
val fileName = relPath.substringAfterLast('/')
val notificationId = CONFLICT_NOTIFICATION_ID_BASE + (relPath.hashCode() and 0xFFFF)
val useLocalIntent = conflictActionIntent(
accountId,
serverUrl,
username,
authHeader,
entry.fileId,
entry.path,
relPath,
"local",
notificationId,
)
val useServerIntent = conflictActionIntent(
accountId,
serverUrl,
username,
authHeader,
entry.fileId,
entry.path,
relPath,
"server",
notificationId,
)
val notification = NotificationCompat.Builder(applicationContext, CHANNEL_ID)
.setSmallIcon(android.R.drawable.stat_notify_error)
.setContentTitle("Sync conflict: $fileName")
.setContentText("Changed both on this device and on the server.")
.setAutoCancel(true)
.addAction(0, "Keep local", useLocalIntent)
.addAction(0, "Use server", useServerIntent)
.build()
manager().notify(notificationId, notification)
}
}
private fun conflictActionIntent(
accountId: String,
serverUrl: String,
username: String,
authHeader: String,
fileId: String,
remotePath: String,
relPath: String,
resolution: String,
notificationId: Int,
): PendingIntent {
val intent = Intent(applicationContext, SyncConflictReceiver::class.java).apply {
action = SyncConflictReceiver.ACTION_RESOLVE
putExtra(SyncConflictReceiver.EXTRA_ACCOUNT_ID, accountId)
putExtra(SyncConflictReceiver.EXTRA_SERVER_URL, serverUrl)
putExtra(SyncConflictReceiver.EXTRA_USERNAME, username)
putExtra(SyncConflictReceiver.EXTRA_AUTH_HEADER, authHeader)
putExtra(SyncConflictReceiver.EXTRA_FILE_ID, fileId)
putExtra(SyncConflictReceiver.EXTRA_REMOTE_PATH, remotePath)
putExtra(SyncConflictReceiver.EXTRA_REL_PATH, relPath)
putExtra(SyncConflictReceiver.EXTRA_RESOLUTION, resolution)
putExtra(SyncConflictReceiver.EXTRA_NOTIFICATION_ID, notificationId)
}
// Request code must be unique per (file, resolution) pair, else the
// two actions' PendingIntents collide and only one survives.
val requestCode = (relPath + resolution).hashCode()
return PendingIntent.getBroadcast(
applicationContext,
requestCode,
intent,
PendingIntent.FLAG_UPDATE_CURRENT or PendingIntent.FLAG_IMMUTABLE,
)
}
private fun manager() =
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)
}
}
}