Skip to content

Commit c4502c6

Browse files
committed
fix(uts): deterministic smoke test, contract-correct mocks, run-to-quiescence FakeClock
Fixes the CI-red UnitInfraSmokeTest race and lands the review/spec-alignment round on the shared UTS infra: - Root cause of the CI flake: FakeClock.waitOn performs a real timed wait, so the disconnected-retry fires on wall-clock regardless of advance() — the "no attempt before advance" assertion was unassertable. The smoke test now owns attempt #2 via the buffered awaitConnectionAttempt() (32/32 green incl. CPU-saturation runs) and README §6.4/§9 teach the true semantics. - FakeClock: advance() now runs due work to quiescence (cascades and timers created mid-advance fire within the same advance — the spec's Fake-time semantics Guarantee); timers/pending hardened against SDK-thread races. The waitOn advisory seam is unchanged. New cascade smoke test covers it. - Mock contract fixes from review triage (verified against the UTS docs): transport cancel() now delivers listener.onClose; respondWith honors the headers param and JSON-serializes non-String bodies; SandboxApp checks HTTP status before parsing; delivery executor shutdown; @volatile channel fields; await helpers unregister listeners on success; AtomicReference for the cross-thread query-params capture. - Docs: uts/README rewritten claims verified against sources; stale "reflection" wording fixed in the skill's objects-mapping notes.
1 parent 2d3128f commit c4502c6

14 files changed

Lines changed: 190 additions & 62 deletions

File tree

