package com.taaazzz.photosync.bg import android.content.Context import android.net.ConnectivityManager import android.net.NetworkCapabilities import android.provider.MediaStore import android.util.Base64 import org.json.JSONObject import java.net.HttpURLConnection import java.net.URL import java.net.URLEncoder import java.util.concurrent.atomic.AtomicBoolean // Upload NATIF des nouveaux medias (le JS en arriere-plan est inutilisable : // moteur en pause -> ne traite pas les reponses reseau). Le natif, lui, marche. // Lit la config ecrite par le JS (SharedPreferences), interroge MediaStore pour // les fichiers ajoutes depuis le dernier passage, calcule le dossier NAS (meme // logique que la sync JS : sous-chemin sous la racine + option Annee/Mois/Jour) // et POST chaque fichier a l'agent. Dedup gere par l'agent (nom+contenu). object Uploader { private val running = AtomicBoolean(false) // Au-dela de cette taille on passe en upload REPRENABLE par morceaux (meme // logique que la sync JS). Une coupure ne reperd qu'une tranche, pas tout. private const val RESUMABLE_MIN = 24L * 1024 * 1024 // 24 Mo private const val CHUNK = 6 * 1024 * 1024 // 6 Mo (binaire ; base64 ~ +33%) // File de reessai des echecs (persistee). Bornee : au-dela on lache les plus // anciens -- le 1er plan JS (verite NAS) reste le filet de securite ultime. private const val RETRY_MAX = 500 private val MOIS = arrayOf( "01-Janvier", "02-Fevrier", "03-Mars", "04-Avril", "05-Mai", "06-Juin", "07-Juillet", "08-Aout", "09-Septembre", "10-Octobre", "11-Novembre", "12-Decembre" ) // Pose la ligne de base "maintenant" si jamais definie (au demarrage du service) // pour que la 1re photo prise apres activation soit bien prise en compte. fun initBaseline(context: Context) { val prefs = context.applicationContext.getSharedPreferences("photosync_bg", Context.MODE_PRIVATE) if (prefs.getLong("lastAdded", 0L) == 0L) { prefs.edit().putLong("lastAdded", System.currentTimeMillis() / 1000 - 5).apply() } } fun run(context: Context) { if (!running.compareAndSet(false, true)) { Diag.ping("up-busy"); return } try { doRun(context.applicationContext) } catch (e: Exception) { Diag.ping("up-EXC-" + (e.message ?: "?")) } finally { running.set(false) } } private fun doRun(ctx: Context) { Diag.ping("up-start") val prefs = ctx.getSharedPreferences("photosync_bg", Context.MODE_PRIVATE) val raw = prefs.getString("config", null) if (raw == null) { Diag.ping("up-noconfig"); return } val cfg = JSONObject(raw) if (cfg.optBoolean("wifiOnly", false) && !isWifi(ctx)) { Diag.ping("up-notwifi"); return } val token = cfg.optString("token", "") val bases = cfg.optJSONArray("bases") ?: return val rules = cfg.optJSONArray("rules") ?: return if (rules.length() == 0) { Diag.ping("up-norules"); return } val base = pickBase(bases, token) ?: run { Diag.ping("up-unreachable"); return } Diag.ping("up-base-" + base.replace(Regex("^https?://"), "")) // Ne traite que les fichiers ajoutes APRES le dernier passage. Au 1er lancement, // on part de "maintenant" (l'historique est gere par la sync JS 15 min / 1er plan). var lastAdded = prefs.getLong("lastAdded", 0L) if (lastAdded == 0L) lastAdded = System.currentTimeMillis() / 1000 - 10 // Nouveaux medias + REESSAIS des echecs precedents. La file de reessai survit // aux passages : un echec (wifi coupe, agent absent, 5xx) est retente jusqu'au // succes -- avant, lastAdded avancait quoi qu'il arrive et l'echec etait perdu // cote natif. On lache un reessai seulement si la source a disparu du telephone // ou ne correspond plus a aucune regle. Le 1er plan JS reste le filet ultime. val files = queryNew(ctx, lastAdded) val retry = loadRetry(prefs) Diag.ping("up-found-" + files.size + "-retry-" + retry.size) // Fusion (dedup par id MediaStore) : nouveaux d'abord, puis reessais restants. val newIds = HashSet() val byId = LinkedHashMap() var maxAdded = lastAdded for (f in files) { val id = android.content.ContentUris.parseId(f.uri) newIds.add(id); byId[id] = f if (f.dateAdded > maxAdded) maxAdded = f.dateAdded } for ((id, f) in retry) if (!byId.containsKey(id)) byId[id] = f var sent = 0 for ((id, f) in byId) { // Un reessai dont la source a disparu du telephone : inutile de le garder. if (!newIds.contains(id) && !stillExists(ctx, f.uri)) { retry.remove(id); continue } val dest = destFor(rules, f) if (dest == null) { retry.remove(id); continue } // hors des regles actuelles -> on lache val ok = try { uploadOne(ctx, base, token, f, dest) } catch (e: Exception) { Diag.ping("up-err-" + f.name + "-" + (e.message ?: "?")); false } if (ok) { sent++; retry.remove(id) } else { retry[id] = f } // garde/repose pour le prochain passage } saveRetry(prefs, retry) prefs.edit().putLong("lastAdded", maxAdded).apply() Diag.ping("up-done-" + sent + "/" + byId.size + "-retryleft-" + retry.size) } // --- File de reessai persistante (SharedPreferences "retryQueue") ------------ // Stocke juste de quoi reconstruire l'Item (l'URI se rebatit a partir de l'id). private fun loadRetry(prefs: android.content.SharedPreferences): LinkedHashMap { val out = LinkedHashMap() val raw = prefs.getString("retryQueue", null) ?: return out try { val arr = org.json.JSONArray(raw) val baseUri = MediaStore.Files.getContentUri("external") for (i in 0 until arr.length()) { val o = arr.optJSONObject(i) ?: continue val id = o.optLong("id", -1L) if (id < 0) continue out[id] = Item( android.content.ContentUris.withAppendedId(baseUri, id), o.optString("name"), o.optLong("size"), o.optLong("dateAdded"), o.optString("folder"), o.optLong("takenMs") ) } } catch (e: Exception) { /* file corrompue -> on repart vide */ } return out } private fun saveRetry(prefs: android.content.SharedPreferences, q: LinkedHashMap) { val items = q.values.toList() val start = if (items.size > RETRY_MAX) items.size - RETRY_MAX else 0 // borne : garde les plus recents val arr = org.json.JSONArray() for (idx in start until items.size) { val f = items[idx] arr.put(JSONObject() .put("id", android.content.ContentUris.parseId(f.uri)) .put("name", f.name).put("size", f.size).put("dateAdded", f.dateAdded) .put("folder", f.folder).put("takenMs", f.takenMs)) } prefs.edit().putString("retryQueue", arr.toString()).apply() } // La source (id MediaStore) existe-t-elle encore sur le telephone ? private fun stillExists(ctx: Context, uri: android.net.Uri): Boolean { return try { ctx.contentResolver.query(uri, arrayOf(MediaStore.Files.FileColumns._ID), null, null, null) ?.use { it.moveToFirst() } ?: false } catch (e: Exception) { false } } // --- Serveur joignable (essaie chaque base, garde la 1re qui repond) --------- private fun pickBase(bases: org.json.JSONArray, token: String): String? { for (i in 0 until bases.length()) { val b = bases.optString(i).trimEnd('/') if (b.isEmpty()) continue try { val c = URL("$b/").openConnection() as HttpURLConnection c.connectTimeout = 4000; c.readTimeout = 4000 if (token.isNotEmpty()) c.setRequestProperty("x-photosync-token", token) val code = c.responseCode c.disconnect() if (code in 200..299) return b } catch (e: Exception) { /* suivant */ } } return null } // --- MediaStore : fichiers (image+video) ajoutes apres `sinceSec` ------------ data class Item(val uri: android.net.Uri, val name: String, val size: Long, val dateAdded: Long, val folder: String, val takenMs: Long) private fun queryNew(ctx: Context, sinceSec: Long): List { val out = ArrayList() val uri = MediaStore.Files.getContentUri("external") val proj = arrayOf( MediaStore.Files.FileColumns._ID, MediaStore.Files.FileColumns.DISPLAY_NAME, MediaStore.Files.FileColumns.SIZE, MediaStore.Files.FileColumns.DATE_ADDED, MediaStore.Files.FileColumns.DATE_MODIFIED, MediaStore.Files.FileColumns.MEDIA_TYPE, MediaStore.Files.FileColumns.RELATIVE_PATH ) val sel = "${MediaStore.Files.FileColumns.DATE_ADDED} > ? AND " + "${MediaStore.Files.FileColumns.MEDIA_TYPE} IN (?, ?)" val args = arrayOf( sinceSec.toString(), MediaStore.Files.FileColumns.MEDIA_TYPE_IMAGE.toString(), MediaStore.Files.FileColumns.MEDIA_TYPE_VIDEO.toString() ) val sort = "${MediaStore.Files.FileColumns.DATE_ADDED} ASC" ctx.contentResolver.query(uri, proj, sel, args, sort)?.use { cur -> val iId = cur.getColumnIndexOrThrow(MediaStore.Files.FileColumns._ID) val iName = cur.getColumnIndexOrThrow(MediaStore.Files.FileColumns.DISPLAY_NAME) val iSize = cur.getColumnIndexOrThrow(MediaStore.Files.FileColumns.SIZE) val iAdded = cur.getColumnIndexOrThrow(MediaStore.Files.FileColumns.DATE_ADDED) val iMod = cur.getColumnIndexOrThrow(MediaStore.Files.FileColumns.DATE_MODIFIED) val iRel = cur.getColumnIndex(MediaStore.Files.FileColumns.RELATIVE_PATH) while (cur.moveToNext()) { val id = cur.getLong(iId) val name = cur.getString(iName) ?: continue val size = cur.getLong(iSize) if (size <= 0) continue val added = cur.getLong(iAdded) val modMs = cur.getLong(iMod) * 1000 val rel = if (iRel >= 0) (cur.getString(iRel) ?: "") else "" // dossier ABSOLU du fichier sur le telephone (sans le nom). val folder = ("/storage/emulated/0/" + rel).trimEnd('/') out.add(Item(android.content.ContentUris.withAppendedId(uri, id), name, size, added, folder, modMs)) } } return out } // --- Dossier NAS pour un fichier (meme logique que la sync JS) --------------- private fun destFor(rules: org.json.JSONArray, f: Item): String? { for (i in 0 until rules.length()) { val r = rules.optJSONObject(i) ?: continue val root = r.optString("rootPath", "").trimEnd('/') val dest = r.optString("dest", "").trimEnd('/') if (root.isEmpty() || dest.isEmpty()) continue val rel: String = when { f.folder == root -> "" f.folder.startsWith("$root/") -> f.folder.substring(root.length + 1) else -> continue } var out = if (rel.isEmpty()) dest else "$dest/$rel" if (r.optBoolean("organizeByDate", false)) out += "/" + datePath(f.name, f.takenMs) return out } return null } private fun datePath(name: String, takenMs: Long): String { val m = Regex("(20\\d{2})[-_]?(\\d{2})[-_]?(\\d{2})").find(name) var y: Int; var mo: Int; var d: Int if (m != null && m.groupValues[2].toInt() in 1..12 && m.groupValues[3].toInt() in 1..31) { y = m.groupValues[1].toInt(); mo = m.groupValues[2].toInt(); d = m.groupValues[3].toInt() } else { val cal = java.util.Calendar.getInstance() cal.timeInMillis = if (takenMs > 0) takenMs else System.currentTimeMillis() y = cal.get(java.util.Calendar.YEAR); mo = cal.get(java.util.Calendar.MONTH) + 1; d = cal.get(java.util.Calendar.DAY_OF_MONTH) } return "$y/${MOIS[mo - 1]}/${d.toString().padStart(2, '0')}" } // --- Upload d'UN fichier : reprenable par morceaux si gros, direct sinon ----- private fun uploadOne(ctx: Context, base: String, token: String, f: Item, dest: String): Boolean { return if (f.size >= RESUMABLE_MIN) uploadChunkedNative(ctx, base, token, f, dest) else uploadDirect(ctx, base, token, f, dest) } // --- Petits fichiers : POST streaming en un coup (l'agent verifie la taille) - private fun uploadDirect(ctx: Context, base: String, token: String, f: Item, dest: String): Boolean { val url = base + "/upload?name=" + enc(f.name) + "&dest=" + enc(dest) + "&total=" + f.size + "&src=" + enc("auto-natif") val conn = URL(url).openConnection() as HttpURLConnection conn.requestMethod = "POST" conn.doOutput = true conn.connectTimeout = 10000 // Reponse attendue APRES l'envoi du corps : l'agent peut hasher un gros // fichier (dedup) avant de repondre -> delai genereux pour les videos. conn.readTimeout = 600000 conn.setFixedLengthStreamingMode(f.size) // streame sans tout charger en RAM if (token.isNotEmpty()) conn.setRequestProperty("x-photosync-token", token) conn.setRequestProperty("Content-Type", "application/octet-stream") ctx.contentResolver.openInputStream(f.uri).use { input -> if (input == null) return false conn.outputStream.use { out -> input.copyTo(out, 64 * 1024) } } val code = conn.responseCode conn.disconnect() return code in 200..299 } // --- Gros fichiers : upload REPRENABLE par morceaux -------------------------- // Demande l'offset deja recu par l'agent (.part), puis envoie le reste tranche // par tranche en base64 vers /upload-chunk. Une coupure ne repart PAS de zero : // on reprend a l'offset reel. Sur 409 (desync) l'agent renvoie son offret reel // et on se repositionne. Meme protocole que la sync JS (lib/sync.js). private fun uploadChunkedNative(ctx: Context, base: String, token: String, f: Item, dest: String): Boolean { val src = "auto-natif" var offset = getOffset(base, token, f.name, dest, src) if (offset < 0 || offset > f.size) offset = 0 Diag.ping("up-chunk-" + f.name + "@" + offset) var input = openAndSkip(ctx, f.uri, offset) ?: return false var fullRestarts = 0 try { val buf = ByteArray(CHUNK) while (offset < f.size) { val want = minOf(CHUNK.toLong(), f.size - offset).toInt() val n = readFully(input, buf, want) if (n <= 0) break // fin inattendue (fichier modifie/tronque) val last = offset + n >= f.size val b64 = Base64.encodeToString(buf, 0, n, Base64.NO_WRAP) val (code, srvOffset) = postChunk(base, token, f.name, dest, src, offset, f.size, last, b64) if (code in 200..299) { offset += n if (last) return true continue } // 409 = desync : l'agent indique ou en est REELLEMENT le .part. if (code == 409 && srvOffset != null && srvOffset >= 0 && srvOffset <= f.size) { if (srvOffset == 0L && offset != 0L && ++fullRestarts > 2) { Diag.ping("up-chunk-give-" + f.name); return false // taille incoherente -> on abandonne } input.close() offset = srvOffset input = openAndSkip(ctx, f.uri, offset) ?: return false continue } return false // autre erreur (auth, reseau, 5xx) -> retentee au prochain passage } return offset >= f.size } finally { try { input.close() } catch (e: Exception) {} } } // Offset deja recu par l'agent pour ce fichier (taille du .part). 0 si rien/erreur. private fun getOffset(base: String, token: String, name: String, dest: String, src: String): Long { return try { val url = base + "/upload-offset?name=" + enc(name) + "&dest=" + enc(dest) + "&src=" + enc(src) val c = URL(url).openConnection() as HttpURLConnection c.connectTimeout = 8000; c.readTimeout = 8000 if (token.isNotEmpty()) c.setRequestProperty("x-photosync-token", token) val j = readJson(c) c.disconnect() j?.optLong("offset", 0L) ?: 0L } catch (e: Exception) { 0L } } // POST d'une tranche base64. Renvoie (code HTTP, offset reel renvoye sur 409). private fun postChunk(base: String, token: String, name: String, dest: String, src: String, offset: Long, total: Long, last: Boolean, b64: String): Pair { val url = base + "/upload-chunk?name=" + enc(name) + "&dest=" + enc(dest) + "&src=" + enc(src) + "&offset=" + offset + "&total=" + total + "&last=" + (if (last) 1 else 0) val conn = URL(url).openConnection() as HttpURLConnection conn.requestMethod = "POST" conn.doOutput = true conn.connectTimeout = 10000 conn.readTimeout = 600000 // l'agent peut hasher (dedup) avant de repondre sur la derniere tranche if (token.isNotEmpty()) conn.setRequestProperty("x-photosync-token", token) conn.setRequestProperty("Content-Type", "text/plain") val bytes = b64.toByteArray(Charsets.US_ASCII) conn.setFixedLengthStreamingMode(bytes.size) conn.outputStream.use { it.write(bytes) } val code = conn.responseCode var srvOffset: Long? = null if (code !in 200..299) { val j = readJson(conn) if (j != null && j.has("offset")) srvOffset = j.optLong("offset", -1L) } conn.disconnect() return Pair(code, srvOffset) } // Lit le corps de la reponse (succes ou erreur) en JSON. null si illisible. private fun readJson(conn: HttpURLConnection): JSONObject? { return try { val stream = if (conn.responseCode in 200..299) conn.inputStream else conn.errorStream val text = stream?.bufferedReader()?.use { it.readText() } ?: return null JSONObject(text) } catch (e: Exception) { null } } // Ouvre le flux du media et avance EXACTEMENT de `offset` octets (reprise). private fun openAndSkip(ctx: Context, uri: android.net.Uri, offset: Long): java.io.InputStream? { val input = ctx.contentResolver.openInputStream(uri) ?: return null return try { var remaining = offset while (remaining > 0) { val s = input.skip(remaining) if (s > 0) { remaining -= s; continue } if (input.read() < 0) break // EOF avant l'offset -> fichier raccourci remaining -= 1 } input } catch (e: Exception) { try { input.close() } catch (e2: Exception) {}; null } } // Remplit `buf` avec jusqu'a `want` octets (read() peut en rendre moins). Renvoie le total lu. private fun readFully(input: java.io.InputStream, buf: ByteArray, want: Int): Int { var off = 0 while (off < want) { val n = input.read(buf, off, want - off) if (n < 0) break off += n } return off } private fun enc(s: String) = URLEncoder.encode(s, "UTF-8") private fun isWifi(ctx: Context): Boolean { return try { val cm = ctx.getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager val n = cm.activeNetwork ?: return false val caps = cm.getNetworkCapabilities(n) ?: return false caps.hasTransport(NetworkCapabilities.TRANSPORT_WIFI) } catch (e: Exception) { true } } }