Demo model-layer e2e harness (Synapse-driven) for the M4 both-strategies criterion #52

Closed
opened 2026-08-06 22:59:45 +02:00 by thecrealm · 1 comment
Owner

While verifying #47, a local-only jvmTest drove AppModel/RoomTimeline headlessly against SynapseContainer: login, E2EE send/echo, reactions, edits, relaunch cold-start-from-cache, offline scrollback with decrypt-on-read (#51). Per D30 the #47 CI gate stays build-only, so the driver was not committed. M4's exit criterion 'the M3 demo runs on both sync strategies with zero app-code changes' wants exactly this harness, parameterized over SyncStrategy (KATRIX_REQUIRE_DOCKER-gated). Resurrect it from the #47 close notes when M4 starts.

While verifying #47, a local-only jvmTest drove AppModel/RoomTimeline headlessly against SynapseContainer: login, E2EE send/echo, reactions, edits, relaunch cold-start-from-cache, offline scrollback with decrypt-on-read (#51). Per D30 the #47 CI gate stays build-only, so the driver was not committed. M4's exit criterion 'the M3 demo runs on both sync strategies with zero app-code changes' wants exactly this harness, parameterized over SyncStrategy (KATRIX_REQUIRE_DOCKER-gated). Resurrect it from the #47 close notes when M4 starts.
Author
Owner

The local-only verification driver used to verify #47 (per D30 the CI gate stays build-only). Needs jvmTest.dependencies { implementation(project(":katrix-testkit")) } in samples/demo and the docker env from the dev-Mac notes. Parameterize over SyncStrategy for the M4 criterion.

package katrix.demo

import de.centoria.katrix.core.id.RoomId
import de.centoria.katrix.testkit.SynapseContainer
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.launch
import kotlinx.coroutines.delay
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeout
import kotlinx.serialization.json.buildJsonObject
import kotlinx.serialization.json.put
import kotlinx.serialization.json.addJsonObject
import kotlinx.serialization.json.putJsonArray
import kotlinx.serialization.json.putJsonObject
import java.nio.file.Files
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.test.fail

/**
 * LOCAL-ONLY verification driver for #47 (not committed; the CI gate is
 * build-only per D30): drives the demo's model layer — the same code the
 * Compose UI calls — against a real Synapse. Covers: login, E2EE send,
 * echo, reactions, edits, relaunch cold-start-from-cache, offline
 * scrollback decrypt-on-read.
 */
class DemoModelEndToEndTest {