‎.claude/skills/uts-to-kotlin/references/objects-mapping.md‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -401,8 +401,9 @@ delivered to subscription listeners) map to ably-java interfaces with getters (p
401401
> obtain an `ObjectMessage` from a subscription event (`event.getMessage()`, §8); there is no public
402402
> factory. The spec's explicit construction-from-wire (`PublicObjectMessage.fromObjectMessage(source,
403403
> channel)` / `PublicObjectOperation.fromObjectOperation(op)`, `PAOM3`/`PAOOP3`, in
404-
> `public_object_message.md`) is `internal` to `:liveobjects` — but the unit helpers expose it **by
405-
> reflection** as `buildPublicObjectMessage(wireJson, channelName)` (§13). So `public_object_message.md` is
404+
> `public_object_message.md`) is `internal` to `:liveobjects` — but the unit helpers expose it as
405+
> `buildPublicObjectMessage(wireMessage, channelName)` (§13), a direct `toPublicMessage` call (no reflection,
406+
> since `:liveobjects`'s internals are visible to its own tests). So `public_object_message.md` is
406407
> translatable: build the source with the op builders (`buildMapSet(...)`, `buildCounterInc(...)`, …) and
407408
> assert the public getters on the result.
408409
@@ -515,7 +516,7 @@ Several **unit** specs assert on the **internal CRDT graph**, not the public API
515516
is public and maps via §2/§12, but `publish` / `publishAndApply` (`RTO15`/`RTO20`, marked `internal` in the
516517
IDL) and the OBJECT/ACK wire assertions are internal.
517518
- `public_object_message.md` — **translatable** via the `buildPublicObjectMessage` helper (below), which
518-
reflectively performs the `PAOM3`/`PAOOP3` construction (`WireObjectMessage` → `DefaultObjectMessage`)
519+
performs the `PAOM3`/`PAOOP3` construction (`WireObjectMessage` → `DefaultObjectMessage`)
519520
that is otherwise `internal`. Build the source with the op builders and assert the public getters (§11).
520521

521522
In ably-java these are **not public**. They live in the `:liveobjects` module as `Internal*` / `Default*` / `Wire*` /

‎java/build.gradle.kts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,9 @@ dependencies {
4040
}
4141

4242
// kotlin-stdlib guardrail (invariant I5): the Kotlin plugin auto-adds kotlin-stdlib to the module's
43-
// main dependency scope, which would leak into :java's published POM/runtime. :java is Kotlin-free at
43+
// main dependency scope, which would leak into :java's published POM/runtime. The leak comes from the
44+
// PLUGIN, NOT from the testImplementation(project(":uts")) dependency — test scopes never enter the
45+
// POM (verified: removing the :uts dep leaves the leak identical). :java is Kotlin-free at
4446
// runtime, so strip it from the main artifact scopes. kotlin-stdlib still reaches the TEST classpath
4547
// transitively (via :uts's kotlin-test-junit5), so the UTS Kotlin suites compile and run.
4648
// Verified empirically on Kotlin 2.1.10: the plugin adds stdlib lazily (it does not appear in any

‎lib/src/test/kotlin/io/ably/lib/uts/unit/realtime/ConnectionRecoveryTest.kt‎

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import io.ably.lib.uts.infra.pollUntil
1212
import kotlinx.coroutines.launch
1313
import kotlinx.coroutines.test.runTest
1414
import java.util.concurrent.CopyOnWriteArrayList
15+
import java.util.concurrent.atomic.AtomicReference
1516
import kotlin.test.*
1617
import kotlin.time.Duration.Companion.seconds
1718

@@ -322,10 +323,12 @@ class ConnectionRecoveryTest {
322323
*/
323324
@Test
324325
fun `RTN16f1 - Malformed recoveryKey logs error and connects normally`() = runTest {
325-
var capturedQueryParams: Map<String, String>? = null
326+
// Written on the SDK transport thread and read from the coroutine dispatcher, so use an
327+
// AtomicReference to publish it safely across threads (avoids a visibility race).
328+
val capturedQueryParams = AtomicReference<Map<String, String>?>(null)
326329
val mock = MockWebSocket {
327330
onConnectionAttempt = { conn ->
328-
capturedQueryParams = conn.queryParams
331+
capturedQueryParams.set(conn.queryParams)
329332
conn.respondWithSuccess(ProtocolMessage().apply {
330333
action = ProtocolMessage.Action.connected
331334
connectionId = "fresh-conn"
@@ -349,8 +352,8 @@ class ConnectionRecoveryTest {
349352
assertEquals(ConnectionState.connected, client.connection.state)
350353
assertEquals("fresh-conn", client.connection.id)
351354
assertEquals("fresh-key", client.connection.key)
352-
assertNull(capturedQueryParams!!["recover"])
353-
assertNull(capturedQueryParams!!["resume"])
355+
assertNull(capturedQueryParams.get()!!["recover"])
356+
assertNull(capturedQueryParams.get()!!["resume"])
354357
assertEquals(1, mock.events.filterIsInstance<MockEvent.ConnectionAttempt>().size)
355358

356359
client.close()

‎liveobjects/src/test/kotlin/io/ably/lib/liveobjects/uts/unit/Helpers.kt‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -252,7 +252,8 @@ internal fun buildPublicObjectMessage(objectMessage: WireObjectMessage, channelN
252252
objectMessage.toPublicMessage(channelName)
253253

254254
// `provision_objects_via_rest(...)` is intentionally not here — it's REST fixture provisioning for
255-
// *integration* tests and belongs with the :uts integration tier.
255+
// *integration* tests and lives with this module's integration tier
256+
// (io.ably.lib.liveobjects.uts.integration.Helpers).
256257

257258
// ---------------------------------------------------------------------------
258259
// STANDARD_POOL_OBJECTS — the fixed tree shared by all objects unit specs

‎uts/README.md‎

Lines changed: 38 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -214,18 +214,18 @@ dependencies {
214214
// :liveobjects.
215215
api(project(":java")) // the SDK + its types (DebugOptions, ProtocolMessage, …)
216216
api(project(":network-client-core")) // HttpEngine / WebSocketEngine SPIs the mocks implement
217-
implementation(libs.coroutine.core)
218-
implementation(libs.ktor.client.core) // sandbox/proxy control HTTP
217+
implementation(libs.ktor.client.core) // proxy infra uses ktor internally — must NOT leak to consumers
219218
implementation(libs.ktor.client.cio)
220219

221-
// :uts's own tests — the tier smoke tests. No :liveobjects edge (I1).
222-
testImplementation(kotlin("test"))
223-
testImplementation(platform(libs.junit.bom))
224-
testImplementation(libs.junit.jupiter.params) // @ParameterizedTest / @ValueSource
225-
testImplementation(libs.coroutine.core)
226-
testImplementation(libs.coroutine.test) // runTest, virtual time
227-
testImplementation(libs.ktor.client.core)
228-
testImplementation(libs.ktor.client.cio)
220+
// The UTS test toolkit — exported (api) so any module consuming the infra via
221+
// testImplementation(project(":uts")) transitively gets JUnit 5, the kotlin.test Jupiter binding,
222+
// and coroutines (runTest etc.). :uts's own smoke tests inherit it from main's api — nothing to declare.
223+
api(platform(libs.junit.bom))
224+
api(libs.junit.jupiter)
225+
api(libs.junit.jupiter.params) // @ParameterizedTest / @ValueSource
226+
api(kotlin("test-junit5"))
227+
api(libs.coroutine.core)
228+
api(libs.coroutine.test) // runTest, virtual time
229229
}
230230

231231
tasks.withType<Test>().configureEach {
@@ -462,9 +462,16 @@ responses back — all without a socket.
462462

463463
### 6.4 `FakeClock` — deterministic time
464464
`FakeClock` implements the SDK's `Clock`. Time is frozen until you call `advance(ms)`; on each
465-
advance it fires any due virtual timers **synchronously**, and wakes any `waitOn` sleepers. This is
465+
advance it fires any due virtual timers **synchronously**, and wakes any `waitOn` sleepers. `advance`
466+
runs due work **to quiescence** — it re-scans until a full pass fires nothing, so cascades due within
467+
the advanced interval (a zero-delay reschedule, or a timer created by fired work) also run in that same
468+
`advance` (the spec run-to-quiescence Guarantee; see §9.3). This is
466469
how the unit test drives reconnection backoff and `connectionStateTtl` expiry **without real
467-
sleeping**:
470+
sleeping**. Caveat: `waitOn(target, timeout)` still performs a real `target.wait(timeout)`, so a
471+
sleeper also wakes once `timeout` ms of wall-clock elapse — `advance()` makes it wake *sooner*, but is
472+
not a hard gate (the **advisory** model of the spec's `mock_websocket.md` §Fake-time semantics).
473+
Drive transitions by owning the resulting attempt (`awaitConnectionAttempt()`), never
474+
by asserting a wait has *not* yet returned.
468475
```kotlin
469476
val fakeClock = FakeClock()
470477
val client = TestRealtimeClient { enableFakeTimers(fakeClock); … }
@@ -595,7 +602,8 @@ future unit-tier UTS test should take.
595602
> (`lib/src/test/kotlin/io/ably/lib/uts/unit/realtime/`, e.g. `ConnectionRecoveryTest`) and
596603
> `:liveobjects` — see §13.
597604
598-
It has **two** `@Test` methods; between them they exercise every teaching point of §5–§8.
605+
It has **three** `@Test` methods: two end-to-end transport tests (§9.1, §9.2) that between them exercise
606+
every teaching point of §5–§8, plus a focused `FakeClock` run-to-quiescence acceptance test (§9.3).
599607

