sync(chunk 4): incremental sync and scheduling

RFC 6578 as an optimisation on top of chunk 3's full path, and sync that runs
by itself.

- :caldav gains sync-collection, with invalidation matched on
  DAV:valid-sync-token in the body on any 4xx rather than on a status code —
  400, 403, 409 and 412 are all used in the wild, and matching the status is
  why Thunderbird never recovers from sabre's 403. No DAV:limit (Nextcloud
  regressed it to an HTML error page); 507 on our own href is truncation.
- The engine persists the token per page, after the bodies. An iteration cap
  and a no-progress guard, since the RFC never requires the token to advance.
  The full path runs anyway every 24h: a token the server accepts over a
  pruned change log returns 207, zero changes and no error, and the protocol
  gives no other way to notice.
- The three membership traps: an unknown removed href is a no-op, a
  delete-then-recreate is re-identified from the UID in the body, and a mass
  removal is refused in favour of a real listing, because ACL churn looks
  exactly like one.
- Scheduling is a PeriodicWorkRequest with a network constraint, expedited
  only for the in-app button. getForegroundInfo is implemented unconditionally
  (setExpedited falls back to a foreground service below API 31 and the
  default throws) but declares no service type, which would have pulled back
  the Android 15 dataSync budget and a Play video-demo requirement.

initialIncomplete turned out not to be needed: adopting a token only after a
full reconciliation completes removes the hazard it guarded, so there is
nothing to persist atomically with anything.

/code-review high raised 8 findings, all fixed. The two that mattered: the
cadence clock sat in the backed-up DataStore, so a restore would have made the
engine trust a stale token for a day — both sync-state stores now have their
own excluded file; and the incremental download path lacked the write-phase
guard, overwriting local edits that had never reached the server. Reasoning in
docs/SYNC-PLAN.md.

