refactor(cordn): move the file stores and EncryptedAppendLog to commonMain on okio

FileCordnStores (CordnStorageLayout and the group, key-package,
coordinator and handoff stores), FileBackedCordnScopeFactory,
CordnMigrationStores and EncryptedAppendLog move from jvmAndroid to
commonMain. They now take an okio Path, plus a FileSystem that defaults
to platformFileSystem, where they used to take a java.io.File.

The on-disk format does not change. These files hold encrypted MLS group
state that cannot be re-derived, so CordnStorageFormatGoldenTest pinned
the exact tree of paths and the bytes of every file against the java.io
implementation first. It passes against this one with its constants
untouched; only its construction glue moved from File to Path.

- Cursor framing: ByteBuffer's default big-endian putLong/getLong is now
  okio Buffer.writeLong/readLong, which is also big-endian.
- Migration base64: java.util.Base64 is now
  kotlin.io.encoding.Base64.Default with PRESENT_OPTIONAL padding, the
  same alphabet and output. A new test checks it against java.util.Base64
  in both directions, padded and unpadded.
- Failure handling: File.delete(), deleteRecursively() and
  renameTo-else-copy returned false where okio throws, so the File
  behaviour is kept by deleteQuietly, deleteRecursivelyQuietly and
  moveOrCopy in commons/util/FileSystemExt.kt.
- Durability: the log's fsyncs are FileHandle.flush(), which is
  FileDescriptor.sync() on the JVM and Android.