600608
### 9.1 `unit infra drives the full mock-WebSocket connection lifecycle` — await style throughout
601609
One long **await-style** test that walks the SDK through the whole transport lifecycle:
@@ -632,12 +640,15 @@ One long **await-style** test that walks the SDK through the whole transport lif
632640
```
633641
3. **Publish**, asserting the full MESSAGE frame (`action`, `channel`, `messages[0].name`/`data`) again
634642
via `awaitNextMessageFromClient()`.
635-
4. **Disconnect + negative check.** `simulateDisconnect()`, await DISCONNECTED, then assert **no**
636-
reconnect has happened yet — exactly one `ConnectionAttempt` is recorded — because the retry is
637-
blocked in `FakeClock.waitOn` until the clock advances.
643+
4. **Disconnect.** `simulateDisconnect()`, await DISCONNECTED, and assert the drop was recorded. Note
644+
we do **not** snapshot the `ConnectionAttempt` count here: `FakeClock.waitOn(target, timeout)` does a
645+
real `target.wait(timeout)`, so the disconnected-retry fires on its own after ~`disconnectedRetryTimeout`
646+
ms of wall-clock even without an `advance()`. `advance()` only wins that race sooner — it is not a
647+
hard gate — so a "still exactly one attempt" assertion would be racy on a loaded runner. Ownership
648+
of attempt #2 belongs to the next step, which gates on it deterministically.
638649
5. **FakeClock-driven reconnect.** A coroutine loops `fakeClock.advance(2.seconds)` then answers the
639-
next attempt with a short-TTL CONNECTED; the test awaits CONNECTED again and asserts a second
640-
`ConnectionAttempt`.
650+
next attempt (received via the buffered `awaitConnectionAttempt()`, so it cannot be missed) with a
651+
short-TTL CONNECTED; the test awaits CONNECTED again and asserts a second `ConnectionAttempt`.
641652
6. **Refuse → SUSPENDED (the centrepiece).** After another `simulateDisconnect()`, a `refuseJob`
642653
coroutine advances the clock and `respondWithRefused()`s every reconnection attempt until the short
643654
`connectionStateTtl` (800 ms, from the short-lived CONNECTED) expires and the client gives up to
@@ -688,6 +699,14 @@ frames, `events` / `awaitNextMessageFromClient` for inspecting client output, an
688699
connect→request two-phase flow. The `RUN_DEVIATIONS` env-gated deviation pattern is **not** here (the
689700
smoke tests carry no deviations) — that teaching lives in §12.
690701

702+
### 9.3 `FakeClock advance runs cascaded work to quiescence in one call` — the run-to-quiescence Guarantee
703+
A focused, SDK-free test that pins the `FakeClock` contract §6.4 depends on: a single `advance(ms)` runs
704+
**all** work due within the advanced interval, including cascades. It schedules a task that reschedules
705+
itself at zero delay and a task that creates a brand-new timer mid-advance, then asserts one `advance`
706+
fires the whole cascade (not just the first pass). This is the infra-level guarantee the reconnect/backoff
707+
walkthroughs in §9.1 rely on; it exercises `FakeClock` directly because the Guarantee is about `advance`
708+
alone reaching quiescence.
709+
691710
---
692711

693712
## 10. Walkthrough: the Direct-Sandbox Smoke Test (`IntegrationInfraSmokeTest`)
@@ -1110,7 +1129,7 @@ implicit.
11101129
| File | Key public surface | Role |
11111130
|------|--------------------|------|
11121131
| `infra/Utils.kt` | `awaitState(client,target,timeout=5s)`, `awaitChannelState(channel,target,timeout=5s)`, `pollUntil(timeout=15s,interval=100ms){ }` | Shared wall-clock coroutine waits (package `io.ably.lib.uts.infra`); listener registered before state check. |
1113-
| `unit/UnitInfraSmokeTest.kt` | 2 `@Test`s: full mock-WS lifecycle (await style), token-auth via mock HTTP (callback WS) | Unit-tier infra acceptance (`io.ably.lib.uts.unit`) — MockWebSocket/MockHttpClient/FakeClock end-to-end. **No** `@UTS`. |
1132+
| `unit/UnitInfraSmokeTest.kt` | 3 `@Test`s: full mock-WS lifecycle (await style), token-auth via mock HTTP (callback WS), FakeClock run-to-quiescence | Unit-tier infra acceptance (`io.ably.lib.uts.unit`) — MockWebSocket/MockHttpClient/FakeClock end-to-end. **No** `@UTS`. |
11141133
| `integration/standard/IntegrationInfraSmokeTest.kt` | 1 `@ParameterizedTest` × {JSON, msgpack} | Direct-sandbox infra acceptance (`io.ably.lib.uts.integration.standard`) — SandboxApp + realtime/REST round-trip, awaited publish + `pollUntil` on `history()`. **No** `@UTS`. |
11151134
| `integration/proxy/ProxyInfraSmokeTest.kt` | 2 `@Test`s: late imperative disconnect, declarative ws-frame rule | Proxy infra acceptance (`io.ably.lib.uts.integration.proxy`) — ProxyManager + ProxySession, both fault-injection styles, proxy-log asserts. **No** `@UTS`. |
11161135

‎uts/src/main/kotlin/io/ably/lib/uts/infra/Utils.kt‎

Lines changed: 20 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -28,11 +28,18 @@ suspend fun awaitState(
2828
withContext(Dispatchers.Default.limitedParallelism(1)) {
2929
withTimeout(timeout) {
3030
suspendCancellableCoroutine { cont ->
31-
val listener = ConnectionStateListener { change ->
32-
if (change.current == target && cont.isActive) cont.resume(Unit)
31+
lateinit var listener: ConnectionStateListener
32+
listener = ConnectionStateListener { change ->
33+
if (change.current == target && cont.isActive) {
34+
client.connection.off(listener)
35+
cont.resume(Unit)
36+
}
3337
}
3438
client.connection.on(listener)
35-
if (client.connection.state == target && cont.isActive) cont.resume(Unit)
39+
if (client.connection.state == target && cont.isActive) {
40+
client.connection.off(listener)
41+
cont.resume(Unit)
42+
}
3643
cont.invokeOnCancellation { client.connection.off(listener) }
3744
}
3845
}
@@ -80,11 +87,18 @@ suspend fun awaitChannelState(
8087
withContext(Dispatchers.Default.limitedParallelism(1)) {
8188
withTimeout(timeout) {
8289
suspendCancellableCoroutine { cont ->
83-
val listener = ChannelStateListener { change ->
84-
if (change.current == target && cont.isActive) cont.resume(Unit)
90+
lateinit var listener: ChannelStateListener
91+
listener = ChannelStateListener { change ->
92+
if (change.current == target && cont.isActive) {
93+
channel.off(listener)
94+
cont.resume(Unit)
95+
}
8596
}
8697
channel.on(listener)
87-
if (channel.state == target && cont.isActive) cont.resume(Unit)
98+
if (channel.state == target && cont.isActive) {
99+
channel.off(listener)
100+
cont.resume(Unit)
101+
}
88102
cont.invokeOnCancellation { channel.off(listener) }
89103
}
90104
}

‎uts/src/main/kotlin/io/ably/lib/uts/infra/integration/SandboxApp.kt‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -93,7 +93,11 @@ class SandboxApp private constructor(
9393
contentType(ContentType.Application.Json)
9494
setBody(loadAppCreationJson().toString())
9595
}
96-
val body = JsonParser.parseString(response.bodyAsText()).asJsonObject
96+
val responseBody = response.bodyAsText()
97+
check(response.status.isSuccess()) {
98+
"sandbox app provisioning failed: HTTP ${response.status} — $responseBody"
99+
}
100+
val body = JsonParser.parseString(responseBody).asJsonObject
97101
val keys = body["keys"].asJsonArray.map { it.asJsonObject["keyStr"].asString }
98102
return SandboxApp(
99103
appId = body["appId"].asString,

‎uts/src/main/kotlin/io/ably/lib/uts/infra/unit/DefaultPendingConnection.kt‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,8 @@ internal class DefaultPendingConnection(
3232
// Async delivery per spec: the library must store the WS reference before processing CONNECTED.
3333
val encoded = Serialisation.gson.toJson(message)
3434
deliveryExecutor.submit { listener.onMessage(encoded) }
35+
// One-shot delivery: release the daemon thread once the single message is queued.
36+
deliveryExecutor.shutdown()
3537
}
3638

3739
override fun respondWithRefused() = listener.onError(IOException("Connection refused to $host:$port"))

‎uts/src/main/kotlin/io/ably/lib/uts/infra/unit/DefaultPendingRequest.kt‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import io.ably.lib.network.FailedConnectionException
44
import io.ably.lib.network.HttpBody
55
import io.ably.lib.network.HttpRequest
66
import io.ably.lib.network.HttpResponse
7+
import io.ably.lib.util.Serialisation
78
import kotlin.time.Duration
89
import kotlinx.coroutines.CompletableDeferred
910
import java.net.SocketTimeoutException
@@ -20,14 +21,15 @@ internal class DefaultPendingRequest(
2021
override fun respondWith(status: Int, body: Any, headers: Map<String, String>) {
2122
val bytes = when (body) {
2223
is ByteArray -> body
23-
else -> body.toString().toByteArray(Charsets.UTF_8)
24+
is String -> body.toByteArray(Charsets.UTF_8)
25+
else -> Serialisation.gson.toJson(body).toByteArray(Charsets.UTF_8)
2426
}
2527
deferred.complete(
2628
HttpResponse.builder()
2729
.code(status)
2830
.message("")
2931
.body(HttpBody("application/json", bytes))
30-
.headers(emptyMap())
32+
.headers(headers.mapValues { listOf(it.value) })
3133
.build()
3234
)
3335
}

0 commit comments

Comments
 (0)