    @Test
    fun demo_model_flow_end_to_end() {
        if (!SynapseContainer.dockerAvailable()) {
            if (System.getenv("KATRIX_REQUIRE_DOCKER") == "1") fail("docker required")
            println("skipping: no docker")
            return
        }
        val home = Files.createTempDirectory("katrix-demo-e2e")
        System.setProperty("katrix.demo.home", home.toString())

        SynapseContainer().use { synapse ->
            synapse.start()
            synapse.registerUser("alice", "alice-pass-1")

            // ── first launch: login through the model ────────────────────
            val scope1 = CoroutineScope(SupervisorJob() + uiThread())
            val model1 = AppModel(scope1)
            model1.start()
            assertEquals(Screen.Login, model1.screen, "fresh profile starts at login")

            model1.login(synapse.baseUrl, "alice", "alice-pass-1")
            awaitTrue("session connects after login") { model1.session != null }
            val session = model1.session!!

            // an encrypted room, created through the session
            val roomId = runBlocking {
                session.createRoom(
                    buildJsonObject {
                        put("name", "demo e2e room")
                        put("preset", "private_chat")
                        putJsonArray("initial_state") {
                            addJsonObject {
                                put("type", "m.room.encryption")
                                put("state_key", "")
                                putJsonObject("content") {
                                    put("algorithm", "m.megolm.v1.aes-sha2")
                                }
                            }
                        }
                    },
                )
            }
            awaitTrue("room list shows the new room") {
                model1.rooms.any { it.roomId == roomId.full }
            }

            // bus tap: everything the change bus delivers for the room
            val busLog = java.util.Collections.synchronizedList(mutableListOf<String>())
            val tap = CoroutineScope(SupervisorJob() + Dispatchers.Default).launch {
                session.events.collect { e ->
                    if (e.roomId == roomId && !e.isState) {
                        val txn = (e.raw.unsigned?.get("transaction_id") as? kotlinx.serialization.json.JsonPrimitive)?.content
                        busLog += "bus type=${e.type} id=${e.raw.eventId?.full} txn=$txn"
                    }
                }
            }

            // ── E2EE chat through the timeline model ─────────────────────
            model1.openRoom(roomId)
            val timeline1 = model1.timeline(roomId)
            timeline1.sendText("hello encrypted")
            awaitTrue("own message echoes decrypted") {
                timeline1.items.any { it.body == "hello encrypted" && !it.pending && it.eventId != null }
            }
            // fail-closed proof: the wire envelope in the cache is megolm
            val echoed = timeline1.items.first { it.body == "hello encrypted" }
            val cachedRaw = model1.store!!.facts.eventCache!!
                .event(model1.principal, roomId, echoed.eventId!!)!!
            assertEquals("m.room.encrypted", cachedRaw.type, "cache stores ciphertext only (§9)")

            // ── reactions and edits ──────────────────────────────────────
            timeline1.sendReaction(echoed.eventId!!, "👍")
            awaitTrue("reaction tally materializes") {
                timeline1.items.first { it.eventId == echoed.eventId }
                    .reactions.any { it.key == "👍" && it.count == 1 && it.mine }
            }
            timeline1.sendEdit(echoed.eventId!!, "hello edited")
            awaitTrue("edit replaces presented content") {
                timeline1.items.first { it.eventId == echoed.eventId }
                    .let { it.body == "hello edited" && it.edited }
            }

            // a few more messages so scrollback has content
            repeat(5) { n -> timeline1.sendText("filler $n") }
            awaitTrue(
                "fillers echo",
                dump = {
                    timeline1.items.joinToString("\n") {
                        "item key=${it.key} eventId=${it.eventId} pending=${it.pending} body=${it.body}"
                    } + "\n" + runBlocking {
                        session.pendingSends().joinToString("\n") { "send ${it.transactionId.value}: ${it.state.value}" }
                    } + "\n" + busLog.joinToString("\n")
                },
            ) {
                timeline1.items.count { it.body.startsWith("filler") && !it.pending } == 5
            }

            // ── relaunch: cold start from cache ──────────────────────────
            runBlocking { delay(500) } // let the last aggregation writes settle
            session.close()
            scope1.cancel()

            val scope2 = CoroutineScope(SupervisorJob() + uiThread())
            val model2 = AppModel(scope2)
            model2.start()
            // start() returns with the cache already rendered — before any
            // network round-trip completed (the M3 cold-start criterion)
            assertTrue(model2.roomsFromCache, "room list renders from cache")
            assertTrue(
                model2.rooms.any { it.roomId == roomId.full },
                "cached room list contains the room",
            )

            // offline-first scrollback: the first page comes from disk; once
            // the session attaches, cached ciphertext decrypts on read (#51)
            model2.openRoom(roomId)
            val timeline2 = model2.timeline(roomId)
            awaitTrue("relaunched timeline decrypts cached history") {
                timeline2.items.any { it.body == "hello edited" && it.edited } &&
                    timeline2.items.count { it.body.startsWith("filler") } == 5
            }
            awaitTrue("cached reaction survives relaunch") {
                timeline2.items.first { it.eventId == echoed.eventId }
                    .reactions.any { it.key == "👍" && it.count == 1 }
            }

            // ── SAS verification: demo model ↔ a second device ───────────
            // The counterparty is a plain engine session (Element's role in
            // the manual interop run); the demo side drives the same
            // surfaces the VerificationDialog uses.
            awaitTrue("relaunched session is live") { model2.session != null }
            val demoSession = model2.session!!
            val other = runBlocking {
                de.centoria.katrix.engine.KatrixSession.connect(
                    baseUrl = synapse.baseUrl,
                    auth = de.centoria.katrix.transport.auth.PasswordAuthProvider(
                        "alice", "alice-pass-1", "second device",
                    ),
                    httpPort = de.centoria.katrix.transport.KtorHttpPort(
                        io.ktor.client.HttpClient(io.ktor.client.engine.cio.CIO),
                    ),
                    cryptoDriver = de.centoria.katrix.crypto.vodozemac.VodozemacCryptoDriver,
                )
            }
            try {
                val inbound = java.util.concurrent.atomic.AtomicReference<
                    de.centoria.katrix.engine.crypto.VerificationHandle?,
                    >(null)
                val collector = CoroutineScope(SupervisorJob() + Dispatchers.Default).launch {
                    other.verificationRequests.collect { inbound.set(it) }
                }
                other.startSync()

                model2.requestVerification("@alice:katrix.test")
                awaitTrue("demo surfaces the outbound handle") { model2.verifications.isNotEmpty() }
                val ours = model2.verifications.first()
                awaitTrue("counterparty receives the request") { inbound.get() != null }
                val theirs = inbound.get()!!
                runBlocking { theirs.accept() }

                val ourSas = runBlocking { ours.sasReady() }
                val theirSas = runBlocking { theirs.sasReady() }
                assertEquals(
                    ourSas.emojis.map { it.symbol },
                    theirSas.emojis.map { it.symbol },
                    "both sides show the same emoji (§10.4.5)",
                )
                runBlocking {
                    ours.confirmSas()
                    theirs.confirmSas()
                    ours.done()
                    theirs.done()
                }
                awaitTrue("demo device sees the other device as verified") {
                    runBlocking {
                        demoSession.isDeviceVerified(demoSession.userId, other.deviceId)
                    }
                }
                collector.cancel()
            } finally {
                other.close()
            }

            model2.session?.close()
            scope2.cancel()
        }
    }