Callers updated: the CLI's CordnContext, CordnRuntime,
Account.cordnFilesDir (now Path?), AccountCacheState, and
EncryptedMarmotMessageStore, which stays in jvmAndroid and passes
toOkioPath() to the log.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01S7FuNBSKiyVecARSoE4B9P
This commit is contained in:
Claude
2026-09-28 03:53:56 +00:00
parent af2404e61c
commit af39263444
17 changed files with 497 additions and 216 deletions
@@ -364,6 +364,7 @@ import kotlinx.coroutines.flow.stateIn
import kotlinx.coroutines.launch import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import okio.Path
import kotlin.coroutines.cancellation.CancellationException import kotlin.coroutines.cancellation.CancellationException
import com.vitorpamplona.quartz.experimental.nip95.header.thumbhash as nip95thumbhash import com.vitorpamplona.quartz.experimental.nip95.header.thumbhash as nip95thumbhash
import com.vitorpamplona.quartz.experimental.profileGallery.thumbhash as galleryThumbhash import com.vitorpamplona.quartz.experimental.profileGallery.thumbhash as galleryThumbhash
@@ -412,7 +413,7 @@ class Account(
* Nothing about it is shared with Marmot's stores above — cordn has its * Nothing about it is shared with Marmot's stores above — cordn has its
* own, by the §3.1 rule in `amethyst/plans/2026-09-19-cordn-ui.md`. * own, by the §3.1 rule in `amethyst/plans/2026-09-19-cordn-ui.md`.
*/ */
val cordnFilesDir: java.io.File? = null, val cordnFilesDir: Path? = null,
val mlsGroupStateStore: MlsGroupStateStore? = null, val mlsGroupStateStore: MlsGroupStateStore? = null,
val marmotMessageStore: com.vitorpamplona.quartz.marmot.groups.MarmotMessageStore? = null, val marmotMessageStore: com.vitorpamplona.quartz.marmot.groups.MarmotMessageStore? = null,
val marmotKeyPackageStore: com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageBundleStore? = null, val marmotKeyPackageStore: com.vitorpamplona.quartz.marmot.mip00KeyPackages.KeyPackageBundleStore? = null,
@@ -380,7 +380,7 @@ class AccountCacheState(
// The same per-account directory the Marmot stores use. cordn // The same per-account directory the Marmot stores use. cordn
// scopes itself further by coordinator underneath it, because a // scopes itself further by coordinator underneath it, because a
// gid is unique only within one (spec/00.md §4). // gid is unique only within one (spec/00.md §4).
cordnFilesDir = accountDir, cordnFilesDir = accountDir.toOkioPath(),
mlsGroupStateStore = mlsStore, mlsGroupStateStore = mlsStore,
marmotMessageStore = marmotMessageStore, marmotMessageStore = marmotMessageStore,
marmotKeyPackageStore = marmotKeyPackageStore, marmotKeyPackageStore = marmotKeyPackageStore,
@@ -46,6 +46,8 @@ import com.vitorpamplona.amethyst.commons.cordn.FileCordnKeyPackageStore
import com.vitorpamplona.amethyst.commons.cordn.KeyStoreCordnBlobCipher import com.vitorpamplona.amethyst.commons.cordn.KeyStoreCordnBlobCipher
import com.vitorpamplona.amethyst.commons.cordn.OpenedWelcome import com.vitorpamplona.amethyst.commons.cordn.OpenedWelcome
import com.vitorpamplona.amethyst.commons.model.cordnGroups.CordnGroupList import com.vitorpamplona.amethyst.commons.model.cordnGroups.CordnGroupList
import com.vitorpamplona.amethyst.commons.util.deleteRecursivelyQuietly
import com.vitorpamplona.amethyst.commons.util.platformFileSystem
import com.vitorpamplona.quartz.contextvm.core.CvmKinds import com.vitorpamplona.quartz.contextvm.core.CvmKinds
import com.vitorpamplona.quartz.cordn.appMultiDevice.CordnHandoffCode import com.vitorpamplona.quartz.cordn.appMultiDevice.CordnHandoffCode
import com.vitorpamplona.quartz.cordn.spec00Coordinator.CoordinatorServerInfo import com.vitorpamplona.quartz.cordn.spec00Coordinator.CoordinatorServerInfo
@@ -71,7 +73,7 @@ import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
import java.io.File import okio.Path
import kotlin.coroutines.cancellation.CancellationException import kotlin.coroutines.cancellation.CancellationException
import kotlin.uuid.ExperimentalUuidApi import kotlin.uuid.ExperimentalUuidApi
import kotlin.uuid.Uuid import kotlin.uuid.Uuid
@@ -98,7 +100,7 @@ import kotlin.uuid.Uuid
class CordnRuntime( class CordnRuntime(
private val accountSigner: NostrSigner, private val accountSigner: NostrSigner,
private val client: INostrClient, private val client: INostrClient,
private val filesDir: File, private val filesDir: Path,
private val scope: CoroutineScope, private val scope: CoroutineScope,
private val cipher: CordnBlobCipher = KeyStoreCordnBlobCipher(), private val cipher: CordnBlobCipher = KeyStoreCordnBlobCipher(),
/** /**
@@ -696,7 +698,7 @@ class CordnRuntime(
forget(coordinatorPubKey) forget(coordinatorPubKey)
groups.forgetCoordinator(coordinatorPubKey) groups.forgetCoordinator(coordinatorPubKey)
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
CordnStorageLayout.directoryFor(filesDir, accountSigner.pubKey, coordinatorPubKey).deleteRecursively() platformFileSystem.deleteRecursivelyQuietly(CordnStorageLayout.directoryFor(filesDir, accountSigner.pubKey, coordinatorPubKey))
} }
} }
@@ -841,7 +843,7 @@ class CordnRuntime(
// The old state goes first. Leaving it would merge two devices' // The old state goes first. Leaving it would merge two devices'
// histories for any gid present in both, which is the one outcome // histories for any gid present in both, which is the one outcome
// this must never produce. // this must never produce.
File(filesDir, "cordn/${accountSigner.pubKey}").deleteRecursively() platformFileSystem.deleteRecursivelyQuietly(filesDir / "cordn" / accountSigner.pubKey)
archive.groups.forEach { group -> archive.groups.forEach { group ->
val store = FileCordnGroupStore(CordnStorageLayout.directoryFor(filesDir, accountSigner.pubKey, group.coordinatorPubKey), cipher) val store = FileCordnGroupStore(CordnStorageLayout.directoryFor(filesDir, accountSigner.pubKey, group.coordinatorPubKey), cipher)
@@ -59,6 +59,7 @@ import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay import kotlinx.coroutines.delay
import kotlinx.coroutines.runBlocking import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeoutOrNull import kotlinx.coroutines.withTimeoutOrNull
import okio.Path.Companion.toOkioPath
import org.junit.After import org.junit.After
import org.junit.Assert.assertEquals import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse import org.junit.Assert.assertFalse
@@ -180,7 +181,7 @@ class CordnRuntimeTest {
) = CordnRuntime( ) = CordnRuntime(
accountSigner = signer, accountSigner = signer,
client = client, client = client,
filesDir = root, filesDir = root.toOkioPath(),
scope = scope, scope = scope,
cipher = XorCipher(), cipher = XorCipher(),
links = links, links = links,
@@ -443,8 +444,8 @@ class CordnRuntimeTest {
// A gid is unique only within one coordinator, so both hold // A gid is unique only within one coordinator, so both hold
// "shared-gid" as unrelated groups. Deleting a level too high takes // "shared-gid" as unrelated groups. Deleting a level too high takes
// both, and the only visible symptom is a room that vanished. // both, and the only visible symptom is a room that vanished.
assertFalse(CordnStorageLayout.directoryFor(root, account, keyA).exists()) assertFalse(CordnStorageLayout.directoryFor(root.toOkioPath(), account, keyA).toFile().exists())
assertTrue(CordnStorageLayout.directoryFor(root, account, keyB).exists()) assertTrue(CordnStorageLayout.directoryFor(root.toOkioPath(), account, keyB).toFile().exists())
assertEquals( assertEquals(
listOf(keyB), listOf(keyB),
runtime.groups.all.value runtime.groups.all.value
@@ -466,7 +467,7 @@ class CordnRuntimeTest {
CordnRuntime( CordnRuntime(
accountSigner = NostrSignerInternal(KeyPair()), accountSigner = NostrSignerInternal(KeyPair()),
client = EmptyNostrClient(), client = EmptyNostrClient(),
filesDir = root, filesDir = root.toOkioPath(),
scope = runtimeScope(), scope = runtimeScope(),
cipher = XorCipher(), cipher = XorCipher(),
links = links, links = links,
@@ -32,7 +32,8 @@ import com.vitorpamplona.amethyst.commons.cordn.KeyedCordnBlobCipher
import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import java.io.File import okio.Path
import okio.Path.Companion.toOkioPath
/** /**
* cordn wiring for the CLI, split out of [Context] the way [CashuContext] is * cordn wiring for the CLI, split out of [Context] the way [CashuContext] is
@@ -72,7 +73,7 @@ class CordnContext(
accountPubKey = accountPubKey, accountPubKey = accountPubKey,
scopes = scopes =
FileBackedCordnScopeFactory( FileBackedCordnScopeFactory(
root = ctx.dataDir.root, root = migrationRoot,
cipher = cipher, cipher = cipher,
links = CordnLinks.over(ctx.signer, ctx.client), links = CordnLinks.over(ctx.signer, ctx.client),
), ),
@@ -81,7 +82,7 @@ class CordnContext(
private val coordinatorStore by lazy { private val coordinatorStore by lazy {
FileCordnCoordinatorStore( FileCordnCoordinatorStore(
CordnStorageLayout.accountDirectoryFor(ctx.dataDir.root, accountPubKey), CordnStorageLayout.accountDirectoryFor(migrationRoot, accountPubKey),
cipher, cipher,
) )
} }
@@ -117,8 +118,8 @@ class CordnContext(
it.keyPackages.restore() it.keyPackages.restore()
} }
/** Where the cordn tree lives, for the migration verbs. */ /** The root the cordn tree is laid out under, for the stores and the migration verbs. */
val migrationRoot: File get() = ctx.dataDir.root val migrationRoot: Path get() = ctx.dataDir.root.toOkioPath()
/** The at-rest cipher, so a migration can read and write the same blobs. */ /** The at-rest cipher, so a migration can read and write the same blobs. */
val blobCipher: CordnBlobCipher get() = cipher val blobCipher: CordnBlobCipher get() = cipher
@@ -1773,3 +1773,29 @@ first. None of the findings below came from the move; all are on `main`.
drops it, so a refresh only spins for a second. drops it, so a refresh only spins for a second.
- `ScheduledFlag` still reads `TimeUtils.now()` once, so a card composed before the start - `ScheduledFlag` still reads `TimeUtils.now()` once, so a card composed before the start
time keeps the date after it passes. time keeps the date after it passes.
## 2026-09-28 — cordn file stores and `EncryptedAppendLog` to `commonMain`
`cordn/FileCordnStores.kt` (`CordnStorageLayout` and the group, key-package,
coordinator and handoff stores), `cordn/FileBackedCordnScopeFactory.kt`,
`cordn/CordnMigrationStores.kt` and `storage/EncryptedAppendLog.kt` moved from
jvmAndroid to commonMain. They take an okio `Path` plus a `FileSystem`
(default `platformFileSystem`), the `ScheduledPostStore` shape.
These files hold encrypted MLS group state that cannot be re-derived, so the
format is pinned first: `CordnStorageFormatGoldenTest` was committed against the
`java.io` implementation, asserting the exact tree of paths and the exact bytes
of every file, and passes unchanged on the okio one. Notes on the port:
- `ByteBuffer` framing (the cursor file) is okio `Buffer.writeLong`/`readLong`;
both are big-endian.
- `java.util.Base64` in the migration snapshot is `kotlin.io.encoding.Base64.Default`
with `PaddingOption.PRESENT_OPTIONAL`, the same shape as the napplet move.
- `File.delete()`/`deleteRecursively()` returned false instead of throwing, and
`renameTo` fell back to copy-and-delete. okio throws; the File behaviour is kept
by `deleteQuietly`, `deleteRecursivelyQuietly` and `moveOrCopy` in
`commons/util/FileSystemExt.kt`. One deliberate difference: the recursive delete
removes symlinks instead of following them.
- The log's `fsync`s are `FileHandle.flush()`, which is `FileDescriptor.sync()`
on the JVM and Android.
- `EncryptedMarmotMessageStore` stays in jvmAndroid and passes `toOkioPath()` to the log.
@@ -20,13 +20,16 @@
*/ */
package com.vitorpamplona.amethyst.commons.cordn package com.vitorpamplona.amethyst.commons.cordn
import com.vitorpamplona.amethyst.commons.util.deleteRecursivelyQuietly
import com.vitorpamplona.amethyst.commons.util.platformFileSystem
import com.vitorpamplona.quartz.cordn.appMultiDevice.CordnCarriedKeyPackage import com.vitorpamplona.quartz.cordn.appMultiDevice.CordnCarriedKeyPackage
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessageCodec import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessageCodec
import com.vitorpamplona.quartz.cordn.sync.GroupCursor import com.vitorpamplona.quartz.cordn.sync.GroupCursor
import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import java.io.File import okio.FileSystem
import java.util.Base64 import okio.Path
import kotlin.io.encoding.Base64
/** /**
* Reading a handoff snapshot off the cordn stores, and writing one back. * Reading a handoff snapshot off the cordn stores, and writing one back.
@@ -46,18 +49,19 @@ object CordnMigrationStores {
* precisely when it matters. * precisely when it matters.
*/ */
suspend fun read( suspend fun read(
root: File, root: Path,
accountPubKey: HexKey, accountPubKey: HexKey,
cipher: CordnBlobCipher, cipher: CordnBlobCipher,
configs: List<CoordinatorConfig>, configs: List<CoordinatorConfig>,
fileSystem: FileSystem = platformFileSystem,
): CordnMigrationSnapshot { ): CordnMigrationSnapshot {
val groups = mutableListOf<CordnMigrationGroup>() val groups = mutableListOf<CordnMigrationGroup>()
val keyPackages = mutableListOf<CordnCarriedKeyPackage>() val keyPackages = mutableListOf<CordnCarriedKeyPackage>()
configs.forEach { config -> configs.forEach { config ->
val dir = CordnStorageLayout.directoryFor(root, accountPubKey, config.pubKey) val dir = CordnStorageLayout.directoryFor(root, accountPubKey, config.pubKey)
val groupStore = FileCordnGroupStore(dir, cipher) val groupStore = FileCordnGroupStore(dir, cipher, fileSystem)
val keyPackageStore = FileCordnKeyPackageStore(dir, cipher) val keyPackageStore = FileCordnKeyPackageStore(dir, cipher, fileSystem)
groupStore.listGroups().forEach { gid -> groupStore.listGroups().forEach { gid ->
val state = groupStore.loadGroup(gid) ?: return@forEach val state = groupStore.loadGroup(gid) ?: return@forEach
@@ -127,19 +131,20 @@ object CordnMigrationStores {
* `gid` present in both cannot end up half from each. * `gid` present in both cannot end up half from each.
*/ */
suspend fun write( suspend fun write(
root: File, root: Path,
accountPubKey: HexKey, accountPubKey: HexKey,
cipher: CordnBlobCipher, cipher: CordnBlobCipher,
snapshot: CordnMigrationSnapshot, snapshot: CordnMigrationSnapshot,
fileSystem: FileSystem = platformFileSystem,
): List<CoordinatorConfig> { ): List<CoordinatorConfig> {
require(snapshot.accountPubKey == accountPubKey) { require(snapshot.accountPubKey == accountPubKey) {
"this migration belongs to a different account" "this migration belongs to a different account"
} }
File(root, "cordn/$accountPubKey").deleteRecursively() fileSystem.deleteRecursivelyQuietly(root / "cordn" / accountPubKey)
snapshot.groups.forEach { group -> snapshot.groups.forEach { group ->
val store = FileCordnGroupStore(CordnStorageLayout.directoryFor(root, accountPubKey, group.coordinatorPubKey), cipher) val store = FileCordnGroupStore(CordnStorageLayout.directoryFor(root, accountPubKey, group.coordinatorPubKey), cipher, fileSystem)
store.saveGroup(group.gid, group.clientStateBase64.fromBase64()) store.saveGroup(group.gid, group.clientStateBase64.fromBase64())
// Before the cursor, for the same reason the live path writes them in // Before the cursor, for the same reason the live path writes them in
// that order: a seeding that wrote the cursor and then failed would // that order: a seeding that wrote the cursor and then failed would
@@ -158,7 +163,7 @@ object CordnMigrationStores {
} }
snapshot.keyPackages.forEach { keyPackage -> snapshot.keyPackages.forEach { keyPackage ->
FileCordnKeyPackageStore(CordnStorageLayout.directoryFor(root, accountPubKey, keyPackage.coordinatorPubKey), cipher) FileCordnKeyPackageStore(CordnStorageLayout.directoryFor(root, accountPubKey, keyPackage.coordinatorPubKey), cipher, fileSystem)
.save(keyPackage.keyPackageRef, keyPackage.bundle.fromBase64()) .save(keyPackage.keyPackageRef, keyPackage.bundle.fromBase64())
} }
@@ -176,7 +181,15 @@ object CordnMigrationStores {
} }
} }
private fun ByteArray.toBase64() = Base64.getEncoder().encodeToString(this) /**
* Standard alphabet, padded on the way out and optional on the way in,
* which is what `java.util.Base64`'s basic encoder and decoder did — so a
* snapshot written by an older build still reads, and one written here
* reads on an older build.
*/
private val base64 = Base64.Default.withPadding(Base64.PaddingOption.PRESENT_OPTIONAL)
private fun String.fromBase64(): ByteArray = Base64.getDecoder().decode(this) private fun ByteArray.toBase64() = base64.encode(this)
private fun String.fromBase64(): ByteArray = base64.decode(this)
} }
@@ -20,8 +20,10 @@
*/ */
package com.vitorpamplona.amethyst.commons.cordn package com.vitorpamplona.amethyst.commons.cordn
import com.vitorpamplona.amethyst.commons.util.platformFileSystem
import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.HexKey
import java.io.File import okio.FileSystem
import okio.Path
/** /**
* The production [CordnCoordinatorScopeFactory]: encrypted files on disk, plus * The production [CordnCoordinatorScopeFactory]: encrypted files on disk, plus
@@ -38,9 +40,10 @@ import java.io.File
* `<root>/cordn/<account>/<coordinator>`; see [CordnStorageLayout]. * `<root>/cordn/<account>/<coordinator>`; see [CordnStorageLayout].
*/ */
class FileBackedCordnScopeFactory( class FileBackedCordnScopeFactory(
private val root: File, private val root: Path,
private val cipher: CordnBlobCipher, private val cipher: CordnBlobCipher,
private val links: CordnCoordinatorLinkFactory, private val links: CordnCoordinatorLinkFactory,
private val fileSystem: FileSystem = platformFileSystem,
) : CordnCoordinatorScopeFactory { ) : CordnCoordinatorScopeFactory {
override suspend fun open( override suspend fun open(
accountPubKey: HexKey, accountPubKey: HexKey,
@@ -54,8 +57,8 @@ class FileBackedCordnScopeFactory(
override suspend fun serverInfo() = link.serverInfo() override suspend fun serverInfo() = link.serverInfo()
override val groupStore = FileCordnGroupStore(dir, cipher) override val groupStore = FileCordnGroupStore(dir, cipher, fileSystem)
override val keyPackageStore = FileCordnKeyPackageStore(dir, cipher) override val keyPackageStore = FileCordnKeyPackageStore(dir, cipher, fileSystem)
// Only the transport closes. The files outlive the session by // Only the transport closes. The files outlive the session by
// design — closing a coordinator is not leaving its groups, and // design — closing a coordinator is not leaving its groups, and
@@ -21,6 +21,11 @@
package com.vitorpamplona.amethyst.commons.cordn package com.vitorpamplona.amethyst.commons.cordn
import com.vitorpamplona.amethyst.commons.storage.EncryptedAppendLog import com.vitorpamplona.amethyst.commons.storage.EncryptedAppendLog
import com.vitorpamplona.amethyst.commons.util.deleteQuietly
import com.vitorpamplona.amethyst.commons.util.deleteRecursivelyQuietly
import com.vitorpamplona.amethyst.commons.util.moveOrCopy
import com.vitorpamplona.amethyst.commons.util.platformFileSystem
import com.vitorpamplona.amethyst.commons.util.sibling
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessage import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessage
import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessageCodec import com.vitorpamplona.quartz.cordn.spec02Envelopes.CordnDeliveredMessageCodec
import com.vitorpamplona.quartz.cordn.sync.EchoState import com.vitorpamplona.quartz.cordn.sync.EchoState
@@ -28,11 +33,13 @@ import com.vitorpamplona.quartz.cordn.sync.GroupCursor
import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.utils.Log import com.vitorpamplona.quartz.utils.Log
import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.IO
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
import java.io.File import okio.Buffer
import java.nio.ByteBuffer import okio.FileSystem
import okio.Path
import kotlin.io.encoding.Base64 import kotlin.io.encoding.Base64
import kotlin.io.encoding.ExperimentalEncodingApi import kotlin.io.encoding.ExperimentalEncodingApi
@@ -56,13 +63,13 @@ object CordnStorageLayout {
* a caller could smuggle `..` into a directory name. * a caller could smuggle `..` into a directory name.
*/ */
fun directoryFor( fun directoryFor(
root: File, root: Path,
accountPubKey: HexKey, accountPubKey: HexKey,
coordinatorPubKey: HexKey, coordinatorPubKey: HexKey,
): File { ): Path {
require(accountPubKey.matches(HEX)) { "account pubkey must be hex" } require(accountPubKey.matches(HEX)) { "account pubkey must be hex" }
require(coordinatorPubKey.matches(HEX)) { "coordinator pubkey must be hex" } require(coordinatorPubKey.matches(HEX)) { "coordinator pubkey must be hex" }
return File(root, "cordn/$accountPubKey/$coordinatorPubKey") return root / "cordn" / accountPubKey / coordinatorPubKey
} }
/** /**
@@ -74,11 +81,11 @@ object CordnStorageLayout {
* make that coordinator's removal delete the record of the others. * make that coordinator's removal delete the record of the others.
*/ */
fun accountDirectoryFor( fun accountDirectoryFor(
root: File, root: Path,
accountPubKey: HexKey, accountPubKey: HexKey,
): File { ): Path {
require(accountPubKey.matches(HEX)) { "account pubkey must be hex" } require(accountPubKey.matches(HEX)) { "account pubkey must be hex" }
return File(root, "cordn/$accountPubKey") return root / "cordn" / accountPubKey
} }
/** /**
@@ -126,19 +133,22 @@ private const val TAG = "CordnStores"
* is not a stale group, it is an unreadable one, and the group cannot be * is not a stale group, it is an unreadable one, and the group cannot be
* re-derived from anywhere else on this device. * re-derived from anywhere else on this device.
*/ */
private fun atomicWrite( private fun FileSystem.atomicWrite(
file: File, file: Path,
data: ByteArray, data: ByteArray,
) { ) {
file.parentFile?.mkdirs() file.parent?.let { createDirectories(it) }
val temp = File(file.parentFile, "${file.name}.tmp") val temp = file.sibling("${file.name}.tmp")
temp.writeBytes(data) write(temp) { write(data) }
if (!temp.renameTo(file)) { moveOrCopy(temp, file)
temp.copyTo(file, overwrite = true)
temp.delete()
}
} }
private fun FileSystem.readBytes(file: Path): ByteArray = read(file) { readByteArray() }
private fun FileSystem.isDirectory(path: Path): Boolean = metadataOrNull(path)?.isDirectory == true
private fun FileSystem.isRegularFile(path: Path): Boolean = metadataOrNull(path)?.isRegularFile == true
/** /**
* A [CordnGroupStore] on the filesystem, encrypted through [cipher]. * A [CordnGroupStore] on the filesystem, encrypted through [cipher].
* *
@@ -158,29 +168,30 @@ private fun atomicWrite(
* that both encrypt blobs, which is not an abstraction worth having. * that both encrypt blobs, which is not an abstraction worth having.
*/ */
class FileCordnGroupStore( class FileCordnGroupStore(
private val dir: File, private val dir: Path,
private val cipher: CordnBlobCipher, private val cipher: CordnBlobCipher,
private val fileSystem: FileSystem = platformFileSystem,
) : CordnGroupStore { ) : CordnGroupStore {
private fun groupDir(gid: String) = File(dir, "groups/${CordnStorageLayout.encodeKey(gid)}") private fun groupDir(gid: String) = dir / "groups" / CordnStorageLayout.encodeKey(gid)
private fun stateFile(gid: String) = File(groupDir(gid), "state") private fun stateFile(gid: String) = groupDir(gid) / "state"
private fun cursorFile(gid: String) = File(groupDir(gid), "cursor") private fun cursorFile(gid: String) = groupDir(gid) / "cursor"
private fun joinOriginFile(gid: String) = File(groupDir(gid), "via-request") private fun joinOriginFile(gid: String) = groupDir(gid) / "via-request"
private fun roomStateFile(gid: String) = File(groupDir(gid), "room") private fun roomStateFile(gid: String) = groupDir(gid) / "room"
private fun echoStateFile(gid: String) = File(groupDir(gid), "echoes") private fun echoStateFile(gid: String) = groupDir(gid) / "echoes"
/** /**
* Inside [groupDir] so [deleteGroup]'s recursive delete already covers it: * Inside [groupDir] so [deleteGroup]'s recursive delete already covers it:
* a group that left its history behind would keep the plaintext of an * a group that left its history behind would keep the plaintext of an
* end-to-end encrypted conversation after the key that read it was gone. * end-to-end encrypted conversation after the key that read it was gone.
*/ */
private fun messagesFile(gid: String) = File(groupDir(gid), "messages") private fun messagesFile(gid: String) = groupDir(gid) / "messages"
private fun summaryFile(gid: String) = File(groupDir(gid), "newest") private fun summaryFile(gid: String) = groupDir(gid) / "newest"
/** /**
* Appending a segment rather than rewriting the conversation, so the cost * Appending a segment rather than rewriting the conversation, so the cost
@@ -192,6 +203,7 @@ class FileCordnGroupStore(
// The log treats null as "this segment is unreadable" and carries on // The log treats null as "this segment is unreadable" and carries on
// with the rest, which is what one corrupt segment should cost. // with the rest, which is what one corrupt segment should cost.
decrypt = { runCatching { cipher.decrypt(it) }.getOrNull() }, decrypt = { runCatching { cipher.decrypt(it) }.getOrNull() },
fileSystem = fileSystem,
) )
private val messageLock = Mutex() private val messageLock = Mutex()
@@ -217,13 +229,13 @@ class FileCordnGroupStore(
gid: String, gid: String,
state: ByteArray, state: ByteArray,
) = withContext(Dispatchers.IO) { ) = withContext(Dispatchers.IO) {
atomicWrite(stateFile(gid), cipher.encrypt(state)) fileSystem.atomicWrite(stateFile(gid), cipher.encrypt(state))
} }
override suspend fun loadGroup(gid: String): ByteArray? = override suspend fun loadGroup(gid: String): ByteArray? =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
val file = stateFile(gid) val file = stateFile(gid)
if (!file.exists()) null else cipher.decrypt(file.readBytes()) if (!fileSystem.exists(file)) null else cipher.decrypt(fileSystem.readBytes(file))
} }
override suspend fun deleteGroup(gid: String) { override suspend fun deleteGroup(gid: String) {
@@ -238,35 +250,39 @@ class FileCordnGroupStore(
// The cursor goes with it. Leaving one behind would mean a later // The cursor goes with it. Leaving one behind would mean a later
// re-join of the same gid resumes from a cursor belonging to a // re-join of the same gid resumes from a cursor belonging to a
// group it is no longer in, skipping everything before it. // group it is no longer in, skipping everything before it.
groupDir(gid).deleteRecursively() fileSystem.deleteRecursivelyQuietly(groupDir(gid))
} }
} }
override suspend fun listGroups(): List<String> = override suspend fun listGroups(): List<String> =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
File(dir, "groups") fileSystem
.listFiles() .listOrNull(dir / "groups")
.orEmpty() .orEmpty()
.filter { it.isDirectory && File(it, "state").exists() } .filter { fileSystem.isDirectory(it) && fileSystem.exists(it / "state") }
.mapNotNull { CordnStorageLayout.decodeKey(it.name) } .mapNotNull { CordnStorageLayout.decodeKey(it.name) }
} }
// Two big-endian int64s, which is what ByteBuffer wrote: its default order
// is BIG_ENDIAN, and so is okio's Buffer.
override suspend fun saveCursor( override suspend fun saveCursor(
gid: String, gid: String,
cursor: GroupCursor, cursor: GroupCursor,
) = withContext(Dispatchers.IO) { ) = withContext(Dispatchers.IO) {
val buffer = ByteBuffer.allocate(16).putLong(cursor.fetchCursor).putLong(cursor.lastCursor) val bytes = Buffer().writeLong(cursor.fetchCursor).writeLong(cursor.lastCursor).readByteArray()
atomicWrite(cursorFile(gid), cipher.encrypt(buffer.array())) fileSystem.atomicWrite(cursorFile(gid), cipher.encrypt(bytes))
} }
override suspend fun loadCursor(gid: String): GroupCursor? = override suspend fun loadCursor(gid: String): GroupCursor? =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
val file = cursorFile(gid) val file = cursorFile(gid)
if (!file.exists()) return@withContext null if (!fileSystem.exists(file)) return@withContext null
val bytes = cipher.decrypt(file.readBytes()) val bytes = cipher.decrypt(fileSystem.readBytes(file))
if (bytes.size < 16) return@withContext null if (bytes.size < 16) return@withContext null
val buffer = ByteBuffer.wrap(bytes) // Only the first 16 bytes are read, as ByteBuffer.wrap did; a longer
GroupCursor(fetchCursor = buffer.long, lastCursor = buffer.long) // blob's tail is ignored rather than refused.
val buffer = Buffer().write(bytes, 0, 16)
GroupCursor(fetchCursor = buffer.readLong(), lastCursor = buffer.readLong())
} }
// Existence IS the flag, so there is nothing to encrypt and nothing to // Existence IS the flag, so there is nothing to encrypt and nothing to
@@ -275,11 +291,11 @@ class FileCordnGroupStore(
// later re-join of the same gid is a different admission. // later re-join of the same gid is a different admission.
override suspend fun saveJoinedViaRequest(gid: String) { override suspend fun saveJoinedViaRequest(gid: String) {
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
atomicWrite(joinOriginFile(gid), ByteArray(0)) fileSystem.atomicWrite(joinOriginFile(gid), ByteArray(0))
} }
} }
override suspend fun loadJoinedViaRequest(gid: String): Boolean = withContext(Dispatchers.IO) { joinOriginFile(gid).exists() } override suspend fun loadJoinedViaRequest(gid: String): Boolean = withContext(Dispatchers.IO) { fileSystem.exists(joinOriginFile(gid)) }
override suspend fun saveRoomState( override suspend fun saveRoomState(
gid: String, gid: String,
@@ -288,18 +304,18 @@ class FileCordnGroupStore(
// Deleted rather than blanked when there is nothing to remember, so an // Deleted rather than blanked when there is nothing to remember, so an
// emptied draft leaves no plaintext behind in an old file. // emptied draft leaves no plaintext behind in an old file.
if (state.isBlank) { if (state.isBlank) {
roomStateFile(gid).delete() fileSystem.deleteQuietly(roomStateFile(gid))
return@withContext return@withContext
} }
atomicWrite(roomStateFile(gid), cipher.encrypt(CordnRoomStateCodec.encode(state))) fileSystem.atomicWrite(roomStateFile(gid), cipher.encrypt(CordnRoomStateCodec.encode(state)))
} }
override suspend fun loadRoomState(gid: String): CordnRoomState = override suspend fun loadRoomState(gid: String): CordnRoomState =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
val file = roomStateFile(gid) val file = roomStateFile(gid)
if (!file.exists()) return@withContext CordnRoomState() if (!fileSystem.exists(file)) return@withContext CordnRoomState()
try { try {
CordnRoomStateCodec.decode(cipher.decrypt(file.readBytes())) CordnRoomStateCodec.decode(cipher.decrypt(fileSystem.readBytes(file)))
} catch (e: Exception) { } catch (e: Exception) {
CordnRoomState() CordnRoomState()
} }
@@ -316,12 +332,12 @@ class FileCordnGroupStore(
// this check that makes re-delivery free rather than duplicating. // this check that makes re-delivery free rather than duplicating.
if (!ids.add(message.envelope.id)) return@withContext if (!ids.add(message.envelope.id)) return@withContext
groupDir(gid).mkdirs() fileSystem.createDirectories(groupDir(gid))
messageLog.append(messagesFile(gid), CordnDeliveredMessageCodec.encode(message)) messageLog.append(messagesFile(gid), CordnDeliveredMessageCodec.encode(message))
// Written on the same beat, so the inbox preview cannot disagree // Written on the same beat, so the inbox preview cannot disagree
// with the room. Whole-blob rather than appended: it is one entry // with the room. Whole-blob rather than appended: it is one entry
// that is always overwritten. // that is always overwritten.
atomicWrite(summaryFile(gid), cipher.encrypt(CordnMessageSummaryCodec.encode(message, ids.size))) fileSystem.atomicWrite(summaryFile(gid), cipher.encrypt(CordnMessageSummaryCodec.encode(message, ids.size)))
} }
} }
@@ -337,9 +353,9 @@ class FileCordnGroupStore(
override suspend fun loadMessageSummary(gid: String): CordnMessageSummary? = override suspend fun loadMessageSummary(gid: String): CordnMessageSummary? =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
val file = summaryFile(gid) val file = summaryFile(gid)
if (!file.exists()) return@withContext null if (!fileSystem.exists(file)) return@withContext null
try { try {
CordnMessageSummaryCodec.decode(cipher.decrypt(file.readBytes())) CordnMessageSummaryCodec.decode(cipher.decrypt(fileSystem.readBytes(file)))
} catch (e: Exception) { } catch (e: Exception) {
// A summary is a derived convenience; losing one costs a preview // A summary is a derived convenience; losing one costs a preview
// line until the next message, not the history it summarises. // line until the next message, not the history it summarises.
@@ -355,18 +371,18 @@ class FileCordnGroupStore(
// Deleted when there is nothing pending, so the common steady state is // Deleted when there is nothing pending, so the common steady state is
// no file rather than an empty one. // no file rather than an empty one.
if (state.isEmpty) { if (state.isEmpty) {
echoStateFile(gid).delete() fileSystem.deleteQuietly(echoStateFile(gid))
return@withContext return@withContext
} }
atomicWrite(echoStateFile(gid), cipher.encrypt(EchoStateCodec.encode(state))) fileSystem.atomicWrite(echoStateFile(gid), cipher.encrypt(EchoStateCodec.encode(state)))
} }
override suspend fun loadEchoState(gid: String): EchoState = override suspend fun loadEchoState(gid: String): EchoState =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
val file = echoStateFile(gid) val file = echoStateFile(gid)
if (!file.exists()) return@withContext EchoState() if (!fileSystem.exists(file)) return@withContext EchoState()
try { try {
EchoStateCodec.decode(cipher.decrypt(file.readBytes())) EchoStateCodec.decode(cipher.decrypt(fileSystem.readBytes(file)))
} catch (e: Exception) { } catch (e: Exception) {
EchoState() EchoState()
} }
@@ -385,36 +401,37 @@ class FileCordnGroupStore(
* publishes different ones to different coordinators. * publishes different ones to different coordinators.
*/ */
class FileCordnKeyPackageStore( class FileCordnKeyPackageStore(
private val dir: File, private val dir: Path,
private val cipher: CordnBlobCipher, private val cipher: CordnBlobCipher,
private val fileSystem: FileSystem = platformFileSystem,
) : CordnKeyPackageStore { ) : CordnKeyPackageStore {
private fun bundleFile(keyPackageRef: String) = File(dir, "keypackages/${CordnStorageLayout.encodeKey(keyPackageRef)}") private fun bundleFile(keyPackageRef: String) = dir / "keypackages" / CordnStorageLayout.encodeKey(keyPackageRef)
override suspend fun save( override suspend fun save(
keyPackageRef: String, keyPackageRef: String,
bundle: ByteArray, bundle: ByteArray,
) = withContext(Dispatchers.IO) { ) = withContext(Dispatchers.IO) {
atomicWrite(bundleFile(keyPackageRef), cipher.encrypt(bundle)) fileSystem.atomicWrite(bundleFile(keyPackageRef), cipher.encrypt(bundle))
} }
override suspend fun load(keyPackageRef: String): ByteArray? = override suspend fun load(keyPackageRef: String): ByteArray? =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
val file = bundleFile(keyPackageRef) val file = bundleFile(keyPackageRef)
if (!file.exists()) null else cipher.decrypt(file.readBytes()) if (!fileSystem.exists(file)) null else cipher.decrypt(fileSystem.readBytes(file))
} }
override suspend fun delete(keyPackageRef: String) { override suspend fun delete(keyPackageRef: String) {
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
bundleFile(keyPackageRef).delete() fileSystem.deleteQuietly(bundleFile(keyPackageRef))
} }
} }
override suspend fun list(): List<String> = override suspend fun list(): List<String> =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
File(dir, "keypackages") fileSystem
.listFiles() .listOrNull(dir / "keypackages")
.orEmpty() .orEmpty()
.filter { it.isFile } .filter { fileSystem.isRegularFile(it) }
// No explicit ".tmp" exclusion, because the encoding already // No explicit ".tmp" exclusion, because the encoding already
// is one: '.' is not in the base64url alphabet, so a temp file // is one: '.' is not in the base64url alphabet, so a temp file
// left behind by a crashed write cannot decode to a ref and // left behind by a crashed write cannot decode to a ref and
@@ -435,22 +452,23 @@ class FileCordnKeyPackageStore(
* account, sitting above the per-coordinator directories it names. * account, sitting above the per-coordinator directories it names.
*/ */
class FileCordnCoordinatorStore( class FileCordnCoordinatorStore(
private val dir: File, private val dir: Path,
private val cipher: CordnBlobCipher, private val cipher: CordnBlobCipher,
private val fileSystem: FileSystem = platformFileSystem,
) : CordnCoordinatorStore { ) : CordnCoordinatorStore {
private val file get() = File(dir, "coordinators") private val file get() = dir / "coordinators"
override suspend fun save(configs: List<CoordinatorConfig>) = override suspend fun save(configs: List<CoordinatorConfig>) =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
atomicWrite(file, cipher.encrypt(CoordinatorListCodec.encode(configs))) fileSystem.atomicWrite(file, cipher.encrypt(CoordinatorListCodec.encode(configs)))
} }
override suspend fun load(): List<CoordinatorConfig> = override suspend fun load(): List<CoordinatorConfig> =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
val stored = file val stored = file
if (!stored.exists()) return@withContext emptyList() if (!fileSystem.exists(stored)) return@withContext emptyList()
try { try {
CoordinatorListCodec.decode(cipher.decrypt(stored.readBytes())) CoordinatorListCodec.decode(cipher.decrypt(fileSystem.readBytes(stored)))
} catch (e: Exception) { } catch (e: Exception) {
// A list written by a future build, or one the keystore can no // A list written by a future build, or one the keystore can no
// longer decrypt. Returning nothing loses the coordinators but // longer decrypt. Returning nothing loses the coordinators but
@@ -463,7 +481,8 @@ class FileCordnCoordinatorStore(
// the rooms are gone from the inbox and the invitations screen // the rooms are gone from the inbox and the invitations screen
// says there is nowhere to look. One line is the difference // says there is nowhere to look. One line is the difference
// between a diagnosable fault and a mystery. // between a diagnosable fault and a mystery.
Log.w(TAG, "could not read ${stored.length()} bytes of coordinators, losing them: ${e.message}", e) val size = fileSystem.metadataOrNull(stored)?.size ?: 0L
Log.w(TAG, "could not read $size bytes of coordinators, losing them: ${e.message}", e)
emptyList() emptyList()
} }
} }
@@ -478,19 +497,20 @@ class FileCordnCoordinatorStore(
* groups to another one, the exact fork the flag exists to prevent. * groups to another one, the exact fork the flag exists to prevent.
*/ */
class FileCordnHandoffStore( class FileCordnHandoffStore(
private val accountDir: File, private val accountDir: Path,
private val fileSystem: FileSystem = platformFileSystem,
) : CordnHandoffStore { ) : CordnHandoffStore {
override suspend fun load(): Boolean = marker().exists() override suspend fun load(): Boolean = fileSystem.exists(marker())
override suspend fun save(handedOff: Boolean) { override suspend fun save(handedOff: Boolean) {
val file = marker() val file = marker()
if (handedOff) { if (handedOff) {
file.parentFile?.mkdirs() file.parent?.let { fileSystem.createDirectories(it) }
file.writeBytes(ByteArray(0)) fileSystem.write(file) { }
} else { } else {
file.delete() fileSystem.deleteQuietly(file)
} }
} }
private fun marker() = File(accountDir, "handed-off") private fun marker() = accountDir / "handed-off"
} }
@@ -20,9 +20,12 @@
*/ */
package com.vitorpamplona.amethyst.commons.storage package com.vitorpamplona.amethyst.commons.storage
import java.io.File import com.vitorpamplona.amethyst.commons.util.moveOrCopy
import java.io.FileOutputStream import com.vitorpamplona.amethyst.commons.util.platformFileSystem
import java.io.RandomAccessFile import com.vitorpamplona.amethyst.commons.util.sibling
import okio.FileSystem
import okio.Path
import okio.use
/** /**
* An encrypted-at-rest log of UTF-8 entries that can be appended to in * An encrypted-at-rest log of UTF-8 entries that can be appended to in
@@ -74,11 +77,14 @@ import java.io.RandomAccessFile
* open. One unreadable segment then costs only its own entries; a throw would * open. One unreadable segment then costs only its own entries; a throw would
* abort the whole read, and a caller that reads an empty log can overwrite a * abort the whole read, and a caller that reads an empty log can overwrite a
* history that was merely unreadable. * history that was merely unreadable.
* @param fileSystem where the logs live. Durability rests on
* [okio.FileHandle.flush], which on the JVM and Android is an `fsync`.
*/ */
class EncryptedAppendLog( class EncryptedAppendLog(
private val encrypt: (ByteArray) -> ByteArray, private val encrypt: (ByteArray) -> ByteArray,
private val decrypt: (ByteArray) -> ByteArray?, private val decrypt: (ByteArray) -> ByteArray?,
private val compactAfterSegments: Int = COMPACT_AFTER_SEGMENTS, private val compactAfterSegments: Int = COMPACT_AFTER_SEGMENTS,
private val fileSystem: FileSystem = platformFileSystem,
) { ) {
/** Entries of one file, plus what it takes to append without re-reading it. */ /** Entries of one file, plus what it takes to append without re-reading it. */
private class LogState( private class LogState(
@@ -94,10 +100,11 @@ class EncryptedAppendLog(
var looseEntries: Int, var looseEntries: Int,
) )
private val logs = mutableMapOf<String, LogState>() /** Keyed by the normalized path, so two spellings of one file share an entry. */
private val logs = mutableMapOf<Path, LogState>()
private fun stateFor(file: File): LogState = private fun stateFor(file: Path): LogState =
logs.getOrPut(file.absolutePath) { logs.getOrPut(file.normalized()) {
val decoded = decodeFile(file) val decoded = decodeFile(file)
// A torn trailing segment is dropped from the FILE, not merely from // A torn trailing segment is dropped from the FILE, not merely from
@@ -105,9 +112,9 @@ class EncryptedAppendLog(
// a record that parsing always stops at, so every later append is // a record that parsing always stops at, so every later append is
// written and then never read back — the log goes silently // written and then never read back — the log goes silently
// write-only, for good. // write-only, for good.
if (decoded.headed && decoded.validLength < file.length()) { if (decoded.headed && decoded.validLength < fileLength(file)) {
try { try {
RandomAccessFile(file, "rw").use { it.setLength(decoded.validLength) } fileSystem.openReadWrite(file).use { it.resize(decoded.validLength) }
} catch (_: Exception) { } catch (_: Exception) {
// Best effort: the next fold rewrites the file wholesale and // Best effort: the next fold rewrites the file wholesale and
// resolves it anyway. // resolves it anyway.
@@ -125,17 +132,17 @@ class EncryptedAppendLog(
} }
/** Every entry in [file], oldest first. */ /** Every entry in [file], oldest first. */
fun readAll(file: File): List<String> = stateFor(file).entries.toList() fun readAll(file: Path): List<String> = stateFor(file).entries.toList()
/** Whether [entry] is already in [file], without reading it back from disk. */ /** Whether [entry] is already in [file], without reading it back from disk. */
fun contains( fun contains(
file: File, file: Path,
entry: String, entry: String,
): Boolean = entry in stateFor(file).seen ): Boolean = entry in stateFor(file).seen
/** Append one entry. Constant time, apart from a periodic fold. */ /** Append one entry. Constant time, apart from a periodic fold. */
fun append( fun append(
file: File, file: Path,
entry: String, entry: String,
) { ) {
val state = stateFor(file) val state = stateFor(file)
@@ -176,13 +183,13 @@ class EncryptedAppendLog(
* between. * between.
*/ */
private fun foldLooseTail( private fun foldLooseTail(
file: File, file: Path,
state: LogState, state: LogState,
) { ) {
if (state.looseEntries == 0) return if (state.looseEntries == 0) return
val prefix = ByteArray(state.foldedLength.toInt()) // Exactly foldedLength bytes, or an EOFException — what readFully did.
RandomAccessFile(file, "r").use { it.readFully(prefix) } val prefix = fileSystem.read(file) { readByteArray(state.foldedLength) }
val folded = encrypt(encodeEntries(state.entries.takeLast(state.looseEntries))) val folded = encrypt(encodeEntries(state.entries.takeLast(state.looseEntries)))
val out = ByteArray(prefix.size + 4 + folded.size) val out = ByteArray(prefix.size + 4 + folded.size)
@@ -200,26 +207,28 @@ class EncryptedAppendLog(
} }
private fun appendSegment( private fun appendSegment(
file: File, file: Path,
entries: List<String>, entries: List<String>,
) { ) {
file.parentFile?.mkdirs() file.parent?.let { fileSystem.createDirectories(it) }
val segment = encrypt(encodeEntries(entries)) val segment = encrypt(encodeEntries(entries))
FileOutputStream(file, true).use { out -> // Opened without truncating and written at its current end: an append.
out.write(lengthPrefix(segment.size)) fileSystem.openReadWrite(file).use { handle ->
out.write(segment) val end = handle.size()
handle.write(end, lengthPrefix(segment.size), 0, 4)
handle.write(end + 4, segment, 0, segment.size)
// An append that survives only in the page cache would lose a // An append that survives only in the page cache would lose a
// message the UI has already shown as sent. // message the UI has already shown as sent.
out.fd.sync() handle.flush()
} }
} }
/** Replace the whole log with [entries], as a single folded segment. */ /** Replace the whole log with [entries], as a single folded segment. */
fun rewrite( fun rewrite(
file: File, file: Path,
entries: List<String>, entries: List<String>,
) { ) {
file.parentFile?.mkdirs() file.parent?.let { fileSystem.createDirectories(it) }
val segment = encrypt(encodeEntries(entries)) val segment = encrypt(encodeEntries(entries))
val out = ByteArray(HEADER_LENGTH + 4 + segment.size) val out = ByteArray(HEADER_LENGTH + 4 + segment.size)
@@ -230,7 +239,7 @@ class EncryptedAppendLog(
atomicWrite(file, out) atomicWrite(file, out)
val state = val state =
logs.getOrPut(file.absolutePath) { logs.getOrPut(file.normalized()) {
LogState(mutableListOf(), mutableSetOf(), headed = true, foldedLength = 0L, looseSegments = 0, looseEntries = 0) LogState(mutableListOf(), mutableSetOf(), headed = true, foldedLength = 0L, looseSegments = 0, looseEntries = 0)
} }
state.entries.clear() state.entries.clear()
@@ -244,8 +253,8 @@ class EncryptedAppendLog(
} }
/** Drop the in-memory cache for [file]; call when the file is deleted. */ /** Drop the in-memory cache for [file]; call when the file is deleted. */
fun forget(file: File) { fun forget(file: Path) {
logs.remove(file.absolutePath) logs.remove(file.normalized())
} }
private class Decoded( private class Decoded(
@@ -258,11 +267,11 @@ class EncryptedAppendLog(
val looseEntries: Int, val looseEntries: Int,
) )
private fun decodeFile(file: File): Decoded { private fun decodeFile(file: Path): Decoded {
// A file that does not exist yet has no header, so the first append has // A file that does not exist yet has no header, so the first append has
// to write one rather than tack a bare segment onto nothing. // to write one rather than tack a bare segment onto nothing.
if (!file.exists()) return Decoded(emptyList(), false, 0L, 0L, 0, 0) if (!fileSystem.exists(file)) return Decoded(emptyList(), false, 0L, 0L, 0, 0)
val bytes = file.readBytes() val bytes = fileSystem.read(file) { readByteArray() }
if (!bytes.startsWithMagic() || bytes.size < HEADER_LENGTH) { if (!bytes.startsWithMagic() || bytes.size < HEADER_LENGTH) {
// The older format: the file is one encrypted blob and nothing else. // The older format: the file is one encrypted blob and nothing else.
@@ -347,20 +356,22 @@ class EncryptedAppendLog(
/** Write via a temp file and rename, so a crash can't leave a half-written log. */ /** Write via a temp file and rename, so a crash can't leave a half-written log. */
private fun atomicWrite( private fun atomicWrite(
target: File, target: Path,
data: ByteArray, data: ByteArray,
) { ) {
val tempFile = File(target.parentFile, "${target.name}.tmp") val tempFile = target.sibling("${target.name}.tmp")
FileOutputStream(tempFile).use { out -> fileSystem.openReadWrite(tempFile).use { handle ->
out.write(data) // Truncated first, as opening a FileOutputStream did: a stale temp
out.fd.sync() // longer than [data] must not leave its tail behind.
} handle.resize(0L)
if (!tempFile.renameTo(target)) { handle.write(0L, data, 0, data.size)
tempFile.copyTo(target, overwrite = true) handle.flush()
tempFile.delete()
} }
fileSystem.moveOrCopy(tempFile, target)
} }
private fun fileLength(file: Path): Long = fileSystem.metadataOrNull(file)?.size ?: 0L
private fun ByteArray.startsWithMagic(): Boolean { private fun ByteArray.startsWithMagic(): Boolean {
if (size < MAGIC.size) return false if (size < MAGIC.size) return false
for (i in MAGIC.indices) if (this[i] != MAGIC[i]) return false for (i in MAGIC.indices) if (this[i] != MAGIC[i]) return false
@@ -0,0 +1,95 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.util
import okio.FileSystem
import okio.IOException
import okio.Path
import okio.Path.Companion.toPath
/*
* The java.io.File idioms that okio spells differently, for stores moving off
* File without changing what they do on failure.
*
* okio reports a failed delete or rename by throwing; File reported it by
* returning false, and the stores were written around that. These keep the
* File behaviour, so a port is a port and not a quiet change in error handling.
*/
/** A file next to this one: `File(parentFile, name)`. */
fun Path.sibling(name: String): Path = parent?.div(name) ?: name.toPath()
/**
* Deletes a file or an empty directory, and never throws: `File.delete()`.
*
* @return true when [path] is gone, including when it never existed.
*/
fun FileSystem.deleteQuietly(path: Path): Boolean =
try {
delete(path)
true
} catch (e: IOException) {
!exists(path)
}
/**
* Deletes [path] and everything under it, carrying on past anything it cannot
* delete: `File.deleteRecursively()`.
*
* One difference, deliberately kept: symlinks are deleted, never followed.
* `File.deleteRecursively` walks into a linked directory and empties it, which
* no store here ever wanted.
*
* @return true when nothing is left, including when [path] never existed.
*/
fun FileSystem.deleteRecursivelyQuietly(path: Path): Boolean {
// Directories come before their children in this listing, so the reverse
// deletes every child before the directory holding it.
val children =
try {
listRecursively(path).toList()
} catch (e: IOException) {
emptyList()
}
var deleted = true
for (child in children.asReversed()) deleted = deleteQuietly(child) && deleted
return deleteQuietly(path) && deleted
}
/**
* Replaces [target] with [source]: an atomic rename, or, where the filesystem
* refuses one, a copy over [target] and a delete of [source].
*
* The fallback is what `if (!temp.renameTo(file)) { temp.copyTo(file, true); temp.delete() }`
* did. It is not atomic, and is only here so that a filesystem without rename
* still gets written.
*/
fun FileSystem.moveOrCopy(
source: Path,
target: Path,
) {
try {
atomicMove(source, target)
} catch (e: IOException) {
copy(source, target)
deleteQuietly(source)
}
}
@@ -29,6 +29,7 @@ import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
import okio.Path.Companion.toOkioPath
import java.io.File import java.io.File
/** /**
@@ -72,13 +73,13 @@ class EncryptedMarmotMessageStore(
logMutex.withLock { logMutex.withLock {
try { try {
val file = messagesFile(nostrGroupId) val file = messagesFile(nostrGroupId)
if (log.contains(file, innerEventJson)) { if (log.contains(file.toOkioPath(), innerEventJson)) {
Log.d(TAG) { "appendMessage($nostrGroupId): duplicate entry skipped" } Log.d(TAG) { "appendMessage($nostrGroupId): duplicate entry skipped" }
return@withLock return@withLock
} }
log.append(file, innerEventJson) log.append(file.toOkioPath(), innerEventJson)
Log.d(TAG) { Log.d(TAG) {
"appendMessage($nostrGroupId): now ${log.readAll(file).size} message(s) persisted" "appendMessage($nostrGroupId): now ${log.readAll(file.toOkioPath()).size} message(s) persisted"
} }
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "appendMessage($nostrGroupId) FAILED: ${e.message}", e) Log.e(TAG, "appendMessage($nostrGroupId) FAILED: ${e.message}", e)
@@ -105,7 +106,7 @@ class EncryptedMarmotMessageStore(
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
logMutex.withLock { logMutex.withLock {
for (file in listOf(messagesFile(nostrGroupId), epochsFile(nostrGroupId), snapshotFile(nostrGroupId), expiriesFile(nostrGroupId), epochRetentionsFile(nostrGroupId))) { for (file in listOf(messagesFile(nostrGroupId), epochsFile(nostrGroupId), snapshotFile(nostrGroupId), expiriesFile(nostrGroupId), epochRetentionsFile(nostrGroupId))) {
log.forget(file) log.forget(file.toOkioPath())
if (file.exists() && !file.delete()) { if (file.exists() && !file.delete()) {
Log.w(TAG) { "delete($nostrGroupId): failed to remove ${file.absolutePath}" } Log.w(TAG) { "delete($nostrGroupId): failed to remove ${file.absolutePath}" }
} }
@@ -340,7 +341,7 @@ class EncryptedMarmotMessageStore(
private fun readAll(nostrGroupId: String): List<String> = readAllFrom(messagesFile(nostrGroupId)) private fun readAll(nostrGroupId: String): List<String> = readAllFrom(messagesFile(nostrGroupId))
private fun readAllFrom(file: File): List<String> = log.readAll(file) private fun readAllFrom(file: File): List<String> = log.readAll(file.toOkioPath())
private fun writeAll( private fun writeAll(
nostrGroupId: String, nostrGroupId: String,
@@ -350,7 +351,7 @@ class EncryptedMarmotMessageStore(
private fun writeAllTo( private fun writeAllTo(
file: File, file: File,
messages: List<String>, messages: List<String>,
) = log.rewrite(file, messages) ) = log.rewrite(file.toOkioPath(), messages)
companion object { companion object {
private const val TAG = "EncryptedMarmotMessageStore" private const val TAG = "EncryptedMarmotMessageStore"
@@ -21,6 +21,7 @@
package com.vitorpamplona.amethyst.commons.cordn package com.vitorpamplona.amethyst.commons.cordn
import kotlinx.coroutines.test.runTest import kotlinx.coroutines.test.runTest
import okio.Path.Companion.toOkioPath
import org.junit.Assert.assertFalse import org.junit.Assert.assertFalse
import org.junit.Assert.assertThrows import org.junit.Assert.assertThrows
import org.junit.Assert.assertTrue import org.junit.Assert.assertTrue
@@ -106,5 +107,5 @@ class CordnHandoffStateTest {
assertTrue(first.handedOff.value) assertTrue(first.handedOff.value)
} }
private fun state() = CordnHandoffState(FileCordnHandoffStore(folder.root)) private fun state() = CordnHandoffState(FileCordnHandoffStore(folder.root.toOkioPath()))
} }
@@ -0,0 +1,95 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.cordn
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import kotlinx.coroutines.test.runTest
import okio.Path.Companion.toOkioPath
import java.util.Base64
import kotlin.io.path.createTempDirectory
import kotlin.random.Random
import kotlin.test.AfterTest
import kotlin.test.Test
import kotlin.test.assertEquals
/**
* A handoff snapshot's base64 must mean the same thing on both sides of an
* upgrade: a phone on an older build (which used `java.util.Base64`) hands off
* to one on this build, and the other way round.
*
* Every length mod 3, so every padding case, and each input both padded and
* unpadded, because `java.util.Base64.getDecoder()` accepted either.
*/
class CordnMigrationStoresBase64Test {
private object Identity : CordnBlobCipher {
override fun encrypt(bytes: ByteArray) = bytes
override fun decrypt(bytes: ByteArray) = bytes
}
private val rootFile = createTempDirectory("cordn-b64").toFile()
private val root = rootFile.toOkioPath()
private val account = "e".repeat(64)
private val coordinator = "b".repeat(64)
@AfterTest
fun cleanUp() {
rootFile.deleteRecursively()
}
@Test
fun `snapshot base64 is java util Base64 in both directions`() =
runTest {
val random = Random(7)
val blobs = (1..13).map { random.nextBytes(it) }
val javaEncoder = Base64.getEncoder()
val unpadded = Base64.getEncoder().withoutPadding()
val snapshot =
CordnMigrationSnapshot(
accountPubKey = account,
groups =
blobs.mapIndexed { i, blob ->
CordnMigrationGroup(
coordinatorPubKey = coordinator,
coordinatorRelays = listOf("wss://coord.example"),
gid = "g" + i.toString().padStart(2, '0'),
clientStateBase64 = if (i % 2 == 0) javaEncoder.encodeToString(blob) else unpadded.encodeToString(blob),
cursor = 0,
)
},
)
CordnMigrationStores.write(root, account, Identity, snapshot)
val read =
CordnMigrationStores.read(
root,
account,
Identity,
listOf(CoordinatorConfig(coordinator, listOf(RelayUrlNormalizer.normalizeOrNull("wss://coord.example")!!))),
)
assertEquals(
blobs.map { javaEncoder.encodeToString(it) },
read.groups.sortedBy { it.gid }.map { it.clientStateBase64 },
)
}
}
@@ -31,6 +31,7 @@ import com.vitorpamplona.quartz.cordn.sync.PendingEpochOperation
import com.vitorpamplona.quartz.nip01Core.core.HexKey import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import kotlinx.coroutines.test.runTest import kotlinx.coroutines.test.runTest
import okio.Path.Companion.toOkioPath
import java.io.ByteArrayOutputStream import java.io.ByteArrayOutputStream
import java.io.File import java.io.File
import java.util.Base64 import java.util.Base64
@@ -92,13 +93,15 @@ class CordnStorageFormatGoldenTest {
// ---- Construction glue: the only part of this file that follows the stores' API. ---- // ---- Construction glue: the only part of this file that follows the stores' API. ----
private fun groupStore(coordinator: HexKey) = FileCordnGroupStore(CordnStorageLayout.directoryFor(rootFile, account, coordinator), cipher) private val root = rootFile.toOkioPath()
private fun keyPackageStore(coordinator: HexKey) = FileCordnKeyPackageStore(CordnStorageLayout.directoryFor(rootFile, account, coordinator), cipher) private fun groupStore(coordinator: HexKey) = FileCordnGroupStore(CordnStorageLayout.directoryFor(root, account, coordinator), cipher)
private fun coordinatorStore() = FileCordnCoordinatorStore(CordnStorageLayout.accountDirectoryFor(rootFile, account), cipher) private fun keyPackageStore(coordinator: HexKey) = FileCordnKeyPackageStore(CordnStorageLayout.directoryFor(root, account, coordinator), cipher)
private fun handoffStore(owner: HexKey) = FileCordnHandoffStore(CordnStorageLayout.accountDirectoryFor(rootFile, owner)) private fun coordinatorStore() = FileCordnCoordinatorStore(CordnStorageLayout.accountDirectoryFor(root, account), cipher)
private fun handoffStore(owner: HexKey) = FileCordnHandoffStore(CordnStorageLayout.accountDirectoryFor(root, owner))
private fun appendLog(compactAfterSegments: Int) = private fun appendLog(compactAfterSegments: Int) =
EncryptedAppendLog( EncryptedAppendLog(
@@ -107,11 +110,11 @@ class CordnStorageFormatGoldenTest {
compactAfterSegments = compactAfterSegments, compactAfterSegments = compactAfterSegments,
) )
private fun logFile(name: String) = File(rootFile, "logs/$name") private fun logFile(name: String) = root / "logs" / name
private suspend fun migrationWrite(snapshot: CordnMigrationSnapshot) = CordnMigrationStores.write(rootFile, migratedAccount, cipher, snapshot) private suspend fun migrationWrite(snapshot: CordnMigrationSnapshot) = CordnMigrationStores.write(root, migratedAccount, cipher, snapshot)
private suspend fun migrationRead(configs: List<CoordinatorConfig>) = CordnMigrationStores.read(rootFile, migratedAccount, cipher, configs) private suspend fun migrationRead(configs: List<CoordinatorConfig>) = CordnMigrationStores.read(root, migratedAccount, cipher, configs)
// ---- End of construction glue. ---- // ---- End of construction glue. ----
@@ -27,6 +27,7 @@ import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.test.runTest import kotlinx.coroutines.test.runTest
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
import okio.Path.Companion.toOkioPath
import java.io.File import java.io.File
import kotlin.test.AfterTest import kotlin.test.AfterTest
import kotlin.test.Test import kotlin.test.Test
@@ -70,16 +71,17 @@ class FileCordnStoresTest {
it.delete() it.delete()
it.mkdirs() it.mkdirs()
} }
private val rootPath = root.toOkioPath()
private val cipher = XorCipher() private val cipher = XorCipher()
private val account = "a".repeat(64) private val account = "a".repeat(64)
private val coordinator = "b".repeat(64) private val coordinator = "b".repeat(64)
private fun dir() = CordnStorageLayout.directoryFor(root, account, coordinator) private fun dir() = CordnStorageLayout.directoryFor(rootPath, account, coordinator)
private fun groups() = FileCordnGroupStore(dir(), cipher) private fun groups() = FileCordnGroupStore(dir(), cipher)
private fun coordinators() = FileCordnCoordinatorStore(CordnStorageLayout.accountDirectoryFor(root, account), cipher) private fun coordinators() = FileCordnCoordinatorStore(CordnStorageLayout.accountDirectoryFor(rootPath, account), cipher)
private fun relay(url: String) = RelayUrlNormalizer.normalizeOrNull(url)!! private fun relay(url: String) = RelayUrlNormalizer.normalizeOrNull(url)!!
@@ -110,7 +112,12 @@ class FileCordnStoresTest {
val state = "ratchet-tree-and-epoch-secrets".encodeToByteArray() val state = "ratchet-tree-and-epoch-secrets".encodeToByteArray()
groups().saveGroup("g1", state) groups().saveGroup("g1", state)
val onDisk = dir().walkTopDown().filter { it.isFile }.toList() val onDisk =
dir()
.toFile()
.walkTopDown()
.filter { it.isFile }
.toList()
assertTrue(onDisk.isNotEmpty(), "nothing was written at all") assertTrue(onDisk.isNotEmpty(), "nothing was written at all")
onDisk.forEach { onDisk.forEach {
val raw = it.readBytes() val raw = it.readBytes()
@@ -138,7 +145,7 @@ class FileCordnStoresTest {
val escaped = File(root, "cordn/etc").exists() || File(root.parentFile, "etc").exists() val escaped = File(root, "cordn/etc").exists() || File(root.parentFile, "etc").exists()
assertTrue(!escaped, "the write escaped the store directory") assertTrue(!escaped, "the write escaped the store directory")
assertTrue(dir().walkTopDown().any { it.isFile }, "it was not written anywhere at all") assertTrue(dir().toFile().walkTopDown().any { it.isFile }, "it was not written anywhere at all")
} }
@Test @Test
@@ -159,7 +166,7 @@ class FileCordnStoresTest {
runTest { runTest {
// The rule the layout exists for. Two coordinators can both serve // The rule the layout exists for. Two coordinators can both serve
// gid "shared" as unrelated groups with different ratchet trees. // gid "shared" as unrelated groups with different ratchet trees.
val other = CordnStorageLayout.directoryFor(root, account, "c".repeat(64)) val other = CordnStorageLayout.directoryFor(rootPath, account, "c".repeat(64))
FileCordnGroupStore(dir(), cipher).saveGroup("shared", byteArrayOf(1)) FileCordnGroupStore(dir(), cipher).saveGroup("shared", byteArrayOf(1))
FileCordnGroupStore(other, cipher).saveGroup("shared", byteArrayOf(2)) FileCordnGroupStore(other, cipher).saveGroup("shared", byteArrayOf(2))
@@ -170,7 +177,7 @@ class FileCordnStoresTest {
@Test @Test
fun `two accounts on one device do not see each other`() = fun `two accounts on one device do not see each other`() =
runTest { runTest {
val theirs = CordnStorageLayout.directoryFor(root, "d".repeat(64), coordinator) val theirs = CordnStorageLayout.directoryFor(rootPath, "d".repeat(64), coordinator)
groups().saveGroup("g1", byteArrayOf(1)) groups().saveGroup("g1", byteArrayOf(1))
assertTrue(FileCordnGroupStore(theirs, cipher).listGroups().isEmpty()) assertTrue(FileCordnGroupStore(theirs, cipher).listGroups().isEmpty())
@@ -179,8 +186,8 @@ class FileCordnStoresTest {
@Test @Test
fun `a non-hex pubkey never becomes a directory`() { fun `a non-hex pubkey never becomes a directory`() {
assertFailsWith<IllegalArgumentException> { CordnStorageLayout.directoryFor(root, "../escape", coordinator) } assertFailsWith<IllegalArgumentException> { CordnStorageLayout.directoryFor(rootPath, "../escape", coordinator) }
assertFailsWith<IllegalArgumentException> { CordnStorageLayout.directoryFor(root, account, "../escape") } assertFailsWith<IllegalArgumentException> { CordnStorageLayout.directoryFor(rootPath, account, "../escape") }
} }
@Test @Test
@@ -248,7 +255,7 @@ class FileCordnStoresTest {
// encoding — swap it for something that passes names through and // encoding — swap it for something that passes names through and
// this fails, which is the regression worth catching. // this fails, which is the regression worth catching.
keyPackages().save("ref1", byteArrayOf(1)) keyPackages().save("ref1", byteArrayOf(1))
File(dir(), "keypackages/${CordnStorageLayout.encodeKey("ref2")}.tmp").writeBytes(byteArrayOf(2)) File(dir().toFile(), "keypackages/${CordnStorageLayout.encodeKey("ref2")}.tmp").writeBytes(byteArrayOf(2))
assertEquals(listOf("ref1"), keyPackages().list()) assertEquals(listOf("ref1"), keyPackages().list())
} }
@@ -259,8 +266,8 @@ class FileCordnStoresTest {
// Ignoring one unreadable name is better than failing the listing // Ignoring one unreadable name is better than failing the listing
// and hiding every real group behind it. // and hiding every real group behind it.
groups().saveGroup("g1", byteArrayOf(1)) groups().saveGroup("g1", byteArrayOf(1))
File(dir(), "groups/not-base64-@@@").mkdirs() File(dir().toFile(), "groups/not-base64-@@@").mkdirs()
File(dir(), "groups/not-base64-@@@/state").writeBytes(byteArrayOf(0)) File(dir().toFile(), "groups/not-base64-@@@/state").writeBytes(byteArrayOf(0))
assertEquals(listOf("g1"), groups().listGroups()) assertEquals(listOf("g1"), groups().listGroups())
} }
@@ -298,7 +305,7 @@ class FileCordnStoresTest {
assertContentEquals(byteArrayOf(2, 2), store.loadGroup("g1")) assertContentEquals(byteArrayOf(2, 2), store.loadGroup("g1"))
assertTrue( assertTrue(
dir().walkTopDown().none { it.name.endsWith(".tmp") }, dir().toFile().walkTopDown().none { it.name.endsWith(".tmp") },
"a temp file survived the rename", "a temp file survived the rename",
) )
} }
@@ -314,7 +321,7 @@ class FileCordnStoresTest {
coordinators().save(configs) coordinators().save(configs)
assertEquals(configs, coordinators().load()) assertEquals(configs, coordinators().load())
val onDisk = File(CordnStorageLayout.accountDirectoryFor(root, account), "coordinators").readBytes() val onDisk = File(CordnStorageLayout.accountDirectoryFor(rootPath, account).toFile(), "coordinators").readBytes()
assertFalse(onDisk.decodeToString().contains(coordinator), "the pubkey went to disk in the clear") assertFalse(onDisk.decodeToString().contains(coordinator), "the pubkey went to disk in the clear")
} }
@@ -340,7 +347,7 @@ class FileCordnStoresTest {
fun `an unreadable list loses the coordinators rather than the login`() = fun `an unreadable list loses the coordinators rather than the login`() =
runTest { runTest {
coordinators().save(listOf(CoordinatorConfig(coordinator, listOf(relay("wss://one.example.com"))))) coordinators().save(listOf(CoordinatorConfig(coordinator, listOf(relay("wss://one.example.com")))))
File(CordnStorageLayout.accountDirectoryFor(root, account), "coordinators").writeBytes(byteArrayOf(9, 9, 9)) File(CordnStorageLayout.accountDirectoryFor(rootPath, account).toFile(), "coordinators").writeBytes(byteArrayOf(9, 9, 9))
assertEquals(emptyList(), coordinators().load()) assertEquals(emptyList(), coordinators().load())
} }
@@ -350,8 +357,8 @@ class FileCordnStoresTest {
runTest { runTest {
// If it lived inside a coordinator's own directory, removing that // If it lived inside a coordinator's own directory, removing that
// coordinator would delete the record of every other one with it. // coordinator would delete the record of every other one with it.
val listFile = File(CordnStorageLayout.accountDirectoryFor(root, account), "coordinators") val listFile = File(CordnStorageLayout.accountDirectoryFor(rootPath, account).toFile(), "coordinators")
assertEquals(dir().parentFile, listFile.parentFile) assertEquals(dir().toFile().parentFile, listFile.parentFile)
} }
@Test @Test
@@ -363,15 +370,15 @@ class FileCordnStoresTest {
// same gid and deleting a level too high takes both. // same gid and deleting a level too high takes both.
val other = "e".repeat(64) val other = "e".repeat(64)
groups().saveGroup("shared-gid", byteArrayOf(1, 2, 3)) groups().saveGroup("shared-gid", byteArrayOf(1, 2, 3))
FileCordnGroupStore(CordnStorageLayout.directoryFor(root, account, other), cipher) FileCordnGroupStore(CordnStorageLayout.directoryFor(rootPath, account, other), cipher)
.saveGroup("shared-gid", byteArrayOf(4, 5, 6)) .saveGroup("shared-gid", byteArrayOf(4, 5, 6))
CordnStorageLayout.directoryFor(root, account, coordinator).deleteRecursively() CordnStorageLayout.directoryFor(rootPath, account, coordinator).toFile().deleteRecursively()
assertNull(groups().loadGroup("shared-gid")) assertNull(groups().loadGroup("shared-gid"))
assertContentEquals( assertContentEquals(
byteArrayOf(4, 5, 6), byteArrayOf(4, 5, 6),
FileCordnGroupStore(CordnStorageLayout.directoryFor(root, account, other), cipher).loadGroup("shared-gid"), FileCordnGroupStore(CordnStorageLayout.directoryFor(rootPath, account, other), cipher).loadGroup("shared-gid"),
) )
} }
@@ -385,7 +392,7 @@ class FileCordnStoresTest {
coordinators().save(listOf(CoordinatorConfig(coordinator, listOf(relay("wss://one.example.com"))))) coordinators().save(listOf(CoordinatorConfig(coordinator, listOf(relay("wss://one.example.com")))))
groups().saveGroup("gid", byteArrayOf(1)) groups().saveGroup("gid", byteArrayOf(1))
CordnStorageLayout.directoryFor(root, account, coordinator).deleteRecursively() CordnStorageLayout.directoryFor(rootPath, account, coordinator).toFile().deleteRecursively()
assertEquals(1, coordinators().load().size) assertEquals(1, coordinators().load().size)
} }
@@ -415,7 +422,7 @@ class FileCordnStoresTest {
store.saveRoomState("gid", CordnRoomState(draft = "the secret plan", lastReadCursor = 7)) store.saveRoomState("gid", CordnRoomState(draft = "the secret plan", lastReadCursor = 7))
assertEquals(CordnRoomState("the secret plan", 7), store.loadRoomState("gid")) assertEquals(CordnRoomState("the secret plan", 7), store.loadRoomState("gid"))
val onDisk = File(dir(), "groups/${CordnStorageLayout.encodeKey("gid")}/room").readBytes() val onDisk = File(dir().toFile(), "groups/${CordnStorageLayout.encodeKey("gid")}/room").readBytes()
assertFalse(onDisk.decodeToString().contains("the secret plan"), "a draft went to disk in the clear") assertFalse(onDisk.decodeToString().contains("the secret plan"), "a draft went to disk in the clear")
} }
@@ -430,7 +437,7 @@ class FileCordnStoresTest {
store.saveRoomState("gid", CordnRoomState(draft = "typed")) store.saveRoomState("gid", CordnRoomState(draft = "typed"))
store.saveRoomState("gid", CordnRoomState()) store.saveRoomState("gid", CordnRoomState())
assertFalse(File(dir(), "groups/${CordnStorageLayout.encodeKey("gid")}/room").exists()) assertFalse(File(dir().toFile(), "groups/${CordnStorageLayout.encodeKey("gid")}/room").exists())
assertEquals(CordnRoomState(), store.loadRoomState("gid")) assertEquals(CordnRoomState(), store.loadRoomState("gid"))
} }
@@ -440,7 +447,7 @@ class FileCordnStoresTest {
val store = groups() val store = groups()
store.saveGroup("gid", byteArrayOf(1)) store.saveGroup("gid", byteArrayOf(1))
store.saveRoomState("gid", CordnRoomState(draft = "typed", lastReadCursor = 3)) store.saveRoomState("gid", CordnRoomState(draft = "typed", lastReadCursor = 3))
File(dir(), "groups/${CordnStorageLayout.encodeKey("gid")}/room").writeBytes(byteArrayOf(7, 7, 7)) File(dir().toFile(), "groups/${CordnStorageLayout.encodeKey("gid")}/room").writeBytes(byteArrayOf(7, 7, 7))
assertEquals(CordnRoomState(), store.loadRoomState("gid")) assertEquals(CordnRoomState(), store.loadRoomState("gid"))
assertNotNull(store.loadGroup("gid")) assertNotNull(store.loadGroup("gid"))
@@ -20,6 +20,7 @@
*/ */
package com.vitorpamplona.amethyst.commons.storage package com.vitorpamplona.amethyst.commons.storage
import okio.Path.Companion.toOkioPath
import java.io.File import java.io.File
import java.nio.file.Files import java.nio.file.Files
import kotlin.test.Test import kotlin.test.Test
@@ -61,39 +62,39 @@ class EncryptedAppendLogTest {
fun `entries survive a round trip through a fresh reader`() { fun `entries survive a round trip through a fresh reader`() {
val file = tempFile() val file = tempFile()
val writer = cipher() val writer = cipher()
writer.append(file, "first") writer.append(file.toOkioPath(), "first")
writer.append(file, "second") writer.append(file.toOkioPath(), "second")
writer.append(file, "third") writer.append(file.toOkioPath(), "third")
// A second instance reads only what reached the disk. // A second instance reads only what reached the disk.
assertContentEquals(listOf("first", "second", "third"), cipher().readAll(file)) assertContentEquals(listOf("first", "second", "third"), cipher().readAll(file.toOkioPath()))
} }
@Test @Test
fun `an empty log reads as empty rather than failing`() { fun `an empty log reads as empty rather than failing`() {
assertContentEquals(emptyList(), cipher().readAll(tempFile())) assertContentEquals(emptyList(), cipher().readAll(tempFile().toOkioPath()))
} }
@Test @Test
fun `contains answers without reading the file back`() { fun `contains answers without reading the file back`() {
val file = tempFile() val file = tempFile()
val log = cipher() val log = cipher()
log.append(file, "hello") log.append(file.toOkioPath(), "hello")
assertTrue(log.contains(file, "hello")) assertTrue(log.contains(file.toOkioPath(), "hello"))
assertFalse(log.contains(file, "goodbye")) assertFalse(log.contains(file.toOkioPath(), "goodbye"))
} }
@Test @Test
fun `a rewrite replaces the whole log`() { fun `a rewrite replaces the whole log`() {
val file = tempFile() val file = tempFile()
val log = cipher() val log = cipher()
log.append(file, "a") log.append(file.toOkioPath(), "a")
log.append(file, "b") log.append(file.toOkioPath(), "b")
log.rewrite(file, listOf("b")) log.rewrite(file.toOkioPath(), listOf("b"))
assertContentEquals(listOf("b"), cipher().readAll(file)) assertContentEquals(listOf("b"), cipher().readAll(file.toOkioPath()))
assertFalse(cipher().contains(file, "a")) assertFalse(cipher().contains(file.toOkioPath(), "a"))
} }
@Test @Test
@@ -101,9 +102,9 @@ class EncryptedAppendLogTest {
val file = tempFile() val file = tempFile()
val log = cipher(compactAfter = 4) val log = cipher(compactAfter = 4)
val written = (1..20).map { "entry-$it" } val written = (1..20).map { "entry-$it" }
written.forEach { log.append(file, it) } written.forEach { log.append(file.toOkioPath(), it) }
assertContentEquals(written, cipher().readAll(file)) assertContentEquals(written, cipher().readAll(file.toOkioPath()))
// Folding every four appends leaves ~5 segments rather than 20, which is // Folding every four appends leaves ~5 segments rather than 20, which is
// the bound on read cost that matters here. // the bound on read cost that matters here.
assertTrue(file.length() < 20 * SEGMENT_OVERHEAD_CEILING, "log should have been compacted, was ${file.length()} bytes") assertTrue(file.length() < 20 * SEGMENT_OVERHEAD_CEILING, "log should have been compacted, was ${file.length()} bytes")
@@ -116,10 +117,10 @@ class EncryptedAppendLogTest {
// The first append lays the header down; the next four are the loose run // The first append lays the header down; the next four are the loose run
// that the fourth of them collapses. // that the fourth of them collapses.
(1..5).forEach { log.append(file, "first-run-$it") } (1..5).forEach { log.append(file.toOkioPath(), "first-run-$it") }
val afterFirstFold = file.readBytes() val afterFirstFold = file.readBytes()
(1..4).forEach { log.append(file, "second-run-$it") } (1..4).forEach { log.append(file.toOkioPath(), "second-run-$it") }
val afterSecondFold = file.readBytes() val afterSecondFold = file.readBytes()
// The SEGMENTS the first fold produced must survive verbatim. If they // The SEGMENTS the first fold produced must survive verbatim. If they
@@ -135,7 +136,7 @@ class EncryptedAppendLogTest {
afterSecondFold.copyOfRange(HEADER_LEN, afterFirstFold.size).toList(), afterSecondFold.copyOfRange(HEADER_LEN, afterFirstFold.size).toList(),
"folding the tail must copy the already-folded segments as ciphertext", "folding the tail must copy the already-folded segments as ciphertext",
) )
assertContentEquals((1..5).map { "first-run-$it" } + (1..4).map { "second-run-$it" }, cipher().readAll(file)) assertContentEquals((1..5).map { "first-run-$it" } + (1..4).map { "second-run-$it" }, cipher().readAll(file.toOkioPath()))
} }
@Test @Test
@@ -145,7 +146,7 @@ class EncryptedAppendLogTest {
// Exactly what the old writer produced: one encrypted blob, no magic. // Exactly what the old writer produced: one encrypted blob, no magic.
file.writeBytes(legacyBlob(listOf("old-one", "old-two"))) file.writeBytes(legacyBlob(listOf("old-one", "old-two")))
assertContentEquals(listOf("old-one", "old-two"), cipher().readAll(file)) assertContentEquals(listOf("old-one", "old-two"), cipher().readAll(file.toOkioPath()))
} }
@Test @Test
@@ -155,9 +156,9 @@ class EncryptedAppendLogTest {
file.writeBytes(legacyBlob(listOf("old-one", "old-two"))) file.writeBytes(legacyBlob(listOf("old-one", "old-two")))
val log = cipher() val log = cipher()
log.append(file, "new-one") log.append(file.toOkioPath(), "new-one")
assertContentEquals(listOf("old-one", "old-two", "new-one"), cipher().readAll(file)) assertContentEquals(listOf("old-one", "old-two", "new-one"), cipher().readAll(file.toOkioPath()))
// and the upgraded file is in the new format, so the next append is cheap // and the upgraded file is in the new format, so the next append is cheap
assertTrue(file.readBytes().decodeToString().startsWith("MRMTLOG3")) assertTrue(file.readBytes().decodeToString().startsWith("MRMTLOG3"))
} }
@@ -166,31 +167,31 @@ class EncryptedAppendLogTest {
fun `a torn final append costs only the torn entry`() { fun `a torn final append costs only the torn entry`() {
val file = tempFile() val file = tempFile()
val log = cipher() val log = cipher()
log.append(file, "kept-one") log.append(file.toOkioPath(), "kept-one")
log.append(file, "kept-two") log.append(file.toOkioPath(), "kept-two")
// Simulate process death partway through writing the third segment. // Simulate process death partway through writing the third segment.
val intact = file.readBytes() val intact = file.readBytes()
log.append(file, "lost") log.append(file.toOkioPath(), "lost")
val torn = file.readBytes() val torn = file.readBytes()
file.writeBytes(torn.copyOfRange(0, intact.size + 6)) file.writeBytes(torn.copyOfRange(0, intact.size + 6))
assertContentEquals(listOf("kept-one", "kept-two"), cipher().readAll(file)) assertContentEquals(listOf("kept-one", "kept-two"), cipher().readAll(file.toOkioPath()))
} }
@Test @Test
fun `a segment that cannot be decrypted does not hide the rest`() { fun `a segment that cannot be decrypted does not hide the rest`() {
val file = tempFile() val file = tempFile()
val log = cipher() val log = cipher()
log.append(file, "before") log.append(file.toOkioPath(), "before")
log.append(file, "after") log.append(file.toOkioPath(), "after")
// Corrupt the first segment's nonce marker so decrypt returns null for it. // Corrupt the first segment's nonce marker so decrypt returns null for it.
val bytes = file.readBytes() val bytes = file.readBytes()
bytes[HEADER_LEN + 4] = 0 bytes[HEADER_LEN + 4] = 0
file.writeBytes(bytes) file.writeBytes(bytes)
assertContentEquals(listOf("after"), cipher().readAll(file)) assertContentEquals(listOf("after"), cipher().readAll(file.toOkioPath()))
} }
@Test @Test
@@ -201,21 +202,21 @@ class EncryptedAppendLogTest {
// writes and never read one back again — silently write-only, for good. // writes and never read one back again — silently write-only, for good.
val file = tempFile() val file = tempFile()
val log = cipher() val log = cipher()
log.append(file, "kept-one") log.append(file.toOkioPath(), "kept-one")
log.append(file, "kept-two") log.append(file.toOkioPath(), "kept-two")
val intact = file.readBytes() val intact = file.readBytes()
log.append(file, "lost") log.append(file.toOkioPath(), "lost")
val torn = file.readBytes() val torn = file.readBytes()
file.writeBytes(torn.copyOfRange(0, intact.size + 6)) file.writeBytes(torn.copyOfRange(0, intact.size + 6))
val reopened = cipher() val reopened = cipher()
assertContentEquals(listOf("kept-one", "kept-two"), reopened.readAll(file)) assertContentEquals(listOf("kept-one", "kept-two"), reopened.readAll(file.toOkioPath()))
reopened.append(file, "three") reopened.append(file.toOkioPath(), "three")
assertContentEquals( assertContentEquals(
listOf("kept-one", "kept-two", "three"), listOf("kept-one", "kept-two", "three"),
cipher().readAll(file), cipher().readAll(file.toOkioPath()),
"an append after a torn tail must be readable by the next reader", "an append after a torn tail must be readable by the next reader",
) )
} }
@@ -230,11 +231,11 @@ class EncryptedAppendLogTest {
repeat(12) { session -> repeat(12) { session ->
// A fresh instance per session, as a cold process gets. // A fresh instance per session, as a cold process gets.
val log = cipher(compactAfter = 4) val log = cipher(compactAfter = 4)
repeat(3) { log.append(file, "s$session-m$it") } repeat(3) { log.append(file.toOkioPath(), "s$session-m$it") }
} }
val expected = (0 until 12).flatMap { session -> (0 until 3).map { "s$session-m$it" } } val expected = (0 until 12).flatMap { session -> (0 until 3).map { "s$session-m$it" } }
assertContentEquals(expected, cipher().readAll(file)) assertContentEquals(expected, cipher().readAll(file.toOkioPath()))
// 36 entries at a fold every 4 appends: a handful of segments, not 36. // 36 entries at a fold every 4 appends: a handful of segments, not 36.
assertTrue( assertTrue(
@@ -256,8 +257,8 @@ class EncryptedAppendLogTest {
decrypt = { null }, decrypt = { null },
) )
runCatching { log.append(file, "never-written") } runCatching { log.append(file.toOkioPath(), "never-written") }
assertFalse(log.contains(file, "never-written"), "a write that failed must not count as persisted") assertFalse(log.contains(file.toOkioPath(), "never-written"), "a write that failed must not count as persisted")
} }
@Test @Test
@@ -265,9 +266,9 @@ class EncryptedAppendLogTest {
val file = tempFile() val file = tempFile()
val log = cipher() val log = cipher()
val awkward = listOf("", "emoji 👩‍👧 here", "a\nb\tc", "\"quoted\": {\"json\": 1}") val awkward = listOf("", "emoji 👩‍👧 here", "a\nb\tc", "\"quoted\": {\"json\": 1}")
awkward.forEach { log.append(file, it) } awkward.forEach { log.append(file.toOkioPath(), it) }
assertEquals(awkward, cipher().readAll(file)) assertEquals(awkward, cipher().readAll(file.toOkioPath()))
} }
/** The pre-segment on-disk shape: `encrypt(uint32 count, (uint32 len, bytes)*)`. */ /** The pre-segment on-disk shape: `encrypt(uint32 count, (uint32 len, bytes)*)`. */