Chunk 2's on-device review is still outstanding; none of this has run on a
device or against a real server.
This commit is contained in:
2026-09-07 16:18:09 +02:00
parent b1189a4884
commit b25f8b231c
24 changed files with 1376 additions and 33 deletions
@@ -20,6 +20,8 @@ import de.jeanlucmakiola.agendula.data.demo.DemoSeeder
import de.jeanlucmakiola.agendula.data.prefs.ThemeMode
import de.jeanlucmakiola.agendula.ui.RootScreen
import de.jeanlucmakiola.agendula.ui.crash.CrashReportActivity
import de.jeanlucmakiola.agendula.data.sync.AccountRepository
import de.jeanlucmakiola.agendula.data.sync.SyncTrigger
import de.jeanlucmakiola.agendula.ui.settings.SettingsViewModel
import de.jeanlucmakiola.agendula.ui.theme.AgendulaTheme
import de.jeanlucmakiola.floret.crash.CrashReportDialog
@@ -38,6 +40,10 @@ class MainActivity : ComponentActivity() {
@Inject lateinit var demoSeeder: DemoSeeder
@Inject lateinit var accounts: AccountRepository
@Inject lateinit var syncTrigger: SyncTrigger
// A captured crash report awaiting the user's decision, surfaced as a dialog
// over the app on the next launch (the single-crash path). A startup
// crash-loop is handled out of band, before setContent — see below.
@@ -60,6 +66,24 @@ class MainActivity : ComponentActivity() {
// Surface a single captured crash as a dialog on the next launch.
if (CrashReporter.shouldPrompt(this)) pendingCrashReport = CrashReporter.pendingReport(this)
// Sync hard on app open: the periodic worker's interval is a floor, and
// in the `rare` and `restricted` App Standby buckets it may not have run
// at all. `KEEP` makes rescheduling idempotent, so this also repairs a
// schedule lost to "clear app data" or to a restore.
//
// ⚠️ Only on a genuine open. This activity declares no `configChanges`,
// so onCreate runs again on every rotation, theme switch, locale change
// and font-scale change — each of which would otherwise start a fresh
// network sync the moment the previous one finished.
if (savedInstanceState == null) {
lifecycleScope.launch {
runCatching {
accounts.rescheduleAll()
accounts.all().forEach { syncTrigger.enqueue(it.displayName) }
}
}
}
// Debug-only sample data: `am start ... --ez agendula_seed true`. Seeds a
// local (non-syncing) demo list once; no-op without the extra.
if (BuildConfig.DEBUG && intent.getBooleanExtra(EXTRA_SEED, false)) {
@@ -42,12 +42,23 @@ private val Context.credentialsDataStore: DataStore<Preferences> by preferencesD
name = CREDENTIALS_DATASTORE,
)
/** See [SyncStateDataStore] for why this is a separate file. */
private val Context.syncStateDataStore: DataStore<Preferences> by preferencesDataStore(
name = SYNC_STATE_DATASTORE,
)
/**
* Named here and in `backup_rules.xml` / `data_extraction_rules.xml`, which
* exclude `datastore/$CREDENTIALS_DATASTORE.preferences_pb` by this name.
*/
const val CREDENTIALS_DATASTORE = "agendula_credentials"
/**
* Named here and in `backup_rules.xml` / `data_extraction_rules.xml`, which
* exclude `datastore/$SYNC_STATE_DATASTORE.preferences_pb` by this name.
*/
const val SYNC_STATE_DATASTORE = "agendula_sync_state"
@Module
@InstallIn(SingletonComponent::class)
abstract class DataBindModule {
@@ -84,6 +95,12 @@ object DataProvideModule {
fun provideCredentialsDataStore(@ApplicationContext context: Context): DataStore<Preferences> =
context.credentialsDataStore
@Provides
@Singleton
@SyncStateDataStore
fun provideSyncStateDataStore(@ApplicationContext context: Context): DataStore<Preferences> =
context.syncStateDataStore
@Provides
@Singleton
fun provideTasksDatabase(@ApplicationContext context: Context): TasksDatabase =
@@ -30,3 +30,21 @@ annotation class ApplicationScope
@Qualifier
@Retention(AnnotationRetention.BINARY)
annotation class CredentialsDataStore
/**
* Marks the DataStore holding per-device **sync bookkeeping** — the quarantine
* counters and the full-reconciliation clock.
*
* Its own file for the same reason the credentials have one: Auto Backup
* includes `datastore/`, and every value in here is a statement about *this*
* device's conversation with a server. Restored onto a new install they are all
* lies, and two of them are dangerous — a restored "reconciled recently" makes
* the engine trust a sync token for another day, which is precisely the silently
* pruned change log the full path exists to catch, and a restored quarantine
* count silently skips resources that were never tried here.
*
* Not user data, so nothing is lost by excluding it.
*/
@Qualifier
@Retention(AnnotationRetention.BINARY)
annotation class SyncStateDataStore
@@ -43,6 +43,8 @@ class AccountRepository @Inject constructor(
private val database: TasksDatabase,
private val credentials: CredentialStore,
private val accounts: CalDavAccounts,
private val syncTrigger: SyncTrigger,
private val cadence: SyncCadenceStore,
@IoDispatcher private val io: CoroutineDispatcher,
) : AccountCreator {
@@ -125,6 +127,12 @@ class AccountRepository @Inject constructor(
return@withContext Outcome.AlreadyExists
}
// On the schedule from the moment it exists, and syncing immediately —
// an account that shows up empty until the first periodic window looks
// broken.
syncTrigger.schedule(displayName)
syncTrigger.enqueue(displayName)
Outcome.Created(accountId)
}
@@ -141,6 +149,18 @@ class AccountRepository @Inject constructor(
database.accounts().delete(accountId)
}
/**
* Puts every existing account back on the periodic schedule.
*
* Cheap and idempotent — `KEEP` means an already-scheduled account is left
* exactly as it is — so calling it on app open costs nothing and repairs the
* one case WorkManager cannot: a schedule lost to "clear app data" or to a
* restore onto a device that never ran the account-add flow.
*/
suspend fun rescheduleAll() = withContext(io) {
database.accounts().all().forEach { syncTrigger.schedule(it.displayName) }
}
/**
* Removes an account and everything that keys off it.
*
@@ -149,6 +169,13 @@ class AccountRepository @Inject constructor(
* instruction to destroy the tasks it held.
*/
suspend fun remove(accountId: Long, displayName: String) = withContext(io) {
syncTrigger.cancel(displayName)
// The lists survive as device-only lists, so their cursors must not: a
// re-added account would otherwise inherit a "reconciled recently" that
// was true of a different account's data.
cadence.forget(
database.taskLists().syncedForAccount(accountId).map { it.id }.toSet(),
)
credentials.clear(accountId)
database.accounts().delete(accountId)
accounts.find(displayName)?.let { accounts.remove(it) }
@@ -5,9 +5,11 @@ import de.jeanlucmakiola.agendula.data.tasks.ical.ResourceValidator
import de.jeanlucmakiola.agendula.data.tasks.ical.VTodoMapper
import de.jeanlucmakiola.agendula.data.tasks.room.TaskEntity
import de.jeanlucmakiola.agendula.data.tasks.room.TaskListEntity
import de.jeanlucmakiola.caldav.ChangeSet
import de.jeanlucmakiola.caldav.DeleteOutcome
import de.jeanlucmakiola.caldav.PutOutcome
import de.jeanlucmakiola.caldav.RemoteCalendar
import de.jeanlucmakiola.caldav.RemoteRef
import de.jeanlucmakiola.caldav.ResourceNames
import okhttp3.HttpUrl
import okhttp3.HttpUrl.Companion.toHttpUrlOrNull
@@ -46,8 +48,9 @@ class CollectionSyncer(
list: TaskListEntity,
remote: RemoteCalendar,
quarantine: MutableMap<String, Int>,
fullReconciliationDue: Boolean = true,
): SyncReport {
val run = Run(list, remote, quarantine)
val run = Run(list, remote, quarantine, fullReconciliationDue)
return try {
run.execute()
} catch (e: Exception) {
@@ -61,6 +64,7 @@ class CollectionSyncer(
val list: TaskListEntity,
val remote: RemoteCalendar,
val quarantine: MutableMap<String, Int>,
val fullReconciliationDue: Boolean,
) {
var report = SyncReport(listId = list.id, listName = list.name)
@@ -93,6 +97,16 @@ class CollectionSyncer(
var writable = true
var shared = false
/**
* The token seen in the PROPFIND at the start of this run.
*
* Captured **before** anything is read, and adopted only after a full
* reconciliation completes. Taking it afterwards would silently swallow
* every change made while we were reading; taking it before merely
* re-reports those next time, which costs one ETag comparison.
*/
var pendingToken: String? = null
fun execute(): SyncReport {
val state = remote.state().getOrElse {
return report.copy(failure = "collection unavailable: $it")
@@ -100,9 +114,56 @@ class CollectionSyncer(
writable = !state.collection.readOnly
shared = state.collection.isShared
persistCollectionState(state.collection.readOnly)
pendingToken = state.syncToken
// Local work first, in both paths. A tombstone must never race the
// download of the resource it is about to remove, and a conflict has
// to be discovered while the local edit still exists.
val pending = localResources()
deletePhase(pending)
uploadPhase(pending)
val cursor = list.syncToken?.takeIf {
// ⚠️ A hint, not a contract — Radicale advertised the report for
// years without implementing it.
state.collection.supportsSyncCollection && !fullReconciliationDue
}
if (cursor == null || !runIncremental(cursor)) {
// ⚠️ Whatever the incremental attempt recorded is a note about a
// path we did not end up using. Leaving it as *the collection's*
// failure refuses the new token, skips the cadence record and
// shows the user "last sync didn't finish" — on a run that
// reconciled the collection perfectly.
val attempt = report.failure
report = report.copy(failure = null)
runFull()
if (report.failure == null && attempt != null) {
report = report.copy(incrementalNote = attempt)
}
} else {
// The write phase asked for fresh copies the change log may never
// mention — a server does not have to echo our own writes back to
// us — so they are fetched explicitly rather than hoped for.
fetchAndApply(refetch.filterNot { isQuarantined(it) })
}
resolveParents()
return report
}
/**
* The full path: a complete listing, an ETag diff, and a sweep.
*
* ⚠️ Not a fallback — a **permanent safety net**. A token the server
* accepts over a change log it has already pruned answers 207 with zero
* changes and no error, and RFC 6578 offers no way to detect it. Running
* this on a slow cadence regardless of the token is the only mitigation
* there is.
*/
private fun runFull() {
val refs = remote.list().getOrElse {
return report.copy(failure = "listing failed: $it")
report = report.copy(failure = "listing failed: $it")
return
}
// ⚠️ Only strong tags are carried forward. A weak one cannot serve
// as `If-Match`, so recording it would make every later write look
@@ -112,14 +173,137 @@ class CollectionSyncer(
}
val locals = localResources()
deletePhase(locals)
uploadPhase(locals)
downloadPhase(locals, remoteETags)
resolveParents()
sweepPhase(locals, remoteETags.keys)
return report
report = report.copy(reconciledInFull = true)
// ⚠️ Adopted only now, and taken from the PROPFIND that *preceded*
// the listing. A token minted after the read would silently swallow
// anything that changed during it; one minted before merely re-reports
// it next time, and a re-report costs an ETag comparison.
if (report.failure == null) store.setSyncToken(list.id, pendingToken)
}
/**
* The incremental path.
*
* @return true when the collection is fully reconciled by change log
* alone; false to fall back to [runFull] this run.
*/
private fun runIncremental(token: String): Boolean {
var cursor = token
repeat(MAX_SYNC_PAGES) {
val page = when (val changes = remote.changes(cursor)) {
is ChangeSet.Page -> changes
ChangeSet.TokenInvalid -> {
// Clear it, so a crash before the full run below does not
// leave a token we already know the server rejects.
store.setSyncToken(list.id, null)
return false
}
ChangeSet.Unsupported -> return false
is ChangeSet.Failed -> {
report = report.copy(failure = "sync-collection failed: ${changes.reason}")
return false
}
}
if (!removalsArePlausible(page.removed)) return false
applyRemovals(page.removed)
downloadChanged(page.changed)
// ⚠️ After the bodies, never before.
store.setSyncToken(list.id, page.token)
val next = page.token
when {
// A 207 with no token at all. The changes were real; the
// cursor is gone, so the next read has to be a full one.
next == null -> return false
!page.truncated -> return true
// ⚠️ The RFC never requires the token to advance, so a server
// that returns the same one forever would loop until the cap.
next == cursor -> {
report = report.copy(
failure = "the server truncated without advancing its sync token",
)
return false
}
else -> cursor = next
}
}
report = report.copy(failure = "sync-collection did not finish in $MAX_SYNC_PAGES pages")
return false
}
/**
* ⚠️ ACL churn can arrive as a mass removal.
*
* A calendar whose share is revoked, or whose permissions change, can be
* reported as every resource in it disappearing at once. Acting on that
* deletes the user's data on the strength of a change log; reconciling
* against a real listing instead costs one PROPFIND.
*/
private fun removalsArePlausible(removed: List<HttpUrl>): Boolean {
if (removed.size < MIN_REMOVALS_TO_QUESTION) return true
val known = store.rowsIn(list.id).mapNotNull { it.href }.distinct().size
if (known == 0 || removed.size * 2 <= known) return true
report = report.copy(
failure = "the change log removed ${removed.size} of $known resources at once — " +
"reconciled against a full listing instead",
)
return false
}
private fun applyRemovals(removed: List<HttpUrl>) {
if (removed.isEmpty()) return
val byHref = localResources().filter { it.href != null }.associateBy { it.href!! }
removed.forEach { url ->
val href = url.toString()
// ⚠️ An unknown href is a no-op, not an error. A resource created
// and deleted between two syncs is reported as removed without
// ever having been reported as added.
val local = byHref[href] ?: return@forEach
// ⚠️ The change log describes the past. If this run has already
// written to that href, the entry predates our write and acting on
// it deletes the row we just uploaded — leaving an orphan on the
// server and nothing here.
if (href in touched) return@forEach
if (local.isDirty) discard(local, DiscardedEdit.Cause.DELETED_ON_SERVER)
purge(local)
report = report.copy(deletedLocally = report.deletedLocally + 1)
}
}
private fun downloadChanged(changed: List<RemoteRef>) {
if (changed.isEmpty()) return
val byHref = localResources().filter { it.href != null }.associateBy { it.href!! }
val wanted = changed.filterNot { ref ->
val href = ref.href.toString()
val local = byHref[href] ?: return@filterNot false
// ⚠️ A row still dirty after the write phase is an edit that never
// reached the server — the collection was demoted to read-only, or
// the upload was refused. Downloading over it destroys the user's
// work silently, with nothing in the report.
//
// Except when the write phase *asked* for the fresh copy: that is
// the 412 path, where the row is deliberately left dirty until its
// replacement lands.
if ((local.isDirty || local.isDeleted) && href !in refetch) return@filterNot true
val eTag = ref.eTag?.takeIf { it.usable }?.value
// Unchanged only when both sides have a strong tag and they agree.
local.eTag != null && eTag != null && local.eTag == eTag
}.map { it.href.toString() }.filterNot { isQuarantined(it) }
fetchAndApply(wanted)
}
// ------------------------------------------------------------ phase 1
@@ -339,22 +523,34 @@ class CollectionSyncer(
wanted += refetch
wanted.removeAll { isQuarantined(it) }
// Built once. Re-reading the whole list per resource turns a first
// sync of a large collection into a quadratic scan on the sync thread.
val byUid = store.rowsIn(list.id).groupBy { it.uid }.toMutableMap()
fetchAndApply(wanted)
}
wanted.chunked(DOWNLOAD_BATCH).forEach { batch ->
/** Shared by both read paths: download these hrefs and write them down. */
private fun fetchAndApply(hrefs: Collection<String>) {
if (hrefs.isEmpty()) return
// Both indexes built once. Re-reading the whole list per resource
// turns a first sync of a large collection into a quadratic scan on
// the sync thread.
val rows = store.rowsIn(list.id)
val index = RowIndex(
byUid = rows.groupBy { it.uid }.toMutableMap(),
byHref = rows.filter { it.href != null }.groupBy { it.href!! }.toMutableMap(),
)
hrefs.chunked(DOWNLOAD_BATCH).forEach { batch ->
val fetched = remote.fetch(batch.map(::url)).getOrElse { error ->
report = report.copy(failure = "download failed: $error")
return
}
fetched.resources.forEach { apply(it, byUid) }
fetched.resources.forEach { apply(it, index) }
}
}
private fun apply(
resource: de.jeanlucmakiola.caldav.RemoteResource,
byUid: MutableMap<String, List<TaskEntity>>,
index: RowIndex,
) {
val href = resource.href.toString()
val todos = runCatching {
@@ -385,7 +581,16 @@ class CollectionSyncer(
// matched pair, and a weak one cannot be used as `If-Match` at all.
val eTag = resource.eTag?.takeIf { it.usable }?.value
val existing = byUid[uid].orEmpty().associateBy { it.recurrenceId }
val existing = index.byUid[uid].orEmpty().associateBy { it.recurrenceId }
// ⚠️ Delete-then-recreate at the same URI is reported as a *change*,
// not as a removal followed by an addition — so the identity in the
// body is the only thing that says the old task is gone. Rows still
// holding this href under a different UID are that old task.
//
// Read before the writes below, which re-point this href's index
// entry at the rows we are about to create.
val displaced = index.byHref[href].orEmpty().filter { it.uid != uid }
// Captured before the overwrite, because that is what the report is
// about: the version the user is losing.
@@ -410,7 +615,17 @@ class CollectionSyncer(
store.deleteAll(stale)
// Keep the hoisted index honest for the resources still to come.
byUid[uid] = written
index.byUid[uid] = written
index.byHref[href] = written
if (displaced.isNotEmpty()) {
val ids = displaced.map { it.id }
store.deleteAll(ids)
displaced.map { it.uid }.distinct().forEach { other ->
index.byUid[other] = index.byUid[other].orEmpty().filterNot { it.id in ids }
}
report = report.copy(deletedLocally = report.deletedLocally + 1)
}
pendingDiscard.remove(href)?.let { cause ->
report = report.copy(
@@ -583,6 +798,18 @@ class CollectionSyncer(
(quarantine[QuarantineStore.key(list.id, href)] ?: 0) >= QuarantineStore.THRESHOLD
}
/**
* The list's rows, indexed both ways [apply] needs them.
*
* Kept for the whole download phase and updated in place, because both
* lookups are per-resource: one to find the rows this UID already has, one to
* find rows that hold this href under a *different* UID.
*/
private class RowIndex(
val byUid: MutableMap<String, List<TaskEntity>>,
val byHref: MutableMap<String, List<TaskEntity>>,
)
/**
* The rows that make up one calendar resource.
*
@@ -626,5 +853,21 @@ class CollectionSyncer(
const val CREATE_ATTEMPTS = 3
const val DOWNLOAD_BATCH = 30
/**
* Pages of `sync-collection` before giving up and reconciling in full.
*
* A cap *and* a no-progress guard, because RFC 6578 never requires the
* token to advance — a server can legitimately truncate forever.
*/
const val MAX_SYNC_PAGES = 50
/**
* Removals in one page below which the sanity threshold does not apply.
*
* Deleting a handful of tasks is ordinary; being told the whole
* collection vanished is what ACL churn looks like.
*/
const val MIN_REMOVALS_TO_QUESTION = 10
}
}
@@ -4,6 +4,7 @@ import androidx.datastore.core.DataStore
import androidx.datastore.preferences.core.Preferences
import androidx.datastore.preferences.core.edit
import androidx.datastore.preferences.core.stringSetPreferencesKey
import de.jeanlucmakiola.agendula.data.di.SyncStateDataStore
import kotlinx.coroutines.flow.first
import javax.inject.Inject
import javax.inject.Singleton
@@ -23,7 +24,7 @@ import javax.inject.Singleton
*/
@Singleton
class QuarantineStore @Inject constructor(
private val dataStore: DataStore<Preferences>,
@SyncStateDataStore private val dataStore: DataStore<Preferences>,
) {
/** Current failure counts, keyed by [key]. */
@@ -0,0 +1,85 @@
package de.jeanlucmakiola.agendula.data.sync
import androidx.datastore.core.DataStore
import androidx.datastore.preferences.core.Preferences
import androidx.datastore.preferences.core.edit
import androidx.datastore.preferences.core.stringSetPreferencesKey
import de.jeanlucmakiola.agendula.data.di.SyncStateDataStore
import kotlinx.coroutines.flow.first
import javax.inject.Inject
import javax.inject.Singleton
import kotlin.time.Duration
import kotlin.time.Duration.Companion.hours
import kotlin.time.Instant
/**
* When each collection was last reconciled against a full listing.
*
* ⚠️ This is the mitigation for the one RFC 6578 failure that has no signal at
* all: a token the server still accepts, over a change log it has already
* pruned, answers `207` with zero changes and no error. Nothing in the protocol
* distinguishes that from "nothing happened". The only defence is to stop
* trusting the token periodically and diff a real listing — so the full path is
* a permanent safety net, not a fallback, and this is its clock.
*
* Kept out of Room deliberately: it is scheduling bookkeeping, not user data,
* and it must never be part of a backup that could restore a stale "we checked
* recently" into a fresh install.
*/
@Singleton
class SyncCadenceStore @Inject constructor(
@SyncStateDataStore private val dataStore: DataStore<Preferences>,
) {
/** Last full reconciliation per list id. */
suspend fun lastFullSync(): Map<Long, Instant> =
dataStore.data.first()[KEY].orEmpty().mapNotNull { entry ->
val separator = entry.lastIndexOf(SEPARATOR)
if (separator <= 0) return@mapNotNull null
val id = entry.substring(0, separator).toLongOrNull() ?: return@mapNotNull null
val at = entry.substring(separator + 1).toLongOrNull() ?: return@mapNotNull null
id to Instant.fromEpochSeconds(at)
}.toMap()
/** Merges, rather than replacing, so concurrent accounts do not erase each other. */
suspend fun record(reconciled: Map<Long, Instant>) {
if (reconciled.isEmpty()) return
dataStore.edit { prefs ->
val current = prefs[KEY].orEmpty()
.mapNotNull { entry ->
val separator = entry.lastIndexOf(SEPARATOR)
if (separator <= 0) null else entry.substring(0, separator) to entry
}
.toMap()
.toMutableMap()
reconciled.forEach { (id, at) ->
current["$id"] = "$id$SEPARATOR${at.epochSeconds}"
}
prefs[KEY] = current.values.toSet()
}
}
/** Forgets a list, so a re-added account starts from a full reconciliation. */
suspend fun forget(listIds: Set<Long>) {
if (listIds.isEmpty()) return
dataStore.edit { prefs ->
prefs[KEY] = prefs[KEY].orEmpty().filterNot { entry ->
entry.substringBefore(SEPARATOR).toLongOrNull() in listIds
}.toSet()
}
}
companion object {
/**
* How long a sync token is trusted before a full listing is diffed anyway.
*
* Long enough that the incremental path still carries almost every sync,
* short enough that a silently pruned change log is a day's divergence
* rather than an indefinite one.
*/
val FULL_RECONCILIATION_INTERVAL: Duration = 24.hours
private const val SEPARATOR = '@'
private val KEY = stringSetPreferencesKey("sync_last_full")
}
}
@@ -27,6 +27,7 @@ class SyncEngine @Inject constructor(
private val store: RoomSyncStore,
private val credentials: CredentialStore,
private val quarantine: QuarantineStore,
private val cadence: SyncCadenceStore,
@IoDispatcher private val io: CoroutineDispatcher,
) {
@@ -45,16 +46,19 @@ class SyncEngine @Inject constructor(
?: return@withContext Result.Misconfigured("no such account: $accountName")
val username = account.username
?: return@withContext Result.Misconfigured("account has no username")
?: return@withContext fatal(account.id, Result.Misconfigured("account has no username"))
val origin = account.principalUrl?.toHttpUrlOrNull()
?: return@withContext Result.Misconfigured("account has no principal URL")
?: return@withContext fatal(
account.id,
Result.Misconfigured("account has no principal URL"),
)
val password = when (val secret = credentials.get(account.id)) {
is CredentialStore.Secret.Present -> secret.value
CredentialStore.Secret.Absent ->
return@withContext Result.NeedsSignIn("no stored password")
return@withContext fatal(account.id, Result.NeedsSignIn("no stored password"))
is CredentialStore.Secret.Unrecoverable ->
return@withContext Result.NeedsSignIn(secret.reason)
return@withContext fatal(account.id, Result.NeedsSignIn(secret.reason))
}
val client = CalDavHttp.authenticated(USER_AGENT, username, password, origin)
@@ -96,12 +100,28 @@ class SyncEngine @Inject constructor(
val counts = before.toMutableMap()
val syncer = CollectionSyncer(store)
val now = kotlin.time.Clock.System.now()
val lastFull = cadence.lastFullSync()
val reports = lists.map { list ->
val url = list.href?.toHttpUrlOrNull()
?: return@map SyncReport(list.id, list.name, failure = "list has no collection URL")
syncer.sync(list, remoteFor(url), counts)
val since = lastFull[list.id]
syncer.sync(
list = list,
remote = remoteFor(url),
quarantine = counts,
// Never reconciled, or the token has been trusted long enough.
fullReconciliationDue = since == null ||
now - since >= SyncCadenceStore.FULL_RECONCILIATION_INTERVAL,
)
}
cadence.record(
reports.filter { it.reconciledInFull && it.failure == null }
.associate { it.listId to now },
)
fun mine(key: String) = key.substringBefore('|').toLongOrNull() in listIds
quarantine.merge(
updates = counts.filterKeys(::mine),
@@ -109,7 +129,7 @@ class SyncEngine @Inject constructor(
)
database.accounts().recordSync(
accountId = account.id,
at = kotlin.time.Clock.System.now(),
at = now,
error = reports.mapNotNull { it.failure }.firstOrNull(),
)
return reports
@@ -26,6 +26,25 @@ data class SyncReport(
* overwritten without us noticing, so it is said out loud.
*/
val unconditionalWrites: Int = 0,
/**
* Whether this run reconciled against a full listing rather than a change log.
*
* ⚠️ Tracked because a token the server accepts over a change log it has
* already pruned returns 207, zero changes and no error — RFC 6578 gives no
* signal for it at all. The only mitigation is to reconcile in full on a slow
* cadence regardless of the token, which means knowing when we last did.
*/
val reconciledInFull: Boolean = false,
/**
* Why the change-log path was abandoned, on a run the full path then
* completed.
*
* Not a [failure]: the collection is reconciled and the user has nothing to
* act on. Kept because a server that rejects `sync-collection` every time
* will do it again, and that is worth seeing in a log without it becoming an
* error in the UI.
*/
val incrementalNote: String? = null,
/** Set when the collection failed as a whole. The account keeps going. */
val failure: String? = null,
) {
@@ -40,6 +40,9 @@ interface SyncStore {
* ran.
*/
fun setListReadOnly(listId: Long, readOnly: Boolean)
/** The RFC 6578 cursor. Null resets the collection to a full reconciliation. */
fun setSyncToken(listId: Long, token: String?)
}
class RoomSyncStore @Inject constructor(
@@ -73,4 +76,8 @@ class RoomSyncStore @Inject constructor(
override fun setListReadOnly(listId: Long, readOnly: Boolean) {
database.taskLists().setReadOnly(listId, readOnly)
}
override fun setSyncToken(listId: Long, token: String?) {
database.taskLists().setSyncToken(listId, token)
}
}
@@ -1,13 +1,21 @@
package de.jeanlucmakiola.agendula.data.sync
import android.content.Context
import androidx.work.Constraints
import androidx.work.Data
import androidx.work.ExistingPeriodicWorkPolicy
import androidx.work.ExistingWorkPolicy
import androidx.work.NetworkType
import androidx.work.OneTimeWorkRequestBuilder
import androidx.work.OutOfQuotaPolicy
import androidx.work.PeriodicWorkRequestBuilder
import androidx.work.WorkManager
import dagger.hilt.android.qualifiers.ApplicationContext
import java.util.concurrent.TimeUnit
import javax.inject.Inject
import javax.inject.Singleton
import kotlin.time.Duration
import kotlin.time.Duration.Companion.hours
/**
* Starts a sync for one account.
@@ -23,13 +31,25 @@ class SyncTrigger @Inject constructor(
@ApplicationContext private val context: Context,
) {
/** @return the unique work name, which the caller may wait on. */
fun enqueue(accountName: String): String {
/**
* Starts a sync now.
*
* @param expedited for a trigger the user is looking at. ⚠️ Paired with
* `RUN_AS_NON_EXPEDITED_WORK_REQUEST`, which is not optional: the expedited
* quota is per-app and exhaustible, and the alternative policy
* (`DROP_WORK_REQUEST`) silently discards the sync the user just asked for.
* Never set from a background trigger — a boot receiver spending the quota
* leaves none for the button.
* @return the unique work name, which the caller may wait on.
*/
fun enqueue(accountName: String, expedited: Boolean = false): String {
val uniqueName = SyncWorker.uniqueNameFor(accountName)
val request = OneTimeWorkRequestBuilder<SyncWorker>()
.setInputData(
Data.Builder().putString(SyncWorker.KEY_ACCOUNT_NAME, accountName).build(),
)
.setInputData(inputFor(accountName))
.setConstraints(NETWORK)
.apply {
if (expedited) setExpedited(OutOfQuotaPolicy.RUN_AS_NON_EXPEDITED_WORK_REQUEST)
}
.build()
WorkManager.getInstance(context).enqueueUniqueWork(
@@ -41,4 +61,63 @@ class SyncTrigger @Inject constructor(
)
return uniqueName
}
/**
* Puts the account on the periodic schedule.
*
* A plain `PeriodicWorkRequest` and no foreground service, deliberately —
* see [SyncWorker]. WorkManager restores its own schedule after a reboot, so
* nothing has to re-arm this from `BOOT_COMPLETED`; that matters because
* Android 15 forbids starting a `dataSync` foreground service from boot, and
* a design that needed one would have no way to run at all.
*
* ⚠️ Be honest about the cadence in the UI. [INTERVAL] is a floor, not a
* promise: in the `rare` and `restricted` App Standby buckets network access
* is off entirely, and the genuine worst case is once overnight.
*/
fun schedule(accountName: String) {
val request = PeriodicWorkRequestBuilder<SyncWorker>(
INTERVAL.inWholeMinutes, TimeUnit.MINUTES,
FLEX.inWholeMinutes, TimeUnit.MINUTES,
)
.setInputData(inputFor(accountName))
.setConstraints(NETWORK)
.build()
WorkManager.getInstance(context).enqueueUniquePeriodicWork(
periodicNameFor(accountName),
// UPDATE would restart the interval on every app launch, so a device
// that is opened often would never reach the end of one.
ExistingPeriodicWorkPolicy.KEEP,
request,
)
}
/** Takes a removed account off the schedule. */
fun cancel(accountName: String) {
WorkManager.getInstance(context).apply {
cancelUniqueWork(periodicNameFor(accountName))
cancelUniqueWork(SyncWorker.uniqueNameFor(accountName))
}
}
private fun inputFor(accountName: String) =
Data.Builder().putString(SyncWorker.KEY_ACCOUNT_NAME, accountName).build()
private companion object {
val INTERVAL: Duration = 4.hours
/** The tail of each interval the system may run us in. */
val FLEX: Duration = 1.hours
/**
* Sync needs a network, and saying so lets WorkManager run us the moment
* connectivity returns rather than on the next interval.
*/
val NETWORK: Constraints = Constraints.Builder()
.setRequiredNetworkType(NetworkType.CONNECTED)
.build()
fun periodicNameFor(accountName: String) = "caldav-sync-periodic:$accountName"
}
}
@@ -1,12 +1,18 @@
package de.jeanlucmakiola.agendula.data.sync
import android.app.Notification
import android.app.NotificationChannel
import android.app.NotificationManager
import android.content.Context
import android.util.Log
import androidx.core.app.NotificationCompat
import androidx.hilt.work.HiltWorker
import androidx.work.CoroutineWorker
import androidx.work.ForegroundInfo
import androidx.work.WorkerParameters
import dagger.assisted.Assisted
import dagger.assisted.AssistedInject
import de.jeanlucmakiola.agendula.R
/**
* Where sync actually happens.
@@ -52,6 +58,43 @@ class SyncWorker @AssistedInject constructor(
}
}
/**
* ⚠️ Implemented **unconditionally**, even though this worker never asks to
* run in the foreground.
*
* `setExpedited` falls back to a foreground service below API 31, and
* WorkManager calls this to build it. The default implementation throws
* `IllegalStateException`, so a worker that only ever runs expedited on
* modern devices crashes on every device running API 29 or 30 — which we
* support. It is never actually shown above API 30.
*/
override suspend fun getForegroundInfo(): ForegroundInfo {
val manager = applicationContext.getSystemService(NotificationManager::class.java)
manager?.createNotificationChannel(
NotificationChannel(
CHANNEL_ID,
applicationContext.getString(R.string.sync_notification_channel),
NotificationManager.IMPORTANCE_LOW,
),
)
val notification: Notification = NotificationCompat.Builder(applicationContext, CHANNEL_ID)
.setContentTitle(applicationContext.getString(R.string.sync_notification_title))
.setSmallIcon(R.drawable.ic_notification)
.setOngoing(true)
.setPriority(NotificationCompat.PRIORITY_LOW)
.build()
// ⚠️ **No `foregroundServiceType`.** Declaring `dataSync` is what drags in
// `FOREGROUND_SERVICE_DATA_SYNC`, the Android 15 six-hours-per-24 budget
// whose failure mode is a fatal `RemoteServiceException`, and a Play
// requirement for a video demo per declared type — the whole tail this
// worker exists to avoid. It is not needed either: above API 30
// `setExpedited` uses an expedited job and never calls this at all, and
// types only became mandatory at API 34.
return ForegroundInfo(NOTIFICATION_ID, notification)
}
companion object {
/** One in-flight sync per account, so a manual trigger cannot pile up. */
fun uniqueNameFor(accountName: String) = "caldav-sync:$accountName"
@@ -59,5 +102,7 @@ class SyncWorker @AssistedInject constructor(
const val KEY_ACCOUNT_NAME = "accountName"
private const val TAG = "SyncWorker"
private const val CHANNEL_ID = "sync"
private const val NOTIFICATION_ID = 4001
}
}
@@ -70,6 +70,19 @@ interface TaskListDao {
@Query("UPDATE task_lists SET is_read_only = :readOnly WHERE id = :listId")
fun setReadOnly(listId: Long, readOnly: Boolean)
/**
* Stores the RFC 6578 cursor.
*
* ⚠️ Written **after** a page's bodies are applied, never before — the RFC's
* own Appendix B has this backwards, and under WorkManager process death
* mid-sync is routine rather than exotic. A token stored ahead of its bodies
* is a permanent hole in the collection.
*
* One column and one statement, so there is nothing to be half-written.
*/
@Query("UPDATE task_lists SET sync_token = :token WHERE id = :listId")
fun setSyncToken(listId: Long, token: String?)
/** The synced collections of one account, in the order sync walks them. */
@Query("SELECT * FROM task_lists WHERE account_id = :accountId AND is_synced = 1 ORDER BY id")
fun syncedForAccount(accountId: Long): List<TaskListEntity>
@@ -39,7 +39,10 @@ class AccountsViewModel @Inject constructor(
* refuses, which would leave the user pressing a button that does nothing.
*/
fun syncNow(account: AccountEntity) {
syncTrigger.enqueue(account.displayName)
// Expedited: the user is looking at the button. Every other trigger is
// ordinary work, so the exhaustible per-app quota is spent here or not
// at all.
syncTrigger.enqueue(account.displayName, expedited = true)
}
fun remove(account: AccountEntity) {
+2
View File
@@ -308,6 +308,8 @@
<string name="accounts_never_synced">Never synced</string>
<string name="accounts_sync_failed">Last sync didn\u2019t finish</string>
<string name="accounts_sync_now">Sync now</string>
<string name="sync_notification_channel">Syncing</string>
<string name="sync_notification_title">Syncing tasks</string>
<string name="add_account_title">Add an account</string>
<string name="add_account_server_label">Email address or server address</string>
+12 -2
View File
@@ -24,10 +24,20 @@
a restored ciphertext can never be decrypted again — it would surface as
an account that silently stops syncing with no way to tell why. Excluding
it means the user signs in again on a new device, which is the honest
outcome. Note this is the *only* exclusion that is right here:
outcome. Note how narrow the exclusions here are, and must stay:
docs/SYNC.md is explicit that dropping the database or all of DataStore
from backup would trade a latent bug for a live one, since Auto Backup is
Local mode's only automatic safety net.
Local mode's only automatic safety net. Only files that are *meaningless
or harmful* on another device come out.
-->
<exclude domain="file" path="datastore/agendula_credentials.preferences_pb" />
<!--
Per-device sync bookkeeping: quarantine counters and the
full-reconciliation clock. Restored onto a new install every value
in it is a claim about a conversation this device never had — and a
restored "reconciled recently" makes the engine trust a stale sync
token for another day, which is the one RFC 6578 failure that has no
signal of its own.
-->
<exclude domain="file" path="datastore/agendula_sync_state.preferences_pb" />
</full-backup-content>
@@ -13,6 +13,15 @@
<include domain="file" path="datastore/" />
<!-- See backup_rules.xml: a restored ciphertext is undecryptable. -->
<exclude domain="file" path="datastore/agendula_credentials.preferences_pb" />
<!--
Per-device sync bookkeeping: quarantine counters and the
full-reconciliation clock. Restored onto a new install every value
in it is a claim about a conversation this device never had — and a
restored "reconciled recently" makes the engine trust a stale sync
token for another day, which is the one RFC 6578 failure that has no
signal of its own.
-->
<exclude domain="file" path="datastore/agendula_sync_state.preferences_pb" />
</cloud-backup>
<device-transfer>
<include domain="database" path="agendula-tasks.db" />
@@ -20,5 +29,14 @@
<include domain="database" path="agendula-tasks.db-shm" />
<include domain="file" path="datastore/" />
<exclude domain="file" path="datastore/agendula_credentials.preferences_pb" />
<!--
Per-device sync bookkeeping: quarantine counters and the
full-reconciliation clock. Restored onto a new install every value
in it is a claim about a conversation this device never had — and a
restored "reconciled recently" makes the engine trust a stale sync
token for another day, which is the one RFC 6578 failure that has no
signal of its own.
-->
<exclude domain="file" path="datastore/agendula_sync_state.preferences_pb" />
</device-transfer>
</data-extraction-rules>
@@ -0,0 +1,298 @@
package de.jeanlucmakiola.agendula.data.sync
import com.google.common.truth.Truth.assertThat
import de.jeanlucmakiola.agendula.data.tasks.room.TaskEntity
import de.jeanlucmakiola.agendula.data.tasks.room.TaskListEntity
import de.jeanlucmakiola.caldav.ChangeSet
import de.jeanlucmakiola.caldav.ETag
import de.jeanlucmakiola.caldav.RemoteRef
import okhttp3.HttpUrl.Companion.toHttpUrl
import org.junit.jupiter.api.Test
import kotlin.time.Instant
/** RFC 6578, from the engine's side. */
class IncrementalSyncTest {
private val store = FakeStore()
private val remote = FakeRemote()
private val quarantine = mutableMapOf<String, Int>()
private var list = TaskListEntity(
id = 1,
name = "Tasks",
color = 0,
accountId = 7,
href = "http://server/dav/tasks/",
syncToken = "urn:x:1",
)
init {
remote.supportsSyncCollection = true
remote.syncToken = "urn:x:from-propfind"
}
private fun sync(fullDue: Boolean = false) =
CollectionSyncer(store) { NOW }.sync(list, remote, quarantine, fullDue)
// ------------------------------------------------------------ the cursor
@Test fun `the token is stored only after its bodies are applied`() {
remote.put("one.ics", vtodo("a", "New"))
remote.changePages += page(changed = listOf("one.ics"), token = "urn:x:2")
sync()
assertThat(store.rows.single().title).isEqualTo("New")
assertThat(store.tokenWrites).containsExactly("urn:x:2")
// The full path was not needed, so no listing was taken.
assertThat(remote.log).doesNotContain("LIST")
}
@Test fun `a truncated page is resumed and each page commits its own token`() {
remote.put("one.ics", vtodo("a", "First"))
remote.put("two.ics", vtodo("b", "Second"))
remote.changePages += page(changed = listOf("one.ics"), token = "urn:x:2", truncated = true)
remote.changePages += page(changed = listOf("two.ics"), token = "urn:x:3")
sync()
assertThat(store.rows.map { it.uid }).containsExactly("a", "b")
// ⚠️ Per page, in order. A token written ahead of its bodies is a
// permanent hole in the collection the moment the process dies.
assertThat(store.tokenWrites).containsExactly("urn:x:2", "urn:x:3").inOrder()
}
@Test fun `a server that truncates without advancing its token is not looped on`() {
remote.changePages += page(changed = emptyList(), token = "urn:x:1", truncated = true)
val report = sync()
// The RFC never requires the token to advance. The full path then
// reconciled the collection, so this is a note rather than a failure.
assertThat(report.incrementalNote).contains("without advancing")
assertThat(report.failure).isNull()
assertThat(remote.log.count { it.startsWith("CHANGES") }).isEqualTo(1)
}
@Test fun `a page with no token at all falls back to a full reconciliation`() {
remote.put("one.ics", vtodo("a", "New"))
remote.changePages += page(changed = listOf("one.ics"), token = null)
val report = sync()
assertThat(store.rows.single().uid).isEqualTo("a")
assertThat(report.reconciledInFull).isTrue()
assertThat(remote.log).contains("LIST")
}
// ------------------------------------------------------- invalidation
@Test fun `an invalidated token is cleared and the run reconciles in full`() {
remote.put("one.ics", vtodo("a", "From the listing"))
remote.changePages += ChangeSet.TokenInvalid
val report = sync()
// Cleared first, so a crash before the full run cannot leave behind a
// token we already know the server rejects.
assertThat(store.tokenWrites.first()).isNull()
assertThat(report.reconciledInFull).isTrue()
assertThat(store.rows.single().uid).isEqualTo("a")
// Adopted from the PROPFIND that preceded the listing, never after it.
assertThat(store.tokenWrites.last()).isEqualTo("urn:x:from-propfind")
}
@Test fun `an unsupported report falls back without clearing the token`() {
remote.put("one.ics", vtodo("a", "From the listing"))
remote.changePages += ChangeSet.Unsupported
val report = sync()
assertThat(report.reconciledInFull).isTrue()
assertThat(store.tokenWrites).doesNotContain(null)
}
// ---------------------------------------------------------- membership
@Test fun `a removal for an href we never had is a no-op`() {
remote.changePages += page(removed = listOf("stranger.ics"), token = "urn:x:2")
val report = sync()
// Created and deleted between two syncs: reported as removed without ever
// having been reported as added.
assertThat(report.failure).isNull()
assertThat(report.deletedLocally).isEqualTo(0)
}
@Test fun `a removal we do know about deletes the rows`() {
store.rows += task(id = 1, uid = "a")
.copy(href = "http://server/dav/tasks/one.ics", etag = "e")
remote.changePages += page(removed = listOf("one.ics"), token = "urn:x:2")
val report = sync()
assertThat(store.rows).isEmpty()
assertThat(report.deletedLocally).isEqualTo(1)
}
@Test fun `delete-then-recreate at the same href replaces the old task`() {
store.rows += task(id = 1, uid = "old")
.copy(href = "http://server/dav/tasks/one.ics", etag = "e-old")
remote.put("one.ics", vtodo("new", "A different task"), eTag = "e-new")
remote.changePages += page(changed = listOf("one.ics"), token = "urn:x:2")
sync()
// ⚠️ Reported as a *change*, not as a removal plus an addition — so the
// UID in the body is the only thing that says the old task is gone.
assertThat(store.rows.map { it.uid }).containsExactly("new")
}
@Test fun `a mass removal is refused and reconciled against a listing instead`() {
(1..20).forEach { n ->
store.rows += task(id = n.toLong(), uid = "u$n")
.copy(href = "http://server/dav/tasks/$n.ics", etag = "e")
remote.put("$n.ics", vtodo("u$n", "Task $n"))
}
remote.changePages += page(removed = (1..20).map { "$it.ics" }, token = "urn:x:2")
val report = sync()
// ⚠️ ACL churn arrives looking exactly like this. Acting on it deletes
// the user's data on the strength of a change log.
assertThat(report.incrementalNote).contains("removed 20 of 20")
assertThat(report.reconciledInFull).isTrue()
assertThat(store.rows).hasSize(20)
}
@Test fun `an un-uploaded local edit is never downloaded over`() {
remote.readOnly = true
store.rows += task(id = 1, uid = "a", title = "Mine, unsent")
.copy(href = "http://server/dav/tasks/one.ics", etag = "e-old", isDirty = true)
remote.put("one.ics", vtodo("a", "Owner's version"), eTag = "e-new")
remote.changePages += page(changed = listOf("one.ics"), token = "urn:x:2")
sync()
// The upload was refused because the share was demoted, so the row is
// still dirty. Downloading over it destroys the edit with nothing in the
// report to say so.
assertThat(store.rows.single().title).isEqualTo("Mine, unsent")
assertThat(store.rows.single().isDirty).isTrue()
}
@Test fun `a stale removal never deletes what this run just wrote`() {
store.rows += task(id = 1, uid = "a", title = "Mine")
.copy(href = "http://server/dav/tasks/one.ics", etag = "e-one.ics", isDirty = true)
remote.put("one.ics", vtodo("a", "Server copy"))
// The change log predates our own upload in this same run.
remote.changePages += page(removed = listOf("one.ics"), token = "urn:x:2")
val report = sync()
assertThat(remote.log).contains("UPDATE one.ics if-match=e-one.ics")
assertThat(store.rows).hasSize(1)
assertThat(report.deletedLocally).isEqualTo(0)
}
@Test fun `a removal that loses a local edit says so`() {
store.rows += task(id = 1, uid = "a", title = "Mine")
.copy(href = "http://server/dav/tasks/gone.ics", etag = "e", isDirty = true)
remote.onPut = { de.jeanlucmakiola.caldav.PutOutcome.Rejected(400, "no") }
remote.changePages += page(removed = listOf("gone.ics"), token = "urn:x:2")
val report = sync()
assertThat(store.rows).isEmpty()
assertThat(report.discardedEdits.single().cause)
.isEqualTo(DiscardedEdit.Cause.DELETED_ON_SERVER)
}
@Test fun `a fallback that succeeds is not reported as a failure`() {
remote.put("one.ics", vtodo("a", "New"))
remote.changePages += ChangeSet.Failed("HTTP 400")
val report = sync()
// The collection is reconciled and the user has nothing to act on, so
// this must not surface as "last sync didn't finish".
assertThat(report.failure).isNull()
assertThat(report.incrementalNote).contains("400")
assertThat(report.reconciledInFull).isTrue()
assertThat(store.tokenWrites.last()).isEqualTo("urn:x:from-propfind")
}
@Test fun `a refetch the change log never mentions still happens`() {
store.rows += task(id = 1, uid = "a", title = "Mine").copy(isDirty = true)
remote.onPut = { href -> de.jeanlucmakiola.caldav.PutOutcome.StoredNeedsRefetch(href) }
remote.put("a.ics", vtodo("a", "Server's canonical form"))
// A server need not echo our own write back to us.
remote.changePages += page(token = "urn:x:2")
sync()
// Otherwise the row keeps a null ETag and never gets the canonical body.
assertThat(store.rows.single().title).isEqualTo("Server's canonical form")
}
// ------------------------------------------------------------- cadence
@Test fun `a due full reconciliation ignores the token entirely`() {
remote.put("one.ics", vtodo("a", "New"))
val report = sync(fullDue = true)
// The pruned-change-log failure has no signal at all, so the full path is
// a permanent safety net rather than a fallback.
assertThat(remote.log).doesNotContain("CHANGES urn:x:1")
assertThat(report.reconciledInFull).isTrue()
}
@Test fun `a collection that does not advertise the report is never asked`() {
remote.supportsSyncCollection = false
remote.put("one.ics", vtodo("a", "New"))
sync()
assertThat(remote.log.none { it.startsWith("CHANGES") }).isTrue()
}
@Test fun `an unchanged etag in the change log is not downloaded`() {
store.rows += task(id = 1, uid = "a", title = "Known")
.copy(href = "http://server/dav/tasks/one.ics", etag = "e-one.ics")
remote.put("one.ics", vtodo("a", "Known"))
remote.changePages += page(changed = listOf("one.ics"), token = "urn:x:2", eTag = "e-one.ics")
sync()
assertThat(remote.log.none { it.startsWith("FETCH") }).isTrue()
}
// --------------------------------------------------------------- setup
private fun page(
changed: List<String> = emptyList(),
removed: List<String> = emptyList(),
token: String?,
truncated: Boolean = false,
eTag: String = "e-changed",
) = ChangeSet.Page(
changed = changed.map {
RemoteRef("http://server/dav/tasks/$it".toHttpUrl(), ETag(eTag, weak = false))
},
removed = removed.map { "http://server/dav/tasks/$it".toHttpUrl() },
token = token,
truncated = truncated,
)
private fun task(id: Long, uid: String, title: String? = null) =
TaskEntity(id = id, listId = list.id, uid = uid, title = title)
private fun vtodo(uid: String, summary: String) =
"BEGIN:VCALENDAR\r\nVERSION:2.0\r\nBEGIN:VTODO\r\nUID:$uid\r\n" +
"SUMMARY:$summary\r\nEND:VTODO\r\nEND:VCALENDAR\r\n"
private companion object {
val NOW: Instant = Instant.parse("2026-03-01T12:00:00Z")
}
}
@@ -2,6 +2,7 @@ package de.jeanlucmakiola.agendula.data.sync
import de.jeanlucmakiola.agendula.data.tasks.room.TaskEntity
import de.jeanlucmakiola.caldav.CalendarCollection
import de.jeanlucmakiola.caldav.ChangeSet
import de.jeanlucmakiola.caldav.DeleteOutcome
import de.jeanlucmakiola.caldav.ETag
import de.jeanlucmakiola.caldav.FetchResult
@@ -18,11 +19,19 @@ class FakeStore(rows: List<TaskEntity> = emptyList()) : SyncStore {
val rows = rows.toMutableList()
val readOnlyWrites = mutableListOf<Pair<Long, Boolean>>()
private var nextId = (rows.maxOfOrNull { it.id } ?: 0L) + 1
/** In order, so a test can prove the token was stored *after* its bodies. */
val tokenWrites = mutableListOf<String?>()
private var nextId = 1L
override fun rowsIn(listId: Long) = this.rows.filter { it.listId == listId }
override fun insert(row: TaskEntity): Long {
// Computed against the rows as they stand, not at construction: tests
// seed `rows` directly after building this, and an id fixed up front
// collides with the first seeded row — which then makes a delete of one
// take the other with it.
nextId = maxOf(nextId, (rows.maxOfOrNull { it.id } ?: 0L) + 1)
val id = nextId++
rows += row.copy(id = id)
return id
@@ -59,6 +68,10 @@ class FakeStore(rows: List<TaskEntity> = emptyList()) : SyncStore {
override fun setListReadOnly(listId: Long, readOnly: Boolean) {
readOnlyWrites += listId to readOnly
}
override fun setSyncToken(listId: Long, token: String?) {
tokenWrites += token
}
}
/** An in-memory CalDAV collection that answers like a server rather than a mock. */
@@ -72,6 +85,11 @@ class FakeRemote(override val url: HttpUrl = "http://server/dav/tasks/".toHttpUr
var readOnly = false
var shared = false
var supportsSyncCollection = false
var syncToken: String? = null
/** Pages the next `changes()` calls will return, in order. */
val changePages = ArrayDeque<ChangeSet>()
var stateFailure: String? = null
var listFailure: String? = null
@@ -94,15 +112,21 @@ class FakeRemote(override val url: HttpUrl = "http://server/dav/tasks/".toHttpUr
color = null,
readOnly = readOnly,
isShared = shared,
supportsSyncCollection = false,
supportsSyncCollection = supportsSyncCollection,
maxResourceSize = null,
),
ctag = "ctag",
syncToken = null,
syncToken = syncToken,
),
)
}
override fun changes(token: String?): ChangeSet {
log += "CHANGES $token"
return changePages.removeFirstOrNull()
?: ChangeSet.Failed("the test queued no page for token $token")
}
override fun list(): Result<List<RemoteRef>> {
listFailure?.let { return Result.failure(IllegalStateException(it)) }
log += "LIST"
@@ -0,0 +1,52 @@
package de.jeanlucmakiola.caldav
import okhttp3.HttpUrl
/** What one `sync-collection` REPORT came back with. */
sealed interface ChangeSet {
/**
* One page of changes.
*
* A page, not a result: RFC 6578 lets a server truncate and hand back a token
* that resumes where it stopped, so the caller applies this page's bodies,
* stores [token], and only then asks for the next.
*/
data class Page(
val changed: List<RemoteRef>,
val removed: List<HttpUrl>,
/**
* The token to store **after** this page is applied.
*
* ⚠️ Null when the server sent a 207 with no `DAV:sync-token` at all,
* which happens. That is not an error and not a reason to throw — the
* changes in this page are still real. It means the next run has no
* cursor and must reconcile in full.
*/
val token: String?,
/**
* The server stopped early and there is more behind [token].
*
* Signalled by a `507` on the collection's **own** href, which is
* truncation rather than failure — a 507 as the *outer* response status
* means something else entirely.
*/
val truncated: Boolean = false,
) : ChangeSet
/**
* The token is no longer usable and the collection must be reconciled in full.
*
* ⚠️ **Invalidation has no status code.** `403` (sabre), `400` (Google),
* `409` (Radicale) and `412` (Evolution) are all observed in the wild, so the
* only reliable signal is `DAV:valid-sync-token` in the body on **any** 4xx.
* Thunderbird matches on 400 alone and therefore never recovers from the 403
* that most of the self-hosted world emits.
*/
data object TokenInvalid : ChangeSet
/** The server will not do `sync-collection` here. Fall back, do not retry. */
data object Unsupported : ChangeSet
data class Failed(val reason: String) : ChangeSet
}
@@ -1,6 +1,8 @@
package de.jeanlucmakiola.caldav
import at.bitfire.dav4jvm.DavCalendar
import at.bitfire.dav4jvm.DavCollection
import at.bitfire.dav4jvm.Error
import at.bitfire.dav4jvm.DavResource
import at.bitfire.dav4jvm.Response
import at.bitfire.dav4jvm.exception.DavException
@@ -99,6 +101,72 @@ class CalendarCollection(
refs
}
/**
* One page of changes since [token] (RFC 6578).
*
* ⚠️ **No `DAV:limit`.** Nextcloud regressed it to a localised HTML error
* page, so asking for a bounded page is how you get an unparseable response
* instead of a bounded one. Truncation is the server's decision, signalled
* back through [ChangeSet.Page.truncated].
*
* ⚠️ **A removal is a `404` at `DAV:response` level**, not inside a
* `propstat`. A client that only reads statuses out of propstats sees every
* deletion as an unremarkable response with no properties and silently keeps
* the row.
*/
override fun changes(token: String?): ChangeSet = try {
val changed = mutableListOf<RemoteRef>()
val removed = mutableListOf<HttpUrl>()
var truncated = false
val properties = DavCollection(httpClient, url).reportChanges(
syncToken = token,
infiniteDepth = false,
limit = null,
GetETag.NAME,
) { response, relation ->
val code = response.status?.code
when {
relation == Response.HrefRelation.SELF ->
// On our own href this is truncation, not a failure.
if (code == INSUFFICIENT_STORAGE) truncated = true
code == NOT_FOUND || code == GONE -> removed += response.href
response.isSuccess() ->
changed += RemoteRef(response.href, ETag.from(response[GetETag::class.java]))
else -> Unit
}
}
ChangeSet.Page(
changed = changed,
removed = removed,
token = properties.filterIsInstance<SyncToken>().firstOrNull()?.token,
truncated = truncated,
)
} catch (e: HttpException) {
when {
e.code / 100 == 4 && mentionsInvalidToken(e) -> ChangeSet.TokenInvalid
e.code in UNSUPPORTED -> ChangeSet.Unsupported
else -> ChangeSet.Failed("HTTP ${e.code}")
}
} catch (e: DavException) {
ChangeSet.Failed(e.toString())
} catch (e: IOException) {
ChangeSet.Failed(e.toString())
}
/**
* ⚠️ The body, not the status. Every observed invalidation status is also a
* legitimate answer to something else, and every server picks a different
* one — so the element is the signal and the code is noise.
*/
private fun mentionsInvalidToken(e: HttpException): Boolean =
Error.VALID_SYNC_TOKEN in e.errors ||
e.responseBody?.contains(VALID_SYNC_TOKEN_ELEMENT) == true
/**
* Downloads [hrefs] with `calendar-multiget`, in batches.
*
@@ -284,5 +352,17 @@ class CalendarCollection(
private const val NOT_FOUND = 404
private const val GONE = 410
private const val PRECONDITION_FAILED = 412
private const val INSUFFICIENT_STORAGE = 507
/** Matched as a local name, so a server's own namespace prefix is irrelevant. */
private const val VALID_SYNC_TOKEN_ELEMENT = "valid-sync-token"
/**
* The server does not implement the report. Deliberately narrow: 400,
* 403, 409 and 412 are all *invalidation* codes on some server, so
* treating them as "unsupported" would abandon `sync-collection` on a
* collection that merely needs a fresh token.
*/
private val UNSUPPORTED = setOf(405, 415, 501)
}
}
@@ -18,6 +18,14 @@ interface RemoteCalendar {
fun list(): Result<List<RemoteRef>>
/**
* One page of changes since [token], or null for a first run.
*
* Never throws for anything a server can legitimately answer — a rejected
* token is [ChangeSet.TokenInvalid], not an exception.
*/
fun changes(token: String?): ChangeSet
fun fetch(hrefs: List<HttpUrl>): Result<FetchResult>
fun create(name: String, iCalendar: String): PutOutcome
@@ -0,0 +1,160 @@
package de.jeanlucmakiola.caldav
import com.google.common.truth.Truth.assertThat
import okhttp3.OkHttpClient
import okhttp3.mockwebserver.MockResponse
import okhttp3.mockwebserver.MockWebServer
import org.junit.After
import org.junit.Before
import org.junit.Test
/** RFC 6578. Every case here is a documented failure of a real client. */
class CalendarChangesTest {
private val server = MockWebServer()
private val httpClient = OkHttpClient.Builder().followRedirects(false).build()
private lateinit var collection: CalendarCollection
@Before fun start() {
server.start()
collection = CalendarCollection(httpClient, server.url("/dav/tasks/"))
}
@After fun stop() = server.shutdown()
@Test fun `changed and removed are told apart`() {
server.enqueue(
multistatus(
"""
<response>
<href>/dav/tasks/kept.ics</href>
<propstat><prop><getetag>"e1"</getetag></prop>
<status>HTTP/1.1 200 OK</status></propstat>
</response>
<response>
<href>/dav/tasks/gone.ics</href>
<status>HTTP/1.1 404 Not Found</status>
</response>
<sync-token>urn:x:2</sync-token>
""",
),
)
val page = collection.changes("urn:x:1") as ChangeSet.Page
// ⚠️ The deletion status is at DAV:response level, not inside a propstat.
assertThat(page.changed.map { it.href.encodedPath })
.containsExactly("/dav/tasks/kept.ics")
assertThat(page.removed.map { it.encodedPath })
.containsExactly("/dav/tasks/gone.ics")
assertThat(page.token).isEqualTo("urn:x:2")
assertThat(page.truncated).isFalse()
}
@Test fun `no DAV limit is ever sent`() {
server.enqueue(multistatus("<sync-token>urn:x:2</sync-token>"))
collection.changes(null)
val body = server.takeRequest().body.readUtf8()
// Nextcloud regressed DAV:limit to a localised HTML error page, so asking
// for a bounded page is how you get an unparseable response.
assertThat(body).doesNotContain("limit")
assertThat(body).contains("sync-level")
}
@Test fun `a 507 on our own href is truncation, not failure`() {
server.enqueue(
multistatus(
"""
<response>
<href>/dav/tasks/one.ics</href>
<propstat><prop><getetag>"e1"</getetag></prop>
<status>HTTP/1.1 200 OK</status></propstat>
</response>
<response>
<href>/dav/tasks/</href>
<status>HTTP/1.1 507 Insufficient Storage</status>
</response>
<sync-token>urn:x:2</sync-token>
""",
),
)
val page = collection.changes("urn:x:1") as ChangeSet.Page
assertThat(page.truncated).isTrue()
assertThat(page.changed).hasSize(1)
assertThat(page.token).isEqualTo("urn:x:2")
}
@Test fun `a 207 with no sync-token is not an error`() {
server.enqueue(
multistatus(
"""
<response>
<href>/dav/tasks/one.ics</href>
<propstat><prop><getetag>"e1"</getetag></prop>
<status>HTTP/1.1 200 OK</status></propstat>
</response>
""",
),
)
val page = collection.changes("urn:x:1") as ChangeSet.Page
// The changes are real; only the cursor is missing.
assertThat(page.changed).hasSize(1)
assertThat(page.token).isNull()
}
// --------------------------------------------------- token invalidation
@Test fun `invalidation is recognised on every status servers use for it`() {
// ⚠️ 403 sabre, 400 Google, 409 Radicale, 412 Evolution. Thunderbird
// matches 400 alone and never recovers from the 403 most self-hosted
// servers emit.
listOf(400, 403, 409, 412).forEach { code ->
server.enqueue(
MockResponse()
.setResponseCode(code)
.setHeader("Content-Type", "application/xml; charset=utf-8")
.setBody(
"<D:error xmlns:D=\"DAV:\"><D:valid-sync-token/></D:error>",
),
)
assertThat(collection.changes("stale")).isEqualTo(ChangeSet.TokenInvalid)
}
}
@Test fun `invalidation is recognised whatever prefix the server uses`() {
server.enqueue(
MockResponse()
.setResponseCode(403)
.setHeader("Content-Type", "application/xml; charset=utf-8")
.setBody("<error xmlns=\"DAV:\"><valid-sync-token/></error>"),
)
assertThat(collection.changes("stale")).isEqualTo(ChangeSet.TokenInvalid)
}
@Test fun `a 4xx without the marker is not an invalidation`() {
server.enqueue(MockResponse().setResponseCode(403).setBody("nope"))
// Treating every 403 as invalidation would discard a working token on a
// permissions error and re-download the whole collection.
assertThat(collection.changes("good")).isInstanceOf(ChangeSet.Failed::class.java)
}
@Test fun `an unimplemented report is reported as unsupported`() {
server.enqueue(MockResponse().setResponseCode(501))
assertThat(collection.changes(null)).isEqualTo(ChangeSet.Unsupported)
}
@Test fun `a 507 outer status is a failure rather than truncation`() {
server.enqueue(MockResponse().setResponseCode(507))
assertThat(collection.changes("t")).isInstanceOf(ChangeSet.Failed::class.java)
}
private fun multistatus(body: String) = MockResponse()
.setResponseCode(207)
.setHeader("Content-Type", "application/xml; charset=utf-8")
.setBody("<multistatus xmlns=\"DAV:\">$body</multistatus>")
}
+90
View File
@@ -600,6 +600,96 @@ genuinely "once overnight".
process kill mid-page loses nothing and duplicates nothing; the periodic worker
runs and the manual trigger is instant.
### Settled while building it
**`initialIncomplete` turned out not to need a column — or to exist.** The plan
called for it to be persisted atomically with the token, to stop a resumed
partial sync sweeping against an incomplete set. But the mark-and-sweep it
protects *is* chunk 3's full path, so the rule that removes the whole hazard is
simpler: **a token is adopted only when a full reconciliation completes.** A run
killed part-way stores no token and the next one starts over. No second column,
nothing to write atomically with anything.
**The token is taken from the PROPFIND that *precedes* the read, not after it.**
Both are imprecise; only one is imprecise safely. A token minted after the
listing silently swallows everything that changed while we were reading it; one
minted before merely re-reports those changes next time, and a re-report costs an
ETag comparison.
**Where the cadence clock lives.** `task_lists.sync_token` is domain state and
stays in Room. "When did this collection last reconcile in full" is scheduling
bookkeeping, so it is in DataStore (`SyncCadenceStore`) — and deliberately
outside any backup, because a restored "checked recently" on a fresh install is
exactly the stale trust the full path exists to break.
**⚠️ `ForegroundInfo` must carry no `foregroundServiceType`.** Implementing
`getForegroundInfo` is mandatory — `setExpedited` falls back to a foreground
service below API 31 and the default implementation throws, so a worker without
it crashes on every API 29/30 device. But passing
`FOREGROUND_SERVICE_TYPE_DATA_SYNC` re-introduces the entire tail the design
rejected: the `FOREGROUND_SERVICE_DATA_SYNC` permission, Android 15's
six-hours-per-24 budget with its fatal `RemoteServiceException`, and a Play video
demo per declared type. Lint caught it, correctly. The type is unnecessary
anyway: above API 30 `setExpedited` uses an expedited job and never calls
`getForegroundInfo`, and types only became mandatory at API 34.
**No boot receiver for sync.** WorkManager restores its own schedule across a
reboot, and Android 15 forbids starting a `dataSync` FGS from `BOOT_COMPLETED` —
so a design that needed one would have had no way to run. App open calls
`rescheduleAll()` instead, which is idempotent under `KEEP` and repairs the one
case WorkManager cannot: a schedule lost to "clear app data" or to a restore.
**Expedited is spent on the button and nowhere else.** The quota is per-app and
exhaustible; a background trigger that spends it leaves none for the trigger the
user is actually watching.
### Found by the review, and worth naming
**⚠️ The cadence clock was in the backed-up DataStore — while its own KDoc, and
this document, claimed it was not.** Auto Backup includes `datastore/`, and only
`agendula_credentials` was excluded. Restored onto a new device it would say
"reconciled minutes ago", so the restored sync token would be trusted for another
day: exactly the silently pruned change log the full path exists to catch. Both
sync-state stores now live in their own `agendula_sync_state` file, excluded
alongside the credentials — a restored quarantine count is the same class of lie,
silently skipping resources that were never tried on this device.
**The incremental download path had no write-phase guard.** `downloadPhase` skips
a row that is still dirty; `downloadChanged` did not. A share demoted to
read-only leaves the upload refused and the row dirty, and the next change-log
entry for it overwrote the user's unsent edit with no `DiscardedEdit` at all. The
guard cannot simply be copied, either — the 412 path *deliberately* leaves the
row dirty and waits for the download — so it is "dirty **and not awaiting a
refetch**".
**A change log describes the past, and `applyRemovals` treated it as the
present.** An entry predating this run's own upload deleted the row we had just
written, leaving an orphan on the server and nothing on the device. It now
respects `touched`, and reports a lost edit the way the sweep does.
**An abandoned incremental attempt poisoned a successful run.** `ChangeSet.Failed`
and the mass-removal refusal both set `report.failure`, which the full fallback
never cleared — so a server answering `sync-collection` with a bare 400 refused
the new token, skipped the cadence record, and showed "last sync didn't finish"
on every other sync while nothing was wrong with the data. A fallback that
succeeds now demotes the attempt to `incrementalNote`.
**`refetch` was silently dropped on a successful incremental run**, so a PUT
accepted without a strong ETag never got its canonical body back on a server that
does not echo our own writes into its change log.
**Two of my own patches were wrong in ways the compiler could not see.** `fatal()`
was written with its KDoc and never called — a `python` patch aborted on its
second assertion after the first had already matched, so nothing was written, and
the function sat there describing a bug it was not preventing. And the
`delete-then-recreate` lookup read the href index *after* the same index had been
re-pointed at the rows just written, so it never found the displaced row. The
first was caught by the review, the second by a test.
**Two smaller ones:** the displaced lookup reintroduced the per-resource
full-table scan the comment twelve lines above it forbids, and `MainActivity`
re-ran the app-open sync on every rotation, theme switch and locale change.
---
## Chunk 5 — hardening, compliance, real servers