    /** The app runs AppModel on the single-threaded Main dispatcher; mirror that. */
    @OptIn(ExperimentalCoroutinesApi::class)
    private fun uiThread() = Dispatchers.Default.limitedParallelism(1)

    private fun awaitTrue(
        what: String,
        timeoutMs: Long = 90_000,
        dump: (() -> String)? = null,
        predicate: () -> Boolean,
    ) {
        runBlocking {
            try {
                withTimeout(timeoutMs) {
                    while (!predicate()) delay(200)
                }
            } catch (e: kotlinx.coroutines.TimeoutCancellationException) {
                fail("timed out: $what" + (dump?.let { "\n${it()}" } ?: ""))
            }
        }
    }
}

The local-only verification driver used to verify #47 (per D30 the CI gate stays build-only). Needs `jvmTest.dependencies { implementation(project(":katrix-testkit")) }` in samples/demo and the docker env from the dev-Mac notes. Parameterize over SyncStrategy for the M4 criterion. ```kotlin package katrix.demo import de.centoria.katrix.core.id.RoomId import de.centoria.katrix.testkit.SynapseContainer import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import kotlinx.coroutines.launch import kotlinx.coroutines.delay import kotlinx.coroutines.runBlocking import kotlinx.coroutines.withTimeout import kotlinx.serialization.json.buildJsonObject import kotlinx.serialization.json.put import kotlinx.serialization.json.addJsonObject import kotlinx.serialization.json.putJsonArray import kotlinx.serialization.json.putJsonObject import java.nio.file.Files import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertTrue import kotlin.test.fail /** * LOCAL-ONLY verification driver for #47 (not committed; the CI gate is * build-only per D30): drives the demo's model layer — the same code the * Compose UI calls — against a real Synapse. Covers: login, E2EE send, * echo, reactions, edits, relaunch cold-start-from-cache, offline * scrollback decrypt-on-read. */ class DemoModelEndToEndTest { @Test fun demo_model_flow_end_to_end() { if (!SynapseContainer.dockerAvailable()) { if (System.getenv("KATRIX_REQUIRE_DOCKER") == "1") fail("docker required") println("skipping: no docker") return } val home = Files.createTempDirectory("katrix-demo-e2e") System.setProperty("katrix.demo.home", home.toString()) SynapseContainer().use { synapse -> synapse.start() synapse.registerUser("alice", "alice-pass-1") // ── first launch: login through the model ──────────────────── val scope1 = CoroutineScope(SupervisorJob() + uiThread()) val model1 = AppModel(scope1) model1.start() assertEquals(Screen.Login, model1.screen, "fresh profile starts at login") model1.login(synapse.baseUrl, "alice", "alice-pass-1") awaitTrue("session connects after login") { model1.session != null } val session = model1.session!! // an encrypted room, created through the session val roomId = runBlocking { session.createRoom( buildJsonObject { put("name", "demo e2e room") put("preset", "private_chat") putJsonArray("initial_state") { addJsonObject { put("type", "m.room.encryption") put("state_key", "") putJsonObject("content") { put("algorithm", "m.megolm.v1.aes-sha2") } } } }, ) } awaitTrue("room list shows the new room") { model1.rooms.any { it.roomId == roomId.full } } // bus tap: everything the change bus delivers for the room val busLog = java.util.Collections.synchronizedList(mutableListOf<String>()) val tap = CoroutineScope(SupervisorJob() + Dispatchers.Default).launch { session.events.collect { e -> if (e.roomId == roomId && !e.isState) { val txn = (e.raw.unsigned?.get("transaction_id") as? kotlinx.serialization.json.JsonPrimitive)?.content busLog += "bus type=${e.type} id=${e.raw.eventId?.full} txn=$txn" } } } // ── E2EE chat through the timeline model ───────────────────── model1.openRoom(roomId) val timeline1 = model1.timeline(roomId) timeline1.sendText("hello encrypted") awaitTrue("own message echoes decrypted") { timeline1.items.any { it.body == "hello encrypted" && !it.pending && it.eventId != null } } // fail-closed proof: the wire envelope in the cache is megolm val echoed = timeline1.items.first { it.body == "hello encrypted" } val cachedRaw = model1.store!!.facts.eventCache!! .event(model1.principal, roomId, echoed.eventId!!)!! assertEquals("m.room.encrypted", cachedRaw.type, "cache stores ciphertext only (§9)") // ── reactions and edits ────────────────────────────────────── timeline1.sendReaction(echoed.eventId!!, "👍") awaitTrue("reaction tally materializes") { timeline1.items.first { it.eventId == echoed.eventId } .reactions.any { it.key == "👍" && it.count == 1 && it.mine } } timeline1.sendEdit(echoed.eventId!!, "hello edited") awaitTrue("edit replaces presented content") { timeline1.items.first { it.eventId == echoed.eventId } .let { it.body == "hello edited" && it.edited } } // a few more messages so scrollback has content repeat(5) { n -> timeline1.sendText("filler $n") } awaitTrue( "fillers echo", dump = { timeline1.items.joinToString("\n") { "item key=${it.key} eventId=${it.eventId} pending=${it.pending} body=${it.body}" } + "\n" + runBlocking { session.pendingSends().joinToString("\n") { "send ${it.transactionId.value}: ${it.state.value}" } } + "\n" + busLog.joinToString("\n") }, ) { timeline1.items.count { it.body.startsWith("filler") && !it.pending } == 5 } // ── relaunch: cold start from cache ────────────────────────── runBlocking { delay(500) } // let the last aggregation writes settle session.close() scope1.cancel() val scope2 = CoroutineScope(SupervisorJob() + uiThread()) val model2 = AppModel(scope2) model2.start() // start() returns with the cache already rendered — before any // network round-trip completed (the M3 cold-start criterion) assertTrue(model2.roomsFromCache, "room list renders from cache") assertTrue( model2.rooms.any { it.roomId == roomId.full }, "cached room list contains the room", ) // offline-first scrollback: the first page comes from disk; once // the session attaches, cached ciphertext decrypts on read (#51) model2.openRoom(roomId) val timeline2 = model2.timeline(roomId) awaitTrue("relaunched timeline decrypts cached history") { timeline2.items.any { it.body == "hello edited" && it.edited } && timeline2.items.count { it.body.startsWith("filler") } == 5 } awaitTrue("cached reaction survives relaunch") { timeline2.items.first { it.eventId == echoed.eventId } .reactions.any { it.key == "👍" && it.count == 1 } } // ── SAS verification: demo model ↔ a second device ─────────── // The counterparty is a plain engine session (Element's role in // the manual interop run); the demo side drives the same // surfaces the VerificationDialog uses. awaitTrue("relaunched session is live") { model2.session != null } val demoSession = model2.session!! val other = runBlocking { de.centoria.katrix.engine.KatrixSession.connect( baseUrl = synapse.baseUrl, auth = de.centoria.katrix.transport.auth.PasswordAuthProvider( "alice", "alice-pass-1", "second device", ), httpPort = de.centoria.katrix.transport.KtorHttpPort( io.ktor.client.HttpClient(io.ktor.client.engine.cio.CIO), ), cryptoDriver = de.centoria.katrix.crypto.vodozemac.VodozemacCryptoDriver, ) } try { val inbound = java.util.concurrent.atomic.AtomicReference< de.centoria.katrix.engine.crypto.VerificationHandle?, >(null) val collector = CoroutineScope(SupervisorJob() + Dispatchers.Default).launch { other.verificationRequests.collect { inbound.set(it) } } other.startSync() model2.requestVerification("@alice:katrix.test") awaitTrue("demo surfaces the outbound handle") { model2.verifications.isNotEmpty() } val ours = model2.verifications.first() awaitTrue("counterparty receives the request") { inbound.get() != null } val theirs = inbound.get()!! runBlocking { theirs.accept() } val ourSas = runBlocking { ours.sasReady() } val theirSas = runBlocking { theirs.sasReady() } assertEquals( ourSas.emojis.map { it.symbol }, theirSas.emojis.map { it.symbol }, "both sides show the same emoji (§10.4.5)", ) runBlocking { ours.confirmSas() theirs.confirmSas() ours.done() theirs.done() } awaitTrue("demo device sees the other device as verified") { runBlocking { demoSession.isDeviceVerified(demoSession.userId, other.deviceId) } } collector.cancel() } finally { other.close() } model2.session?.close() scope2.cancel() } } /** The app runs AppModel on the single-threaded Main dispatcher; mirror that. */ @OptIn(ExperimentalCoroutinesApi::class) private fun uiThread() = Dispatchers.Default.limitedParallelism(1) private fun awaitTrue( what: String, timeoutMs: Long = 90_000, dump: (() -> String)? = null, predicate: () -> Boolean, ) { runBlocking { try { withTimeout(timeoutMs) { while (!predicate()) delay(200) } } catch (e: kotlinx.coroutines.TimeoutCancellationException) { fail("timed out: $what" + (dump?.let { "\n${it()}" } ?: "")) } } } } ```
thecrealm added
now
and removed
next
labels 2026-08-11 13:51:54 +02:00
Sign in to join this conversation.
No milestone
No project
No assignees
1 participant
Notifications
Due date
The due date is invalid or out of range. Please use the format "yyyy-mm-dd".

No due date set.

Dependencies

No dependencies set.

Reference
thecrealm/katrix#52
No description provided.