diff --git a/README.md b/README.md index c0d8032..3b5f171 100644 --- a/README.md +++ b/README.md @@ -13,8 +13,8 @@ and `@zennotes/shared-domain` packages. The exact archives are vendored under `vendor/zennotes/` with their source identity and checksums (`manifest.json`), and `package-lock.json` pins the complete install. No source checkout is used. The vendored set is the core release -[core-2.60.0-core.hb0d0b54f320a8e3f](https://github.com/ZenNotes/zennotes/releases/tag/core-2.60.0-core.hb0d0b54f320a8e3f) -(desktop commit `15829394`, clean tree). Run `npm run +[core-2.60.1-core.h05ebb55c14afffb2](https://github.com/ZenNotes/zennotes/releases/tag/core-2.60.1-core.h05ebb55c14afffb2) +(desktop commit `4c74b478`, clean tree). Run `npm run boundaries:check` to verify archives, installed versions, singleton editor/React peers, and imports; it refuses an archive built from a dirty upstream tree unless `ZEN_ALLOW_DIRTY_CORE=1` is set for a local try-out. diff --git a/android/app/src/androidTest/AndroidManifest.xml b/android/app/src/androidTest/AndroidManifest.xml index 7d82873..0341a71 100644 --- a/android/app/src/androidTest/AndroidManifest.xml +++ b/android/app/src/androidTest/AndroidManifest.xml @@ -4,5 +4,9 @@ + diff --git a/android/app/src/androidTest/java/md/zennotes/CloudUploadFixtureProvider.java b/android/app/src/androidTest/java/md/zennotes/CloudUploadFixtureProvider.java new file mode 100644 index 0000000..04d0d1a --- /dev/null +++ b/android/app/src/androidTest/java/md/zennotes/CloudUploadFixtureProvider.java @@ -0,0 +1,15 @@ +package md.zennotes; + +import android.content.Intent; +import android.net.Uri; +import android.os.Bundle; + +/** Protected like a document provider; the test must explicitly receive a URI grant. */ +public class CloudUploadFixtureProvider extends SafFixtureProvider { + @Override public Bundle call(String method, String arg, Bundle extras) { + if (!"grant".equals(method) || !"md.zennotes".equals(arg)) throw new SecurityException("Unsupported fixture call"); + Uri uri = Uri.parse("content://md.zennotes.test.cloud-files/tree/root/document/large"); + getContext().grantUriPermission(arg, uri, Intent.FLAG_GRANT_READ_URI_PERMISSION); + return Bundle.EMPTY; + } +} diff --git a/android/app/src/androidTest/java/md/zennotes/DirectUploadInstrumentedTest.java b/android/app/src/androidTest/java/md/zennotes/DirectUploadInstrumentedTest.java new file mode 100644 index 0000000..94d6811 --- /dev/null +++ b/android/app/src/androidTest/java/md/zennotes/DirectUploadInstrumentedTest.java @@ -0,0 +1,139 @@ +package md.zennotes; + +import static org.junit.Assert.*; + +import android.content.Context; +import android.content.Intent; +import android.net.Uri; +import android.provider.DocumentsContract; +import androidx.test.ext.junit.runners.AndroidJUnit4; +import androidx.test.platform.app.InstrumentationRegistry; +import com.getcapacitor.JSObject; +import com.getcapacitor.PluginCall; +import java.io.ByteArrayOutputStream; +import java.io.File; +import java.io.FileOutputStream; +import java.io.InputStream; +import java.net.ServerSocket; +import java.net.Socket; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.util.Arrays; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import org.junit.Test; +import org.junit.runner.RunWith; + +@RunWith(AndroidJUnit4.class) +public class DirectUploadInstrumentedTest { + private static class Result extends PluginCall { + String message; + JSObject value; + Result(JSObject data) { super(null, "ZenDirectUpload", "fixture", "test", data); } + @Override public void resolve(JSObject data) { value = data; } + @Override public void resolve() { value = new JSObject(); } + @Override public void reject(String message, String code, Exception error, JSObject data) { this.message = message + " " + error; } + } + + private DirectUploadPlugin plugin(Context context) { + return new DirectUploadPlugin() { + @Override public Context getContext() { return context; } + @Override public void execute(Runnable work) { work.run(); } + }; + } + + @Test public void inspectsAndUploadsLocalAndDocumentProviderFilesWithoutBase64() throws Exception { + Context context = InstrumentationRegistry.getInstrumentation().getTargetContext(); + File folder = new File(context.getFilesDir(), "ZenNotes/CapacityFixture"); + assertTrue(folder.isDirectory() || folder.mkdirs()); + File file = new File(folder, "large.bin"); + try { + byte[] chunk = new byte[64 * 1024]; + Arrays.fill(chunk, (byte) 197); + try (FileOutputStream output = new FileOutputStream(file)) { + for (int left = 6_000_000; left > 0; left -= Math.min(left, chunk.length)) { + output.write(chunk, 0, Math.min(left, chunk.length)); + } + } + Uri tree = DocumentsContract.buildTreeDocumentUri("md.zennotes.test.cloud-files", "root"); + for (Uri uri : new Uri[] { Uri.fromFile(file), DocumentsContract.buildDocumentUriUsingTree(tree, "large") }) { + DirectUploadPlugin plugin = plugin(context); + if ("content".equals(uri.getScheme())) { + context.revokeUriPermission(uri, Intent.FLAG_GRANT_READ_URI_PERMISSION); + Result denied = new Result(new JSObject().put("uri", uri.toString())); + plugin.inspect(denied); + assertNotNull(denied.message); + context.getContentResolver().call(uri, "grant", context.getPackageName(), null); + } + Result inspected = new Result(new JSObject().put("uri", uri.toString()).put("textCandidate", false)); + plugin.inspect(inspected); + assertNull(uri.toString(), inspected.message); + assertEquals(6_000_000L, inspected.value.getLong("byteLength")); + assertFalse(inspected.value.has("inlineBase64")); + File copy = new File(folder, "copy.bin"); + Result copied = new Result(new JSObject().put("from", uri.toString()).put("to", Uri.fromFile(copy).toString()) + .put("byteLength", 6_000_000).put("sha256", inspected.value.getString("sha256"))); + plugin.copy(copied); + assertNull(copied.message); + Result copyCheck = new Result(new JSObject().put("uri", Uri.fromFile(copy).toString())); + plugin.inspect(copyCheck); + assertEquals(inspected.value.getString("sha256"), copyCheck.value.getString("sha256")); + assertTrue(copy.delete()); + try (ServerSocket server = new ServerSocket(0, 1, java.net.InetAddress.getByName("127.0.0.1"))) { + CompletableFuture received = CompletableFuture.supplyAsync(() -> receive(server)); + Result uploaded = new Result(new JSObject() + .put("url", "http://127.0.0.1:" + server.getLocalPort() + "/object") + .put("uri", uri.toString()).put("byteLength", 6_000_000) + .put("sha256", inspected.value.getString("sha256"))); + plugin.put(uploaded); + assertNull(uploaded.message); + assertEquals(200, uploaded.value.getInteger("status").intValue()); + assertEquals(inspected.value.getString("sha256"), received.get(10, TimeUnit.SECONDS)); + } + } + } finally { + file.delete(); + new File(context.getCacheDir(), "saf-capacity-fixture.bin").delete(); + } + } + + @Test public void refusesFilesOutsideVaultStorage() { + Context context = InstrumentationRegistry.getInstrumentation().getTargetContext(); + Result request = new Result(new JSObject().put("uri", Uri.fromFile(new File(context.getFilesDir(), "credential.json")).toString())); + plugin(context).inspect(request); + assertNotNull(request.message); + assertNull(request.value); + } + + private static String receive(ServerSocket server) { + try (Socket socket = server.accept()) { + socket.setSoTimeout(10_000); + InputStream input = socket.getInputStream(); + ByteArrayOutputStream header = new ByteArrayOutputStream(); + int ending = 0; + while (ending != 0x0d0a0d0a && header.size() < 16_384) { + int value = input.read(); + if (value < 0) throw new IllegalStateException("Incomplete upload headers"); + header.write(value); + ending = (ending << 8) | value; + } + String headers = header.toString("UTF-8").toLowerCase(); + assertTrue(headers.contains("content-length: 6000000")); + assertFalse(headers.contains("transfer-encoding:")); + assertFalse(headers.contains("authorization:")); + MessageDigest hash = MessageDigest.getInstance("SHA-256"); + byte[] chunk = new byte[64 * 1024]; + int remaining = 6_000_000; + while (remaining > 0) { + int count = input.read(chunk, 0, Math.min(chunk.length, remaining)); + if (count < 0) throw new IllegalStateException("Truncated upload"); + hash.update(chunk, 0, count); + remaining -= count; + } + socket.getOutputStream().write("HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n".getBytes(StandardCharsets.US_ASCII)); + StringBuilder digest = new StringBuilder(); + for (byte value : hash.digest()) digest.append(String.format("%02x", value)); + return digest.toString(); + } catch (Exception error) { throw new IllegalStateException(error); } + } +} diff --git a/android/app/src/androidTest/java/md/zennotes/SafFixtureProvider.java b/android/app/src/androidTest/java/md/zennotes/SafFixtureProvider.java index a828b2f..3dfc681 100644 --- a/android/app/src/androidTest/java/md/zennotes/SafFixtureProvider.java +++ b/android/app/src/androidTest/java/md/zennotes/SafFixtureProvider.java @@ -10,6 +10,7 @@ import android.provider.DocumentsContract; import java.io.File; import java.io.FileNotFoundException; +import java.io.FileOutputStream; import java.nio.charset.StandardCharsets; import java.nio.file.Files; @@ -35,8 +36,19 @@ public class SafFixtureProvider extends ContentProvider { } @Override public ParcelFileDescriptor openFile(Uri uri, String mode) throws FileNotFoundException { String id = DocumentsContract.getDocumentId(uri); - if (!id.equals("present")) throw new FileNotFoundException("Provider cannot open this document"); + if (!id.equals("present") && !id.equals("large")) throw new FileNotFoundException("Provider cannot open this document"); try { + if (id.equals("large")) { + File large = new File(getContext().getCacheDir(), "saf-capacity-fixture.bin"); + try (FileOutputStream output = new FileOutputStream(large)) { + byte[] chunk = new byte[64 * 1024]; + java.util.Arrays.fill(chunk, (byte) 197); + for (int left = 6_000_000; left > 0; left -= Math.min(left, chunk.length)) { + output.write(chunk, 0, Math.min(left, chunk.length)); + } + } + return ParcelFileDescriptor.open(large, ParcelFileDescriptor.MODE_READ_ONLY); + } File file = new File(getContext().getCacheDir(), "saf-fixture.md"); Files.write(file.toPath(), "Exact café 日本語. \n".getBytes(StandardCharsets.UTF_8)); return ParcelFileDescriptor.open(file, ParcelFileDescriptor.MODE_READ_ONLY); diff --git a/android/app/src/main/java/md/zennotes/CloudFileStream.java b/android/app/src/main/java/md/zennotes/CloudFileStream.java new file mode 100644 index 0000000..1191613 --- /dev/null +++ b/android/app/src/main/java/md/zennotes/CloudFileStream.java @@ -0,0 +1,114 @@ +package md.zennotes; + +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.nio.ByteBuffer; +import java.nio.Buffer; +import java.nio.CharBuffer; +import java.nio.charset.CharsetDecoder; +import java.nio.charset.CoderResult; +import java.nio.charset.CodingErrorAction; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; + +/** Bounded file inspection and copying, including UTF-8 sequences split across reads. */ +final class CloudFileStream { + static final int BUFFER_SIZE = 64 * 1024; + static final int INLINE_LIMIT = 5 * 1024 * 1024; + + static final class Fingerprint { + final long byteLength; + final String sha256; + final boolean utf8; + final byte[] inlineBytes; + + Fingerprint(long byteLength, String sha256, boolean utf8, byte[] inlineBytes) { + this.byteLength = byteLength; + this.sha256 = sha256; + this.utf8 = utf8; + this.inlineBytes = inlineBytes; + } + } + + static Fingerprint inspect(InputStream input, boolean textCandidate) throws IOException { + MessageDigest digest = digest(); + byte[] chunk = new byte[BUFFER_SIZE]; + ByteArrayOutputStream inline = new ByteArrayOutputStream(); + Utf8Check text = textCandidate ? new Utf8Check() : null; + long length = 0; + int count; + while ((count = input.read(chunk)) != -1) { + if (Thread.currentThread().isInterrupted()) throw new IOException("File inspection cancelled."); + digest.update(chunk, 0, count); + length += count; + if (inline != null) { + if (length <= INLINE_LIMIT) inline.write(chunk, 0, count); + else inline = null; + } + if (text != null) text.accept(chunk, count, false); + } + if (text != null) text.accept(chunk, 0, true); + return new Fingerprint(length, hex(digest.digest()), text != null && text.valid, + inline == null ? null : inline.toByteArray()); + } + + static void copyVerified(InputStream input, OutputStream output, long expectedBytes, String expectedHash) throws IOException { + MessageDigest digest = digest(); + byte[] chunk = new byte[BUFFER_SIZE]; + long written = 0; + int count; + while ((count = input.read(chunk)) != -1) { + if (Thread.currentThread().isInterrupted()) throw new IOException("File upload cancelled."); + written += count; + if (written > expectedBytes) throw new IOException("The upload file changed size."); + digest.update(chunk, 0, count); + output.write(chunk, 0, count); + } + if (written != expectedBytes || !hex(digest.digest()).equals(expectedHash)) { + throw new IOException("The upload file changed after its scan."); + } + output.flush(); + } + + private static MessageDigest digest() { + try { + return MessageDigest.getInstance("SHA-256"); + } catch (NoSuchAlgorithmException impossible) { + throw new IllegalStateException(impossible); + } + } + + private static String hex(byte[] bytes) { + char[] alphabet = "0123456789abcdef".toCharArray(); + char[] out = new char[bytes.length * 2]; + for (int i = 0; i < bytes.length; i++) { + out[i * 2] = alphabet[(bytes[i] & 255) >>> 4]; + out[i * 2 + 1] = alphabet[bytes[i] & 15]; + } + return new String(out); + } + + private static final class Utf8Check { + final CharsetDecoder decoder = StandardCharsets.UTF_8.newDecoder() + .onMalformedInput(CodingErrorAction.REPORT).onUnmappableCharacter(CodingErrorAction.REPORT); + final ByteBuffer pending = ByteBuffer.allocate(BUFFER_SIZE + 4); + final CharBuffer characters = CharBuffer.allocate(8192); + boolean valid = true; + + void accept(byte[] chunk, int count, boolean last) { + if (!valid) return; + pending.put(chunk, 0, count); + ((Buffer) pending).flip(); + CoderResult result; + do { + ((Buffer) characters).clear(); + result = decoder.decode(pending, characters, last); + } while (result.isOverflow()); + valid = !result.isError(); + pending.compact(); + } + } +} diff --git a/android/app/src/main/java/md/zennotes/DirectUploadPlugin.java b/android/app/src/main/java/md/zennotes/DirectUploadPlugin.java index 028f812..ee7440e 100644 --- a/android/app/src/main/java/md/zennotes/DirectUploadPlugin.java +++ b/android/app/src/main/java/md/zennotes/DirectUploadPlugin.java @@ -2,6 +2,10 @@ import android.util.Base64; import android.util.Base64InputStream; +import android.content.Intent; +import android.content.pm.PackageManager; +import android.net.Uri; +import android.os.Process; import com.getcapacitor.JSObject; import com.getcapacitor.Plugin; @@ -10,6 +14,10 @@ import com.getcapacitor.annotation.CapacitorPlugin; import java.io.ByteArrayInputStream; +import java.io.File; +import java.io.FileInputStream; +import java.io.FileOutputStream; +import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; import java.net.HttpURLConnection; @@ -18,7 +26,7 @@ import java.util.Iterator; /** - * Streams a base64 Cloud item to its short-lived signed object-storage URL. + * Streams a vault file to its short-lived signed object-storage URL. * * Capacitor's Android `dataType: file` implementation only decodes the body * on API 26+, while ZenNotes supports API 24. Keeping this tiny uploader in @@ -32,6 +40,111 @@ public class DirectUploadPlugin extends Plugin { private static final int READ_TIMEOUT_MS = 300_000; private static final int BUFFER_SIZE = 64 * 1024; + @PluginMethod + public void inspect(PluginCall call) { + execute(() -> { + try (InputStream input = openVaultFile(call.getString("uri"))) { + CloudFileStream.Fingerprint fingerprint = CloudFileStream.inspect(input, call.getBoolean("textCandidate", false)); + JSObject result = new JSObject(); + result.put("uri", call.getString("uri")); + result.put("byteLength", fingerprint.byteLength); + result.put("sha256", fingerprint.sha256); + result.put("utf8", fingerprint.utf8); + if (fingerprint.inlineBytes != null) { + result.put("inlineBase64", Base64.encodeToString(fingerprint.inlineBytes, Base64.NO_WRAP)); + } + call.resolve(result); + } catch (Exception error) { + call.reject("Cloud file inspection failed.", "CLOUD_FILE_READ_FAILED", error); + } + }); + } + + @PluginMethod + public void copy(PluginCall call) { + execute(() -> { + Long bytes = call.getLong("byteLength"); + if (bytes == null && call.getInt("byteLength") != null) bytes = call.getInt("byteLength").longValue(); + String hash = call.getString("sha256"); + if (bytes == null || bytes < 0 || hash == null || !hash.matches("[0-9a-f]{64}")) { + call.reject("Invalid Cloud copy request.", "INVALID_COPY_REQUEST"); + return; + } + try { + Uri target = vaultFileUri(call.getString("to"), Intent.FLAG_GRANT_WRITE_URI_PERMISSION); + Uri source = vaultFileUri(call.getString("from"), Intent.FLAG_GRANT_READ_URI_PERMISSION); + if (source.equals(target)) throw new IOException("Cloud copy source and destination must differ."); + try (InputStream input = openVaultFile(source.toString()); + OutputStream output = "content".equals(target.getScheme()) + ? getContext().getContentResolver().openOutputStream(target, "wt") + : new FileOutputStream(target.getPath())) { + if (output == null) throw new IOException("Cloud copy destination is unavailable."); + CloudFileStream.copyVerified(input, output, bytes, hash); + } + call.resolve(); + } catch (Exception error) { + call.reject("Cloud file copy failed.", "CLOUD_FILE_COPY_FAILED", error); + } + }); + } + + /** + * Stream a signed Cloud revision GET into a vault staging file. The body + * never enters the WebView; a length or SHA-256 mismatch deletes the + * partial file and rejects, so a truncated download can never be published. + */ + @PluginMethod + public void download(PluginCall call) { + execute(() -> { + Long bytes = call.getLong("byteLength"); + if (bytes == null && call.getInt("byteLength") != null) bytes = call.getInt("byteLength").longValue(); + String hash = call.getString("sha256"); + String value = call.getString("url"); + JSObject headers = call.getObject("headers", new JSObject()); + if (bytes == null || bytes < 0 || hash == null || !hash.matches("[0-9a-f]{64}") || value == null) { + call.reject("Invalid Cloud download request.", "INVALID_DOWNLOAD_REQUEST"); + return; + } + HttpURLConnection connection = null; + Uri target = null; + try { + target = vaultFileUri(call.getString("to"), Intent.FLAG_GRANT_WRITE_URI_PERMISSION); + URL url = new URL(value); + validateUrl(url); + connection = (HttpURLConnection) url.openConnection(); + connection.setRequestMethod("GET"); + connection.setInstanceFollowRedirects(false); + connection.setConnectTimeout(CONNECT_TIMEOUT_MS); + connection.setReadTimeout(READ_TIMEOUT_MS); + connection.setRequestProperty("Accept-Encoding", "identity"); + applyHeaders(connection, headers); + int status = connection.getResponseCode(); + if (status < 200 || status >= 300) { + drain(connection.getErrorStream()); + JSObject data = new JSObject(); + data.put("status", status); + call.reject("Cloud object download failed (" + status + ").", "DIRECT_DOWNLOAD_FAILED", null, data); + return; + } + long declared = connection.getContentLengthLong(); + if (declared >= 0 && declared != bytes) throw new IOException("Cloud download length did not match its reference."); + try (InputStream input = connection.getInputStream(); + OutputStream output = "content".equals(target.getScheme()) + ? getContext().getContentResolver().openOutputStream(target, "wt") + : new FileOutputStream(target.getPath())) { + if (output == null) throw new IOException("Cloud download destination is unavailable."); + CloudFileStream.copyVerified(input, output, bytes, hash); + } + call.resolve(); + } catch (Exception error) { + if (target != null && "file".equals(target.getScheme())) new File(target.getPath()).delete(); + call.reject("Cloud object download failed.", "DIRECT_DOWNLOAD_FAILED", error); + } finally { + if (connection != null) connection.disconnect(); + } + }); + } + @PluginMethod public void put(PluginCall call) { execute(() -> upload(call)); @@ -42,9 +155,12 @@ private void upload(PluginCall call) { try { String value = call.getString("url"); String base64 = call.getString("base64"); + String sourceUri = call.getString("uri"); + String expectedHash = call.getString("sha256"); Integer expectedBytes = call.getInt("byteLength"); JSObject headers = call.getObject("headers", new JSObject()); - if (value == null || base64 == null || expectedBytes == null || expectedBytes < 0) { + if (value == null || (base64 == null && sourceUri == null) || expectedBytes == null || expectedBytes < 0 + || (sourceUri != null && (expectedHash == null || !expectedHash.matches("[0-9a-f]{64}")))) { call.reject("Invalid direct-upload request.", "INVALID_DIRECT_UPLOAD_REQUEST"); return; } @@ -60,22 +176,24 @@ private void upload(PluginCall call) { connection.setFixedLengthStreamingMode(expectedBytes); applyHeaders(connection, headers); - long written = 0; - byte[] buffer = new byte[BUFFER_SIZE]; try ( - InputStream encoded = new ByteArrayInputStream(base64.getBytes(StandardCharsets.US_ASCII)); - InputStream decoded = new Base64InputStream(encoded, Base64.DEFAULT); + InputStream decoded = sourceUri != null ? openVaultFile(sourceUri) : new Base64InputStream( + new ByteArrayInputStream(base64.getBytes(StandardCharsets.US_ASCII)), Base64.DEFAULT); OutputStream output = connection.getOutputStream() ) { - int count; - while ((count = decoded.read(buffer)) != -1) { - output.write(buffer, 0, count); - written += count; + if (expectedHash != null) { + CloudFileStream.copyVerified(decoded, output, expectedBytes, expectedHash); + } else { + long written = 0; + byte[] buffer = new byte[BUFFER_SIZE]; + int count; + while ((count = decoded.read(buffer)) != -1) { + output.write(buffer, 0, count); + written += count; + } + if (written != expectedBytes) throw new IOException("Decoded upload size changed before transmission."); + output.flush(); } - output.flush(); - } - if (written != expectedBytes) { - throw new IllegalArgumentException("Decoded upload size changed before transmission."); } int status = connection.getResponseCode(); @@ -90,6 +208,37 @@ private void upload(PluginCall call) { } } + private InputStream openVaultFile(String value) throws IOException { + Uri uri = vaultFileUri(value, Intent.FLAG_GRANT_READ_URI_PERMISSION); + if ("content".equals(uri.getScheme())) { + InputStream input = getContext().getContentResolver().openInputStream(uri); + if (input == null) throw new IOException("Vault file is unavailable."); + return input; + } + return new FileInputStream(uri.getPath()); + } + + private Uri vaultFileUri(String value, int access) throws IOException { + if (value == null) throw new IOException("Missing vault file."); + Uri uri = Uri.parse(value); + if ("content".equals(uri.getScheme())) { + if (getContext().checkUriPermission(uri, Process.myPid(), Process.myUid(), access) + != PackageManager.PERMISSION_GRANTED) throw new IOException("Vault file access is not granted."); + return uri; + } + if (!"file".equals(uri.getScheme()) || uri.getPath() == null || + (uri.getAuthority() != null && !uri.getAuthority().isEmpty())) throw new IOException("Invalid vault file URI."); + File file = new File(uri.getPath()).getCanonicalFile(); + if (!insideVaults(file, getContext().getFilesDir()) && !insideVaults(file, getContext().getExternalFilesDir(null))) { + throw new IOException("File is outside the local vault storage."); + } + return Uri.fromFile(file); + } + + private static boolean insideVaults(File file, File storage) throws IOException { + return storage != null && file.getPath().startsWith(new File(storage, "ZenNotes").getCanonicalPath() + File.separator); + } + private static void validateUrl(URL url) { String protocol = url.getProtocol(); String host = url.getHost() == null ? "" : url.getHost().toLowerCase(); diff --git a/android/app/src/test/java/md/zennotes/CloudFileStreamTest.java b/android/app/src/test/java/md/zennotes/CloudFileStreamTest.java new file mode 100644 index 0000000..d07c1e4 --- /dev/null +++ b/android/app/src/test/java/md/zennotes/CloudFileStreamTest.java @@ -0,0 +1,70 @@ +package md.zennotes; + +import static org.junit.Assert.*; + +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.util.Arrays; +import org.junit.Test; + +public class CloudFileStreamTest { + @Test + public void fingerprintsLargeInputWithBoundedReadsAndNoInlineBody() throws Exception { + byte[] bytes = new byte[8_000_000]; + Arrays.fill(bytes, (byte) 129); + InputStream input = new ByteArrayInputStream(bytes) { + @Override public synchronized int read(byte[] target, int offset, int length) { + assertTrue("Native reads must stay bounded", length <= 65_536); + return super.read(target, offset, length); + } + }; + CloudFileStream.Fingerprint result = CloudFileStream.inspect(input, false); + assertEquals(bytes.length, result.byteLength); + assertEquals(hash(bytes), result.sha256); + assertNull(result.inlineBytes); + } + + @Test + public void validatesUtf8AcrossChunkBoundariesWithoutChangingRawBytes() throws Exception { + byte[] bytes = ("a".repeat(65_535) + "日本語 café").getBytes(StandardCharsets.UTF_8); + CloudFileStream.Fingerprint result = CloudFileStream.inspect(new ByteArrayInputStream(bytes), true); + assertTrue(result.utf8); + assertArrayEquals(bytes, result.inlineBytes); + assertEquals(hash(bytes), result.sha256); + bytes[bytes.length - 1] = (byte) 255; + assertFalse(CloudFileStream.inspect(new ByteArrayInputStream(bytes), true).utf8); + } + + @Test + public void detectsTruncatedUtf8AndStillHashesTheEntireFile() throws Exception { + byte[] bytes = new byte[] { 97, (byte) 0xe6, (byte) 0x97 }; + CloudFileStream.Fingerprint result = CloudFileStream.inspect(new ByteArrayInputStream(bytes), true); + assertFalse(result.utf8); + assertEquals(hash(bytes), result.sha256); + } + + @Test + public void streamsAnExactUploadAndRejectsChangedContentOrLength() throws Exception { + byte[] bytes = "exact café bytes".getBytes(StandardCharsets.UTF_8); + ByteArrayOutputStream output = new ByteArrayOutputStream(); + CloudFileStream.copyVerified(new ByteArrayInputStream(bytes), output, bytes.length, hash(bytes)); + assertArrayEquals(bytes, output.toByteArray()); + for (int expected : new int[] { bytes.length - 1, bytes.length + 1 }) { + assertThrows(IOException.class, () -> CloudFileStream.copyVerified( + new ByteArrayInputStream(bytes), OutputStream.nullOutputStream(), expected, hash(bytes))); + } + assertThrows(IOException.class, () -> CloudFileStream.copyVerified( + new ByteArrayInputStream(bytes), OutputStream.nullOutputStream(), bytes.length, hash(new byte[] { 1 }))); + } + + private static String hash(byte[] bytes) throws Exception { + StringBuilder hex = new StringBuilder(); + for (byte value : MessageDigest.getInstance("SHA-256").digest(bytes)) hex.append(String.format("%02x", value)); + return hex.toString(); + } +} diff --git a/package-lock.json b/package-lock.json index e358302..4132ecf 100644 --- a/package-lock.json +++ b/package-lock.json @@ -32,9 +32,9 @@ "@lezer/highlight": "^1.2.5", "@replit/codemirror-vim": "^6.3.0", "@xyflow/react": "^12.12.0", - "@zennotes/app-core": "file:vendor/zennotes/zennotes-app-core-2.60.0-core.hb0d0b54f320a8e3f.tgz", - "@zennotes/bridge-contract": "file:vendor/zennotes/zennotes-bridge-contract-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", - "@zennotes/shared-domain": "file:vendor/zennotes/zennotes-shared-domain-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", + "@zennotes/app-core": "file:vendor/zennotes/zennotes-app-core-2.60.1-core.h05ebb55c14afffb2.tgz", + "@zennotes/bridge-contract": "file:vendor/zennotes/zennotes-bridge-contract-2.60.1-boundaries.ha881a2f8575a38bf.tgz", + "@zennotes/shared-domain": "file:vendor/zennotes/zennotes-shared-domain-2.60.1-boundaries.ha881a2f8575a38bf.tgz", "codemirror": "^6.0.1", "dompurify": "^3.4.16", "function-plot": "^1.25.3", @@ -698,6 +698,7 @@ "cpu": [ "ppc64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -714,6 +715,7 @@ "cpu": [ "arm" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -730,6 +732,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -746,6 +749,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -762,6 +766,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -778,6 +783,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -794,6 +800,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -810,6 +817,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -826,6 +834,7 @@ "cpu": [ "arm" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -842,6 +851,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -858,6 +868,7 @@ "cpu": [ "ia32" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -874,6 +885,7 @@ "cpu": [ "loong64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -890,6 +902,7 @@ "cpu": [ "mips64el" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -906,6 +919,7 @@ "cpu": [ "ppc64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -922,6 +936,7 @@ "cpu": [ "riscv64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -938,6 +953,7 @@ "cpu": [ "s390x" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -954,6 +970,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -970,6 +987,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -986,6 +1004,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -1002,6 +1021,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -1018,6 +1038,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -1034,6 +1055,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -1050,6 +1072,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -1066,6 +1089,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -1082,6 +1106,7 @@ "cpu": [ "ia32" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -1098,6 +1123,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2579,6 +2605,7 @@ "cpu": [ "arm" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2595,6 +2622,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2611,6 +2639,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2627,6 +2656,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2643,6 +2673,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2659,6 +2690,7 @@ "cpu": [ "arm" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2675,9 +2707,7 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2694,9 +2724,7 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2713,9 +2741,7 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2732,9 +2758,7 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2751,9 +2775,7 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2770,9 +2792,7 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2789,6 +2809,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2805,6 +2826,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -2821,6 +2843,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MIT", "optional": true, "os": [ @@ -3149,7 +3172,7 @@ "version": "22.20.4", "resolved": "https://registry.npmjs.org/@types/node/-/node-22.20.4.tgz", "integrity": "sha512-zJRE40jpHtKqE/C4fgHrAKQLJuSpzEnP9ff9Y7YtoR3Wd2pwqzlekDeEuUQXjRd+QCYnVnNwuJYmhdk9XV8gvA==", - "devOptional": true, + "dev": true, "license": "MIT", "dependencies": { "undici-types": "~6.21.0" @@ -3330,9 +3353,9 @@ } }, "node_modules/@zennotes/app-core": { - "version": "2.60.0-core.hb0d0b54f320a8e3f", - "resolved": "file:vendor/zennotes/zennotes-app-core-2.60.0-core.hb0d0b54f320a8e3f.tgz", - "integrity": "sha512-aMZ8DegEHQTwT4tf6kdvAHEUMfhi4nayN81lHKw4PSZZfDr4UKvhpufIZEm5xO4UR4gyKEJSwd4/1QNBB+aC+Q==", + "version": "2.60.1-core.h05ebb55c14afffb2", + "resolved": "file:vendor/zennotes/zennotes-app-core-2.60.1-core.h05ebb55c14afffb2.tgz", + "integrity": "sha512-pTn7+3VmHG1le5pzBruIAzlHg4Nc+CRucL/u4FLbktdPDeIQ/PWECeS+fkia7JM5w6UVhqFxY5ZP24gj77vfFg==", "license": "MIT", "dependencies": { "@codemirror/autocomplete": "^6.18.3", @@ -3360,8 +3383,8 @@ "@myriaddreamin/typst.ts": "^0.7.0", "@replit/codemirror-vim": "^6.3.0", "@xyflow/react": "^12.11.2", - "@zennotes/bridge-contract": "2.60.0-boundaries.hbdc4a2a12368fec1", - "@zennotes/shared-domain": "2.60.0-boundaries.hbdc4a2a12368fec1", + "@zennotes/bridge-contract": "2.60.1-boundaries.ha881a2f8575a38bf", + "@zennotes/shared-domain": "2.60.1-boundaries.ha881a2f8575a38bf", "dompurify": "^3.3.4", "function-plot": "^1.25.3", "gray-matter": "^4.0.3", @@ -3405,18 +3428,18 @@ } }, "node_modules/@zennotes/bridge-contract": { - "version": "2.60.0-boundaries.hbdc4a2a12368fec1", - "resolved": "file:vendor/zennotes/zennotes-bridge-contract-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", - "integrity": "sha512-5uh7dQjtrWNds26/4FH3wcrEQ0cF6SgUYd0+bGHKtl8GOLlayJWmDaJPXgm7FQ97Gc18KNM47k6k2rxYuECx3A==", + "version": "2.60.1-boundaries.ha881a2f8575a38bf", + "resolved": "file:vendor/zennotes/zennotes-bridge-contract-2.60.1-boundaries.ha881a2f8575a38bf.tgz", + "integrity": "sha512-/BjWqbzR4+S8Bg8gOwT6OZO8Ikr8qO8zNmVVscjr45vIWasAykfaLS9H6IuFfIE48/wbAXb3W8J6aTHNEdk0xA==", "license": "MIT" }, "node_modules/@zennotes/shared-domain": { - "version": "2.60.0-boundaries.hbdc4a2a12368fec1", - "resolved": "file:vendor/zennotes/zennotes-shared-domain-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", - "integrity": "sha512-Gr/F9IkAnX24T9apcPWLhErAtQ4N/QUr/rIdHJqNOTxERVSjrvIqP+wGCzNyuVRDujX9eXOfAt0y4Ir6Ol0TwQ==", + "version": "2.60.1-boundaries.ha881a2f8575a38bf", + "resolved": "file:vendor/zennotes/zennotes-shared-domain-2.60.1-boundaries.ha881a2f8575a38bf.tgz", + "integrity": "sha512-5Sb02aV5yUy3MkqR9rkrx7eHQTRmurivhMVpeVa4kU0S07dm08N6kpstB8uIeKVlxhDvLILQHhBw606N3wx+aA==", "license": "MIT", "dependencies": { - "@zennotes/bridge-contract": "2.60.0-boundaries.hbdc4a2a12368fec1", + "@zennotes/bridge-contract": "2.60.1-boundaries.ha881a2f8575a38bf", "lz-string": "^1.5.0" } }, @@ -4716,7 +4739,7 @@ "version": "0.28.2", "resolved": "https://registry.npmjs.org/esbuild/-/esbuild-0.28.2.tgz", "integrity": "sha512-HKVLS8dvII+xoKW9kmqxbRKrnWEXfJJr/FZhhJmiqIB0e053QNYFqOBouTMO/k5sID4MvCiUCvv8b9M4h32wIA==", - "devOptional": true, + "dev": true, "hasInstallScript": true, "license": "MIT", "bin": { @@ -5542,7 +5565,7 @@ "version": "1.21.7", "resolved": "https://registry.npmjs.org/jiti/-/jiti-1.21.7.tgz", "integrity": "sha512-/imKNG4EbWNrVjoNC/1H5/9GFy+tqjGBHCaSsN+P2RnPqjsLmv6UD3Ej+Kj8nBWaRAwyk7kK5ZUc+OEatnTR3A==", - "devOptional": true, + "dev": true, "license": "MIT", "bin": { "jiti": "bin/jiti.js" @@ -5715,6 +5738,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5735,6 +5759,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5755,6 +5780,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5775,6 +5801,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5795,6 +5822,7 @@ "cpu": [ "arm" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5815,6 +5843,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5835,6 +5864,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5855,6 +5885,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5875,6 +5906,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5895,6 +5927,7 @@ "cpu": [ "arm64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -5915,6 +5948,7 @@ "cpu": [ "x64" ], + "dev": true, "license": "MPL-2.0", "optional": true, "os": [ @@ -8644,7 +8678,7 @@ "version": "6.21.0", "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz", "integrity": "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ==", - "devOptional": true, + "dev": true, "license": "MIT" }, "node_modules/unified": { diff --git a/package.json b/package.json index 4086073..804ec7a 100644 --- a/package.json +++ b/package.json @@ -43,9 +43,9 @@ "@lezer/highlight": "^1.2.5", "@replit/codemirror-vim": "^6.3.0", "@xyflow/react": "^12.12.0", - "@zennotes/app-core": "file:vendor/zennotes/zennotes-app-core-2.60.0-core.hb0d0b54f320a8e3f.tgz", - "@zennotes/bridge-contract": "file:vendor/zennotes/zennotes-bridge-contract-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", - "@zennotes/shared-domain": "file:vendor/zennotes/zennotes-shared-domain-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", + "@zennotes/app-core": "file:vendor/zennotes/zennotes-app-core-2.60.1-core.h05ebb55c14afffb2.tgz", + "@zennotes/bridge-contract": "file:vendor/zennotes/zennotes-bridge-contract-2.60.1-boundaries.ha881a2f8575a38bf.tgz", + "@zennotes/shared-domain": "file:vendor/zennotes/zennotes-shared-domain-2.60.1-boundaries.ha881a2f8575a38bf.tgz", "codemirror": "^6.0.1", "dompurify": "^3.4.16", "function-plot": "^1.25.3", diff --git a/src/bridge/cloud-rate-limit.test.ts b/src/bridge/cloud-rate-limit.test.ts new file mode 100644 index 0000000..a9eaaaa --- /dev/null +++ b/src/bridge/cloud-rate-limit.test.ts @@ -0,0 +1,84 @@ +import assert from 'node:assert/strict' +import { it } from 'node:test' +import { loadMobileModule } from '../../tooling/load-mobile-module.ts' + +const flush = async () => { for (let i = 0; i < 20; i++) await Promise.resolve() } + +it('retries the same native manifest page after the complete Retry-After delay', async (test) => { + const requests: string[] = [] + const { createCloudSyncClient } = await loadMobileModule('./src/bridge/cloud-sync-client', { + '@capacitor/core': { + registerPlugin: () => ({}), + CapacitorHttp: { request: async ({ url }: { url: string }) => { + requests.push(url) + return requests.length === 1 + ? { status: 429, data: 'Slow down', headers: { 'rEtRy-AfTeR': '2' } } + : { status: 200, data: { data: [], cursor: 500, next_page: null } } + } } + } + }) + test.mock.timers.enable({ apis: ['Date', 'setTimeout'], now: 1_000_000 }) + const client = createCloudSyncClient('https://retry.example.test', 'test-token', { accountId: 'account' }) + let settled = false + const pending = client.manifest('vault', { includeContent: false, page: 3, perPage: 250 }) + .then((data: unknown) => ({ data }), (error: unknown) => ({ error })).finally(() => { settled = true }) + await flush() + assert.equal(settled, false) + test.mock.timers.tick(1999) + await flush() + assert.equal(requests.length, 1) + test.mock.timers.tick(1) + assert.deepEqual(await pending, { data: { data: [], cursor: 500, next_page: null } }) + assert.equal(requests.length, 2) + assert.equal(requests[0], requests[1]) + assert.equal(new URL(requests[1]).searchParams.get('page'), '3') +}) + +it('cancels a native retry wait and preserves its cooldown across recreated clients', async (test) => { + const requests: string[] = [] + const api = await loadMobileModule('./src/bridge/cloud-sync-client', { + '@capacitor/core': { + registerPlugin: () => ({}), + CapacitorHttp: { request: async ({ url }: { url: string }) => { + requests.push(url) + return requests.length === 1 + ? { status: 429, data: {}, headers: { 'Retry-After': '2' } } + : { status: 200, data: { data: [] } } + } } + } + }) + test.mock.timers.enable({ apis: ['Date', 'setTimeout'], now: 1_000_000 }) + const first = api.createCloudSyncClient('https://cancel.example.test', 'old-token', { accountId: 'account' }) + .listVaults().catch((error: unknown) => error) + await flush() + api.stopMobileCloudRequests() + assert.equal((await first).name, 'AbortError') + api.resumeMobileCloudRequests() + const pending = api.createCloudSyncClient('https://cancel.example.test', 'new-token', { accountId: 'account' }).listVaults() + await flush() + assert.equal(requests.length, 1) + test.mock.timers.tick(2000) + assert.deepEqual(await pending, { data: [] }) + assert.equal(requests.length, 2) +}) + +it('preserves explicitly requested content-page offsets after negotiating references', async () => { + const requests: URL[] = [] + const { createCloudSyncClient } = await loadMobileModule('./src/bridge/cloud-sync-client', { + '@capacitor/core': { registerPlugin: () => ({}), CapacitorHttp: { request: async ({ url }: { url: string }) => { + const request = new URL(url) + requests.push(request) + if (request.pathname.endsWith('/account')) { + return { status: 200, data: { data: { capabilities: { content_references: true }, user: {}, device: {}, features: {}, usage: {} } } } + } + return { status: 200, data: { data: [], cursor: 0, next_page: null } } + } } } + }) + const client = createCloudSyncClient('https://pages.example.test', 'test-token', { accountId: 'account' }) + await client.negotiateContentReferences() + await client.manifest('vault', { includeContent: true, page: 3, perPage: 25 }) + const manifest = requests.find((url) => url.pathname.endsWith('/manifest'))! + assert.equal(manifest.searchParams.get('page'), '3') + assert.equal(manifest.searchParams.get('per_page'), '25') + assert.equal(manifest.searchParams.get('content_mode'), 'references') +}) diff --git a/src/bridge/cloud-sync-client.test.ts b/src/bridge/cloud-sync-client.test.ts index 2e9df0b..42bd7d0 100644 --- a/src/bridge/cloud-sync-client.test.ts +++ b/src/bridge/cloud-sync-client.test.ts @@ -1,53 +1,78 @@ import assert from 'node:assert/strict' +import { createHash } from 'node:crypto' import { it } from 'node:test' import { loadMobileModule } from '../../tooling/load-mobile-module.ts' -it('downloads every manifest and change page without grouping large attachments in one native response', async () => { - const requests: URL[] = [] - const items = [1, 2, 3].map((id) => ({ - item_id: `asset-${id}`, path: `assets/${id}.jpg`, kind: 'binary', revision: 1, - sha256: `hash-${id}`, byte_length: 8_700_000, media_type: 'image/jpeg', - content: { - encoding: 'base64', data: `fixture-${id}`, sha256: `hash-${id}`, - byte_length: 8_700_000, media_type: 'image/jpeg' +/** Account capability + reference-mode responses, as the backend now serves them. */ +function referenceServer(items: any[], feedRef: { feed: any[] }, requests: URL[], downloads: string[]) { + return async ({ url, method }: { url: string; method?: string }) => { + const request = new URL(url) + requests.push(request) + if (request.pathname.endsWith('/account')) { + return { status: 200, data: { data: { capabilities: { content_references: true }, user: {}, device: {}, features: {}, usage: {} } } } } - })) - let feed: any[] = [] - const { createCloudSyncClient, CloudSyncCoordinator } = await loadMobileModule([ - './src/bridge/cloud-sync-client.ts', '@zennotes/shared-domain/cloud-sync-coordinator' - ], { - '@capacitor/core': { - registerPlugin: () => ({}), - CapacitorHttp: { - request: async ({ url }: { url: string }) => { - const request = new URL(url) - requests.push(request) - if (request.pathname.endsWith('/manifest')) { - const page = Number(request.searchParams.get('page') ?? 1) - const size = Number(request.searchParams.get('per_page') ?? 100) - const data = items.slice((page - 1) * size, page * size) - assert.ok(data.length <= 1, 'content responses must fit one attachment through the native bridge') - return { status: 200, data: { data, cursor: 3, next_page: page * size < items.length ? page + 1 : null } } - } - assert.ok(request.pathname.endsWith('/changes')) - const after = Number(request.searchParams.get('after')) - const size = Number(request.searchParams.get('limit')) - const remaining = feed.filter((change) => change.sequence > after) - const data = remaining.slice(0, size) - assert.ok(data.length <= 1, 'change responses must fit one attachment through the native bridge') - return { status: 200, data: { data, cursor: feed.at(-1)?.sequence ?? 3, has_more: remaining.length > size } } - } - } + if (request.pathname.endsWith('/manifest')) { + assert.equal(request.searchParams.get('content_mode'), 'references') + const page = Number(request.searchParams.get('page') ?? 1) + const size = Number(request.searchParams.get('per_page') ?? 100) + const data = items.slice((page - 1) * size, page * size).map(({ content, ...metadata }) => ({ + ...metadata, content_ref: { item_id: metadata.item_id, revision: metadata.revision, ...content, data: undefined } + })).map((row) => ({ ...row, content_ref: Object.fromEntries(Object.entries(row.content_ref).filter(([, v]) => v !== undefined)) })) + return { status: 200, data: { data, cursor: 3, next_page: page * size < items.length ? page + 1 : null } } + } + if (request.pathname.endsWith('/download')) { + const [, itemId, revision] = request.pathname.match(/items\/([^/]+)\/revisions\/(\d+)\/download/)! + const item = items.find((candidate) => candidate.item_id === itemId)! + const content = item.versions?.[Number(revision)] ?? item.content + downloads.push(`${itemId}:${revision}`) + return { status: 200, data: { data: { + item_id: itemId, revision: Number(revision), + content: { encoding: content.encoding, sha256: content.sha256, byte_length: content.byte_length, media_type: content.media_type }, + download: { url: `https://objects.example.test/${itemId}/${revision}`, method: 'GET', headers: {}, expires_at: new Date(Date.now() + 300_000).toISOString() } + } } } } + assert.ok(request.pathname.endsWith('/changes'), `unexpected ${method} ${request.pathname}`) + assert.equal(request.searchParams.get('content_mode'), 'references') + const after = Number(request.searchParams.get('after')) + const size = Number(request.searchParams.get('limit')) + const remaining = feedRef.feed.filter((change) => change.sequence > after) + const data = remaining.slice(0, size) + return { status: 200, data: { data, cursor: feedRef.feed.at(-1)?.sequence ?? 3, has_more: remaining.length > size } } + } +} + +it('syncs large attachments as native-staged references, never as base64 through the bridge', async () => { + const requests: URL[] = [] + const downloads: string[] = [] + const items = [1, 2, 3].map((id) => { + const bytes = Buffer.alloc(1_048_577, id) + const hash = createHash('sha256').update(bytes).digest('hex') + const content = { encoding: 'base64', sha256: hash, byte_length: bytes.length, media_type: 'image/jpeg' } + return { item_id: `asset-${id}`, path: `assets/${id}.jpg`, kind: 'binary', revision: 1, sha256: hash, byte_length: bytes.length, media_type: 'image/jpeg', content, versions: { 1: content } as Record } + }) + const feedRef = { feed: [] as any[] } + const { createCloudSyncClient, CloudSyncCoordinator, registerCloudSyncStagedFile } = await loadMobileModule([ + './src/bridge/cloud-sync-client.ts', '@zennotes/shared-domain/cloud-sync-coordinator', '@zennotes/shared-domain/cloud-sync-content' + ], { + '@capacitor/core': { registerPlugin: () => ({}), CapacitorHttp: { request: referenceServer(items, feedRef, requests, downloads) } } }) const files = new Map() let state: any = null - const sync = new CloudSyncCoordinator('vault-1', createCloudSyncClient('https://example.test', 'test-only'), { + const staged: string[] = [] + const repository = { scan: async () => [...files.values()], - apply: async (change: any) => { - files.set(change.path, { path: change.path, kind: 'binary', content: change.content }) - } - }, { + apply: async () => assert.fail('reference-mode hosts must receive staged files, not inline apply'), + stageCloudContent: async (source: any) => { + await source.getInstruction() + staged.push(`${source.reference.item_id}:${source.reference.revision}`) + return registerCloudSyncStagedFile(source.reference, { native: true }, async () => {}) + }, + applyStagedCloudContent: async (change: any) => { + files.set(change.path, { path: change.path, kind: 'binary', content: { ...change.content_ref, data: '' } }) + }, + resolveStagedCloudConflict: async () => assert.fail('no conflicts expected') + } + const sync = new CloudSyncCoordinator('vault-1', createCloudSyncClient('https://example.test', 'test-only', { accountId: 'account' }), repository, { load: async () => state, save: async (next: any) => { state = structuredClone(next) } }, { @@ -58,18 +83,25 @@ it('downloads every manifest and change page without grouping large attachments await sync.sync() assert.equal(state.cursor, 3) assert.equal(files.size, 3) - assert.deepEqual(requests.filter((url) => url.pathname.endsWith('/manifest')).map((url) => url.searchParams.get('page')), ['1', '2', '3']) + assert.deepEqual(staged.sort(), ['asset-1:1', 'asset-2:1', 'asset-3:1']) + assert.deepEqual(downloads.sort(), staged.sort()) + // No request ever carried a body larger than metadata. + assert.ok(requests.every((url) => !url.pathname.includes('/revisions/') || url.pathname.endsWith('/download'))) - feed = items.map((item, index) => ({ - sequence: 4 + index, item_id: item.item_id, path: item.path, previous_path: item.path, - type: 'upsert', revision: 2, - content: { ...item.content, data: `updated-${index}`, sha256: `updated-hash-${index}` } - })) + feedRef.feed = items.map((item, index) => { + const bytes = Buffer.alloc(2_000_000, 100 + index) + const hash = createHash('sha256').update(bytes).digest('hex') + const content = { encoding: 'base64', sha256: hash, byte_length: bytes.length, media_type: 'image/jpeg' } + item.versions[2] = content + return { sequence: 4 + index, item_id: item.item_id, path: item.path, previous_path: item.path, type: 'upsert', revision: 2, + content_ref: { item_id: item.item_id, revision: 2, ...content } } + }) requests.length = 0 + staged.length = 0 await sync.sync() assert.equal(state.cursor, 6) - assert.deepEqual([...files.values()].map((item) => item.content.sha256), feed.map((change) => change.content.sha256)) - assert.deepEqual(requests.slice(0, 3).map((url) => url.searchParams.get('after')), ['3', '4', '5']) + assert.deepEqual([...files.values()].map((item) => item.content.sha256), feedRef.feed.map((change) => change.content_ref.sha256)) + assert.deepEqual(staged.sort(), ['asset-1:2', 'asset-2:2', 'asset-3:2']) }) it('keeps metadata-only manifest pagination intact', async () => { @@ -83,7 +115,7 @@ it('keeps metadata-only manifest pagination intact', async () => { } } } }) - const client = createCloudSyncClient('https://example.test', 'test-only') + const client = createCloudSyncClient('https://example.test', 'test-only', { accountId: 'account' }) await client.manifest('vault-1', { includeContent: false, page: 2, perPage: 250 }) await client.manifest('vault-1', { includeContent: false, perPage: 1 }) assert.equal(requests[0].searchParams.get('per_page'), '250') @@ -91,3 +123,21 @@ it('keeps metadata-only manifest pagination intact', async () => { assert.equal(requests[1].searchParams.get('per_page'), '1') assert.ok(requests.every((url) => url.searchParams.get('include_content') === 'false')) }) + +it('preserves the requested change page size and metadata-only content budget', async () => { + const requests: URL[] = [] + const { createCloudSyncClient } = await loadMobileModule('./src/bridge/cloud-sync-client.ts', { + '@capacitor/core': { + registerPlugin: () => ({}), + CapacitorHttp: { request: referenceServer([], { feed: [] }, requests, []) } + } + }) + const client = createCloudSyncClient('https://pagination.example.test', 'test-only', { accountId: 'account' }) + await client.negotiateContentReferences() + await client.changes('vault-1', 10, 250, { contentMode: 'references', maxInlineBytes: 0 }) + const request = requests.at(-1)! + assert.equal(request.searchParams.get('limit'), '250') + assert.equal(request.searchParams.get('after'), '10') + assert.equal(request.searchParams.get('max_inline_bytes'), '0') + assert.equal(request.searchParams.get('max_response_bytes'), '1048576') +}) diff --git a/src/bridge/cloud-sync-client.ts b/src/bridge/cloud-sync-client.ts index bdc2600..38b1a5e 100644 --- a/src/bridge/cloud-sync-client.ts +++ b/src/bridge/cloud-sync-client.ts @@ -1,6 +1,8 @@ import { CapacitorHttp, registerPlugin } from '@capacitor/core' import { CloudSyncApiClient, + cloudSyncRateLimits, + type CloudSyncResponseHeaders, type CloudSyncHttpRequest, type CloudSyncHttpTransport } from '@zennotes/shared-domain/cloud-sync-api' @@ -14,6 +16,21 @@ import { type MobileObjectUpload } from './mobile-direct-upload' +let requestLifetime = new AbortController() + +export function mobileCloudRequestSignal(): AbortSignal { + return requestLifetime.signal +} + +export function stopMobileCloudRequests(): void { + requestLifetime.abort() + cloudSyncRateLimits.cancelAll() +} + +export function resumeMobileCloudRequests(): void { + if (requestLifetime.signal.aborted) requestLifetime = new AbortController() +} + export class CloudServiceRequestError extends Error { readonly status: number readonly code: string | null @@ -23,7 +40,8 @@ export class CloudServiceRequestError extends Error { message: string, status: number, code: string | null, - details: Record | null = null + details: Record | null = null, + readonly headers: CloudSyncResponseHeaders = {} ) { super(message) this.name = 'CloudServiceRequestError' @@ -33,8 +51,9 @@ export class CloudServiceRequestError extends Error { } } -export function createCloudSyncClient(baseUrl: string, token: string): CloudSyncApiClient { +export function createCloudSyncClient(baseUrl: string, token: string, options: { accountId: string; signal?: AbortSignal }): CloudSyncApiClient { const normalizedBaseUrl = baseUrl.trim().replace(/\/+$/, '') + const lifetime = options.signal ?? requestLifetime.signal const transport: CloudSyncHttpTransport = { async request(request: CloudSyncHttpRequest): Promise { const multipart = request.body instanceof FormData @@ -71,7 +90,8 @@ export function createCloudSyncClient(baseUrl: string, token: string): CloudSync : `ZenNotes Cloud request failed (${response.status}).`), response.status, typeof error?.code === 'string' ? error.code : null, - isRecord(error?.details) ? error.details : null + isRecord(error?.details) ? error.details : null, + response.headers ?? {} ) } @@ -91,28 +111,31 @@ export function createCloudSyncClient(baseUrl: string, token: string): CloudSync } } - return new MobileCloudSyncApiClient(transport, uploadObject) + const loopback = /^http:\/\/(localhost|127\.\d+\.\d+\.\d+|\[::1\])(:\d+)?$/.test(normalizedBaseUrl) + return new MobileCloudSyncApiClient(cloudSyncRateLimits.wrap(transport, { + baseUrl: normalizedBaseUrl, accountId: options.accountId, signal: lifetime + }), async (request) => { + if (lifetime.aborted) throw new DOMException('Cloud request cancelled.', 'AbortError') + await uploadObject(request) + if (lifetime.aborted) throw new DOMException('Cloud request cancelled.', 'AbortError') + }, { + // Streaming host: large revisions arrive as references and are downloaded + // natively into staging, never as base64 through the WebView bridge. + contentReferences: true, + accountScope: { baseUrl: normalizedBaseUrl, accountId: options.accountId }, + signal: lifetime, + bootstrapContentPageBytes: 1024 * 1024, + allowInsecureLoopbackDownloads: loopback + }) } class MobileCloudSyncApiClient extends CloudSyncApiClient { constructor( http: CloudSyncHttpTransport, - private readonly uploadObject: MobileObjectUpload + private readonly uploadObject: MobileObjectUpload, + options: ConstructorParameters[1] ) { - super(http) - } - - // Capacitor copies JSON through Java and the WebView. A page containing - // several near-limit attachments can exhaust Android's native heap. - override manifest( - vaultId: string, - options: { includeContent?: boolean; page?: number; perPage?: number } = {} - ) { - return super.manifest(vaultId, options.includeContent ? { ...options, perPage: 1 } : options) - } - - override changes(vaultId: string, after: number, _limit = 100) { - return super.changes(vaultId, after, 1) + super(http, options) } override async mutate( @@ -137,7 +160,9 @@ const DirectUpload = registerPlugin<{ put(options: { url: string headers: Record - base64: string + base64?: string + uri?: string + sha256?: string byteLength: number }): Promise<{ status: number }> }>('ZenDirectUpload') @@ -149,6 +174,8 @@ const uploadObject: MobileObjectUpload = async (request) => { url: request.url, headers: request.headers, base64: request.base64, + uri: request.uri, + sha256: request.sha256, byteLength: request.byteLength }) } catch { diff --git a/src/bridge/cloud-sync-repository.test.ts b/src/bridge/cloud-sync-repository.test.ts index 292c82e..a343387 100644 --- a/src/bridge/cloud-sync-repository.test.ts +++ b/src/bridge/cloud-sync-repository.test.ts @@ -9,7 +9,9 @@ import type { CloudSyncRepository } from '@zennotes/shared-domain/cloud-sync-coo import { loadMobileModule } from '../../tooling/load-mobile-module.ts' -const { CachedCloudSyncRepository } = await loadMobileModule('./src/bridge/cloud-sync-repository') +const { CachedCloudSyncRepository, mutateWithMobileDirectUploads } = await loadMobileModule([ + './src/bridge/cloud-sync-repository', './src/bridge/mobile-direct-upload' +]) const { CloudSyncCoordinator } = await loadMobileModule('@zennotes/shared-domain/cloud-sync-coordinator') type StoredFile = { bytes: Buffer; mtime: number } @@ -21,8 +23,13 @@ function harness(initial: Record = { 'note.md': 'Hello' let state: CloudSyncState | null = null let clock = 1000 const reads: string[] = [] + const writes: string[] = [] + const copies: string[] = [] const failures = { cacheRead: false, cacheWrite: false, stateRead: false, directory: false, file: false } let onRead: ((path: string) => void) | undefined + let onWrite: ((path: string) => void) | undefined + let onRename: ((from: string, to: string) => void) | undefined + let onCopy: ((from: string, to: string) => void) | undefined const put = (path: string, body: string | Buffer) => { files.set(path, { bytes: Buffer.from(body), mtime: ++clock }) } @@ -55,23 +62,56 @@ function harness(initial: Record = { 'note.md': 'Hello' onRead?.(path) const file = files.get(path) if (!file) throw new Error('File missing') + if (file.bytes.length > 5 * 1024 * 1024) throw new Error('Whole-file bridge read exceeded the inline limit') return file.bytes.toString('base64') } const fs = { readdir, stat: async (path: string) => (await stat(path))?.type ?? null, readBase64, - writeText: async (path: string, data: string) => put(path, data), - writeBase64: async (path: string, data: string) => put(path, Buffer.from(data, 'base64')), - deleteFile: async (path: string) => { files.delete(path) }, + writeText: async (path: string, data: string) => { writes.push(path); put(path, data); onWrite?.(path) }, + writeBase64: async (path: string, data: string) => { writes.push(path); put(path, Buffer.from(data, 'base64')); onWrite?.(path) }, + deleteFile: async (path: string) => { writes.push(path); files.delete(path) }, rename: async (from: string, to: string) => { const file = files.get(from) if (!file) throw new Error('File missing') + writes.push(to) files.set(to, file) files.delete(from) + onRename?.(from, to) + } + } + const native = { + readdirStrict: readdir, readBase64, statOrNull: stat, stat, + async copyForSync(from: string, to: string) { + copies.push(from) + const file = files.get(from) + if (!file) throw new Error('File missing') + writes.push(to) + put(to, file.bytes) + onCopy?.(from, to) + }, + async readForSync(path: string, textCandidate: boolean) { + reads.push(path) + if (failures.file) throw new Error('File unavailable') + onRead?.(path) + const file = files.get(path) + if (!file) throw new Error('File missing') + let utf8 = false + try { + if (textCandidate) { + new TextDecoder('utf-8', { fatal: true }).decode(file.bytes) + utf8 = true + } + } catch {} + return { + uri: `file:///vault/${path}`, + sha256: createHash('sha256').update(file.bytes).digest('hex'), + byteLength: file.bytes.length, utf8, + ...(file.bytes.length <= 5 * 1024 * 1024 ? { inlineBase64: file.bytes.toString('base64') } : {}) + } } } - const native = { readdirStrict: readdir, readBase64, statOrNull: stat, stat } const store = { loadTracked: async () => { if (failures.stateRead) throw new Error('State unavailable') @@ -97,10 +137,10 @@ function harness(initial: Record = { 'note.md': 'Hello' }])) } } - const coordinator = () => { + const coordinator = (manifestItems: unknown[] = [], syncRepository = repository) => { const mutations: CloudSyncMutation[] = [] const remote = { - manifest: async () => ({ data: [], cursor: state?.cursor ?? 0, next_page: null }), + manifest: async () => ({ data: manifestItems, cursor: state?.cursor ?? 0, next_page: null }), changes: async () => ({ data: [], cursor: state?.cursor ?? 0, has_more: false }), mutate: async (_vaultId: string, body: { mutations: CloudSyncMutation[] }) => { // Serialization is deliberately real: a cache placeholder must never be uploaded. @@ -116,15 +156,18 @@ function harness(initial: Record = { 'note.md': 'Hello' let id = 0 return { mutations, - service: new CloudSyncCoordinator('vault-1', remote, repository, { + service: new CloudSyncCoordinator('vault-1', remote, syncRepository, { load: async () => state, save: async (next: CloudSyncState) => { state = structuredClone(next) } }, { itemId: () => `new-${++id}`, operationId: () => `op-${++id}` }) } } return { - files, reads, failures, repository, acknowledge, coordinator, put, + files, reads, writes, copies, failures, repository, acknowledge, coordinator, put, fs, native, store, setReadHook: (hook: typeof onRead) => { onRead = hook }, + setWriteHook: (hook: typeof onWrite) => { onWrite = hook }, + setRenameHook: (hook: typeof onRename) => { onRename = hook }, + setCopyHook: (hook: typeof onCopy) => { onCopy = hook }, get cache() { return cache }, set cache(next: unknown) { cache = next }, get state() { return state }, set state(next: CloudSyncState | null) { state = next } } @@ -147,6 +190,53 @@ function pending(path: string, local: CloudSyncContent, cloud = content('Other d } describe('cached mobile Cloud scan', () => { + it('scans a large attachment without materializing its contents across the bridge', async () => { + const bytes = Buffer.alloc(8_000_000, 129) + const h = harness({ 'attachements/large.bin': bytes }) + const [item] = await h.repository.scan() + assert.equal(item.content.byte_length, bytes.length) + assert.equal(item.content.sha256, createHash('sha256').update(bytes).digest('hex')) + assert.equal(item.content.data, '') + assert.equal(item.kind, 'binary') + assert.ok(JSON.stringify(item).length < 500) + }) + + it('preserves UTF-8 classification and raw-byte hashes for file-backed text', async () => { + const bytes = Buffer.from('日本語 café\n'.repeat(400_000)) + const h = harness({ 'large.md': bytes }) + const [item] = await h.repository.scan() + assert.equal(item.kind, 'text') + assert.equal(item.content.encoding, 'utf8') + assert.equal(item.content.sha256, createHash('sha256').update(bytes).digest('hex')) + assert.equal(item.content.data, '') + }) + + it('passes a scanned file to the uploader by URI and completes only after the native transfer', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'attachements/large.bin': bytes }) + const [item] = await h.repository.scan() + let uploaded = false + const result = await mutateWithMobileDirectUploads({ + mutate: async () => { throw new Error('Unexpected inline upload') }, + initiateUpload: async (_vault: string, request: any) => ({ data: { + id: 'upload', operation_id: request.operation_id, expected_bytes: bytes.length, + upload: { method: 'PUT', url: 'https://storage.example.test/object', headers: {} } + } }), + completeUpload: async () => { + assert.equal(uploaded, true) + return { data: { result: { acknowledged: [{ item_id: 'item' }], conflicts: [], cursor: 1 } } } + }, + abortUpload: async () => { throw new Error('Unexpected abort') } + }, 'vault', { mutations: [{ ...item, type: 'upsert', item_id: 'item', operation_id: 'operation', base_revision: null }] }, async (request: any) => { + assert.equal(request.uri, 'file:///vault/attachements/large.bin') + assert.equal(request.base64, undefined) + assert.equal(request.sha256, createHash('sha256').update(bytes).digest('hex')) + assert.equal(request.byteLength, bytes.length) + uploaded = true + }) + assert.equal(result.acknowledged.length, 1) + }) + it('reads and hashes new text and binary files with the same portable semantics', async () => { const bytes = Buffer.from([0, 255, 1, 128]) const h = harness({ 'note.md': 'Hello', 'assets/photo.png': bytes, '.zennotes/cache.json': '{}' }) @@ -294,6 +384,220 @@ describe('cached mobile Cloud scan', () => { }) describe('cached scan with the pinned conflict coordinator', () => { + it('keeps all 6 MB of the local version while replacing the original with Cloud bytes', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + const [local] = await h.repository.scan() + h.acknowledge([local]) + const cloud = content('Cloud replacement') + h.state!.pending_conflicts = { 'conflict-1': pending('asset.bin', local.content, cloud) } + const coordinator = h.coordinator([{ + item_id: 'item-0', path: 'asset.bin', kind: 'text', revision: 2, + sha256: cloud.sha256, byte_length: cloud.byte_length, media_type: cloud.media_type + }]) + + await coordinator.service.resolveConflict({ + conflict_id: 'conflict-1', choice: 'both', keep_both_path: 'copies/local.bin', + expected_local_sha256: local.content.sha256, expected_cloud_revision: 2 + }) + + assert.deepEqual(h.files.get('copies/local.bin')?.bytes, bytes) + assert.equal(h.files.get('asset.bin')?.bytes.toString(), cloud.data) + assert.equal(h.state!.pending_conflicts?.['conflict-1'], undefined) + assert.deepEqual([...h.files.keys()].sort(), ['asset.bin', 'copies/local.bin']) + assert.ok(h.copies.length > 0) + }) + + it('rejects a changed file-backed source before writing any resolution files', async () => { + const h = harness({ 'source.bin': Buffer.alloc(6_000_000, 197), 'note.md': 'Keep me' }) + const source = (await h.repository.scan()).find((item) => item.path === 'source.bin')! + h.put('source.bin', Buffer.alloc(6_000_000, 198)) + await assert.rejects(h.repository.applyConflictResolutionFiles!({ + expected_path: 'note.md', expected_sha256: content('Keep me').sha256, + files: [{ path: 'first.md', content: content('First') }, { path: 'copy.bin', content: source.content }] + }), /changed|source/i) + assert.deepEqual(h.writes, []) + assert.equal(h.files.get('note.md')?.bytes.toString(), 'Keep me') + }) + + it('preserves raw UTF-8 bytes when a large local text version is copied', async () => { + const bytes = Buffer.from('日本語é\n'.repeat(500_000)) + assert.equal(bytes.length, 6_000_000) + const h = harness({ 'large.md': bytes }) + const [local] = await h.repository.scan() + assert.equal(local.content.encoding, 'utf8') + await h.repository.applyConflictResolutionFiles!({ + expected_path: 'large.md', expected_sha256: local.content.sha256, + files: [{ path: 'large.md', content: content('Cloud') }, { path: 'local.md', content: local.content }] + }) + assert.deepEqual(h.files.get('local.md')?.bytes, bytes) + assert.equal(h.files.get('large.md')?.bytes.toString(), 'Cloud') + }) + + it('does not upload a deletion when an interrupted replacement left a rollback file', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + h.acknowledge(await h.repository.scan()) + const rollback = '.zennotes/sync/rollback-interrupted.bin' + h.files.delete('asset.bin') + h.put(rollback, bytes) + const restarted = new CachedCloudSyncRepository(h.fs, h.native, h.store) + const coordinator = h.coordinator([], restarted) + await assert.rejects(coordinator.service.sync(), /needs recovery/) + assert.deepEqual(coordinator.mutations, []) + assert.deepEqual(h.files.get(rollback)?.bytes, bytes) + }) + + it('rechecks recovery files on the next scan of an already-running repository', async () => { + const h = harness({ 'note.md': 'original' }) + h.acknowledge(await h.repository.scan()) + h.files.delete('note.md') + h.put('.zennotes/sync/rollback-failed.md', 'original') + const coordinator = h.coordinator() + await assert.rejects(coordinator.service.sync(), /needs recovery/) + assert.deepEqual(coordinator.mutations, []) + }) + + it('rejects unknown metadata-only content before a bootstrap rename or a multi-file write', async () => { + const h = harness({ 'note.md': 'Keep me' }) + const missing = { ...content('Missing bytes'), data: '' } + await assert.rejects(h.repository.resolveBootstrapConflict!({ + path: 'note.md', expectedLocalSha256: content('Keep me').sha256, cloudContent: missing, + resolution: { choice: 'both', keep_both_path: 'copy.md', conflict: { + code: 'BOOTSTRAP_CONTENT_CONFLICT', item_id: 'item', path: 'note.md', + local_sha256: content('Keep me').sha256, remote_sha256: missing.sha256 + } } + }), /source|bytes|content/i) + await assert.rejects(h.repository.applyConflictResolutionFiles!({ + expected_path: 'note.md', expected_sha256: content('Keep me').sha256, + files: [{ path: 'first.md', content: content('First') }, { path: 'note.md', content: missing }] + }), /source|bytes|content/i) + assert.deepEqual(h.writes, []) + assert.deepEqual([...h.files.keys()], ['note.md']) + }) + + it('rolls back a failed Cloud replacement after creating the large local copy', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + const [local] = await h.repository.scan() + let failed = false + h.setRenameHook((_from, to) => { + if (to === 'asset.bin' && !failed) { + failed = true + h.put(to, 'Partial write') + throw new Error('Native replacement failed after modifying the target') + } + }) + await assert.rejects(h.repository.applyConflictResolutionFiles!({ + expected_path: 'asset.bin', expected_sha256: local.content.sha256, + files: [{ path: 'asset.bin', content: content('Cloud') }, { path: 'copy.bin', content: local.content }] + }), /failed/) + assert.equal(failed, true) + assert.deepEqual(h.files.get('asset.bin')?.bytes, bytes) + assert.deepEqual([...h.files.keys()], ['asset.bin']) + }) + + it('rejects a corrupt native copy before replacing the original', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + const [local] = await h.repository.scan() + h.setCopyHook((_from, to) => h.put(to, 'Truncated copy')) + await assert.rejects(h.repository.applyConflictResolutionFiles!({ + expected_path: 'asset.bin', expected_sha256: local.content.sha256, + files: [{ path: 'copy.bin', content: local.content }, { path: 'asset.bin', content: content('Cloud') }] + }), /bytes|hash|changed|verification/i) + assert.deepEqual(h.files.get('asset.bin')?.bytes, bytes) + assert.deepEqual([...h.files.keys()], ['asset.bin']) + }) + + it('retains a source edited during copying and does not publish the stale copy', async () => { + const h = harness({ 'source.bin': Buffer.alloc(6_000_000, 197), 'note.md': 'Keep me' }) + const source = (await h.repository.scan()).find((item) => item.path === 'source.bin')! + const changed = Buffer.alloc(6_000_000, 198) + h.setCopyHook((from) => h.put(from, changed)) + await assert.rejects(h.repository.applyConflictResolutionFiles!({ + expected_path: 'note.md', expected_sha256: content('Keep me').sha256, + files: [{ path: 'copy.bin', content: source.content }, { path: 'note.md', content: content('Cloud') }] + }), /changed/) + assert.deepEqual(h.files.get('source.bin')?.bytes, changed) + assert.equal(h.files.get('note.md')?.bytes.toString(), 'Keep me') + assert.deepEqual([...h.files.keys()].sort(), ['note.md', 'source.bin']) + }) + + it('preserves the large original when staging a replacement fails after a partial write', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + const [local] = await h.repository.scan() + h.setWriteHook((path) => { + h.put(path, 'Partial') + throw new Error('Staging failed') + }) + await assert.rejects(h.repository.replaceConflictFile!({ + path: local.path, expectedSha256: local.content.sha256, content: content('Cloud') + }), /Staging failed/) + assert.deepEqual(h.files.get('asset.bin')?.bytes, bytes) + assert.deepEqual([...h.files.keys()], ['asset.bin']) + }) + + it('keeps a large bootstrap local copy without reading its body across the bridge', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + const [local] = await h.repository.scan() + const cloud = content('Cloud') + await h.repository.resolveBootstrapConflict!({ + path: local.path, expectedLocalSha256: local.content.sha256, cloudContent: cloud, + resolution: { choice: 'both', keep_both_path: 'local.bin', conflict: { + code: 'BOOTSTRAP_CONTENT_CONFLICT', item_id: 'item', path: local.path, + local_sha256: local.content.sha256, remote_sha256: cloud.sha256 + } } + }) + assert.deepEqual(h.files.get('local.bin')?.bytes, bytes) + assert.equal(h.files.get('asset.bin')?.bytes.toString(), 'Cloud') + }) + + it('rolls back the bootstrap rename when the Cloud write fails', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + const [local] = await h.repository.scan() + const cloud = content('Cloud') + h.setWriteHook(() => { throw new Error('Write failed') }) + await assert.rejects(h.repository.resolveBootstrapConflict!({ + path: local.path, expectedLocalSha256: local.content.sha256, cloudContent: cloud, + resolution: { choice: 'both', keep_both_path: 'local.bin', conflict: { + code: 'BOOTSTRAP_CONTENT_CONFLICT', item_id: 'item', path: local.path, + local_sha256: local.content.sha256, remote_sha256: cloud.sha256 + } } + }), /Write failed/) + assert.deepEqual(h.files.get('asset.bin')?.bytes, bytes) + assert.deepEqual([...h.files.keys()], ['asset.bin']) + }) + + it('rejects unknown metadata on inherited apply and replace before changing files', async () => { + const h = harness({ 'note.md': 'Keep me' }) + const missing = { ...content('Unavailable'), data: '' } + await assert.rejects(h.repository.apply({ sequence: 2, revision: 2, item_id: 'item', + path: 'note.md', previous_path: null, type: 'upsert', content: missing }, undefined), /source/) + await assert.rejects(h.repository.replaceConflictFile!({ + path: 'note.md', expectedSha256: content('Keep me').sha256, content: missing + }), /source/) + assert.deepEqual(h.writes, []) + assert.equal(h.files.get('note.md')?.bytes.toString(), 'Keep me') + }) + + it('uses streaming reads for inherited apply and preserves an unsynced large edit', async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const h = harness({ 'asset.bin': bytes }) + const [local] = await h.repository.scan() + const change = { sequence: 2, item_id: 'item', revision: 2, type: 'upsert' as const, + path: 'asset.bin', previous_path: null, content: content('Cloud') } + const conflict = await h.repository.apply(change, undefined) + assert.equal(conflict?.code, 'LOCAL_EDIT_CONFLICT') + assert.deepEqual(h.files.get('asset.bin')?.bytes, bytes) + await h.repository.apply(change, { item_id: 'item', path: local.path, kind: local.kind, revision: 1, + sha256: local.content.sha256, byte_length: bytes.length, media_type: local.content.media_type }) + assert.equal(h.files.get('asset.bin')?.bytes.toString(), 'Cloud') + }) + it('returns real bytes for review when pending local content matches acknowledged content', async () => { const h = harness() h.acknowledge(await h.repository.scan()) diff --git a/src/bridge/cloud-sync-repository.ts b/src/bridge/cloud-sync-repository.ts index cf5063f..981aa92 100644 --- a/src/bridge/cloud-sync-repository.ts +++ b/src/bridge/cloud-sync-repository.ts @@ -37,6 +37,9 @@ import { import type { CloudSyncLocalItem, CloudSyncState } from '@zennotes/shared-domain/cloud-sync-engine' import type { NativeFs } from './native-fs' import { cloudSyncWorkBudget, decodeCloudSyncBase64 } from './cloud-sync-work' +import { CLOUD_SYNC_INLINE_UPLOAD_LIMIT_BYTES, rememberMobileUploadSource } from './mobile-direct-upload' +import { isNotFoundError } from './fs-errors' +import { NativeCloudStaging, type CloudStagingNative } from './cloud-sync-staging' export interface ScanCacheEntry { mtime: number @@ -49,6 +52,8 @@ export interface ScanCacheEntry { export type ScanCache = Record +type CloudSyncNativeFiles = Pick & Partial + export interface ScanCacheStore { loadTracked(): Promise loadCache(): Promise @@ -56,16 +61,51 @@ export interface ScanCacheStore { } export class CachedCloudSyncRepository extends PortableCloudSyncRepository { + private readonly contentFiles: NativeCloudSyncContent + constructor( fs: PortableCloudSyncFileSystem, - private readonly native: Pick, + private readonly native: CloudSyncNativeFiles, private readonly store: ScanCacheStore, private readonly onChanged: () => void = () => {} ) { - super(fs) + const contentFiles = new NativeCloudSyncContent(fs, native) + const staging = isStagingNative(native) + ? new NativeCloudStaging(native, (path) => TEXT_EXTENSIONS.has(extension(path)), (path) => contentFiles.read(path)) + : null + const sourceAwareFs = { + ...(staging ? { + stageCloudContent: async (source: Parameters[0]) => { + await contentFiles.assertReady(true) + return staging.stage(source) + }, + applyStagedCloudContent: (...args: Parameters) => staging.apply(...args), + resolveStagedCloudConflict: (input: Parameters[0]) => staging.resolve(input) + } : {}), + readdir: (path: string) => fs.readdir(path), + stat: (path: string) => fs.stat(path), + readBase64: (path: string) => fs.readBase64(path), + writeText: (path: string, value: string) => fs.writeText(path, value), + writeBase64: (path: string, value: string) => fs.writeBase64(path, value), + deleteFile: (path: string) => fs.deleteFile(path), + rename: (from: string, to: string) => fs.rename(from, to), + readItem: (path: string) => contentFiles.read(path), + validateContent: (content: CloudSyncContent) => contentFiles.validate(content), + writeContent: (path: string, content: CloudSyncContent) => contentFiles.write(path, content) + } + super(sourceAwareFs) + this.contentFiles = contentFiles + this.staging = staging + } + + private readonly staging: NativeCloudStaging | null + + override async matchesCloudContent(path: string, reference: Parameters[1]): Promise { + return this.staging ? this.staging.matches(path, reference) : super.matchesCloudContent(path, reference) } override async scan(): Promise { + await this.contentFiles.assertReady(true) const trackedSha = trackedShaByPath(await this.store.loadTracked().catch(() => null)) const cache = normalizeScanCache(await this.store.loadCache().catch(() => null)) const nextCache: ScanCache = {} @@ -113,7 +153,7 @@ export class CachedCloudSyncRepository extends PortableCloudSyncRepository { continue } - const item = await this.readItemFresh(path) + const item = await this.readItemFresh(path, entry.uri) if (!cached || cached.mtime !== entry.mtime || cached.size !== entry.size || cached.sha256 !== item.content.sha256) { this.onChanged() } @@ -135,12 +175,53 @@ export class CachedCloudSyncRepository extends PortableCloudSyncRepository { } // --------------------------------------------------------------------- - // Preserve upstream readItem's encoding/hash semantics while yielding - // during large base64 decoding and avoiding a binary re-encode. + // Only inline-sized bodies cross the native bridge. Large files retain a + // host-only source reference, like the desktop's disk-backed uploader. // --------------------------------------------------------------------- - private async readItemFresh(path: string): Promise { - const { bytes, base64 } = await decodeCloudSyncBase64(await this.native.readBase64(path)) + private async readItemFresh(path: string, uri?: string): Promise { + return this.contentFiles.read(path, uri) + } +} + +/** Large local content remains an identity-bound source, never an empty payload. */ +class NativeCloudSyncContent { + private readonly sources = new WeakMap() + private recoveryChecked = false + + constructor(private readonly fs: PortableCloudSyncFileSystem, private readonly native: CloudSyncNativeFiles) {} + + async assertReady(force = false): Promise { + if (this.recoveryChecked && !force) return + const entries = await this.native.readdirStrict('.zennotes/sync').catch((error) => { + if (!isNotFoundError(error)) throw error + return [] + }) + const recovery = entries.find((entry) => entry.name.startsWith('rollback-')) + if (recovery) { + // Never turn a process-interrupted rename into a new local deletion. + throw new Error(`Cloud sync needs recovery of .zennotes/sync/${recovery.name} before it can continue.`) + } + this.recoveryChecked = true + } + + async read(path: string, uri?: string): Promise { + await this.assertReady() + const file = await this.native.readForSync(path, TEXT_EXTENSIONS.has(extension(path)), uri) + if (file.byteLength > CLOUD_SYNC_INLINE_UPLOAD_LIMIT_BYTES) { + const content = rememberMobileUploadSource({ + encoding: file.utf8 ? 'utf8' : 'base64', data: '', sha256: file.sha256, + byte_length: file.byteLength, media_type: mediaType(path, file.utf8) + }, file.uri) + this.sources.set(content, { path, hash: file.sha256, bytes: file.byteLength }) + return { + path, + kind: file.utf8 ? 'text' : 'binary', + content + } + } + if (file.inlineBase64 === undefined) throw new Error('Native file inspection omitted inline content.') + const { bytes, base64 } = await decodeCloudSyncBase64(file.inlineBase64) const text = decodeText(path, bytes) return { path, @@ -148,12 +229,89 @@ export class CachedCloudSyncRepository extends PortableCloudSyncRepository { content: { encoding: text === null ? 'base64' : 'utf8', data: text === null ? base64 : text, - sha256: await sha256(bytes), + sha256: file.sha256, byte_length: bytes.byteLength, media_type: mediaType(path, text !== null) } } } + + async validate(content: CloudSyncContent): Promise { + if (content.encoding !== 'utf8' && content.encoding !== 'base64') { + throw new Error('Encrypted cloud sync content must be decrypted before filesystem apply') + } + if (content.data !== '' || content.byte_length === 0) return + const source = this.sources.get(content) + if (!source || source.hash !== content.sha256 || source.bytes !== content.byte_length) { + throw new Error('This file content has no recognized source bytes. Sync again and retry.') + } + await this.verify(source.path, content) + } + + private async verify(path: string, content: CloudSyncContent): Promise { + // Resolve the current path again: a saved document URI may name a replaced file. + const file = await this.native.readForSync(path, false) + if (file.sha256 !== content.sha256 || file.byteLength !== content.byte_length) { + throw new Error('The file source changed or its copy failed byte verification. Sync again and retry.') + } + } + + async write(path: string, content: CloudSyncContent): Promise { + await this.assertReady() + await this.validate(content) + const kind = await this.fs.stat(path) + if (kind === 'directory') throw new Error('The destination is a directory.') + const before = kind === 'file' ? await this.read(path) : null + const temporary = `.zennotes/sync/write-${crypto.randomUUID()}${extension(path)}` + const backup = `.zennotes/sync/rollback-${crypto.randomUUID()}${extension(path)}` + let publishing = false + let completed = false + try { + const source = this.sources.get(content) + if (content.data === '' && content.byte_length > 0 && source) { + await this.native.copyForSync(source.path, temporary, content.byte_length) + await this.validate(content) + } else if (content.encoding === 'utf8') { + await this.fs.writeText(temporary, content.data) + } else { + await this.fs.writeBase64(temporary, content.data) + } + await this.verify(temporary, content) + if (before) { + await this.verify(path, before.content) + await this.fs.rename(path, backup) + } else if (await this.fs.stat(path) !== null) { + throw new Error('The destination changed before its Cloud write.') + } + publishing = true + await this.fs.rename(temporary, path) + await this.verify(path, content) + completed = true + } catch (error) { + // Native operations may modify the destination and then reject. The + // original is kept separately until the published bytes are verified. + if (await this.fs.stat(backup) === 'file') { + try { + if (await this.fs.stat(path) !== null) await this.fs.deleteFile(path) + await this.fs.rename(backup, path) + } catch (rollbackError) { + throw new Error(`Cloud write rollback failed; the original is preserved at ${backup}.`, { cause: rollbackError }) + } + } else if (!before && publishing && await this.fs.stat(path) !== null) { + await this.fs.deleteFile(path) + } + throw error + } finally { + this.recoveryChecked = false + if (await this.fs.stat(temporary) === 'file') await this.fs.deleteFile(temporary).catch(() => {}) + if (completed) await this.fs.deleteFile(backup).catch(() => {}) + } + } +} + +function isStagingNative(native: CloudSyncNativeFiles): native is CloudSyncNativeFiles & CloudStagingNative { + return ['statVerified', 'download', 'rename', 'deleteFile', 'mkdir', 'readText'] + .every((method) => typeof (native as Record)[method] === 'function') } function itemFromCache(path: string, cached: ScanCacheEntry): CloudSyncLocalItem { @@ -288,10 +446,3 @@ function extension(path: string): string { function mediaType(path: string, text: boolean): string { return MEDIA_TYPES[extension(path)] ?? (text ? 'text/plain' : 'application/octet-stream') } - -async function sha256(bytes: Uint8Array): Promise { - const digest = await crypto.subtle.digest('SHA-256', bytes.buffer) - return [...new Uint8Array(digest)] - .map((byte) => byte.toString(16).padStart(2, '0')) - .join('') -} diff --git a/src/bridge/cloud-sync-staging.test.ts b/src/bridge/cloud-sync-staging.test.ts new file mode 100644 index 0000000..6054bce --- /dev/null +++ b/src/bridge/cloud-sync-staging.test.ts @@ -0,0 +1,247 @@ +import assert from 'node:assert/strict' +import { createHash } from 'node:crypto' +import { describe, it } from 'node:test' +import { loadMobileModule } from '../../tooling/load-mobile-module.ts' + +// One bundle: the staging token registry is a module-private WeakMap, so the +// helper that reads a handle must share the module instance that wrote it. +const { NativeCloudStaging, cloudSyncStagedHandle } = await loadMobileModule([ + './src/bridge/cloud-sync-staging', '@zennotes/shared-domain/cloud-sync-content' +]) + +const hash = (bytes: Buffer) => createHash('sha256').update(bytes).digest('hex') + +/** In-memory vault with a fake native downloader; bytes never pass through JS strings. */ +function harness(initial: Record = {}) { + const files = new Map(Object.entries(initial)) + const objects = new Map() + const log: string[] = [] + const native = { + async statVerified(path: string) { return files.has(path) ? 'file' as const : null }, + async readForSync(path: string) { + const bytes = files.get(path) + if (!bytes) throw new Error(`missing ${path}`) + let utf8 = true + try { new TextDecoder('utf-8', { fatal: true }).decode(bytes) } catch { utf8 = false } + return { uri: `file:///${path}`, sha256: hash(bytes), byteLength: bytes.length, utf8 } + }, + async download(options: { url: string; to: string; byteLength: number; sha256: string }) { + log.push(`download ${options.url}`) + const bytes = objects.get(options.url) + if (!bytes) throw Object.assign(new Error('not found'), { status: 404 }) + if (bytes.length !== options.byteLength || hash(bytes) !== options.sha256) throw new Error('verification failed') + files.set(options.to, bytes) + }, + async copyForSync(from: string, to: string) { + log.push(`copy ${from} -> ${to}`) + files.set(to, Buffer.from(files.get(from)!)) + }, + async rename(from: string, to: string) { + log.push(`rename ${from} -> ${to}`) + const bytes = files.get(from) + if (!bytes) throw new Error(`rename missing ${from}`) + files.set(to, bytes) + files.delete(from) + }, + async deleteFile(path: string) { files.delete(path) }, + async mkdir() {}, + async readText(path: string) { return files.get(path)!.toString('utf8') } + } + const staging = new NativeCloudStaging(native, (path: string) => path.endsWith('.md'), async (path: string) => { + const bytes = files.get(path)! + return { path, kind: 'text', content: { encoding: 'utf8', data: bytes.toString('utf8'), + sha256: hash(bytes), byte_length: bytes.length, media_type: 'text/markdown' } } + }) + const reference = (itemId: string, revision: number, bytes: Buffer, encoding = 'base64') => ({ + item_id: itemId, revision, encoding, sha256: hash(bytes), byte_length: bytes.length, media_type: 'application/octet-stream' + }) + const source = (ref: any, url: string) => ({ + reference: ref, previewLimitBytes: 262_144, allowInsecureLoopback: false, + getInstruction: async () => ({ item_id: ref.item_id, revision: ref.revision, + content: { encoding: ref.encoding, sha256: ref.sha256, byte_length: ref.byte_length, media_type: ref.media_type }, + download: { url, method: 'GET', headers: {}, expires_at: new Date(Date.now() + 300_000).toISOString() } }) + }) + return { files, objects, log, staging, reference, source, native } +} + +describe('native Cloud staging', () => { + for (const failure of ['copy', 'publish', 'restore', 'cancel'] as const) { + it(`preserves the original and cleans partial writes after ${failure} failure`, async () => { + const local = Buffer.from('local') + const h = harness({ 'note.md': local }) + const cloud = Buffer.from('cloud') + const ref = h.reference('item', 2, cloud, 'utf8') + const url = 'https://objects.example.test/item/2' + h.objects.set(url, cloud) + const controller = new AbortController() + const file = await h.staging.stage({ ...h.source(ref, url), signal: controller.signal }) + const rename = h.native.rename + if (failure === 'copy') h.native.copyForSync = async (_from, to) => { + h.files.set(to, Buffer.from('partial')) + throw new Error('copy failed') + } + if (failure === 'publish' || failure === 'restore') h.native.rename = async (from, to) => { + if (failure === 'restore' && from.startsWith('.zennotes/sync/rollback-')) throw new Error('restore failed') + await rename(from, to) + if (from.startsWith('.zennotes/sync/downloads/')) { + h.files.set(to, Buffer.from('corrupt')) + if (failure === 'restore') throw new Error('publish failed') + } + } + if (failure === 'cancel') controller.abort() + await assert.rejects(h.staging.resolve({ + expected_path: 'note.md', expected_sha256: hash(local), cloud_path: 'note.md', file, + ...(failure === 'copy' ? { keep_both_path: 'local.md' } : {}) + })) + assert.equal(h.files.has('local.md'), false) + if (failure === 'restore') { + const recovery = [...h.files.keys()].find((path) => /^\.zennotes\/sync\/rollback-[^/]+$/.test(path)) + assert.ok(recovery, 'stranded original must be visible to the scan recovery guard') + assert.deepEqual(h.files.get(recovery), local) + } else { + assert.deepEqual(h.files.get('note.md'), local) + } + }) + } + + it('removes a newly created destination when native rename partially publishes then rejects', async () => { + const h = harness() + const bytes = Buffer.from('cloud') + const ref = h.reference('item', 1, bytes) + const url = 'https://objects.example.test/item/1' + h.objects.set(url, bytes) + const file = await h.staging.stage(h.source(ref, url)) + h.native.rename = async (_from, to) => { + h.files.set(to, Buffer.from('partial')) + throw new Error('rename failed') + } + await assert.rejects(h.staging.apply({ sequence: 1, item_id: 'item', type: 'upsert', + path: 'new.bin', previous_path: null, revision: 1, content_ref: ref }, undefined, file)) + assert.equal(h.files.has('new.bin'), false) + }) + + for (const failure of ['publish', 'edited-copy', 'rollback'] as const) { + it(`keeps a failed keep-both resolution safe and retryable after ${failure}`, async () => { + const local = Buffer.from('local') + const cloud = Buffer.from('cloud') + const h = harness({ 'note.md': local }) + const ref = h.reference('item', 2, cloud, 'utf8') + const url = 'https://objects.example.test/item/2' + h.objects.set(url, cloud) + const file = await h.staging.stage(h.source(ref, url)) + const rename = h.native.rename + h.native.rename = async (from, to) => { + if (from.startsWith('.zennotes/sync/downloads/') && to === 'note.md') { + if (failure === 'edited-copy') h.files.set('local.md', Buffer.from('later user edit')) + throw new Error('disk full') + } + if (failure === 'rollback' && from.startsWith('.zennotes/sync/rollback-')) throw new Error('rollback failed') + await rename(from, to) + } + const input = { expected_path: 'note.md', expected_sha256: hash(local), cloud_path: 'note.md', file, keep_both_path: 'local.md' } + await assert.rejects(h.staging.resolve(input)) + if (failure === 'publish') { + assert.deepEqual(h.files.get('note.md'), local) + assert.equal(h.files.has('local.md'), false, 'the unfinished decision must not block its own retry') + h.native.rename = rename + await h.staging.resolve(input) + assert.deepEqual(h.files.get('note.md'), cloud) + assert.deepEqual(h.files.get('local.md'), local) + } else { + assert.deepEqual(h.files.get('local.md'), failure === 'edited-copy' ? Buffer.from('later user edit') : local) + } + }) + } + + it('refreshes an expired signed URL reported as a native Capacitor rejection', async () => { + const h = harness() + const bytes = Buffer.from('cloud') + const ref = h.reference('item', 1, bytes) + const url = 'https://objects.example.test/item/1' + h.objects.set(url, bytes) + const download = h.native.download + let attempts = 0 + h.native.download = async (options) => { + if (++attempts === 1) { + h.files.set(options.to, Buffer.from('partial')) + throw Object.assign(new Error('expired'), { data: { status: 403 } }) + } + assert.equal(h.files.has(options.to), false) + await download(options) + } + await h.staging.stage(h.source(ref, url)) + assert.equal(attempts, 2) + }) + + it('downloads into staging, verifies, then publishes atomically', async () => { + const h = harness() + const bytes = Buffer.alloc(6_000_000, 42) + const ref = h.reference('item', 3, bytes) + h.objects.set('https://objects.example.test/item/3', bytes) + const file = await h.staging.stage(h.source(ref, 'https://objects.example.test/item/3')) + assert.equal(cloudSyncStagedHandle(file).owner, h.staging) + const staged = [...h.files.keys()].find((path) => path.startsWith('.zennotes/sync/downloads/'))! + assert.ok(staged, 'file staged in vault-private directory') + const conflict = await h.staging.apply( + { sequence: 9, item_id: 'item', type: 'upsert', path: 'attachements/big.bin', previous_path: null, revision: 3, content_ref: ref }, + undefined, file + ) + assert.equal(conflict, undefined) + assert.deepEqual(h.files.get('attachements/big.bin'), bytes) + assert.equal(h.files.has(staged), false, 'staging file was moved, not copied') + }) + + it('rejects a download whose bytes do not match the reference and leaves nothing behind', async () => { + const h = harness() + const bytes = Buffer.alloc(1000, 1) + const ref = h.reference('item', 1, bytes) + h.objects.set('https://objects.example.test/item/1', Buffer.alloc(1000, 2)) + await assert.rejects(h.staging.stage(h.source(ref, 'https://objects.example.test/item/1')), /verification/) + assert.equal([...h.files.keys()].some((path) => path.startsWith('.zennotes/sync/downloads/')), false) + }) + + it('preserves a local edit as a conflict instead of overwriting it', async () => { + const localEdit = Buffer.from('my local edit') + const h = harness({ 'note.md': localEdit }) + const cloud = Buffer.from('cloud version') + const ref = h.reference('item', 2, cloud, 'utf8') + h.objects.set('https://objects.example.test/item/2', cloud) + const file = await h.staging.stage(h.source(ref, 'https://objects.example.test/item/2')) + const tracked = { item_id: 'item', path: 'note.md', kind: 'text', revision: 1, sha256: hash(Buffer.from('original')), byte_length: 8, media_type: 'text/markdown' } + const conflict = await h.staging.apply( + { sequence: 5, item_id: 'item', type: 'upsert', path: 'note.md', previous_path: null, revision: 2, content_ref: ref }, + tracked, file + ) + assert.equal(conflict?.code, 'LOCAL_EDIT_CONFLICT') + assert.equal(conflict?.local?.content.data, 'my local edit') + assert.deepEqual(h.files.get('note.md'), localEdit, 'local bytes untouched') + }) + + it('keeps both versions with a byte-verified copy of the local file', async () => { + const local = Buffer.alloc(6_000_000, 7) + const h = harness({ 'attachements/big.bin': local }) + const cloud = Buffer.alloc(6_000_000, 9) + const ref = h.reference('item', 4, cloud) + h.objects.set('https://objects.example.test/item/4', cloud) + const file = await h.staging.stage(h.source(ref, 'https://objects.example.test/item/4')) + await h.staging.resolve({ + expected_path: 'attachements/big.bin', expected_sha256: hash(local), + cloud_path: 'attachements/big.bin', file, keep_both_path: 'attachements/big (local).bin' + }) + assert.deepEqual(h.files.get('attachements/big.bin'), cloud) + assert.deepEqual(h.files.get('attachements/big (local).bin'), local) + assert.ok(h.log.some((entry) => entry.startsWith('copy attachements/big.bin'))) + }) + + it('refuses to resolve when the local file changed since the conflict was recorded', async () => { + const h = harness({ 'note.md': Buffer.from('changed again') }) + const cloud = Buffer.from('cloud') + const ref = h.reference('item', 1, cloud, 'utf8') + h.objects.set('https://objects.example.test/item/1', cloud) + const file = await h.staging.stage(h.source(ref, 'https://objects.example.test/item/1')) + await assert.rejects(h.staging.resolve({ + expected_path: 'note.md', expected_sha256: hash(Buffer.from('stale')), cloud_path: 'note.md', file + }), /changed on this device/) + assert.equal(h.files.get('note.md')!.toString(), 'changed again') + }) +}) diff --git a/src/bridge/cloud-sync-staging.ts b/src/bridge/cloud-sync-staging.ts new file mode 100644 index 0000000..57e64bd --- /dev/null +++ b/src/bridge/cloud-sync-staging.ts @@ -0,0 +1,263 @@ +import type { CloudSyncChange, CloudSyncContentReference } from '@zennotes/bridge-contract/cloud-sync' +import { cloudSyncPathKey, normalizeCloudSyncPath, shouldSyncVaultPath } from '@zennotes/shared-domain/cloud-sync' +import { + cloudSyncStagedHandle, + registerCloudSyncStagedFile, + throwIfCloudSyncCancelled, + validateCloudSyncContentReference, + validateCloudSyncDownloadInstruction, + type CloudSyncDownloadSource, + type CloudSyncStagedConflict, + type CloudSyncStagedFile +} from '@zennotes/shared-domain/cloud-sync-content' +import type { CloudSyncRepositoryConflict } from '@zennotes/shared-domain/cloud-sync-coordinator' +import type { CloudSyncLocalItem, CloudSyncTrackedItem } from '@zennotes/shared-domain/cloud-sync-engine' +import type { CloudFileFingerprint } from './native-fs' + +const ROLLBACK_DIRECTORY = '.zennotes/sync' +const STAGING_DIRECTORY = `${ROLLBACK_DIRECTORY}/downloads` + +/** Native file operations the staging layer needs; bytes never cross the bridge. */ +export interface CloudStagingNative { + statVerified(path: string): Promise<'file' | 'directory' | null> + readForSync(path: string, textCandidate: boolean): Promise + download(options: { url: string; headers: Record; to: string; byteLength: number; sha256: string }): Promise + copyForSync(from: string, to: string, byteLength: number): Promise + rename(from: string, to: string): Promise + deleteFile(path: string): Promise + mkdir(path: string): Promise + readText(path: string): Promise +} + +interface StagedHandle { + owner: object + path: string + signal?: AbortSignal +} + +/** Downloads a referenced revision into vault-private staging, then publishes + * it only after the local destination is re-verified against sync state. */ +export class NativeCloudStaging { + constructor( + private readonly native: CloudStagingNative, + private readonly textCandidate: (path: string) => boolean, + private readonly readLocal: (path: string) => Promise + ) {} + + async stage(source: CloudSyncDownloadSource): Promise { + const reference = validateCloudSyncContentReference(source.reference) + throwIfCloudSyncCancelled(source.signal) + await this.native.mkdir(STAGING_DIRECTORY) + const path = `${STAGING_DIRECTORY}/${crypto.randomUUID()}` + const discard = async () => { await this.native.deleteFile(path).catch(() => {}) } + try { + for (let attempt = 0; attempt < 2; attempt++) { + throwIfCloudSyncCancelled(source.signal) + const instruction = validateCloudSyncDownloadInstruction( + { data: await source.getInstruction() }, reference, source.allowInsecureLoopback + ) + if (Date.parse(instruction.download.expires_at) <= Date.now()) { + if (attempt === 0) continue + throw new Error('Cloud returned an expired download URL.') + } + try { + await this.native.download({ + url: instruction.download.url, headers: instruction.download.headers, to: path, + byteLength: reference.byte_length, sha256: reference.sha256 + }) + } catch (error) { + if (attempt === 0 && isExpiredUrlError(error)) { + await discard() + continue + } + throw error + } + throwIfCloudSyncCancelled(source.signal) + const staged = await this.native.readForSync(path, false) + if (staged.sha256 !== reference.sha256 || staged.byteLength !== reference.byte_length) { + throw new Error('Cloud download failed length or SHA-256 verification.') + } + const preview = reference.encoding === 'utf8' && reference.byte_length <= Math.min(source.previewLimitBytes, 262_144) + ? { encoding: 'utf8' as const, data: await this.native.readText(path), sha256: reference.sha256, + byte_length: reference.byte_length, media_type: reference.media_type } + : undefined + if (preview && new TextEncoder().encode(preview.data).byteLength !== reference.byte_length) { + throw new Error('Cloud preview did not match its downloaded bytes.') + } + throwIfCloudSyncCancelled(source.signal) + return registerCloudSyncStagedFile(reference, { owner: this, path, signal: source.signal } satisfies StagedHandle, discard, preview) + } + throw new Error('Cloud download could not be refreshed.') + } catch (error) { + await discard() + throw error + } + } + + async matches(path: string, reference: CloudSyncContentReference): Promise { + if (await this.native.statVerified(path) !== 'file') return false + const local = await this.native.readForSync(path, false) + return local.sha256 === reference.sha256 && local.byteLength === reference.byte_length + } + + async apply(change: CloudSyncChange, previous: CloudSyncTrackedItem | undefined, file: CloudSyncStagedFile): Promise { + const handle = this.handle(file) + if (change.type !== 'upsert' || !change.content_ref || change.item_id !== file.reference.item_id || change.revision < file.reference.revision) { + throw new Error('Invalid staged Cloud change.') + } + cloudSyncStagedHandle(file, change.content_ref) + const target = normalizeCloudSyncPath(change.path) + const sourcePath = previous?.path ? normalizeCloudSyncPath(previous.path) : target + if (![target, sourcePath].every(shouldSyncVaultPath)) return + if (await this.matches(target, file.reference)) { + if (sourcePath !== target && await this.native.statVerified(sourcePath) === 'file') { + if (!await this.vouched(sourcePath, previous)) return this.conflict(sourcePath) + await this.native.deleteFile(sourcePath) + } + return + } + for (const candidate of new Set([sourcePath, target])) { + if (await this.native.statVerified(candidate) === 'file' && !await this.vouched(candidate, previous)) { + return this.conflict(candidate) + } + } + const targetExists = await this.native.statVerified(target) === 'file' + await this.publish(target, handle, file.reference, targetExists ? previous?.sha256 ?? null : null, previous) + if (sourcePath !== target && await this.native.statVerified(sourcePath) === 'file') { + if (!await this.vouched(sourcePath, previous)) return this.conflict(sourcePath) + throwIfCloudSyncCancelled(handle.signal) + await this.native.deleteFile(sourcePath) + } + } + + async resolve(input: CloudSyncStagedConflict): Promise { + const handle = this.handle(input.file) + const expected = input.expected_path === null ? null : normalizeCloudSyncPath(input.expected_path) + const cloud = normalizeCloudSyncPath(input.cloud_path) + const keep = input.keep_both_path === undefined ? null : normalizeCloudSyncPath(input.keep_both_path) + if (!shouldSyncVaultPath(cloud) || (keep && (!shouldSyncVaultPath(keep) || cloudSyncPathKey(keep) === cloudSyncPathKey(cloud) || + (expected && cloudSyncPathKey(keep) === cloudSyncPathKey(expected))))) { + throw new Error('Choose a different filename inside the synced vault.') + } + const local = expected ? await this.fingerprint(expected) : null + if ((local?.sha256 ?? null) !== input.expected_sha256) throw new Error('This file changed on this device.') + const replacing = expected !== null && cloudSyncPathKey(expected) === cloudSyncPathKey(cloud) + if (!replacing && await this.native.statVerified(cloud) !== null) throw new Error(`${cloud} already exists.`) + if (keep) { + if (!expected || !local) throw new Error('Both versions are no longer available.') + if (await this.native.statVerified(keep) !== null) throw new Error(`${keep} already exists.`) + // keep_both_path is user-visible and scanned: a partial copy left behind + // would be uploaded as the "conflict copy". Any failure removes it. + try { + await this.native.copyForSync(expected, keep, local.byteLength) + const copy = await this.native.readForSync(keep, false) + if (copy.sha256 !== local.sha256 || copy.byteLength !== local.byteLength) { + throw new Error('The local copy failed byte verification.') + } + } catch (error) { + await this.native.deleteFile(keep).catch(() => {}) + throw error + } + } + try { + if (expected && (await this.fingerprint(expected))?.sha256 !== input.expected_sha256) { + throw new Error('This file changed on this device.') + } + await this.publish(cloud, handle, input.file.reference, replacing ? input.expected_sha256 : null) + } catch (error) { + // Remove only our unchanged duplicate after the original is confirmed + // safe. A later edit or incomplete rollback must retain both copies. + if (keep && expected && local) { + const original = await this.fingerprint(expected).catch(() => null) + const duplicate = await this.fingerprint(keep).catch(() => null) + if (original?.sha256 === local.sha256 && original.byteLength === local.byteLength && + duplicate?.sha256 === local.sha256 && duplicate.byteLength === local.byteLength) { + await this.native.deleteFile(keep).catch(() => {}) + } + } + throw error + } + if (expected && !replacing) { + if ((await this.fingerprint(expected))?.sha256 !== input.expected_sha256) throw new Error('This file changed on this device.') + throwIfCloudSyncCancelled(handle.signal) + await this.native.deleteFile(expected) + } + } + + private async publish(path: string, handle: StagedHandle, reference: CloudSyncContentReference, expectedHash: string | null, previous?: CloudSyncTrackedItem): Promise { + throwIfCloudSyncCancelled(handle.signal) + if (!await this.matches(handle.path, reference)) throw new Error('The staged Cloud file changed before publication.') + const parent = path.includes('/') ? path.slice(0, path.lastIndexOf('/')) : '' + if (parent) await this.native.mkdir(parent) + const current = await this.fingerprint(path) + if ((current?.sha256 ?? null) !== expectedHash && !(current && previous && await this.vouched(path, previous))) { + throw new Error('This file changed on this device.') + } + // The backup lives directly under .zennotes/sync so the repository's + // fail-closed recovery check (which lists that directory) sees it if the + // process dies mid-publish. A stranded original must never read as a + // local deletion on the next scan. + const backup = `${ROLLBACK_DIRECTORY}/rollback-${crypto.randomUUID()}` + let publishing = false + try { + throwIfCloudSyncCancelled(handle.signal) + if (current) await this.native.rename(path, backup) + throwIfCloudSyncCancelled(handle.signal) + publishing = true + await this.native.rename(handle.path, path) + const published = await this.native.readForSync(path, false) + if (published.sha256 !== reference.sha256 || published.byteLength !== reference.byte_length) { + throw new Error('Published Cloud file failed byte verification.') + } + throwIfCloudSyncCancelled(handle.signal) + } catch (error) { + if (await this.native.statVerified(backup) === 'file') { + try { + if (await this.native.statVerified(path) !== null) await this.native.deleteFile(path) + await this.native.rename(backup, path) + } catch (restoreError) { + throw new Error(`Cloud publish rollback failed; the original is preserved at ${backup}.`, { cause: restoreError }) + } + } else if (!current && publishing && await this.native.statVerified(path) !== null) { + await this.native.deleteFile(path) + } + throw error + } + if (current) await this.native.deleteFile(backup).catch(() => {}) + } + + private async vouched(path: string, previous: CloudSyncTrackedItem | undefined): Promise { + if (!previous) return false + const local = await this.fingerprint(path) + return local !== null && local.sha256 === previous.sha256 && local.byteLength === previous.byte_length + } + + private async conflict(path: string): Promise { + const local = await this.localItem(path) + return { code: 'LOCAL_EDIT_CONFLICT', path, conflict_copy_path: null, local } + } + + private async localItem(path: string): Promise { + if (await this.native.statVerified(path) !== 'file') return null + return this.readLocal(path) + } + + private async fingerprint(path: string): Promise { + if (await this.native.statVerified(path) !== 'file') return null + return this.native.readForSync(path, this.textCandidate(path)) + } + + private handle(file: CloudSyncStagedFile): StagedHandle { + const handle = cloudSyncStagedHandle(file) + if (handle.owner !== this) throw new Error('Cloud staging belongs to another repository.') + throwIfCloudSyncCancelled(handle.signal) + return handle + } +} + +function isExpiredUrlError(error: unknown): boolean { + // Capacitor plugin rejections carry the HTTP status under `data`. + const candidate = error as { status?: unknown; data?: { status?: unknown } } + const status = candidate?.status ?? candidate?.data?.status + return status === 401 || status === 403 +} diff --git a/src/bridge/mobile-cloud-auth.test.ts b/src/bridge/mobile-cloud-auth.test.ts index c18ff28..4371250 100644 --- a/src/bridge/mobile-cloud-auth.test.ts +++ b/src/bridge/mobile-cloud-auth.test.ts @@ -12,6 +12,10 @@ const account = { const credential = JSON.stringify({ base_url: account.base_url, token: 'test-only', account }) async function coldLaunch(saved: Map) { + const lifecycle: string[] = [] + const listeners = new Map void>() + let lifetime = new AbortController() + let delayRead: { entered(): void; wait: Promise } | null = null // Exercise the real Capacitor lazy proxy: concurrent first calls can // instantiate separate implementations, each with its own key prefix. const secureStorage = registerPlugin(`TestCloudStorage${randomUUID()}`, { @@ -22,7 +26,16 @@ async function coldLaunch(saved: Map) { async setKeyPrefix(prefix: string) { this.prefix = prefix } async setSynchronize(_value: boolean) {} async setDefaultKeychainAccess(_value: unknown) {} - async getItem(key: string) { return saved.get(this.prefix + key) ?? null } + async getItem(key: string) { + const value = saved.get(this.prefix + key) ?? null + if (key === 'credential' && delayRead) { + const delayed = delayRead + delayRead = null + delayed.entered() + await delayed.wait + } + return value + } async setItem(key: string, value: string) { saved.set(this.prefix + key, value) } async removeItem(key: string) { saved.delete(this.prefix + key) } }() @@ -30,14 +43,32 @@ async function coldLaunch(saved: Map) { }) const api = await loadMobileModule('./src/bridge/mobile-cloud-auth.ts', { '@capacitor/core': { Capacitor: { isNativePlatform: () => true }, CapacitorHttp: {} }, - '@capacitor/app': { App: { addListener: async () => ({}), getLaunchUrl: async () => null } }, + '@capacitor/app': { App: { addListener: async (name: string, listener: (event: any) => void) => { + listeners.set(name, listener) + return {} + }, getLaunchUrl: async () => null } }, '@aparajita/capacitor-secure-storage': { SecureStorage: secureStorage, KeychainAccess: { whenUnlockedThisDeviceOnly: 'device-only' } }, - './cloud-sync-client': { createCloudSyncClient: () => assert.fail('status must only read storage') } + './cloud-sync-client': { + createCloudSyncClient: () => assert.fail('status must only read storage'), + stopMobileCloudRequests: () => { lifecycle.push('stop'); lifetime.abort() }, + resumeMobileCloudRequests: () => { lifecycle.push('resume'); if (lifetime.signal.aborted) lifetime = new AbortController() }, + mobileCloudRequestSignal: () => lifetime.signal + } }) await api.configureMobileCloudAuth('test-version') + api.lifecycle = lifecycle + api.appState = (isActive: boolean) => listeners.get('appStateChange')!({ isActive }) + api.pauseNextCredentialRead = () => { + let entered!: () => void + let release!: () => void + const started = new Promise(resolve => { entered = resolve }) + const wait = new Promise(resolve => { release = resolve }) + delayRead = { entered, wait } + return { started, release } + } return api } @@ -61,10 +92,32 @@ it('preserves the canonical account and prevents a legacy credential from return const api = await coldLaunch(saved) assert.equal((await api.authenticatedCredential()).token, 'test-only') await api.logoutMobileCloudAccount() + assert.ok(api.lifecycle.includes('stop')) assert.equal(saved.size, 0) assert.deepEqual(await (await coldLaunch(saved)).getMobileCloudAccountStatus(), { state: 'disconnected', account: null }) }) +it('stops pending Cloud requests in the background and resumes on foreground', async () => { + const api = await coldLaunch(new Map([['zennotes.cloud.credential', credential]])) + api.appState(false) + api.appState(true) + assert.deepEqual(api.lifecycle, ['stop', 'resume']) +}) + +it('rejects a credential read spanning logout even after foreground creates a new lifetime', async () => { + const saved = new Map([['zennotes.cloud.credential', credential]]) + const api = await coldLaunch(saved) + await api.getMobileCloudAccountStatus() + const paused = api.pauseNextCredentialRead() + const pending = api.authenticatedClient().catch((error: unknown) => error) + await paused.started + await api.logoutMobileCloudAccount() + api.appState(true) + paused.release() + assert.equal((await pending).name, 'AbortError') + assert.equal(saved.size, 0) +}) + it('rejects invalid recovered credentials through the shared auth validator', async () => { const invalid = JSON.stringify({ base_url: 'https://wrong.example.test', token: 'test-only', account }) const saved = new Map([['capacitor-storage_credential', invalid]]) diff --git a/src/bridge/mobile-cloud-auth.ts b/src/bridge/mobile-cloud-auth.ts index a5c62f4..8c4fe27 100644 --- a/src/bridge/mobile-cloud-auth.ts +++ b/src/bridge/mobile-cloud-auth.ts @@ -21,7 +21,7 @@ import { type CloudAuthPending, type CloudAuthStorage } from '@zennotes/shared-domain/cloud-auth-flow' -import { createCloudSyncClient } from './cloud-sync-client' +import { createCloudSyncClient, stopMobileCloudRequests, resumeMobileCloudRequests, mobileCloudRequestSignal } from './cloud-sync-client' const DEVELOPMENT_CLOUD_BASE_URL = import.meta.env.VITE_ZENNOTES_CLOUD_DEV_URL?.trim() const PRODUCTION_CLOUD_BASE_URL = 'https://zennotes.org' @@ -33,6 +33,7 @@ const accountListeners = new Set<(status: CloudAccountStatus) => void>() let authFlow: CloudAuthFlow | null = null let callbackQueue = Promise.resolve() +let appActive = true // Lazy and retryable so a transient native storage failure does not poison // every later account read for the session. @@ -113,12 +114,15 @@ const storage: CloudAuthStorage = { return credential.value }, async saveCredential(credential: CloudAuthCredential): Promise { + stopMobileCloudRequests() assertNativeCloudAuth() await secureStorageReady() const canonicalCredential = migrateLegacyCloudCredential(credential).value await SecureStorage.setItem(CREDENTIAL_KEY, JSON.stringify(canonicalCredential)) + if (appActive) resumeMobileCloudRequests() }, async deleteCredential(): Promise { + stopMobileCloudRequests() if (!Capacitor.isNativePlatform()) return await secureStorageReady() await SecureStorage.removeItem(CREDENTIAL_KEY) @@ -144,6 +148,11 @@ export async function configureMobileCloudAuth(appVersion: string): Promise { if (isCloudAuthUrl(url)) scheduleAuthCallback(url) }) + await CapApp.addListener('appStateChange', ({ isActive }) => { + appActive = isActive + if (isActive) resumeMobileCloudRequests() + else stopMobileCloudRequests() + }) const launch = await CapApp.getLaunchUrl() if (launch?.url && isCloudAuthUrl(launch.url)) scheduleAuthCallback(launch.url) } @@ -176,6 +185,7 @@ export async function connectMobileCloudAccount( } export async function logoutMobileCloudAccount(): Promise { + stopMobileCloudRequests() const status = await requireAuthFlow().logout() notify(status) return status @@ -218,8 +228,10 @@ export async function listMobileCloudVaults(): Promise { } export async function authenticatedClient() { + const signal = mobileCloudRequestSignal() const credential = await authenticatedCredential() - return createCloudSyncClient(credential.base_url, credential.token) + if (signal.aborted) throw new DOMException('Cloud account changed while loading credentials.', 'AbortError') + return createCloudSyncClient(credential.base_url, credential.token, { accountId: credential.account.user.email, signal }) } export async function authenticatedCredential(): Promise { diff --git a/src/bridge/mobile-cloud-sync.integration.test.ts b/src/bridge/mobile-cloud-sync.integration.test.ts index 28fcc8e..d21d7c6 100644 --- a/src/bridge/mobile-cloud-sync.integration.test.ts +++ b/src/bridge/mobile-cloud-sync.integration.test.ts @@ -128,6 +128,16 @@ async function fixture(initial: Record = { 'note.md': 'Original' assert.ok(file) return file.bytes.toString('base64') }, + async readForSync(path: string, textCandidate: boolean) { + reads.push(path) + const file = files.get(path) + assert.ok(file) + return { + uri: `file:///vault/${path}`, byteLength: file.bytes.length, + sha256: createHash('sha256').update(file.bytes).digest('hex'), utf8: textCandidate, + inlineBase64: file.bytes.toString('base64') + } + }, async writeText(path: string, data: string) { put(path, data) if (path === failWritePath) throw new Error('Native write failed after writing') @@ -143,6 +153,7 @@ async function fixture(initial: Record = { 'note.md': 'Original' assert.ok(file) files.set(to, file) files.delete(from) + if (to === failWritePath) throw new Error('Native write failed after publishing') } } const vault = { @@ -272,7 +283,7 @@ describe('mobile Cloud adapter wiring', () => { if (uploaded?.type === 'upsert') assert.equal(uploaded.content.data, 'Local changes') }) - it('exposes partial native writes when a pull fails, and can retry safely', async () => { + it('rolls back a partial native publish while preserving earlier completed pulls', async () => { const h = await fixture() await h.sync() h.refreshes.length = 0 @@ -280,7 +291,8 @@ describe('mobile Cloud adapter wiring', () => { h.remoteText('second.md', 'Partially written') h.setFailWrite('second.md') await assert.rejects(h.sync(), /Native write failed/) - assert.deepEqual(h.refreshes, [{ 'note.md': 'Remote update', 'second.md': 'Partially written' }]) + assert.deepEqual(h.refreshes, [{ 'note.md': 'Remote update' }]) + assert.equal(h.files.has('second.md'), false) h.setFailWrite(null) await h.sync() assert.equal(h.files.get('note.md')?.bytes.toString(), 'Remote update') diff --git a/src/bridge/mobile-direct-upload.test.ts b/src/bridge/mobile-direct-upload.test.ts index 2a99f5b..604a8e7 100644 --- a/src/bridge/mobile-direct-upload.test.ts +++ b/src/bridge/mobile-direct-upload.test.ts @@ -15,6 +15,30 @@ import { } from './mobile-direct-upload.ts' describe('mutateWithMobileDirectUploads', () => { + it('adds a content type to host-only signed uploads so Android writes the binary body', () => { + const headers = { Host: 'objects.example.test' } + const options = mobileObjectUploadOptions({ + url: 'https://objects.example.test/upload', + method: 'PUT', + headers, + base64: 'AQID', + byteLength: 3 + }) + assert.equal(options.headers['Content-Type'], 'application/octet-stream') + assert.deepEqual(headers, { Host: 'objects.example.test' }) + }) + + it('preserves an existing signed content type regardless of casing', () => { + const options = mobileObjectUploadOptions({ + url: 'https://objects.example.test/upload', + method: 'PUT', + headers: { 'content-type': 'image/png' }, + base64: 'AQID', + byteLength: 3 + }) + assert.deepEqual(options.headers, { 'content-type': 'image/png' }) + }) + it('builds a native binary PUT without an account bearer token or redirects', () => { const options = mobileObjectUploadOptions({ url: 'https://objects.example.test/upload?signature=signed', diff --git a/src/bridge/mobile-direct-upload.ts b/src/bridge/mobile-direct-upload.ts index 22aa87b..e82f62d 100644 --- a/src/bridge/mobile-direct-upload.ts +++ b/src/bridge/mobile-direct-upload.ts @@ -2,6 +2,7 @@ import type { CloudSyncCapacityConflict, CloudSyncConflict, CloudSyncConflictCode, + CloudSyncContent, CloudSyncMutation, CloudSyncMutationRequest, CloudSyncMutationResponse, @@ -13,6 +14,13 @@ import type { export const CLOUD_SYNC_INLINE_UPLOAD_LIMIT_BYTES = 5 * 1024 * 1024 +const uploadSources = new WeakMap() + +export function rememberMobileUploadSource(content: CloudSyncContent, uri: string): CloudSyncContent { + uploadSources.set(content, uri) + return content +} + const DIRECT_UPLOAD_COMPLETION_ATTEMPTS = 3 const SYNC_CONFLICT_CODES = new Set([ 'REVISION_CONFLICT', @@ -33,13 +41,12 @@ export interface MobileDirectUploadApi { abortUpload(vaultId: string, uploadId: string): Promise } -export interface MobileObjectUploadRequest { +export type MobileObjectUploadRequest = { url: string method: 'PUT' headers: Record - base64: string byteLength: number -} +} & ({ uri: string; sha256: string; base64?: never } | { base64: string; uri?: never; sha256?: string }) export type MobileObjectUpload = (request: MobileObjectUploadRequest) => Promise @@ -57,10 +64,16 @@ export interface MobileObjectUploadOptions { export function mobileObjectUploadOptions( request: MobileObjectUploadRequest ): MobileObjectUploadOptions { + if (request.uri !== undefined) throw new Error('File-backed uploads require the native file uploader.') + // Capacitor skips the request body entirely when Content-Type is absent. + const headers = { ...request.headers } + if (!Object.keys(headers).some((key) => key.toLowerCase() === 'content-type')) { + headers['Content-Type'] = 'application/octet-stream' + } return { url: request.url, method: request.method, - headers: request.headers, + headers, data: request.base64, dataType: 'file', connectTimeout: 30_000, @@ -132,7 +145,8 @@ async function directUpload( mutation: CloudSyncUpsertMutation, uploadObject: MobileObjectUpload ): Promise { - const base64 = uploadBase64(mutation) + const uri = uploadSources.get(mutation.content) + const source = uri === undefined ? { base64: uploadBase64(mutation) } : { uri } let initiation: CloudSyncUploadInitiationResponse try { @@ -157,7 +171,8 @@ async function directUpload( url: secureDirectUploadUrl(instruction.upload.url), method: instruction.upload.method, headers: instruction.upload.headers, - base64, + ...source, + sha256: mutation.content.sha256, byteLength: mutation.content.byte_length }) } catch (error) { diff --git a/src/bridge/native-fs-cloud-copy.test.ts b/src/bridge/native-fs-cloud-copy.test.ts new file mode 100644 index 0000000..d28f387 --- /dev/null +++ b/src/bridge/native-fs-cloud-copy.test.ts @@ -0,0 +1,75 @@ +import assert from 'node:assert/strict' +import { createHash } from 'node:crypto' +import { it } from 'node:test' +import { loadMobileModule } from '../../tooling/load-mobile-module.ts' + +for (const provider of ['file', 'content']) { + it(`copies a 6 MB ${provider} source using a verified native file handle`, async () => { + const bytes = Buffer.alloc(6_000_000, 197) + const uri = (path: string) => `${provider}:///${path.replace('ZenNotes/Test/', '')}` + const files = new Map([[uri('source.bin'), [bytes]]]) + const copies: unknown[] = [] + const missing = () => Object.assign(new Error('Not found'), { code: 'OS-PLUG-FILE-0008' }) + const stat = (path: string) => { + const chunks = files.get(uri(path)) + if (!chunks) throw missing() + return { type: 'file', uri: uri(path), mtime: 1, size: chunks.reduce((sum, chunk) => sum + chunk.length, 0) } + } + const storage = Object.getOwnPropertyDescriptor(globalThis, 'localStorage') + Object.defineProperty(globalThis, 'localStorage', { configurable: true, value: { getItem: () => null } }) + try { + const { NativeFs } = await loadMobileModule('./src/bridge/native-fs', { + './icloud': { ensureDownloaded: async () => {} }, + '@capacitor/core': { + Capacitor: {}, + registerPlugin: (name: string) => name === 'SafFs' ? { + stat: async ({ path }: { path: string }) => stat(path), + writeBase64: async ({ path, data }: { path: string; data: string }) => { + assert.equal(data, '') + files.set(uri(path), []) + }, + copy: () => assert.fail('SafFs.copy buffers the whole source') + } : { + inspect: async ({ uri: path }: { uri: string }) => { + const chunks = files.get(path) + assert.ok(chunks) + const hash = createHash('sha256') + for (const chunk of chunks) hash.update(chunk) + return { uri: path, byteLength: chunks.reduce((sum, chunk) => sum + chunk.length, 0), + sha256: hash.digest('hex'), utf8: false } + }, + copy: async (request: { from: string; to: string; byteLength: number; sha256: string }) => { + assert.deepEqual(request, { from: uri('source.bin'), to: uri('copy.bin'), byteLength: bytes.length, + sha256: createHash('sha256').update(bytes).digest('hex') }) + assert.ok(JSON.stringify(request).length < 400) + copies.push(request) + files.set(request.to, [Buffer.from(bytes)]) + } + } + }, + '@capacitor/filesystem': { + Directory: { Data: 'DATA', External: 'EXTERNAL' }, Encoding: { UTF8: 'utf8' }, + Filesystem: { + getUri: async ({ path }: { path: string }) => ({ uri: uri(path) }), + stat: async ({ path }: { path: string }) => stat(path), + writeFile: async ({ path, data }: { path: string; data: string }) => { + assert.equal(data, '') + files.set(uri(path), []) + }, + readFile: () => assert.fail('File bytes must not cross the bridge'), + appendFile: () => assert.fail('File bytes must not cross the bridge'), + copy: () => assert.fail('Whole-file native copy is not used') + } + } + }) + const fs = new NativeFs('Test', provider === 'content' ? 'content:///tree/root' : null) + await fs.copyForSync('source.bin', 'copy.bin', bytes.length) + assert.deepEqual(Buffer.concat(files.get(uri('copy.bin'))!), bytes) + assert.equal(copies.length, 1) + await assert.rejects(fs.copyForSync('source.bin', 'copy.bin', bytes.length), /already exists/) + } finally { + if (storage) Object.defineProperty(globalThis, 'localStorage', storage) + else Reflect.deleteProperty(globalThis, 'localStorage') + } + }) +} diff --git a/src/bridge/native-fs.ts b/src/bridge/native-fs.ts index 3f1b9e1..507e4ef 100644 --- a/src/bridge/native-fs.ts +++ b/src/bridge/native-fs.ts @@ -95,6 +95,20 @@ const SafFs = registerPlugin<{ delete(o: { root: string; path: string }): Promise }>('SafFs') +export interface CloudFileFingerprint { + uri: string + sha256: string + byteLength: number + utf8: boolean + inlineBase64?: string +} + +const CloudFiles = registerPlugin<{ + inspect(options: { uri: string; textCandidate: boolean }): Promise + copy(options: { from: string; to: string; byteLength: number; sha256: string }): Promise + download(options: { url: string; headers: Record; to: string; byteLength: number; sha256: string }): Promise +}>('ZenDirectUpload') + export function isSafRoot(uri: string | null | undefined): boolean { return typeof uri === 'string' && uri.startsWith('content://') } @@ -227,6 +241,15 @@ export class NativeFs { return bytesToBase64(new Uint8Array(await res.data.arrayBuffer())) } + async readForSync(relPath: string, textCandidate: boolean, knownUri?: string): Promise { + const location = this.loc(relPath) + const uri = knownUri ?? (this.saf + ? (await SafFs.stat({ root: this.cloudRootUri!, path: relPath })).uri + : location.directory === undefined ? location.path + : (await Filesystem.getUri({ path: location.path, directory: location.directory })).uri) + return CloudFiles.inspect({ uri, textCandidate }) + } + async readTextOrNull(relPath: string): Promise { try { return await this.readText(relPath) @@ -369,6 +392,41 @@ export class NativeFs { }) } + /** Copy into staging without carrying file bytes across the bridge or + * reopening a remote document provider for every chunk. */ + async copyForSync(fromRel: string, toRel: string, byteLength: number): Promise { + if (fromRel === toRel || !Number.isSafeInteger(byteLength) || byteLength < 0) { + throw new Error('Invalid Cloud file copy.') + } + if (await this.statVerified(toRel) !== null) throw new Error('Cloud copy destination already exists.') + const source = await this.readForSync(fromRel, false) + if (source.byteLength !== byteLength) throw new Error('The Cloud copy source changed size.') + const parent = toRel.includes('/') ? toRel.slice(0, toRel.lastIndexOf('/')) : '' + if (parent) await this.mkdir(parent) + await this.writeBase64(toRel, '') + const target = await this.statOrNull(toRel) + if (!target) throw new Error('Cloud copy staging file is unavailable.') + await CloudFiles.copy({ from: source.uri, to: target.uri, byteLength, sha256: source.sha256 }) + } + + /** Stream a signed Cloud revision into a fresh vault-private staging file. + * The native side verifies length and SHA-256 before resolving, so the + * WebView never holds the body and a truncated download never lands. */ + async download(options: { url: string; headers: Record; to: string; byteLength: number; sha256: string }): Promise { + if (await this.statVerified(options.to) !== null) throw new Error('Cloud download destination already exists.') + const parent = options.to.includes('/') ? options.to.slice(0, options.to.lastIndexOf('/')) : '' + if (parent) await this.mkdir(parent) + await this.writeBase64(options.to, '') + const target = await this.statOrNull(options.to) + if (!target) throw new Error('Cloud download staging file is unavailable.') + try { + await CloudFiles.download({ ...options, to: target.uri }) + } catch (error) { + await this.deleteFile(options.to).catch(() => {}) + throw error + } + } + async deleteFile(relPath: string): Promise { if (this.saf) { await SafFs.delete({ root: this.cloudRootUri!, path: relPath }) diff --git a/vendor/zennotes/manifest.json b/vendor/zennotes/manifest.json index 51a076d..65e306d 100644 --- a/vendor/zennotes/manifest.json +++ b/vendor/zennotes/manifest.json @@ -1,12 +1,12 @@ { "name": "@zennotes/app-core", - "version": "2.60.0-core.hb0d0b54f320a8e3f", - "file": "zennotes-app-core-2.60.0-core.hb0d0b54f320a8e3f.tgz", - "sha256": "bef83d372725736b26b62724e90126dfb474ad4a3ed88205b806cf5a3be2f97f", - "integrity": "sha512-aMZ8DegEHQTwT4tf6kdvAHEUMfhi4nayN81lHKw4PSZZfDr4UKvhpufIZEm5xO4UR4gyKEJSwd4/1QNBB+aC+Q==", - "sourceCommit": "158293947e451f2a88011ddf647517c847cc0a76", + "version": "2.60.1-core.h05ebb55c14afffb2", + "file": "zennotes-app-core-2.60.1-core.h05ebb55c14afffb2.tgz", + "sha256": "2a1781c1121f37cf494b5f34d67ea8dce25b5bdfe6dee0d7a9837ceaec8cfd75", + "integrity": "sha512-pTn7+3VmHG1le5pzBruIAzlHg4Nc+CRucL/u4FLbktdPDeIQ/PWECeS+fkia7JM5w6UVhqFxY5ZP24gj77vfFg==", + "sourceCommit": "4c74b478148a68adc8d77cc5c0d54a46ab469342", "workingTreeDirty": false, - "sourceLockSha256": "34573e8de5f27f233eef597e71c61a5d1e161d5679469b1049dd9b5c708a255c", + "sourceLockSha256": "503c4639d5dbcb1a909ad122882c66aea751747f1a2a8fbcb8bdf8077eb8d871", "toolchain": { "node": "v22.23.3", "typescript": "5.9.3", @@ -15,20 +15,20 @@ "dependencies": [ { "name": "@zennotes/bridge-contract", - "version": "2.60.0-boundaries.hbdc4a2a12368fec1", - "file": "zennotes-bridge-contract-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", - "sha256": "02f9beedd7d6d9454571bcaeb158d5194615c9f0b9dad42a428757aa774d068c", - "integrity": "sha512-5uh7dQjtrWNds26/4FH3wcrEQ0cF6SgUYd0+bGHKtl8GOLlayJWmDaJPXgm7FQ97Gc18KNM47k6k2rxYuECx3A==", - "sourceCommit": "158293947e451f2a88011ddf647517c847cc0a76", + "version": "2.60.1-boundaries.ha881a2f8575a38bf", + "file": "zennotes-bridge-contract-2.60.1-boundaries.ha881a2f8575a38bf.tgz", + "sha256": "79f327d998e999bbeaaa12f25272be3f4c2276fc36e4a1470d252586b6e40e3f", + "integrity": "sha512-/BjWqbzR4+S8Bg8gOwT6OZO8Ikr8qO8zNmVVscjr45vIWasAykfaLS9H6IuFfIE48/wbAXb3W8J6aTHNEdk0xA==", + "sourceCommit": "4c74b478148a68adc8d77cc5c0d54a46ab469342", "workingTreeDirty": false }, { "name": "@zennotes/shared-domain", - "version": "2.60.0-boundaries.hbdc4a2a12368fec1", - "file": "zennotes-shared-domain-2.60.0-boundaries.hbdc4a2a12368fec1.tgz", - "sha256": "5bec93dfe153ad0687c5e003f3db3bf53cae33f5ba37deefb4327fba96f77711", - "integrity": "sha512-Gr/F9IkAnX24T9apcPWLhErAtQ4N/QUr/rIdHJqNOTxERVSjrvIqP+wGCzNyuVRDujX9eXOfAt0y4Ir6Ol0TwQ==", - "sourceCommit": "158293947e451f2a88011ddf647517c847cc0a76", + "version": "2.60.1-boundaries.ha881a2f8575a38bf", + "file": "zennotes-shared-domain-2.60.1-boundaries.ha881a2f8575a38bf.tgz", + "sha256": "4e97e45e1f319091c85ead57ebaf4b0e64928d4d2d6d814af00cae2ca13535a2", + "integrity": "sha512-5Sb02aV5yUy3MkqR9rkrx7eHQTRmurivhMVpeVa4kU0S07dm08N6kpstB8uIeKVlxhDvLILQHhBw606N3wx+aA==", + "sourceCommit": "4c74b478148a68adc8d77cc5c0d54a46ab469342", "workingTreeDirty": false } ] diff --git a/vendor/zennotes/zennotes-app-core-2.60.0-core.hb0d0b54f320a8e3f.tgz b/vendor/zennotes/zennotes-app-core-2.60.1-core.h05ebb55c14afffb2.tgz similarity index 51% rename from vendor/zennotes/zennotes-app-core-2.60.0-core.hb0d0b54f320a8e3f.tgz rename to vendor/zennotes/zennotes-app-core-2.60.1-core.h05ebb55c14afffb2.tgz index 13db6b2..0124ea0 100644 Binary files a/vendor/zennotes/zennotes-app-core-2.60.0-core.hb0d0b54f320a8e3f.tgz and b/vendor/zennotes/zennotes-app-core-2.60.1-core.h05ebb55c14afffb2.tgz differ diff --git a/vendor/zennotes/zennotes-bridge-contract-2.60.0-boundaries.hbdc4a2a12368fec1.tgz b/vendor/zennotes/zennotes-bridge-contract-2.60.0-boundaries.hbdc4a2a12368fec1.tgz deleted file mode 100644 index c8afe37..0000000 Binary files a/vendor/zennotes/zennotes-bridge-contract-2.60.0-boundaries.hbdc4a2a12368fec1.tgz and /dev/null differ diff --git a/vendor/zennotes/zennotes-bridge-contract-2.60.1-boundaries.ha881a2f8575a38bf.tgz b/vendor/zennotes/zennotes-bridge-contract-2.60.1-boundaries.ha881a2f8575a38bf.tgz new file mode 100644 index 0000000..c934918 Binary files /dev/null and b/vendor/zennotes/zennotes-bridge-contract-2.60.1-boundaries.ha881a2f8575a38bf.tgz differ diff --git a/vendor/zennotes/zennotes-shared-domain-2.60.0-boundaries.hbdc4a2a12368fec1.tgz b/vendor/zennotes/zennotes-shared-domain-2.60.0-boundaries.hbdc4a2a12368fec1.tgz deleted file mode 100644 index fce5729..0000000 Binary files a/vendor/zennotes/zennotes-shared-domain-2.60.0-boundaries.hbdc4a2a12368fec1.tgz and /dev/null differ diff --git a/vendor/zennotes/zennotes-shared-domain-2.60.1-boundaries.ha881a2f8575a38bf.tgz b/vendor/zennotes/zennotes-shared-domain-2.60.1-boundaries.ha881a2f8575a38bf.tgz new file mode 100644 index 0000000..1c73e99 Binary files /dev/null and b/vendor/zennotes/zennotes-shared-domain-2.60.1-boundaries.ha881a2f8575a38bf.tgz differ