Merge pull request #3995 from vitorpamplona/fix/tor-bootstrap-stall-and-ondemand

fix(tor): stop fresh installs stranding on a Tor bootstrap, and start them on clearnet defaults
This commit is contained in:
Vitor Pamplona
2026-08-26 20:46:13 -04:00
committed by GitHub
35 changed files with 1864 additions and 291 deletions
@@ -0,0 +1,182 @@
# Defaults stand in for the user's relay lists only while we have no event
**Status:** proposal — not implemented
**Goal:** first-login startup on a Tor-enabled install
**Related:** `fix/tor-bootstrap-stall-and-ondemand`, `[[fresh-install-routes-everything-via-tor]]`
## The rule
Three states, currently collapsed into two:
| we have | effective list | today |
|---|---|---|
| **no event** for the user | app defaults | defaults ✅ |
| event, **empty** list | **empty** — the user chose nothing | defaults ❌ |
| event with relays | those relays | those relays ✅ |
Everything below follows from separating "we don't know" from "we know, and it's nothing".
## Why the first login is slow
On a fresh install **100% of relay traffic is Tor-routed by construction**.
`TorRelayState.trustedRelays` is empty, so `TorRelayEvaluation.useTor()` falls through to
`newRelaysViaTor` (**default true**) for every URL — and the kind-10002 that would populate it can
only be fetched over Tor. Measured (SM-T220, same account, same ~app+8-10s login, fresh install
each; the Tor-OFF arm sets the pref, force-stops, then starts the timed run so Arti never boots):
| @20s census | Tor ON | Tor OFF |
|---|---|---|
| feed on screen | login+18s | **login+11s** |
| relays opened | 18/40 | **32/41** |
| relays serving events | 9 | **22** |
| events ingested | 2,830 | **6,134 / 7,641** |
≈7s of first paint and half the relay coverage.
## Finding 1 — every `WithBackup` helper keys on emptiness, not absence
This is a pre-existing bug against the rule above, and it must be fixed first because the whole
feature depends on the distinction being real.
```kotlin
// AdvertisedRelayListEvent
fun relays() = tags.mapNotNull(AdvertisedRelayInfo::parse) // [] when none
fun readRelaysNorm() = tags.mapNotNull(AdvertisedRelayInfo::parseReadNorm).ifEmpty { null } // null!
fun writeRelaysNorm()= tags.mapNotNull(AdvertisedRelayInfo::parseWriteNorm).ifEmpty { null } // null!
```
| helper | fallback fires when | correct |
|---|---|---|
| `normalizeNIP65AllRelayListWithBackup` | event absent only | ✅ (by accident — `relays()` has no `ifEmpty`) |
| `normalizeNIP65Read/WriteRelayListWithBackup` | event absent **or list empty** | ❌ |
| `normalizeIndexerRelayListWithBackup` | `?.ifEmpty { null } ?: DefaultIndexerRelayList` | ❌ |
| `normalizeSearchRelayListWithBackup` | `?.ifEmpty { null } ?: DefaultSearchRelayList` | ❌ |
Consequence today: **a user who publishes a kind-10002 with only write relays gets
`Constants.bootstrapInbox` silently substituted as their inbox list.** Same for a deliberately empty
search or indexer list. The app overrides an explicit choice.
The mirror problem sinks the obvious implementation: the `NoDefaults` variants return `emptySet()`
for *both* "no event" and "empty event", so `trustedRelays.isEmpty()` cannot be used as the
"do we have data yet" signal.
**Fix:** make presence explicit, and never infer it from emptiness.
```kotlin
// absent -> defaults; present -> whatever it says, including nothing
fun readRelayList(note: Note): Set<NormalizedRelayUrl> =
nip65Event(note)?.let { it.readRelaysNorm()?.toSet() ?: emptySet() } ?: Constants.bootstrapInbox
```
Same shape for write/all, and drop the `?.ifEmpty { null }` from the indexer and search helpers.
Worth doing on its own merits even if the rest of this plan is dropped.
**This removes the need for any window or timeout.** The fallback becomes a pure function of "do we
have the event", so it ends the instant one arrives — even an empty one. No per-account bookkeeping,
no 30s backstop, no race to close.
## Finding 2 — do NOT put defaults into `TrustedRelayListsState`
Tempting (it already merges all nine lists) but wrong: `account.trustedRelays.flow` feeds
`Account.kt:454`
```kotlin
isInMyRelayList = { relayUrl -> ... it in trustedRelays.flow.value }
```
which feeds `RelayAuthPermissionLedger` -> `RelayAuthResolver` -> **the NIP-42 AUTH decision**.
Adding defaults there would make the app **auto-AUTH to the six hardcoded bootstrap relays as if
they were the user's own** — signing a challenge with the user's key and revealing the pubkey — at
exactly the moment we are also going clearnet. That converts a modest timing leak into a signed
identity assertion. See `[[relay-auth-always-was-gated]]` and `[[inbox-wine-notify-auth-billing]]`
for why AUTH is the sensitive edge.
(The `saveTrustedRelayList(trustedRelays + relay)` write path in `RelayGroupChannelListScreen:449`
is **not** a hazard — it reads `account.trustedRelayList` (the NIP-51 list), not the merged
`trustedRelays`. Checked.)
**Instead:** add a separate, purpose-named flow consumed only by Tor evaluation, e.g.
`Account.relaysAssumedWhileUnknown` — the union of the with-defaults views, non-empty only while the
corresponding events are absent. `AccountsTorStateConnector` feeds it into a new
`TorRelayState.assumedRelays`. Nothing else reads it.
## Where the check goes in `useTor()`
```
torType == OFF -> false
isLocalHost -> false
isOverlayNetwork -> false
isOnion -> onionRelaysViaTor
in moneyOpRelayList -> moneyOperationsViaTor
in dmRelayList -> dmRelaysViaTor
in trustedRelayList -> trustedRelaysViaTor
in assumedRelayList -> trustedRelaysViaTor <-- new, immediately above the fallback
else -> newRelaysViaTor
```
Landing immediately above the fallback means **.onion, money-operation and DM relays keep their own
policy for free** — the change can only ever affect URLs that would have been treated as "new".
Resolve to `trustedRelaysViaTor`, **not** a hardcoded `false`:
- default user (`false`) -> clearnet -> fast start;
- hardened user (`true`) -> stays on Tor, automatically, with no new setting to discover.
That is the difference between "the app overrides you" and "the app treats its stand-in list the way
you asked your own list to be treated".
## Privacy, for the PR body
The window correlates the user's **IP with their pubkey** at ~6 hardcoded relays, because the REQ
asks those relays for that pubkey's events. A first login is the most sensitive moment there is.
What makes it defensible: **`trustedRelaysViaTor` already defaults to false**, so the moment
kind-10002 lands the user's own relays are dialled over clearnet anyway. This moves an existing
disclosure slightly earlier, to a different well-known set. It is not a new class of exposure for
the default configuration — and it is *not* an AUTH disclosure, provided Finding 2 is respected.
If `trustedRelaysViaTor` ever becomes default-true, **this feature must be revisited in the same
commit** — its justification disappears. Leave a comment at the default linking the two.
Residual, worth verifying rather than assuming: `useTor()` is keyed by relay **URL**, and the pool
multiplexes every subscription for a URL over one socket. During the window, anything addressed to a
default relay rides that clearnet socket — including a kind-1059 giftwrap subscription, since the DM
list is also absent. Measure it (below) before deciding it is acceptable.
## Testing
Unit — the rule itself, per list type: absent event -> defaults; present-but-empty -> **empty**;
present-with-values -> values. The middle case is the regression guard and the one that fails today.
Unit (`TorRelayEvaluationTest`): an assumed relay resolves to `trustedRelaysViaTor` (both values);
.onion / money-op / DM keep their own policy while also listed as assumed; a non-assumed "new" relay
still resolves to `newRelaysViaTor`; an empty assumed set is byte-for-byte today's behaviour.
Unit: `isInMyRelayList` does **not** see assumed relays (guards Finding 2 permanently).
Device — the number that justifies the change. `relaytiming.sh` + `BootRelayDiag` census,
`VERBOSE_LOGS=true` benchmark build, fresh install each, counterbalanced, n>=3:
- primary: login -> first note; login -> own profile + follow list;
- secondary: relays opened / serving / events at the 20s census;
- guard: grep the verbose log for any request to a default relay during the window that is not for
the account's own pubkey, and for any AUTH sent to one.
Harness traps (all in `[[fresh-install-routes-everything-via-tor]]`): the tablet raises its lock
screen during long waits (`wm dismiss-keyguard`, not just `KEYCODE_WAKEUP`); the login layout shifts
when the IME opens, so dismiss it before tapping fixed coordinates; `BACK` on the home screen exits
the app; always assert the run left the login screen before trusting its timing.
## Expected outcome
Approach the Tor-OFF column: ≈**-7s to first paint, ~2x relay coverage** in the first 20s, with
everything after the first event behaving exactly as today.
If the gain is materially smaller, the likely cause is that the feed is gated on outbox-discovered
relays (which stay "new", hence Tor) rather than the user's own list — in which case the win is
limited to profile and follows, and may not be worth the privacy cost. Decide on the numbers.
## Order of work
1. Fix the absent-vs-empty bug in the four helpers + tests. Independently correct; ship separately.
2. Add `relaysAssumedWhileUnknown` + `TorRelayState.assumedRelays` + the `useTor()` branch.
3. Device A/B. Keep only if it earns its keep.
@@ -25,7 +25,10 @@ import androidx.test.filters.LargeTest
import androidx.test.platform.app.InstrumentationRegistry
import com.vitorpamplona.amethyst.ui.tor.TorService
import com.vitorpamplona.amethyst.ui.tor.TorServiceStatus
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeout
@@ -58,7 +61,7 @@ import kotlin.system.measureTimeMillis
* 3. `./gradlew :amethyst:connectedPlayDebugAndroidTest -P android.testInstrumentationRunnerArguments.class=com.vitorpamplona.amethyst.tor.TorBootstrapInstrumentedTest`
*
* **What it covers that [TorManagerTest] does not:**
* - Real `ArtiNative.initialize` → `create_bootstrapped` → SOCKS listener bind.
* - Real `ArtiNative.initialize` → `create_unbootstrapped_async` → SOCKS listener bind.
* - Real rustls `CryptoProvider` install (regression check after the arti-v2.3.0 bump).
* - Real `destroy()` releasing the state file lock so a second `initialize()` succeeds.
* - OkHttp routing traffic through the SOCKS port and Arti exiting through the
@@ -73,7 +76,14 @@ import kotlin.system.measureTimeMillis
@Ignore("Tier-3 integration test — requires on-device network access to Tor. See class kdoc to enable.")
class TorBootstrapInstrumentedTest {
private val context = InstrumentationRegistry.getInstrumentation().targetContext
private val torService = TorService(context)
/**
* [TorService] promotes Bootstrapping -> Active from a coroutine on this scope, so the test
* must own one and cancel it — without a live scope `status` would never reach Active and every
* assertion below would hang until its timeout.
*/
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
private val torService = TorService(context, scope)
@After
fun tearDown() =
@@ -81,11 +91,12 @@ class TorBootstrapInstrumentedTest {
// Drop the native client so this test's state file lock doesn't bleed into
// the next instrumented run on the same device.
torService.reset()
scope.cancel()
}
/**
* Cold-start bootstrap. The whole point of the custom Arti build is that this
* works at all — if create_bootstrapped panics (e.g., because we forgot to install
* works at all — if client creation panics (e.g., because we forgot to install
* a rustls CryptoProvider after an arti bump) the test catches it.
*/
@Test
@@ -82,9 +82,14 @@ class Amethyst : Application() {
*/
val DEFAULT_LOG_LEVEL: LogLevel =
when {
!BuildConfig.DEBUG -> LogLevel.WARN
VERBOSE_LOGS -> LogLevel.DEBUG
else -> LogLevel.INFO
// `isDebug` also covers the `benchmark` build type — a release build (R8 + AOT)
// that exists purely to be measured and is never shipped. Treating it as a release
// build left it at WARN, which drops every INFO milestone the boot narrative is
// made of (account load timings, Tor status transitions, the relay census), so the
// one variant whose numbers are trustworthy was also the one we could not read.
VERBOSE_LOGS && isDebug -> LogLevel.DEBUG
isDebug -> LogLevel.INFO
else -> LogLevel.WARN
}
lateinit var instance: AppModules
@@ -130,7 +130,6 @@ import com.vitorpamplona.amethyst.ui.screen.AccountState
import com.vitorpamplona.amethyst.ui.screen.UiSettingsState
import com.vitorpamplona.amethyst.ui.tor.TorManager
import com.vitorpamplona.amethyst.ui.tor.TorService
import com.vitorpamplona.amethyst.ui.tor.TorServiceStatus
import com.vitorpamplona.quartz.nip01Core.core.Address
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
@@ -282,7 +281,7 @@ class AppModules(
UiSettingsState(uiPrefs.value, connManager.isMobileOrFalse, applicationIOScope)
}
private val torService = TorService(appContext)
private val torService = TorService(appContext, applicationIOScope)
val torManager = TorManager(torPrefs, torService, applicationIOScope)
// Network identity change (wifi↔cellular, regained from offline, captive portal
@@ -414,7 +413,11 @@ class AppModules(
init {
applicationIOScope.launch {
torService.status
.map { it is TorServiceStatus.Active }
// Battery ledger: Tor is doing work from the moment the client exists — the
// directory download is the most expensive part of a launch — so this tracks
// "running", not "bootstrapped". Keying it on Active alone would silently omit the
// 12-34s download from every cold start.
.map { it.socksPort != null }
.distinctUntilChanged()
.collect { torSession.setActive(it) }
}
@@ -647,7 +650,7 @@ class AppModules(
// proxy during bootstrap. RelayProxyClientConnector reconnects them (with
// ignoreRetryDelays=true) the instant Tor flips to Active.
canDial = { url ->
!torEvaluatorFlow.shouldUseTorForRelay(url) || torManager.isSocksReady()
!torEvaluatorFlow.shouldUseTorForRelay(url) || torManager.isTorReady()
},
)
@@ -716,7 +719,7 @@ class AppModules(
TorCircuitHealthTracker(
client = client,
isTorRouted = { torEvaluatorFlow.shouldUseTorForRelay(it) },
isTorActive = { torManager.isSocksReady() },
isTorActive = { torManager.isTorReady() },
isConnectivityActive = { connManager.status.value is ConnectivityStatus.Active },
onCircuitsDead = { torManager.onTorCircuitsDead() },
).also { it.register() }
@@ -31,7 +31,6 @@ import com.vitorpamplona.amethyst.commons.connectedApps.signers.InMemoryNostrSig
import com.vitorpamplona.amethyst.commons.connectedApps.signers.NostrSignerPermissionLedger
import com.vitorpamplona.amethyst.commons.connectedApps.signers.NostrSignerPermissionStore
import com.vitorpamplona.amethyst.commons.defaults.Constants
import com.vitorpamplona.amethyst.commons.defaults.DefaultIndexerRelayList
import com.vitorpamplona.amethyst.commons.marmot.MarmotManager
import com.vitorpamplona.amethyst.commons.model.IAccount
import com.vitorpamplona.amethyst.commons.model.buzz.BuzzChannelStars
@@ -138,6 +137,7 @@ import com.vitorpamplona.amethyst.model.nip78AppSpecific.AppSpecificState
import com.vitorpamplona.amethyst.model.nip89AppHandlers.AppRecommendationsState
import com.vitorpamplona.amethyst.model.nipA3PaymentTargets.NipA3PaymentTargetsState
import com.vitorpamplona.amethyst.model.nipB7Blossom.BlossomServerListState
import com.vitorpamplona.amethyst.model.serverList.AssumedRelayListsState
import com.vitorpamplona.amethyst.model.serverList.MergedFollowListsState
import com.vitorpamplona.amethyst.model.serverList.MergedFollowPlusMineRelayListsState
import com.vitorpamplona.amethyst.model.serverList.MergedFollowPlusMineWithIndexRelayListsState
@@ -384,12 +384,16 @@ class Account(
// doubles as the attribution pubkey for ExplainedFilter.accountPubKeys.
override val userFinderPubkeyHex: HexKey get() = userProfile().pubkeyHex
override fun indexRelays(): Set<NormalizedRelayUrl> = indexerRelayList.flow.value.ifEmpty { DefaultIndexerRelayList }
// No ifEmpty here on purpose: an empty kind:10086 is the user asking for no indexers, and
// IndexerRelayListState already substitutes the defaults for the only case we may override —
// never having seen the event. Re-substituting here would undo that choice.
override fun indexRelays(): Set<NormalizedRelayUrl> = indexerRelayList.flow.value
override fun outboxHomeRelays(): Set<NormalizedRelayUrl> = nip65RelayList.allFlowNoDefaults.value + privateStorageRelayList.flow.value + localRelayList.flow.value
// searchRelayList.flow already applies the DefaultSearchRelayList fallback internally
// (SearchRelayListState.normalizeSearchRelayListWithBackup), so no ifEmpty needed here.
// searchRelayList.flow applies DefaultSearchRelayList internally when no kind:10007 has ever
// been seen (SearchRelayListState.normalizeSearchRelayListWithBackup); an empty published list
// stays empty. No ifEmpty here either way.
override fun searchRelays(): Set<NormalizedRelayUrl> = (trustedRelayList.flow.value + searchRelayList.flow.value).toSet()
override fun searchOnlyRelays(): Set<NormalizedRelayUrl> = searchRelayList.flow.value
@@ -809,6 +813,9 @@ class Account(
val trustedRelays = TrustedRelayListsState(nip65RelayList, privateStorageRelayList, localRelayList, dmRelayList, searchRelayList, indexerRelayList, proxyRelayList, trustedRelayList, broadcastRelayList, scope)
/** Relays guessed on the user's behalf until their own lists arrive. Read only by Tor routing. */
val assumedRelays = AssumedRelayListsState(nip65RelayList, searchRelayList, indexerRelayList, scope)
// Follows Relays
val followOutboxesOrProxy = FollowListOutboxOrProxyRelays(kind3FollowList, blockedRelayList, proxyRelayList, cache, scope)
@@ -58,7 +58,11 @@ class IndexerRelayListState(
fun indexListEvent(note: Note) = note.event as? IndexerRelayListEvent ?: settings.backupIndexRelayList
suspend fun normalizeIndexerRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> = indexListEvent(note)?.let { decryptionCache.relays(it) }?.ifEmpty { null } ?: DefaultIndexerRelayList
suspend fun normalizeIndexerRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> {
val event = indexListEvent(note) ?: return DefaultIndexerRelayList
// Fully decrypted here, so empty means the user listed nothing — not "not decrypted yet".
return decryptionCache.relays(event)
}
suspend fun normalizeIndexerRelayListWithBackupNoDefaults(note: Note): Set<NormalizedRelayUrl> = indexListEvent(note)?.let { decryptionCache.relays(it) } ?: emptySet()
@@ -73,12 +77,26 @@ class IndexerRelayListState(
*/
fun normalizeIndexerRelayListPrecached(note: Note): Set<NormalizedRelayUrl> = indexListEvent(note)?.let { decryptionCache.cachedRelays(it) }?.ifEmpty { null } ?: DefaultIndexerRelayList
/** See `Nip65RelayListState.assumedDefaults`. Empty as soon as any kind:10086 exists. */
fun assumedDefaults(note: Note): Set<NormalizedRelayUrl> = if (indexListEvent(note) == null) DefaultIndexerRelayList else emptySet()
val assumedDefaultsFlow =
getIndexerRelayListFlow()
.map { assumedDefaults(it.note) }
.onStart { emit(assumedDefaults(indexerListNote)) }
.flowOn(Dispatchers.IO)
.stateIn(
scope,
SharingStarted.Eagerly,
assumedDefaults(indexerListNote),
)
/**
* The account's indexer relays, **never empty** — [normalizeIndexerRelayListWithBackup]
* substitutes [DefaultIndexerRelayList] both when there is no kind:10086 and when the
* one we have decodes to zero relays. Callers assembling metadata / relay-list REQs read
* this and can rely on getting a usable set; use [flowNoDefaults] instead to show or diff
* what the user actually configured.
* The account's indexer relays. [normalizeIndexerRelayListWithBackup] substitutes
* [DefaultIndexerRelayList] when there is no kind:10086 at all — but **not** when the one we
* have decodes to zero relays, which is the user saying "no indexers" and is honored. Callers
* assembling metadata / relay-list REQs must therefore tolerate an empty set; use
* [flowNoDefaults] to show or diff what the user actually configured.
*
* Seeded via [normalizeIndexerRelayListPrecached] rather than `emptySet()`, for the same
* reason as the search list: `flowOn(IO)` makes the first real emission asynchronous, so an
@@ -58,7 +58,11 @@ class SearchRelayListState(
fun searchListEvent(note: Note) = note.event as? SearchRelayListEvent ?: settings.backupSearchRelayList
suspend fun normalizeSearchRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> = searchListEvent(note)?.let { decryptionCache.relays(it) }?.ifEmpty { null } ?: DefaultSearchRelayList
suspend fun normalizeSearchRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> {
val event = searchListEvent(note) ?: return DefaultSearchRelayList
// Fully decrypted here, so empty means the user listed nothing — not "not decrypted yet".
return decryptionCache.relays(event)
}
suspend fun normalizeSearchRelayListWithBackupNoDefaults(note: Note): Set<NormalizedRelayUrl> = searchListEvent(note)?.let { decryptionCache.relays(it) } ?: emptySet()
@@ -74,16 +78,31 @@ class SearchRelayListState(
*/
fun normalizeSearchRelayListPrecached(note: Note): Set<NormalizedRelayUrl> = searchListEvent(note)?.let { decryptionCache.cachedRelays(it) }?.ifEmpty { null } ?: DefaultSearchRelayList
/** See `Nip65RelayListState.assumedDefaults`. Empty as soon as any kind:10007 exists. */
fun assumedDefaults(note: Note): Set<NormalizedRelayUrl> = if (searchListEvent(note) == null) DefaultSearchRelayList else emptySet()
val assumedDefaultsFlow =
getSearchRelayListFlow()
.map { assumedDefaults(it.note) }
.onStart { emit(assumedDefaults(searchListNote)) }
.flowOn(Dispatchers.IO)
.stateIn(
scope,
SharingStarted.Eagerly,
assumedDefaults(searchListNote),
)
/**
* The account's search relays, **never empty** — [normalizeSearchRelayListWithBackup]
* substitutes [DefaultSearchRelayList] both when there is no kind:10007 and when the
* one we have decodes to zero relays. Callers assembling NIP-50 REQs read this and can
* The account's search relays. [normalizeSearchRelayListWithBackup] substitutes
* [DefaultSearchRelayList] when there is no kind:10007 at all — but **not** when the one we
* have decodes to zero relays, which is the user saying "no search relays" and is honored.
* Callers assembling NIP-50 REQs must tolerate an empty set, and can
* rely on getting a usable set; use [flowNoDefaults] instead to show or diff what the
* user actually configured.
*
* Seeded via [normalizeSearchRelayListPrecached] rather than `emptySet()`: `flowOn(IO)` means
* the first real emission can never be synchronous with `stateIn`, so an `emptySet()` seed
* left a window where `.value` contradicted the "never empty" contract above and search
* left a window where `.value` reported nothing before the event had been read at all, so search
* silently queried nothing. That window is unbounded for a NIP-46 signer whose list has
* private entries, since the first emission waits on a remote decrypt.
*/
@@ -21,6 +21,7 @@
package com.vitorpamplona.amethyst.model.nip65RelayList
import com.vitorpamplona.amethyst.commons.defaults.Constants
import com.vitorpamplona.amethyst.commons.defaults.relayListOrDefaultsWhenUnknown
import com.vitorpamplona.amethyst.model.AccountSettings
import com.vitorpamplona.amethyst.model.LocalCache
import com.vitorpamplona.amethyst.model.Note
@@ -58,9 +59,9 @@ class Nip65RelayListState(
fun nip65Event(note: Note) = note.event as? AdvertisedRelayListEvent ?: settings.backupNIP65RelayList
fun normalizeNIP65WriteRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> = nip65Event(note)?.writeRelaysNorm()?.toSet() ?: Constants.eventFinderRelays
fun normalizeNIP65WriteRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> = relayListOrDefaultsWhenUnknown(nip65Event(note), Constants.eventFinderRelays) { it.writeRelaysNorm()?.toSet() }
fun normalizeNIP65ReadRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> = nip65Event(note)?.readRelaysNorm()?.toSet() ?: Constants.bootstrapInbox
fun normalizeNIP65ReadRelayListWithBackup(note: Note): Set<NormalizedRelayUrl> = relayListOrDefaultsWhenUnknown(nip65Event(note), Constants.bootstrapInbox) { it.readRelaysNorm()?.toSet() }
fun normalizeNIP65WriteRelayListNoDefaults(note: Note): Set<NormalizedRelayUrl> = nip65Event(note)?.writeRelaysNorm()?.toSet() ?: emptySet()
@@ -70,6 +71,27 @@ class Nip65RelayListState(
fun normalizeNIP65AllRelayListWithBackupNoDefaults(note: Note): Set<NormalizedRelayUrl> = nip65Event(note)?.relays()?.map { it.relayUrl }?.toSet() ?: emptySet()
/**
* The app defaults currently standing in for a user we have no kind:10002 for — empty as soon
* as one exists, including an empty one.
*
* Uses the same `nip65Event(note) == null` predicate the substitution itself uses, so the two
* cannot drift: whatever is listed here is exactly what the app is guessing on the user's
* behalf. See [relayListOrDefaultsWhenUnknown].
*/
fun assumedDefaults(note: Note): Set<NormalizedRelayUrl> = if (nip65Event(note) == null) Constants.bootstrapInbox + Constants.eventFinderRelays else emptySet()
val assumedDefaultsFlow =
getNIP65RelayListFlow()
.map { assumedDefaults(it.note) }
.onStart { emit(assumedDefaults(nip65ListNote)) }
.flowOn(Dispatchers.IO)
.stateIn(
scope,
SharingStarted.Eagerly,
assumedDefaults(nip65ListNote),
)
val outboxFlow =
getNIP65RelayListFlow()
.map { normalizeNIP65WriteRelayListWithBackup(it.note) }
@@ -0,0 +1,68 @@
/*
* 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.model.serverList
import com.vitorpamplona.amethyst.model.nip51Lists.indexerRelays.IndexerRelayListState
import com.vitorpamplona.amethyst.model.nip51Lists.searchRelays.SearchRelayListState
import com.vitorpamplona.amethyst.model.nip65RelayList.Nip65RelayListState
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.flowOn
import kotlinx.coroutines.flow.stateIn
/**
* The relays the app is **guessing** on the user's behalf because it has not seen their lists yet.
*
* Non-empty only while the corresponding event is absent — never because a list is empty, which is
* a choice we honor (see `relayListOrDefaultsWhenUnknown`). It therefore empties itself, per list,
* the moment the user's own data lands; no window, no timeout, no bookkeeping.
*
* **Deliberately NOT merged into [TrustedRelayListsState].** That one feeds `Account.isInMyRelayList`
* -> `RelayAuthPermissionLedger` -> `RelayAuthResolver`, i.e. the NIP-42 AUTH decision. Guessed
* relays must never make the app sign an AUTH challenge as though they were the user's own — that
* would turn a timing signal into a signed identity assertion. The single consumer of this flow is
* Tor routing.
*/
class AssumedRelayListsState(
val nip65RelayList: Nip65RelayListState,
val searchRelayList: SearchRelayListState,
val indexerRelayList: IndexerRelayListState,
val scope: CoroutineScope,
) {
val flow: StateFlow<Set<NormalizedRelayUrl>> =
combine(
nip65RelayList.assumedDefaultsFlow,
searchRelayList.assumedDefaultsFlow,
indexerRelayList.assumedDefaultsFlow,
) { nip65, search, indexer ->
nip65 + search + indexer
}.flowOn(Dispatchers.IO)
.stateIn(
scope,
kotlinx.coroutines.flow.SharingStarted.Eagerly,
nip65RelayList.assumedDefaultsFlow.value +
searchRelayList.assumedDefaultsFlow.value +
indexerRelayList.assumedDefaultsFlow.value,
)
}
@@ -20,14 +20,17 @@
*/
package com.vitorpamplona.amethyst.model.torState
import com.vitorpamplona.amethyst.model.Account
import com.vitorpamplona.amethyst.model.accountsCache.AccountCacheState
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.utils.Log
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.FlowPreview
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharingStarted
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.debounce
import kotlinx.coroutines.flow.emitAll
@@ -35,112 +38,132 @@ import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.stateIn
import kotlinx.coroutines.flow.transformLatest
/**
* Pushes the relay classifications [TorRelayState] needs — which relays are DM, trusted, guessed, or
* money-operation relays — as a union across every logged-in account.
*
* All four are the same fold: pick one set per account, union them, publish. It used to be written
* out four times at ~30 lines each, and the copies had already drifted apart in trivial ways (an
* `if (isEmpty)` guard that could never fire, differently-named accumulators). Sharing one
* implementation is what keeps a fifth classification from being another 30 lines of the same
* thing — and, more importantly, from being 30 lines that quietly forget a step.
*/
class AccountsTorStateConnector(
accountsCache: AccountCacheState,
torEvaluatorFlow: TorRelayState,
scope: CoroutineScope,
) {
@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class)
val allDmRelayFlows: Flow<Set<NormalizedRelayUrl>> =
accountsCache.accounts
.debounce(200)
.transformLatest { snapshot ->
val dmFlows = snapshot.map { it.value.dmRelayList.flow }
val dmFlowReady =
dmFlows.ifEmpty {
listOf(MutableStateFlow(emptySet()))
}
if (dmFlowReady.isEmpty()) {
emit(emptySet())
} else {
emitAll(
combine(dmFlowReady) {
val dmRelays = mutableSetOf<NormalizedRelayUrl>()
it.forEach {
dmRelays.addAll(it)
}
dmRelays.toSet()
},
)
}
}.onEach {
torEvaluatorFlow.dmRelays.tryEmit(it)
}.stateIn(
scope,
SharingStarted.Eagerly,
emptySet(),
)
/**
* Union of one relay set across all logged-in accounts, republished into [TorRelayState].
*
* `debounce(200)` rides out the burst of account churn at login; `transformLatest` drops the
* previous fan-in when the account set changes so a logged-out account cannot keep contributing.
* The seed is `emptySet()` for every classification: before any account exists, nothing is
* classified.
*
* Takes its collaborators as parameters rather than reading constructor properties because the
* call sites are property initializers, where non-`val` constructor parameters are in scope but
* member functions cannot see them.
*/
@OptIn(FlowPreview::class, ExperimentalCoroutinesApi::class)
val allTrustedRelaysFlow: Flow<Set<NormalizedRelayUrl>> =
private fun unionAcrossAccounts(
accountsCache: AccountCacheState,
scope: CoroutineScope,
select: (Account) -> Flow<Set<NormalizedRelayUrl>>,
publish: (Set<NormalizedRelayUrl>) -> Unit,
): StateFlow<Set<NormalizedRelayUrl>> =
accountsCache.accounts
.debounce(200)
.transformLatest { snapshot ->
val trustedRelayFlows = snapshot.map { it.value.trustedRelays.flow }
val trustedRelayFlowReady =
trustedRelayFlows.ifEmpty {
listOf(MutableStateFlow(emptySet()))
}
if (trustedRelayFlowReady.isEmpty()) {
emit(emptySet())
} else {
emitAll(
combine(trustedRelayFlowReady) {
val trustedRelays = mutableSetOf<NormalizedRelayUrl>()
it.forEach {
trustedRelays.addAll(it)
}
trustedRelays.toSet()
},
)
}
}.onEach {
torEvaluatorFlow.trustedRelays.tryEmit(it)
}.stateIn(
scope,
SharingStarted.Eagerly,
emptySet(),
)
// Persistent money-operation relays across all accounts: NIP-47 wallet relays and saved CLINK
// Debits service relays. Feeds TorRelayState.moneyOpRelays so these connections honor the
// money-operations Tor preference instead of being classified as generic "new" relays.
@OptIn(FlowPreview::class, ExperimentalCoroutinesApi::class)
val allMoneyOpRelaysFlow: Flow<Set<NormalizedRelayUrl>> =
accountsCache.accounts
.debounce(200)
.transformLatest { snapshot ->
val perAccountFlows =
snapshot.map { (_, account) ->
combine(
account.settings.nwcWallets,
account.settings.clinkDebitWallets,
) { nwcWallets, clinkDebitWallets ->
val relays = mutableSetOf<NormalizedRelayUrl>()
nwcWallets.forEach { relays.add(it.uri.relayUri) }
clinkDebitWallets.forEach { relays.addAll(it.pointer.relays) }
relays.toSet()
}
}
val ready = perAccountFlows.ifEmpty { listOf(MutableStateFlow(emptySet())) }
val perAccount =
snapshot
.map { select(it.value) }
.ifEmpty { listOf(MutableStateFlow(emptySet())) }
emitAll(
combine(ready) { perAccount ->
val moneyOpRelays = mutableSetOf<NormalizedRelayUrl>()
perAccount.forEach { moneyOpRelays.addAll(it) }
moneyOpRelays.toSet()
combine(perAccount) { sets ->
sets.flatMapTo(mutableSetOf()) { it }
},
)
}.onEach {
torEvaluatorFlow.moneyOpRelays.tryEmit(it)
}.stateIn(
}.onEach(publish)
.stateIn(
scope,
SharingStarted.Eagerly,
emptySet(),
)
/** NIP-17 DM relays: these follow the dedicated DM preference, never the generic "new" one. */
val allDmRelayFlows: StateFlow<Set<NormalizedRelayUrl>> =
unionAcrossAccounts(
accountsCache,
scope,
select = { it.dmRelayList.flow },
publish = { torEvaluatorFlow.dmRelays.tryEmit(it) },
)
/** Everything the user actually put in one of their own relay lists. */
val allTrustedRelaysFlow: StateFlow<Set<NormalizedRelayUrl>> =
unionAcrossAccounts(
accountsCache,
scope,
select = { it.trustedRelays.flow },
publish = { torEvaluatorFlow.trustedRelays.tryEmit(it) },
)
/**
* Relays the app is *guessing* while an account's own lists are unknown. Feeds
* [TorRelayState.assumedRelays] and nothing else — see `AssumedRelayListsState` for why these
* must never reach the AUTH decision.
*
* Per account, so a second login cannot re-open the guess for an established one; each
* account's contribution empties itself as soon as that account's own lists land.
*/
val allAssumedRelaysFlow: StateFlow<Set<NormalizedRelayUrl>> =
unionAcrossAccounts(
accountsCache,
scope,
select = { it.assumedRelays.flow },
publish = {
logHandover(it)
torEvaluatorFlow.assumedRelays.tryEmit(it)
},
)
/**
* Persistent money-operation relays: NIP-47 wallet relays and saved CLINK Debits service
* relays, so these connections honor the money-operations preference rather than being
* classified as generic "new" relays.
*/
val allMoneyOpRelaysFlow: StateFlow<Set<NormalizedRelayUrl>> =
unionAcrossAccounts(
accountsCache,
scope,
select = { account ->
combine(
account.settings.nwcWallets,
account.settings.clinkDebitWallets,
) { nwcWallets, clinkDebitWallets ->
val relays = mutableSetOf<NormalizedRelayUrl>()
nwcWallets.forEach { relays.add(it.uri.relayUri) }
clinkDebitWallets.forEach { relays.addAll(it.pointer.relays) }
relays.toSet()
}
},
publish = { torEvaluatorFlow.moneyOpRelays.tryEmit(it) },
)
@Volatile private var lastAssumedCount: Int = -1
/**
* The handover is the whole contract of the guessed-relay feature: the moment a user's own
* lists arrive, every relay we were guessing about goes back to the policy they actually asked
* for. Logged at INFO because "did it hand over, and when" is not answerable from any other
* line — the reconnect that follows looks identical to an ordinary one.
*/
private fun logHandover(relays: Set<NormalizedRelayUrl>) {
if (relays.size == lastAssumedCount) return
val released = if (relays.isEmpty()) " (own lists arrived; released to their real Tor policy)" else ""
Log.i("AccountsTorState") { "Guessed relays: $lastAssumedCount -> ${relays.size}$released" }
lastAssumedCount = relays.size
}
}
@@ -22,3 +22,5 @@ package com.vitorpamplona.amethyst.model.torState
// Canonical type now lives in commons
typealias TorRelayEvaluation = com.vitorpamplona.amethyst.commons.tor.TorRelayEvaluation
typealias RelayClassification = com.vitorpamplona.amethyst.commons.tor.RelayClassification
@@ -46,6 +46,13 @@ class TorRelayState(
val dmRelays = MutableStateFlow<Set<NormalizedRelayUrl>>(emptySet())
val trustedRelays = MutableStateFlow<Set<NormalizedRelayUrl>>(emptySet())
/**
* Relays guessed on the user's behalf while their own lists are unknown. Fed by
* [AccountsTorStateConnector]; see `AssumedRelayListsState` for why this is separate from
* [trustedRelays] rather than merged into it.
*/
val assumedRelays = MutableStateFlow<Set<NormalizedRelayUrl>>(emptySet())
/**
* Relays known to be used for money operations from persistent configuration: NIP-47 wallet
* relays and saved CLINK Debits service relays. Fed by [AccountsTorStateConnector] across all
@@ -130,47 +137,49 @@ class TorRelayState(
currentSettings(),
)
val flow =
combineTransform(
torSettings,
private fun currentClassification() =
RelayClassification(
trusted = trustedRelays.value,
dm = dmRelays.value,
moneyOp = currentMoneyOpRelays(),
assumed = assumedRelays.value,
)
/**
* The four category sets as one value. Folding them here also keeps the evaluation flow below
* at two sources instead of six — `combineTransform`'s typed overloads stop at five.
*/
private val classification =
combine(
trustedRelays,
dmRelays,
moneyOpRelays,
adHocMoneyOpCounts,
) {
torSettings: TorRelaySettings,
trustedRelayList: Set<NormalizedRelayUrl>,
dmRelayList: Set<NormalizedRelayUrl>,
moneyOpRelayList: Set<NormalizedRelayUrl>,
adHocMoneyOps: Map<NormalizedRelayUrl, Int>,
->
emit(
TorRelayEvaluation(
torSettings = torSettings,
trustedRelayList = trustedRelayList,
dmRelayList = dmRelayList,
moneyOpRelayList = moneyOpRelayList + adHocMoneyOps.keys,
),
assumedRelays,
) { trusted, dm, moneyOp, adHocMoneyOps, assumed ->
RelayClassification(
trusted = trusted,
dm = dm,
moneyOp = moneyOp + adHocMoneyOps.keys,
assumed = assumed,
)
}
val flow =
combineTransform(
torSettings,
classification,
) { torSettings: TorRelaySettings, classification: RelayClassification ->
emit(TorRelayEvaluation(torSettings, classification))
}.onStart {
emit(
TorRelayEvaluation(
torSettings = torSettings.value,
trustedRelayList = trustedRelays.value,
dmRelayList = dmRelays.value,
moneyOpRelayList = currentMoneyOpRelays(),
),
TorRelayEvaluation(torSettings.value, currentClassification()),
)
}.flowOn(Dispatchers.IO)
.stateIn(
scope,
SharingStarted.Eagerly,
TorRelayEvaluation(
torSettings = torSettings.value,
trustedRelayList = trustedRelays.value,
dmRelayList = dmRelays.value,
moneyOpRelayList = currentMoneyOpRelays(),
),
TorRelayEvaluation(torSettings.value, currentClassification()),
)
/**
@@ -178,13 +187,7 @@ class TorRelayState(
* snapshot. This makes ad-hoc money-op registration ([registerMoneyOpRelays]) take effect on the
* very next connection attempt, with no dependency on the combine pipeline having propagated yet.
*/
fun shouldUseTorForRelay(relay: NormalizedRelayUrl) =
TorRelayEvaluation(
torSettings = currentSettings(),
trustedRelayList = trustedRelays.value,
dmRelayList = dmRelays.value,
moneyOpRelayList = currentMoneyOpRelays(),
).useTor(relay)
fun shouldUseTorForRelay(relay: NormalizedRelayUrl) = TorRelayEvaluation(currentSettings(), currentClassification()).useTor(relay)
fun okHttpClientForRelay(url: NormalizedRelayUrl): OkHttpClient = okHttpClient.getHttpClient(shouldUseTorForRelay(url))
}
@@ -20,13 +20,13 @@
*/
package com.vitorpamplona.amethyst.service.relayClient
import com.vitorpamplona.amethyst.commons.tor.RelayClassification
import com.vitorpamplona.amethyst.commons.tor.TorRelaySettings
import com.vitorpamplona.amethyst.model.torState.TorRelayEvaluation
import com.vitorpamplona.amethyst.service.connectivity.ConnectivityStatus
import com.vitorpamplona.amethyst.service.resourceusage.UsageKeys
import com.vitorpamplona.amethyst.ui.tor.TorServiceStatus
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.utils.Log
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
@@ -101,9 +101,7 @@ class RelayProxyClientConnector(
// flipped relay would sit out its (now-irrelevant) backoff. We track these so such a relay can
// skip its retry delay on the next reconnect — scoped to onlyIfChanged, so only the relays that
// actually flipped re-dial and the rest of the pool's backoff is left untouched.
private var lastTrustedRelays: Set<NormalizedRelayUrl>? = null
private var lastDmRelays: Set<NormalizedRelayUrl>? = null
private var lastMoneyOpRelays: Set<NormalizedRelayUrl>? = null
private var lastClassification: RelayClassification? = null
@OptIn(FlowPreview::class)
val relayServices =
@@ -152,7 +150,7 @@ class RelayProxyClientConnector(
onTrigger(UsageKeys.TRIGGER_OFF)
client.disconnect()
}
if (infra.torStatus is TorServiceStatus.Active) {
if (infra.torStatus.isFullyBootstrapped) {
Log.d("ManageRelayServices", "Connectivity off, Tor idle")
}
// disconnect() already cleared every relay's backoff. Forget the network
@@ -163,7 +161,7 @@ class RelayProxyClientConnector(
infra.connectivity is ConnectivityStatus.Active && !client.isActive() -> {
Log.d("ManageRelayServices", "Connectivity On: Resuming Relay Services")
if (infra.torStatus is TorServiceStatus.Active) {
if (infra.torStatus.isFullyBootstrapped) {
Log.d("ManageRelayServices", "Connectivity resumed, Tor active")
}
@@ -174,9 +172,7 @@ class RelayProxyClientConnector(
lastTorSettings = torSettings
lastTorConnection = infra.torConnection
lastClearConnection = infra.clearConnection
lastTrustedRelays = infra.evaluator.trustedRelayList
lastDmRelays = infra.evaluator.dmRelayList
lastMoneyOpRelays = infra.evaluator.moneyOpRelayList
lastClassification = infra.evaluator.classification
}
else -> {
@@ -202,13 +198,13 @@ class RelayProxyClientConnector(
// so let onlyIfChanged pick out the flipped relay(s) and skip THEIR retry delay —
// without resetBackoff(), so the rest of the pool's backoff is untouched (these sets
// churn while relay lists load, and forgiving the whole pool then would be too much).
//
// One comparison over the whole classification, not one per category: this used to
// be a four-way `||` and adding a category meant remembering to extend it. Missing
// a term fails silently — the affected relays keep a socket on a transport the
// policy has already moved them off.
val classificationChanged =
lastTrustedRelays != null &&
(
infra.evaluator.trustedRelayList != lastTrustedRelays ||
infra.evaluator.dmRelayList != lastDmRelays ||
infra.evaluator.moneyOpRelayList != lastMoneyOpRelays
)
lastClassification != null && infra.evaluator.classification != lastClassification
val previousNetworkId = lastNetworkId
@@ -216,9 +212,7 @@ class RelayProxyClientConnector(
lastClearConnection = infra.clearConnection
lastNetworkId = networkId ?: lastNetworkId
lastTorSettings = torSettings
lastTrustedRelays = infra.evaluator.trustedRelayList
lastDmRelays = infra.evaluator.dmRelayList
lastMoneyOpRelays = infra.evaluator.moneyOpRelayList
lastClassification = infra.evaluator.classification
if (networkChanged) {
Log.d("ManageRelayServices") {
@@ -98,7 +98,6 @@ import com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.publicChannels.relayG
import com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.publicChannels.relayGroup.datasource.RelayGroupsOnRelaySubscription
import com.vitorpamplona.amethyst.ui.stringRes
import com.vitorpamplona.amethyst.ui.theme.warningColor
import com.vitorpamplona.amethyst.ui.tor.TorServiceStatus
import com.vitorpamplona.quartz.buzz.workspace.BUZZ_CHANNEL_TYPE_DM
import com.vitorpamplona.quartz.buzz.workspace.BUZZ_CHANNEL_TYPE_FORUM
import com.vitorpamplona.quartz.nip01Core.core.HexKey
@@ -351,7 +350,7 @@ fun RelayGroupChannelListScreen(
// own dialog), so don't let it read as "this relay blocks Tor exits".
val torStatus by Amethyst.instance.torManager.status
.collectAsStateWithLifecycle()
val torIsUp = torStatus is TorServiceStatus.Active
val torIsUp = torStatus.isFullyBootstrapped
// The offer adds the relay to the kind-10089 Trusted list, which only moves it to clearnet while
// trusted relays are *off* Tor. Under the Small-Payloads / Full-Privacy presets they are on Tor,
@@ -49,6 +49,30 @@ object ArtiNative {
*/
external fun initialize(dataDir: String): Int
/**
* Whether Tor can carry traffic right now.
*
* [initialize] returns as soon as the client exists (the proxy is routable immediately and each
* stream waits for its own circuit), so this is the separate signal for "circuits can be built
* now". Polled rather than pushed: driving state off Arti's log strings is what caused the
* Connecting->Active race this wrapper already had to fix once.
*
* Reports Arti's *live* readiness rather than the outcome of the initial download, so a
* bootstrap that failed once and then succeeded on a later stream is picked up.
*
* @return 1 when ready for traffic, 0 when not yet, -1 when there is no client.
*/
external fun isBootstrapped(): Int
/**
* Directory-download progress in permille (0..1000), or -1 when there is no client.
*
* Lets the lifecycle tell a slow download from a stalled one, which a timeout cannot: measured
* cold downloads ran 12.6-34.4s on the same hardware, so any fixed patience is either short
* enough to kill healthy ones or long enough to sit on a dead one.
*/
external fun bootstrapProgressPermille(): Int
/**
* Start the SOCKS5 proxy on the given port.
* Can be called multiple times — stops any existing listener first.
@@ -31,6 +31,27 @@ import kotlinx.coroutines.flow.StateFlow
interface TorBackend {
val status: StateFlow<TorServiceStatus>
/**
* True while a native bootstrap attempt is actually running.
*
* [TorServiceStatus.Connecting] conflates two states that need opposite responses: a bootstrap
* that is working through a cold consensus download (leave it alone) and a lifecycle that has
* stopped trying (reset it). The stuck-Connecting watchdog cannot tell them apart from status
* alone, and a fresh install spends its first minute in the first one — so the watchdog used to
* queue a reset behind the in-flight attempt's lifecycle lock and tear the client down the
* moment it succeeded. This flag is the missing half of the signal.
*/
val bootstrapInFlight: StateFlow<Boolean>
/**
* Directory-download progress in permille while [TorServiceStatus.Bootstrapping]; -1 when there
* is no client or nothing is being downloaded.
*
* Emits only on change (it is a [StateFlow]), which is exactly the signal the stall detector
* needs: the timestamp of the last distinct value is the last time Tor made forward progress.
*/
val bootstrapProgress: StateFlow<Int>
suspend fun start()
suspend fun stop()
@@ -26,6 +26,7 @@ import kotlinx.coroutines.CoroutineDispatcher
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharingStarted
@@ -34,6 +35,7 @@ import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.drop
import kotlinx.coroutines.flow.emitAll
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.flowOn
import kotlinx.coroutines.flow.launchIn
import kotlinx.coroutines.flow.map
@@ -103,6 +105,36 @@ class TorManager(
*/
@Volatile private var hasEverBootstrapped: Boolean = false
/**
* Epoch-millis of the first moment Tor was expected to work and didn't, spanning the self-heal
* retries in between. 0 while Tor is working or off. See [connectionFailure].
*/
@Volatile private var tryingSinceMs: Long = 0L
/**
* Epoch-millis of the most recent transition INTO a trying state. Distinct from
* [tryingSinceMs], which deliberately spans self-heal retries: this one restarts on every
* attempt, because the patience owed to a directory download is per-attempt.
*/
@Volatile private var lastTryingTransitionMs: Long = 0L
/**
* Epoch-millis of the last time the directory download moved. Seeded when a download starts so
* a fresh attempt is never mistaken for a stalled one, then stamped by the collector below on
* every distinct progress value.
*/
@Volatile private var lastProgressAtMs: Long = 0L
/**
* Whether the bypass prompt is currently raised. Survives the transient [TorServiceStatus.Off]
* a self-heal reset passes through, so the dialog stays up instead of blinking. Cleared
* wherever [tryingSinceMs] is.
*/
@Volatile private var failureRaised: Boolean = false
/** Consecutive gentle (state-preserving) self-heals with no successful bootstrap in between. */
@Volatile private var consecutiveGentleResets: Int = 0
init {
// Seed hasEverBootstrapped from persisted on-disk evidence before the watchdog can fire
// (well under SELF_HEAL_AFTER_MS), so a stuck bootstrap on a previously-working install
@@ -129,6 +161,8 @@ class TorManager(
.onEach {
sessionBypass.value = false
lastBypassApprovalMs = 0L
tryingSinceMs = 0L
failureRaised = false
torPrefs.saveLastBypassApprovalMs(0L)
}.launchIn(scope)
}
@@ -150,8 +184,19 @@ class TorManager(
}
when (torType) {
TorType.INTERNAL -> {
service.start()
emitAll(service.status)
// Subscribe to the backend's status BEFORE start(), not after.
//
// start() awaits a blocking JNI bootstrap that runs to its own timeout, so
// awaiting it first meant nothing observed Connecting until that whole attempt
// had already finished. Every timer keyed on the Connecting span therefore
// started one full attempt late: on device the stuck-watchdog fired at 105s
// instead of 45s, and the connection-failure dialog measured its 60s from the
// wrong instant. Running start() alongside the emitAll makes the status the
// app reacts to the status the service is actually in.
coroutineScope {
launch { service.start() }
emitAll(service.status)
}
}
TorType.OFF -> {
@@ -181,11 +226,11 @@ class TorManager(
val activePortOrNull: StateFlow<Int?> =
status
.map {
(it as? TorServiceStatus.Active)?.port
it.socksPort
}.stateIn(
scope,
SharingStarted.WhileSubscribed(2000),
(status.value as? TorServiceStatus.Active)?.port,
status.value.socksPort,
)
/**
@@ -198,17 +243,45 @@ class TorManager(
val connectionFailure: StateFlow<Boolean> =
status
.transformLatest { s ->
if (s is TorServiceStatus.Connecting) {
if (!s.isTryingToConnect()) {
// Deliberately does NOT clear [tryingSinceMs]: the self-heal watchdog's own
// reset passes through Off on its way back to Bootstrapping, and clearing here
// would let the retry cycle rearm the timer forever. It is cleared where the
// outage genuinely ends — on a bootstrapped Tor, or on a user intent change.
//
// For the same reason this re-emits [failureRaised] rather than a flat false:
// that transient Off would otherwise dismiss the dialog, and the following
// Bootstrapping would re-raise it immediately (its deadline has already
// passed), so a stuck Tor blinked a modal at the user on every watchdog tick.
emit(failureRaised)
return@transformLatest
}
// Measure from when Tor STOPPED WORKING, not from this status span.
//
// `transformLatest` restarts on every status change, and the self-heal watchdog
// deliberately bounces Off -> Bootstrapping every SELF_HEAL_AFTER_MS (45s) while
// stuck — less than this 60s timeout. Keyed on the span, the timer was reset by its
// own watchdog before it could ever expire, so the user was never offered the
// bypass no matter how long Tor stayed broken. The question being asked is "has Tor
// been down for a minute", which spans those retries.
if (tryingSinceMs == 0L) tryingSinceMs = nowMs()
emit(false)
val remaining = BOOTSTRAP_TIMEOUT_MS - (nowMs() - tryingSinceMs)
if (remaining > 0) delay(remaining)
// Never offer to give up on a bootstrap attempt that is still running: a cold
// install legitimately outlasts this timeout, and prompting mid-attempt asks the
// user to abandon something that is working.
service.bootstrapInFlight.first { !it }
if (rememberedApprovalActive()) {
sessionBypass.value = true
emit(false)
delay(BOOTSTRAP_TIMEOUT_MS)
if (rememberedApprovalActive()) {
sessionBypass.value = true
emit(false)
} else {
emit(true)
}
} else {
emit(false)
failureRaised = true
emit(true)
}
}.stateIn(
scope,
@@ -217,16 +290,26 @@ class TorManager(
)
/**
* Fires once after [SELF_HEAL_AFTER_MS] of continuous [TorServiceStatus.Connecting].
* Drives the watchdog wired up below. `transformLatest` cancels the pending delay
* whenever the status changes, so a brief Connecting blip never fires.
* Fires every [SELF_HEAL_AFTER_MS] for as long as status stays [TorServiceStatus.Connecting].
* Drives the watchdog wired up below. `transformLatest` cancels the pending delay whenever the
* status changes, so a brief Connecting blip never fires.
*
* It **repeats** rather than firing once per Connecting span, and that is the whole point. The
* retry loop is driven by status transitions, but the failure it has to recover from produces
* no transition: when the native bootstrap hits its own timeout, [TorService.start] gives up
* and deliberately leaves status at Connecting for this watchdog to retry. A one-shot signal
* has already been consumed by then, so nothing ever re-armed and Tor sat at Connecting
* forever — no retry, no dialog change, no recovery short of a network-identity change or a
* process restart. Repeating means every stuck span is re-examined until it stops being stuck.
*/
@OptIn(ExperimentalCoroutinesApi::class)
private val selfHealSignal =
status.transformLatest { s ->
if (s is TorServiceStatus.Connecting) {
delay(SELF_HEAL_AFTER_MS)
emit(Unit)
if (s.isTryingToConnect()) {
while (true) {
delay(SELF_HEAL_AFTER_MS)
emit(Unit)
}
}
}
@@ -247,7 +330,21 @@ class TorManager(
// state from a different network needs to go.
status
.onEach {
if (it is TorServiceStatus.Active) hasEverBootstrapped = true
if (it.isFullyBootstrapped) {
hasEverBootstrapped = true
tryingSinceMs = 0L
lastTryingTransitionMs = 0L
failureRaised = false
consecutiveGentleResets = 0
} else if (it.isTryingToConnect() && lastTryingTransitionMs == 0L) {
lastTryingTransitionMs = nowMs()
// A download that has just begun has not stalled, whatever the last attempt did.
lastProgressAtMs = nowMs()
} else if (it == TorServiceStatus.Off) {
// A reset passes through Off on its way to a fresh attempt; the next
// Bootstrapping earns a full patience window of its own.
lastTryingTransitionMs = 0L
}
}.launchIn(scope)
// Rotten guard sample while Tor is otherwise UP. The watchdog above only fires on a status
@@ -258,6 +355,13 @@ class TorManager(
// AllGuardsDown log is the only reliable signal, so route it through the same rate-limited
// wipe. Always a clean-state reset: the whole point is that the persisted sample is the
// problem.
// Stamps the moment Tor last moved. A StateFlow only emits distinct values, so this fires
// exactly when progress changes — making `lastProgressAtMs` the age of the last real
// advance rather than the age of the attempt.
service.bootstrapProgress
.onEach { lastProgressAtMs = nowMs() }
.launchIn(scope)
service.guardsDownSignal
.onEach {
val now = nowMs()
@@ -270,20 +374,91 @@ class TorManager(
selfHealSignal
.onEach {
// Re-check: the signal repeats, so by the time it lands the status may have moved
// on. Resetting a client that just reached Active is the opposite of self-healing.
if (!status.value.isTryingToConnect()) return@onEach
// A native attempt that is still running is not stuck — it is working. (Under
// on-demand bootstrap `initialize` returns in ~130ms, so this only covers client
// creation; the directory download is covered by the patience window below.)
if (service.bootstrapInFlight.value) return@onEach
val downloading = status.value is TorServiceStatus.Bootstrapping
// A running directory download is judged on forward progress, not elapsed time.
//
// A timer cannot tell slow from stalled: measured cold downloads ran 12.6-34.4s on
// this same hardware and network, so any fixed patience is either short enough to
// kill healthy ones — and a reset discards the partial consensus, so firing early
// can stop the download EVER completing — or long enough to sit uselessly on a dead
// one. Progress separates them exactly: a download that is still advancing is left
// alone indefinitely, and one that has not moved at all is reset promptly.
if (downloading && nowMs() - lastProgressAtMs < BOOTSTRAP_STALL_MS) return@onEach
val now = nowMs()
if (now - lastSelfHealAtMs < SELF_HEAL_COOLDOWN_MS) return@onEach
if (now - lastSelfHealAtMs < selfHealCooldownMs()) return@onEach
lastSelfHealAtMs = now
if (hasEverBootstrapped) {
Log.w("TorManager") { "Tor stuck Connecting >${SELF_HEAL_AFTER_MS}ms — self-healing (drop client + wipe state)" }
// Never wipe while downloading. `resetWithCleanState` deletes `arti/cache`, which
// is precisely the consensus this state is in the middle of fetching: wiping it
// guarantees the next attempt restarts from zero, and on a slow link that loops
// forever. The clean-state hammer is for a lifecycle that cannot even get a client
// up, where the persisted state is the prime suspect.
// Escalate a fresh install that cannot even get a client up.
//
// The inline `clearAllArtiData()` retry used to cover corrupt on-disk state; it was
// removed because it fired on every failure, including "no network". But with
// `hasEverBootstrapped` false there is no confirmed guard on disk, so the gentle
// branch below would drop the client forever without ever wiping — and
// `noUsableGuards()` only inspects `guards.json`, so a corrupt `cache/` or the rest
// of `state/` is invisible to it. Escalate once the gentle path has demonstrably
// failed several times in a row.
val exhaustedGentleRetries = !downloading && consecutiveGentleResets >= GENTLE_RESETS_BEFORE_WIPE
if ((hasEverBootstrapped || exhaustedGentleRetries) && !downloading) {
Log.w("TorManager") { "Tor stuck with no client for >${SELF_HEAL_AFTER_MS}ms — self-healing (drop client + wipe state)" }
consecutiveGentleResets = 0
service.resetWithCleanState()
} else {
Log.w("TorManager") { "Tor stuck Connecting >${SELF_HEAL_AFTER_MS}ms on first bootstrap — self-healing (drop client only)" }
consecutiveGentleResets++
val what =
if (downloading) {
"directory download stuck at ${service.bootstrapProgress.value}/1000 for >${BOOTSTRAP_STALL_MS}ms"
} else {
"stuck with no client >${SELF_HEAL_AFTER_MS}ms"
}
Log.w("TorManager") { "Tor $what — self-healing (drop client only, keeping the consensus cache)" }
service.reset()
}
resetEpoch.update { it + 1 }
}.launchIn(scope)
}
/**
* How long to wait between self-heals.
*
* Once Tor has bootstrapped on this install, a reset is expensive and rarely the answer, so the
* full [SELF_HEAL_COOLDOWN_MS] applies — a permanently broken network must not put us in a
* reset loop. Before the first successful bootstrap the trade is reversed: there is no working
* state to protect, retrying is nearly free (Arti's directory cache persists across attempts,
* so each retry resumes rather than restarts), and the alternative is a brand-new install
* sitting on a dead Tor for five minutes at a time. So a fresh install retries on
* [FIRST_BOOTSTRAP_RETRY_COOLDOWN_MS] instead.
*/
private fun selfHealCooldownMs(): Long = if (hasEverBootstrapped) SELF_HEAL_COOLDOWN_MS else FIRST_BOOTSTRAP_RETRY_COOLDOWN_MS
/**
* Tor is meant to be working and isn't yet — the span both the stuck watchdog and the
* connection-failure dialog exist to bound.
*
* It is deliberately NOT `is Connecting`. Under on-demand bootstrap the client is created in
* ~130ms, so status leaves Connecting almost immediately and spends the entire 12-34s directory
* download in [TorServiceStatus.Bootstrapping]. Keying on Connecting alone would have made both
* safety nets unreachable: a download that never completes would sit at Bootstrapping forever
* with nothing watching it.
*/
private fun TorServiceStatus.isTryingToConnect() = this != TorServiceStatus.Off && !isFullyBootstrapped
fun rememberedApprovalActive(): Boolean {
val ts = lastBypassApprovalMs
return ts > 0 && (nowMs() - ts) < APPROVAL_REMEMBER_MS
@@ -292,6 +467,8 @@ class TorManager(
/** Called when the user picks "Use regular connection". Starts a fresh 1-hour window. */
fun approveBypassForOneHour() {
val now = nowMs()
tryingSinceMs = 0L
failureRaised = false
lastBypassApprovalMs = now
sessionBypass.value = true
scope.launch(ioDispatcher) {
@@ -311,6 +488,8 @@ class TorManager(
fun onNetworkChange() {
sessionBypass.value = false
lastBypassApprovalMs = 0L
tryingSinceMs = 0L
failureRaised = false
// Prevent the stuck-Connecting watchdog from firing a second reset while the
// network-change bootstrap is still legitimately in progress (initial bootstrap
// on a new network can take ~10–30s, sometimes longer).
@@ -349,7 +528,7 @@ class TorManager(
*/
fun onTorCircuitsDead() {
if (sessionBypass.value) return
if (status.value !is TorServiceStatus.Active) return
if (!status.value.isFullyBootstrapped) return
val now = nowMs()
if (now - lastSelfHealAtMs < SELF_HEAL_COOLDOWN_MS) return
lastSelfHealAtMs = now
@@ -360,9 +539,28 @@ class TorManager(
}
}
fun isSocksReady() = status.value is TorServiceStatus.Active
/**
* Whether a Tor-routed dial has somewhere to go. Both callers
* (`AppModules`' relay gate and the media-http `isTorActive` probe) are asking "can I send this
* through Tor", not "is the directory ready" — a dial during the download queues on its own
* circuit, which is strictly better than the alternative of refusing it or sending it in clear.
*/
fun isSocksReady() = status.value.socksPort != null
fun socksPort(): Int = (status.value as? TorServiceStatus.Active)?.port ?: 17392
/**
* Tor can carry traffic now — the gate for "start using Tor", as opposed to [isSocksReady]'s
* "route through it if you do".
*
* Dialling merely because the port exists costs more than it saves: measured, relays dialled
* during the download simply time out (Tor connect timeout is 30s, inside the 12-34s window)
* and enter exponential backoff, and because the port is identical either side of
* Bootstrapping -> Active the transport never "changes", so `RelayProxyClientConnector` never
* calls `resetBackoff()` to forgive them. Time-to-first-socket was unchanged by dialling early
* (n=3), so the backoff is pure loss.
*/
fun isTorReady() = status.value.isFullyBootstrapped
fun socksPort(): Int = status.value.socksPort ?: 17392
companion object {
const val BOOTSTRAP_TIMEOUT_MS: Long = 60_000L
@@ -371,5 +569,22 @@ class TorManager(
/** Self-heal kicks in BEFORE the 60s [BOOTSTRAP_TIMEOUT_MS] dialog so most users never see it. */
const val SELF_HEAL_AFTER_MS: Long = 45_000L
const val SELF_HEAL_COOLDOWN_MS: Long = 5L * 60L * 1000L
/** Cooldown before Tor has ever bootstrapped on this install. See [selfHealCooldownMs]. */
const val FIRST_BOOTSTRAP_RETRY_COOLDOWN_MS: Long = 30_000L
/**
* How long a directory download may make **no forward progress at all** before it counts as
* stalled. Time spent downloading does not count against it — only time spent not moving.
*/
const val BOOTSTRAP_STALL_MS: Long = 60_000L
/**
* Gentle self-heals to try before wiping on-disk state on an install that has never
* bootstrapped. Corrupt `cache/`/`state/` is invisible to [ArtiGuardState.hasNoUsableGuards]
* (it only reads `guards.json`), so without this a fresh install with a bad cache would
* drop-and-retry the client forever and never clear the thing actually blocking it.
*/
const val GENTLE_RESETS_BEFORE_WIPE: Int = 3
}
}
@@ -24,14 +24,18 @@ import android.content.Context
import com.fasterxml.jackson.databind.JsonNode
import com.fasterxml.jackson.module.kotlin.jacksonObjectMapper
import com.vitorpamplona.quartz.utils.Log
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
@@ -47,16 +51,12 @@ private const val GUARDS_DOWN_THRESHOLD = 40
/** Window for [GUARDS_DOWN_THRESHOLD]. Wide enough that ordinary transient churn never trips it. */
private const val GUARDS_DOWN_WINDOW_MS = 60_000L
/** Cheap: an atomic read against a download that takes 12-34s. */
private const val BOOTSTRAP_POLL_MS = 500L
private const val DEFAULT_SOCKS_PORT = 17392
private const val MAX_PORT_RETRIES = 10
/**
* Return code from [ArtiNative.initialize] when the native bootstrap exceeds its
* internal timeout. The native side has already torn down the half-built client;
* we treat this differently from a hard failure (see [TorService.start]).
*/
private const val ARTI_ERROR_BOOTSTRAP_TIMEOUT = -4
/**
* Manages the Arti Tor client via custom JNI bindings.
*
@@ -70,6 +70,7 @@ private const val ARTI_ERROR_BOOTSTRAP_TIMEOUT = -4
*/
class TorService(
val context: Context,
private val scope: CoroutineScope,
) : TorBackend {
private var socksPort = DEFAULT_SOCKS_PORT
private val initialized = AtomicBoolean(false)
@@ -82,6 +83,16 @@ class TorService(
*/
@Volatile private var bootstrapStartedAtMs: Long = -1L
/**
* How many native bootstrap attempts this process has made. Logged at INFO on every attempt so
* a boot log answers "did Tor retry, and how often" — the question that separates a genuinely
* slow first bootstrap from a lifecycle that stopped retrying altogether.
*/
@Volatile private var bootstrapAttempts: Int = 0
/** Poller promoting Bootstrapping -> Active. Cancelled by every reset so it can't outlive its client. */
private var bootstrapWatcher: Job? = null
/**
* Serializes every native lifecycle transition ([start], [stop], [reset],
* [resetWithCleanState]). [ArtiNative] is a process-global singleton over a
@@ -98,6 +109,17 @@ class TorService(
private val _status = MutableStateFlow<TorServiceStatus>(TorServiceStatus.Off)
override val status: StateFlow<TorServiceStatus> = _status.asStateFlow()
/**
* True for exactly as long as [ArtiNative.initialize] is running. See
* [TorBackend.bootstrapInFlight] — this is what lets the stuck-Connecting watchdog tell a
* bootstrap that is still working from a lifecycle that has given up.
*/
private val _bootstrapInFlight = MutableStateFlow(false)
override val bootstrapInFlight: StateFlow<Boolean> = _bootstrapInFlight.asStateFlow()
private val _bootstrapProgress = MutableStateFlow(-1)
override val bootstrapProgress: StateFlow<Int> = _bootstrapProgress.asStateFlow()
/**
* Every status change goes through here so the transition is logged exactly once, at INFO, with
* the time since bootstrap started.
@@ -159,7 +181,14 @@ class TorService(
private fun artiDataDir() = File(context.filesDir, "arti")
/** Diagnostic: total bytes of the consensus/descriptor cache, to correlate with bootstrap time. */
/**
* Diagnostic: total bytes of the consensus/descriptor cache, to correlate with bootstrap time —
* a cold 0-byte cache costs 12-34s where a warm one costs ~6s, so it is the first thing you
* want beside a slow bootstrap.
*
* Walks the whole cache directory, so it is called exactly once per attempt, from the INFO line
* below. Read it there rather than adding another call.
*/
private fun cacheSizeBytes(): Long {
val cacheDir = File(artiDataDir(), "cache")
if (!cacheDir.exists()) return 0
@@ -235,8 +264,20 @@ class TorService(
override suspend fun start() =
lifecycleMutex.withLock {
if (proxyRunning.get()) {
if (_status.value is TorServiceStatus.Active) return@withLock
setStatus(TorServiceStatus.Connecting)
// The proxy is already bound, so re-assert the state that matches reality rather
// than falling back to Connecting.
//
// Connecting reports no port. Downgrading to it here would strand a perfectly good
// listener: `activePortOrNull` goes null, every Tor-routed dial drops to the Orbot
// default 9050 where nothing listens, and — because this returns without arming
// [watchBootstrap] — nothing would ever promote it back. `TorManager` re-enters
// this branch on any combine re-fire (torType/port/bypass change, resetEpoch bump,
// or the status flow restarting after WhileSubscribed's 30s timeout), so it is very
// much reachable.
if (_status.value !is TorServiceStatus.Active) {
setStatus(TorServiceStatus.Bootstrapping(socksPort))
watchBootstrap(socksPort)
}
return@withLock
}
@@ -272,7 +313,7 @@ class TorService(
// in state/ — not cache/ — and Arti already validates consensus freshness
// and refetches whatever has expired. The reset/clean-state paths still call
// clearAllArtiData() for genuine corruption recovery.
Log.d("TorService") { "Preserving Arti cache for warm bootstrap (cache size: ${cacheSizeBytes()} bytes)" }
Log.d("TorService") { "Preserving Arti cache for warm bootstrap" }
// Self-heal the wedged guard sample (see [noUsableGuards]): if
// the persisted sample has no usable guard left, Arti can
@@ -288,29 +329,37 @@ class TorService(
Log.d("TorService") { "Initializing Arti with data dir: $dataDir" }
bootstrapStartedAtMs = System.currentTimeMillis()
var initResult = ArtiNative.initialize(dataDir)
if (initResult == ARTI_ERROR_BOOTSTRAP_TIMEOUT) {
// The native bootstrap hit its timeout (hostile network) and
// already tore down the half-built client. Don't wipe state or
// retry inline — that would hold lifecycleMutex for another full
// timeout. Drop the init flag and leave status at Connecting so
// TorManager's self-heal watchdog resets and retries on its own
// cadence (and the connection-failure dialog can still surface).
Log.w("TorService") { "Arti bootstrap timed out — leaving Connecting for the self-heal watchdog to retry" }
initialized.set(false)
return@withContext
}
bootstrapAttempts++
Log.i("TorService") { "Bootstrapping Arti (attempt $bootstrapAttempts, cache ${cacheSizeBytes()} bytes)" }
_bootstrapInFlight.value = true
val initResult =
try {
ArtiNative.initialize(dataDir)
} finally {
_bootstrapInFlight.value = false
}
Log.i("TorService") { "Arti bootstrap attempt $bootstrapAttempts returned $initResult after ${System.currentTimeMillis() - bootstrapStartedAtMs}ms" }
if (initResult != 0) {
Log.e("TorService") { "Failed to initialize Arti: error $initResult, clearing data and retrying" }
clearAllArtiData()
initResult = ArtiNative.initialize(dataDir)
}
if (initResult != 0) {
Log.e("TorService") { "Failed to initialize Arti on retry: error $initResult" }
// Every failure mode ends the same way: don't decide recovery here.
//
// A timeout has already torn down its half-built client natively, and
// retrying inline would hold lifecycleMutex for another full timeout. A
// hard failure used to be treated differently — wipe all Arti data, retry
// inline, then fall back to status Off — and both halves of that were
// wrong. The wipe treated every failure as corruption, so a bootstrap that
// failed for the most ordinary reason there is (no network) threw away a
// guard sample that was working fine. And Off is a terminal state for this
// lifecycle: the watchdog and the connection-failure dialog both only arm
// while status is Connecting, so a hard failure left Tor switched off with
// nothing ever retrying it.
//
// Leaving Connecting hands the decision to [TorManager], which already
// owns the escalation: a gentle client-drop before Tor has ever
// bootstrapped here, a clean-state wipe once a confirmed guard on disk
// proves the persisted state used to work and is therefore suspect.
Log.w("TorService") { "Arti client creation failed (error $initResult) — leaving Connecting for the self-heal watchdog to retry" }
initialized.set(false)
setStatus(TorServiceStatus.Off)
return@withContext
}
}
@@ -330,8 +379,12 @@ class TorService(
}
if (!started) {
Log.e("TorService") { "Failed to start SOCKS proxy after $MAX_PORT_RETRIES attempts" }
setStatus(TorServiceStatus.Off)
// Same reasoning as the init-failure branch: Off is terminal for this
// lifecycle — neither the watchdog nor the connection-failure dialog arms on
// it — so reporting Off here would leave Tor silently disabled with nothing
// retrying and no way for the user to find out. Stay Connecting and let the
// watchdog retry; a port collision is usually transient.
Log.w("TorService") { "Failed to bind a SOCKS port after $MAX_PORT_RETRIES attempts — leaving Connecting for the self-heal watchdog to retry" }
return@withContext
}
@@ -350,11 +403,54 @@ class TorService(
// reset/stop can't clobber it.
val startedAt = bootstrapStartedAtMs
val elapsed = if (startedAt > 0) System.currentTimeMillis() - startedAt else -1
setStatus(TorServiceStatus.Active(socksPort))
Log.d("TorService") { "Arti SOCKS proxy active on port $socksPort (bootstrap took ${elapsed}ms)" }
// Routable, not yet bootstrapped. The proxy is bound so dials belong here rather
// than at the dead 9050 fallback, but circuits cannot be built until the directory
// download lands — which [watchBootstrap] turns into Active.
setStatus(TorServiceStatus.Bootstrapping(socksPort))
Log.i("TorService") { "Arti SOCKS proxy routable on port $socksPort after ${elapsed}ms (directory still downloading)" }
watchBootstrap(socksPort)
}
}
/**
* Polls the native directory-download result and promotes [TorServiceStatus.Bootstrapping] to
* [TorServiceStatus.Active] once circuits can actually be built.
*
* Polling rather than reacting to a log line is deliberate: the previous log-callback-driven
* transition raced `startSocksProxy` returning and silently dropped the Active transition,
* stranding Tor at Connecting until the 60s dialog. The poll reads Arti's live readiness, so a
* download that fails and is retried by a later stream still promotes.
*/
private fun watchBootstrap(port: Int) {
bootstrapWatcher?.cancel()
bootstrapWatcher =
scope.launch(Dispatchers.IO) {
while (true) {
// Publish progress on the same tick we check readiness — one extra cheap JNI
// read, and it is what lets TorManager distinguish slow from stalled.
_bootstrapProgress.value = ArtiNative.bootstrapProgressPermille()
when (ArtiNative.isBootstrapped()) {
1 -> {
// Only promote if this is still the run we started watching for: a
// reset in between will have moved us to Off/Connecting already.
if (_status.value == TorServiceStatus.Bootstrapping(port)) {
setStatus(TorServiceStatus.Active(port))
}
return@launch
}
-1 -> {
// No native client behind the proxy — a reset is in flight, or init
// failed. Nothing to promote; leave the status where it is so
// TorManager's watchdog treats it as stuck and retries.
Log.w("TorService") { "No Arti client while watching bootstrap — leaving status for the self-heal watchdog" }
return@launch
}
else -> delay(BOOTSTRAP_POLL_MS)
}
}
}
}
/**
* Stop the SOCKS proxy and release the port.
* The TorClient stays alive — no file lock issues on restart.
@@ -363,6 +459,10 @@ class TorService(
lifecycleMutex.withLock {
if (!proxyRunning.compareAndSet(true, false)) return@withLock
bootstrapWatcher?.cancel()
bootstrapWatcher = null
_bootstrapProgress.value = -1
withContext(Dispatchers.IO) {
ArtiNative.stopSocksProxy()
Log.d("TorService") { "SOCKS proxy stopped" }
@@ -409,6 +509,9 @@ class TorService(
*/
private suspend fun resetLocked() =
withContext(Dispatchers.IO) {
bootstrapWatcher?.cancel()
bootstrapWatcher = null
_bootstrapProgress.value = -1
if (proxyRunning.compareAndSet(true, false)) {
ArtiNative.stopSocksProxy()
}
@@ -20,12 +20,52 @@
*/
package com.vitorpamplona.amethyst.ui.tor
/**
* Two independent facts about Tor that used to coincide and no longer do.
*
* The SOCKS proxy binds within ~130ms of process start, but the directory download it needs before
* it can build a circuit takes 12-34s on a cold install. While `create_bootstrapped` blocked until
* both were true, one "Active" could honestly mean both. Under `BootstrapBehavior::OnDemand` the
* proxy is usable immediately and streams wait for their own circuits, so the two facts diverge by
* that whole window — and callers want different ones. Read [socksPort] to route bytes and
* [isFullyBootstrapped] to tell a user (or a watchdog) whether Tor is actually working; matching on
* the variants directly is how you end up with the wrong one, because both look like "Active".
*/
sealed class TorServiceStatus {
/** Proxy bound and the directory is ready: circuits build immediately. */
data class Active(
val port: Int,
) : TorServiceStatus()
/**
* Proxy bound and routable, directory still downloading. Dials sent here are not lost — each
* stream waits for its own circuit — but they will not complete until the download lands.
*/
data class Bootstrapping(
val port: Int,
) : TorServiceStatus()
object Off : TorServiceStatus()
/** No proxy yet: the native client is still being created. Nothing is routable. */
object Connecting : TorServiceStatus()
/**
* Where to send bytes, or null if there is nowhere to send them.
*
* [Bootstrapping] counts. Routing to it queues the dial behind the directory download, which is
* what we want; treating it as "no proxy" is what made every dial fall back to the Orbot
* default port 9050, where nothing listens, and fail instantly into backoff.
*/
val socksPort: Int?
get() =
when (this) {
is Active -> port
is Bootstrapping -> port
else -> null
}
/** Tor can build circuits right now. What a user is told, and what the watchdogs judge. */
val isFullyBootstrapped: Boolean
get() = this is Active
}
Binary file not shown.
@@ -53,8 +53,11 @@ class TorRelayEvaluationTest {
newRelaysViaTor = newViaTor,
trustedRelaysViaTor = trustedViaTor,
),
trustedRelayList = trustedRelays,
dmRelayList = dmRelays,
classification =
RelayClassification(
trusted = trustedRelays,
dm = dmRelays,
),
)
// --- Tor OFF: always false ---
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.amethyst.service.relayClient
import com.vitorpamplona.amethyst.commons.tor.RelayClassification
import com.vitorpamplona.amethyst.commons.tor.TorRelaySettings
import com.vitorpamplona.amethyst.commons.tor.TorType
import com.vitorpamplona.amethyst.model.torState.TorRelayEvaluation
@@ -28,6 +29,7 @@ import com.vitorpamplona.amethyst.service.relayClient.RelayProxyClientConnector.
import com.vitorpamplona.amethyst.ui.tor.TorServiceStatus
import com.vitorpamplona.quartz.nip01Core.relay.client.EmptyNostrClient
import com.vitorpamplona.quartz.nip01Core.relay.client.INostrClient
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import io.mockk.mockk
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
@@ -112,8 +114,11 @@ class RelayProxyClientConnectorTest {
private fun evaluation(settings: TorRelaySettings = TorRelaySettings()) =
TorRelayEvaluation(
torSettings = settings,
trustedRelayList = emptySet(),
dmRelayList = emptySet(),
classification =
RelayClassification(
trusted = emptySet(),
dm = emptySet(),
),
)
private fun infra(
@@ -179,6 +184,43 @@ class RelayProxyClientConnectorTest {
assertEquals(listOf(true to true), client.reconnects)
}
/**
* The handover the guessed-relay feature exists to perform: when a user's own lists arrive, the
* relays we were guessing about must re-dial onto whatever policy the user actually chose.
*
* The case that matters is an **empty** arriving list. Normally `trusted` grows at the same
* moment and would have flagged the change on its own — but an account that publishes an empty
* relay list leaves `trusted` untouched while `assumed` empties, and before the classification
* was compared as one value that combination produced no reconnect at all, stranding those
* relays on clearnet against a policy that had already moved them to Tor.
*/
@Test
fun `guessed relays re-dial when an empty list arrives and trusted does not change`() {
settleOnFirstNetwork()
val guessing =
TorRelayEvaluation(
torSettings = TorRelaySettings(),
classification = RelayClassification(assumed = setOf(NormalizedRelayUrl("wss://guessed.example/"))),
)
connector.apply(infra(networkId = 1L, evaluation = guessing))
client.reconnects.clear()
// The user's own list arrives and is empty: `assumed` empties, `trusted` stays empty.
val released =
TorRelayEvaluation(
torSettings = TorRelaySettings(),
classification = RelayClassification(),
)
connector.apply(infra(networkId = 1L, evaluation = released))
assertEquals(
"The relays we stopped guessing about must be asked to re-dial onto their real policy",
listOf(true to true),
client.reconnects,
)
}
@Test
fun `unrelated churn on the same network leaves the backoff alone`() {
settleOnFirstNetwork()
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.amethyst.ui.screen.loggedIn.chats.publicChannels.relayGroup
import com.vitorpamplona.amethyst.commons.tor.RelayClassification
import com.vitorpamplona.amethyst.commons.tor.TorRelayEvaluation
import com.vitorpamplona.amethyst.commons.tor.TorRelaySettings
import com.vitorpamplona.amethyst.commons.tor.TorSettings
@@ -88,8 +89,11 @@ class TorClearnetFallbackTest {
trustedRelaysViaTor = preset.trustedRelaysViaTor,
moneyOperationsViaTor = preset.moneyOperationsViaTor,
),
trustedRelayList = emptySet(),
dmRelayList = emptySet(),
classification =
RelayClassification(
trusted = emptySet(),
dm = emptySet(),
),
)
@Test
@@ -22,6 +22,7 @@ package com.vitorpamplona.amethyst.ui.tor
import com.vitorpamplona.amethyst.commons.tor.TorType
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.awaitCancellation
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
@@ -199,8 +200,8 @@ class TorManagerTest {
val backend = FakeTorBackend()
val manager = buildManager(backend = backend, clock = { 1_000_000_000_000L })
advanceUntilIdle()
// Still Connecting (never reached Active).
assertEquals(TorServiceStatus.Connecting, manager.status.value)
// Started but never reached a working Tor.
assertFalse(manager.status.value.isFullyBootstrapped)
val resetCountBefore = backend.resetCount
manager.onTorCircuitsDead()
@@ -270,11 +271,14 @@ class TorManagerTest {
fun `watchdog uses gentle reset before first Active`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
// No client at all: this is the fast 45s no-client cadence, not a running download.
backend.startFailsToConnect = true
// Big constant clock so (now - lastSelfHealAtMs=0) is well past cooldown.
val manager = buildManager(backend = backend, clock = { 1_000_000_000_000L })
advanceUntilIdle()
assertEquals(TorServiceStatus.Connecting, manager.status.value)
// Big constant clock so (now - lastSelfHealAtMs=0) is well past cooldown.
assertFalse(manager.status.value.isFullyBootstrapped)
assertEquals(0, backend.resetCount)
advanceTimeBy(TorManager.SELF_HEAL_AFTER_MS + 1_000L)
@@ -288,11 +292,14 @@ class TorManagerTest {
fun `watchdog wipes state on first stuck-Connecting when guards prove prior bootstrap`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend().apply { bootstrappedBefore = true }
// No client at all: this is the fast 45s no-client cadence, not a running download.
backend.startFailsToConnect = true
// Never reaches Active in this session, but on-disk state proves a prior bootstrap.
val manager = buildManager(backend = backend, clock = { 1_000_000_000_000L })
advanceUntilIdle()
assertEquals(TorServiceStatus.Connecting, manager.status.value)
// Never reaches Active in this session, but on-disk state proves a prior bootstrap.
assertFalse(manager.status.value.isFullyBootstrapped)
advanceTimeBy(TorManager.SELF_HEAL_AFTER_MS + 1_000L)
runCurrent()
@@ -348,6 +355,8 @@ class TorManagerTest {
fun `watchdog cooldown blocks a second fire within the window`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
// No client at all: this is the fast 45s no-client cadence, not a running download.
backend.startFailsToConnect = true
var clockNow = 1_000_000_000_000L
val manager = buildManager(backend = backend, clock = { clockNow })
advanceUntilIdle()
@@ -368,6 +377,8 @@ class TorManagerTest {
fun `watchdog can fire again once the cooldown elapses`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
// No client at all: this is the fast 45s no-client cadence, not a running download.
backend.startFailsToConnect = true
var clockNow = 1_000_000_000_000L
val manager = buildManager(backend = backend, clock = { clockNow })
advanceUntilIdle()
@@ -430,7 +441,7 @@ class TorManagerTest {
val manager = buildManager(backend = backend)
advanceUntilIdle()
assertEquals(TorServiceStatus.Connecting, manager.status.value)
assertEquals(TorServiceStatus.Bootstrapping(17392), manager.status.value)
backend.setActive(17392)
advanceUntilIdle()
@@ -470,6 +481,396 @@ class TorManagerTest {
assertTrue(backend.stopCount >= 1)
}
// ------------------------------------------------------------------
// fresh install: the bootstrap-timeout retry loop
// ------------------------------------------------------------------
/**
* The regression that stranded brand-new installs on "Connecting" indefinitely.
*
* On a native bootstrap timeout `TorService.start()` deliberately leaves status at Connecting
* and delegates the retry to this watchdog. The watchdog used to fire once per Connecting
* *span* — and a timeout produces no status change, so no new span ever began. Two attempts
* were made and then the app stopped trying, permanently: no retry, no recovery short of a
* network-identity change or a process restart.
*/
@Test
fun `keeps retrying when the bootstrap keeps timing out`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
// No client at all: this is the fast 45s no-client cadence, not a running download.
backend.startFailsToConnect = true
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
advanceTimeBy(10L * 60_000L)
runCurrent()
assertTrue(
"expected repeated bootstrap retries over 10 stuck minutes, got ${backend.startCount}",
backend.startCount >= 4,
)
sub.cancel()
}
/**
* A bootstrap that is still running is not stuck. The native call can hold its lifecycle lock
* for its full timeout, so a reset issued while it runs queues behind it and lands the instant
* the attempt finishes — tearing down a client that may have just succeeded. On a fresh install
* with a cold consensus cache that window is the common case, not a corner case.
*/
@Test
fun `watchdog leaves an in-flight bootstrap alone`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
backend.holdBootstrapInFlight = true
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
// Well past SELF_HEAL_AFTER_MS, but the attempt is still running.
advanceTimeBy(3L * TorManager.SELF_HEAL_AFTER_MS)
runCurrent()
assertEquals(0, backend.resetCount)
assertEquals(0, backend.resetWithCleanStateCount)
assertEquals(1, backend.startCount)
// The attempt returns without reaching Active — now it is genuinely stuck.
backend.finishBootstrapAttempt()
advanceTimeBy(TorManager.SELF_HEAL_AFTER_MS + 1_000L)
runCurrent()
assertTrue("watchdog should fire once the attempt returned", backend.resetCount >= 1)
sub.cancel()
}
/** A fresh install has no working state to protect, so it must not wait out the 5-min cooldown. */
@Test
fun `first-bootstrap retries use the short cooldown`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
// No client at all: this is the fast 45s no-client cadence, not a running download.
backend.startFailsToConnect = true
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
// Two watchdog windows: with the 5-min cooldown only one reset could land.
advanceTimeBy(3L * TorManager.SELF_HEAL_AFTER_MS)
runCurrent()
assertTrue(
"expected more than one retry inside 3 watchdog windows, got ${backend.resetCount}",
backend.resetCount >= 2,
)
sub.cancel()
}
/** Once Tor has worked, resets stay rate-limited — a broken network must not cause a reset loop. */
@Test
fun `after a successful bootstrap the long cooldown still applies`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
backend.setActive(17392)
advanceUntilIdle()
backend.setConnecting()
advanceUntilIdle()
val before = backend.resetWithCleanStateCount
advanceTimeBy(3L * TorManager.SELF_HEAL_AFTER_MS)
runCurrent()
assertEquals(
"post-bootstrap self-heal must stay on the long cooldown",
1,
backend.resetWithCleanStateCount - before,
)
sub.cancel()
}
/**
* `start()` blocks for a whole native bootstrap attempt. Awaiting it before subscribing to the
* backend's status meant nothing observed Connecting until that attempt was already over, so
* every timer keyed on the Connecting span — the stuck watchdog, the connection-failure
* dialog — started one full attempt late.
*/
@Test
fun `status is observable while the first bootstrap is still running`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
backend.startNeverReturns = true
val manager = buildManager(backend = backend)
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
assertFalse(manager.status.value.isFullyBootstrapped)
assertNotEquals(TorServiceStatus.Off, manager.status.value)
sub.cancel()
}
/**
* The bypass prompt asks the user to give up on Tor. Asking that while the first bootstrap
* attempt is still downloading a cold consensus offers to abandon something that is working.
*/
@Test
fun `connection-failure prompt waits for the bootstrap attempt to return`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
backend.holdBootstrapInFlight = true
val manager = buildManager(backend = backend)
val sub = launch { manager.status.collect { } }
val fail = mutableListOf<Boolean>()
val subFail = launch { manager.connectionFailure.collect { fail.add(it) } }
advanceUntilIdle()
advanceTimeBy(TorManager.BOOTSTRAP_TIMEOUT_MS * 2)
runCurrent()
assertFalse("must not prompt while the attempt is still running", fail.contains(true))
backend.finishBootstrapAttempt()
advanceTimeBy(1_000L)
runCurrent()
assertTrue("must prompt once the attempt returned without connecting", fail.contains(true))
sub.cancel()
subFail.cancel()
}
// ------------------------------------------------------------------
// Bootstrapping: routable, not yet ready
// ------------------------------------------------------------------
/**
* The regression this state exists to prevent. On-demand bootstrap leaves Connecting in ~130ms
* and spends the whole 12-34s directory download in Bootstrapping, so a watchdog keyed on
* Connecting would never fire again — a download that never completes would sit there forever
* with nothing watching it.
*/
@Test
fun `watchdog still fires while stuck Bootstrapping`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
assertTrue(
"start() should leave us routable but not bootstrapped",
manager.status.value is TorServiceStatus.Bootstrapping,
)
advanceTimeBy(TorManager.BOOTSTRAP_STALL_MS + 2L * TorManager.SELF_HEAL_AFTER_MS)
runCurrent()
assertTrue("watchdog must arm on Bootstrapping, not just Connecting", backend.resetCount >= 1)
sub.cancel()
}
/** The bypass prompt must also survive the state it now spends its time in. */
@Test
fun `connection-failure prompt fires from a stuck Bootstrapping span`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
// Virtual clock, not the default constant one: the timer now measures elapsed
// wall-clock across the watchdog's retries, so a frozen clock makes it un-expirable.
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
val fail = mutableListOf<Boolean>()
val subFail = launch { manager.connectionFailure.collect { fail.add(it) } }
advanceUntilIdle()
advanceTimeBy(TorManager.BOOTSTRAP_TIMEOUT_MS + 1_000L)
runCurrent()
assertTrue(
"the prompt must survive the self-heal watchdog restarting the status span at 45s",
fail.contains(true),
)
sub.cancel()
subFail.cancel()
}
/**
* The whole point of the split: a dial issued during the download must be routed through the
* proxy, not dropped to the Orbot default port where nothing listens.
*/
@Test
fun `port is routable while still bootstrapping`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
val manager = buildManager(backend = backend)
val sub = launch { manager.status.collect { } }
val port = launch { manager.activePortOrNull.collect { } }
advanceUntilIdle()
backend.setBootstrapping(17392)
advanceUntilIdle()
assertEquals(17392, manager.activePortOrNull.value)
assertTrue(manager.isSocksReady())
sub.cancel()
port.cancel()
}
/** ...but it must not be reported to the user, or to the exit-rotation path, as working Tor. */
@Test
fun `bootstrapping does not count as fully bootstrapped`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
val manager = buildManager(backend = backend)
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
backend.setBootstrapping(17392)
advanceUntilIdle()
assertFalse(manager.status.value.isFullyBootstrapped)
// onTorCircuitsDead is about dead exits behind a working Tor; a download in progress is
// not that, and resetting here would restart the very download we are waiting on.
manager.onTorCircuitsDead()
advanceUntilIdle()
assertEquals(0, backend.resetCount)
backend.setActive(17392)
advanceUntilIdle()
assertTrue(manager.status.value.isFullyBootstrapped)
sub.cancel()
}
// ------------------------------------------------------------------
// audit regressions
// ------------------------------------------------------------------
/**
* The worst bug the audit found, now guarded by progress rather than a timer.
*
* A cold directory download legitimately takes 12.6-34.4s (measured), and a reset discards the
* partial consensus — so a watchdog firing on the 45s no-client cadence could stop the download
* ever completing. A download that is still advancing must be left alone no matter how long it
* takes.
*/
@Test
fun `a download that keeps making progress is never reset`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
backend.bootstrappedBefore = true
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
assertTrue(manager.status.value is TorServiceStatus.Bootstrapping)
// Ten minutes of slow-but-real progress — far past any fixed patience window.
repeat(20) { step ->
advanceTimeBy(30_000L)
backend.advanceBootstrapProgress(step * 50)
runCurrent()
}
assertEquals("a progressing download must never be reset", 0, backend.resetCount)
assertEquals("and never wiped", 0, backend.resetWithCleanStateCount)
sub.cancel()
}
/** A download that stops moving is reset promptly — and still never has its cache wiped. */
@Test
fun `a download that stops progressing is reset but never wiped`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
backend.bootstrappedBefore = true
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
backend.advanceBootstrapProgress(120)
runCurrent()
// Nothing moves from here on.
advanceTimeBy(TorManager.BOOTSTRAP_STALL_MS + TorManager.SELF_HEAL_AFTER_MS)
runCurrent()
assertTrue("a stalled download must be escaped", backend.resetCount >= 1)
assertEquals(
"but the consensus cache is what it is trying to fetch — never wipe it",
0,
backend.resetWithCleanStateCount,
)
sub.cancel()
}
/**
* A fresh install with corrupt on-disk state can never produce a client, and `guards.json`
* (the only thing [ArtiGuardState] inspects) may look fine. Without an escalation the gentle
* branch drop-and-retries forever and never clears what is actually blocking it.
*/
@Test
fun `a fresh install that never gets a client eventually wipes state`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
// Never bootstrapped here, and start() cannot even bind a proxy.
backend.startFailsToConnect = true
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
advanceUntilIdle()
advanceTimeBy(TorManager.SELF_HEAL_AFTER_MS * 12)
runCurrent()
assertTrue(
"gentle resets alone never clear corrupt state; expected an escalation",
backend.resetWithCleanStateCount >= 1,
)
sub.cancel()
}
/**
* Each self-heal drives Bootstrapping -> Off -> Bootstrapping. If the prompt drops on the
* transient Off and re-raises immediately after, the user gets a modal blinking at them on
* every watchdog tick instead of a stable choice.
*/
@Test
fun `the bypass prompt stays up across a self-heal reset`() =
runTest(UnconfinedTestDispatcher()) {
val backend = FakeTorBackend()
val epoch = 1_700_000_000_000L
val manager = buildManager(backend = backend, clock = { epoch + testScheduler.currentTime })
val sub = launch { manager.status.collect { } }
val seen = mutableListOf<Boolean>()
val subFail = launch { manager.connectionFailure.collect { seen.add(it) } }
advanceUntilIdle()
advanceTimeBy(TorManager.BOOTSTRAP_TIMEOUT_MS + 1_000L)
runCurrent()
assertTrue("prompt should be up", manager.connectionFailure.value)
// Drive several watchdog cycles; the prompt must not drop back to false.
val raisedAt = seen.size
advanceTimeBy(TorManager.BOOTSTRAP_STALL_MS * 3)
runCurrent()
assertFalse(
"prompt blinked off during a self-heal cycle",
seen.drop(raisedAt).contains(false),
)
sub.cancel()
subFail.cancel()
}
// ------------------------------------------------------------------
// helpers
// ------------------------------------------------------------------
@@ -496,6 +897,28 @@ private class FakeTorBackend : TorBackend {
private val _status = MutableStateFlow<TorServiceStatus>(TorServiceStatus.Off)
override val status: StateFlow<TorServiceStatus> = _status.asStateFlow()
private val _bootstrapInFlight = MutableStateFlow(false)
override val bootstrapInFlight: StateFlow<Boolean> = _bootstrapInFlight.asStateFlow()
private val _bootstrapProgress = MutableStateFlow(-1)
override val bootstrapProgress: StateFlow<Int> = _bootstrapProgress.asStateFlow()
/** Models the directory download advancing. Only distinct values count as progress. */
fun advanceBootstrapProgress(permille: Int) {
_bootstrapProgress.value = permille
}
/**
* When true, [start] models a native bootstrap that is still running: status goes Connecting
* and [bootstrapInFlight] stays true until [finishBootstrapAttempt] is called.
*/
var holdBootstrapInFlight = false
/** Models the native bootstrap attempt returning without reaching Active (Arti's own timeout). */
fun finishBootstrapAttempt() {
_bootstrapInFlight.value = false
}
var startCount = 0
private set
var stopCount = 0
@@ -515,9 +938,30 @@ private class FakeTorBackend : TorBackend {
override suspend fun hasBootstrappedBefore(): Boolean = bootstrappedBefore
/**
* Set to model the real backend, whose `start()` does not return until the blocking native
* bootstrap attempt has finished. The status is published as soon as the attempt begins, the
* same as [TorService] does.
*/
var startNeverReturns = false
/** When true, start() models a client that cannot be created at all: status stays Connecting. */
var startFailsToConnect = false
/**
* Mirrors [TorService]: the native client is created in ~130ms and the proxy binds, so a real
* start lands in [TorServiceStatus.Bootstrapping] — routable, directory still downloading — not
* in Connecting. Tests that want the pre-proxy state call [setConnecting].
*/
override suspend fun start() {
startCount++
_status.value = TorServiceStatus.Connecting
if (startFailsToConnect) {
_status.value = TorServiceStatus.Connecting
return
}
_status.value = TorServiceStatus.Bootstrapping(17392)
if (holdBootstrapInFlight) _bootstrapInFlight.value = true
if (startNeverReturns) awaitCancellation()
}
override suspend fun stop() {
@@ -527,21 +971,28 @@ private class FakeTorBackend : TorBackend {
override suspend fun reset() {
resetCount++
_bootstrapInFlight.value = false
_status.value = TorServiceStatus.Off
}
override suspend fun resetWithCleanState() {
resetWithCleanStateCount++
_bootstrapInFlight.value = false
_status.value = TorServiceStatus.Off
}
fun setActive(port: Int) {
_bootstrapInFlight.value = false
_status.value = TorServiceStatus.Active(port)
}
fun setConnecting() {
_status.value = TorServiceStatus.Connecting
}
fun setBootstrapping(port: Int = 17392) {
_status.value = TorServiceStatus.Bootstrapping(port)
}
}
/** In-memory [TorPreferencesPort] driven by tests. */
@@ -0,0 +1,52 @@
/*
* 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.defaults
/**
* Substitute app defaults only when we have **never seen** the user's list — never when they
* published an empty one.
*
* There are three states, and collapsing the last two is how the app ends up overriding an explicit
* choice:
*
* | we have | effective list |
* |---|---|
* | no event | [defaults] — we do not know what they want |
* | an event, empty list | **empty** — they told us: nothing |
* | an event with relays | those relays |
*
* Written as one named primitive because the rule was open-coded at four call sites and three of
* them got it wrong the same way: `event?.relays()?.ifEmpty { null } ?: DEFAULTS` reads naturally
* but folds "published nothing" into "published nothing we know of", so a kind:10002 carrying only
* write relays silently acquired a default *inbox* list.
*
* [read] may return null — several event accessors end in `.ifEmpty { null }` — and null from a
* present event means the same thing as an empty set: the user listed nothing.
*
* **Not for partially-resolved sources.** A reader that can only see *already decrypted* private
* tags returns empty both for "the user listed nothing" and for "we have not decrypted it yet",
* which this cannot distinguish; those callers legitimately want defaults until the decrypt lands.
*/
inline fun <E : Any, T> relayListOrDefaultsWhenUnknown(
event: E?,
defaults: Set<T>,
read: (E) -> Set<T>?,
): Set<T> = if (event == null) defaults else read(event) ?: emptySet()
@@ -25,11 +25,32 @@ import com.vitorpamplona.quartz.nip01Core.relay.normalizer.isLocalHost
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.isOnion
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.isOverlayNetwork
/**
* Which relays fall into each Tor-routing category, as one value.
*
* Grouped deliberately rather than passed as four loose sets. Consumers that must react when the
* categories change — `RelayProxyClientConnector` re-dials the relays whose transport flipped —
* previously compared the sets field by field, so adding a fifth category meant remembering to add
* a fifth `||`. Forgetting it fails silently: relays keep a socket on a transport the policy has
* already moved them off. Structural equality on one object makes that impossible to forget.
*/
data class RelayClassification(
/** Relays the user actually put in one of their own lists. */
val trusted: Set<NormalizedRelayUrl> = emptySet(),
/** NIP-17 DM relays. */
val dm: Set<NormalizedRelayUrl> = emptySet(),
/** NIP-47 wallet and CLINK debit relays, including ad-hoc registrations. */
val moneyOp: Set<NormalizedRelayUrl> = emptySet(),
/**
* Relays the app is guessing on the user's behalf while their own lists are unknown. Empties
* itself as soon as any of their events arrive; see `AssumedRelayListsState`.
*/
val assumed: Set<NormalizedRelayUrl> = emptySet(),
)
class TorRelayEvaluation(
val torSettings: TorRelaySettings,
val trustedRelayList: Set<NormalizedRelayUrl>,
val dmRelayList: Set<NormalizedRelayUrl>,
val moneyOpRelayList: Set<NormalizedRelayUrl> = emptySet(),
val classification: RelayClassification = RelayClassification(),
) {
fun useTor(relay: NormalizedRelayUrl): Boolean =
if (torSettings.torType == TorType.OFF) {
@@ -45,15 +66,21 @@ class TorRelayEvaluation(
} else if (relay.isOnion()) {
// .onion is only reachable over Tor regardless of any other classification.
torSettings.onionRelaysViaTor
} else if (relay in moneyOpRelayList) {
} else if (relay in classification.moneyOp) {
// Relays used for money operations (NIP-47 wallets, CLINK offer/debit services)
// follow the dedicated money-operations preference, taking precedence over the
// generic DM/trusted/new classification so a payment never silently inherits a
// different Tor policy than the one the user set for money.
torSettings.moneyOperationsViaTor
} else if (relay in dmRelayList) {
} else if (relay in classification.dm) {
torSettings.dmRelaysViaTor
} else if (relay in trustedRelayList) {
} else if (relay in classification.trusted) {
torSettings.trustedRelaysViaTor
} else if (relay in classification.assumed) {
// Last resort before treating it as a stranger. Sits below every other
// classification on purpose: .onion, money-operation and DM relays keep their own
// policy even while we are guessing, because this branch can only ever capture
// relays that would otherwise have fallen through to `newRelaysViaTor`.
torSettings.trustedRelaysViaTor
} else {
torSettings.newRelaysViaTor
@@ -32,4 +32,18 @@ sealed class TorServiceStatus {
data class Error(
val message: String,
) : TorServiceStatus()
/**
* Where to send bytes, or null. Mirrors the Android status class so callers express intent
* rather than matching variants. No `Bootstrapping` here on purpose: the desktop backend drives
* an external Tor, so it never observes the routable-but-not-yet-bootstrapped window that the
* in-process Arti client has, and a variant nothing emits is dead weight (see [Error], which is
* only ever constructed by `DesktopTorManager`).
*/
val socksPort: Int?
get() = (this as? Active)?.port
/** Tor can build circuits right now. */
val isFullyBootstrapped: Boolean
get() = this is Active
}
@@ -0,0 +1,69 @@
/*
* 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.defaults
import kotlin.test.Test
import kotlin.test.assertEquals
class RelayListDefaultsTest {
private val defaults = setOf("wss://default.one", "wss://default.two")
/** No event: we genuinely do not know what the user wants, so the app's defaults stand in. */
@Test
fun `absent event yields the defaults`() {
assertEquals(defaults, relayListOrDefaultsWhenUnknown<String, String>(null, defaults) { setOf("wss://ignored") })
}
/**
* The case every open-coded copy of this rule got wrong: a published-but-empty list is the user
* saying "nothing", and must not acquire the defaults.
*/
@Test
fun `present event with an empty list stays empty`() {
assertEquals(emptySet(), relayListOrDefaultsWhenUnknown("event", defaults) { emptySet<String>() })
}
/**
* Same, via null: several event accessors end in `.ifEmpty { null }`, so null from a *present*
* event means the user listed nothing — not that the event is missing.
*/
@Test
fun `present event whose reader returns null stays empty`() {
assertEquals(emptySet(), relayListOrDefaultsWhenUnknown<String, String>("event", defaults) { null })
}
@Test
fun `present event with relays yields those relays`() {
val mine = setOf("wss://mine.example")
assertEquals(mine, relayListOrDefaultsWhenUnknown("event", defaults) { mine })
}
/** The reader must not even be consulted when there is no event to read. */
@Test
fun `reader is not invoked for an absent event`() {
var invoked = false
relayListOrDefaultsWhenUnknown<String, String>(null, defaults) {
invoked = true
emptySet()
}
assertEquals(false, invoked)
}
}
@@ -34,6 +34,7 @@ class TorRelayEvaluationTest {
private val dmRelay = NormalizedRelayUrl("wss://dm.relay.com/")
private val trustedRelay = NormalizedRelayUrl("wss://trusted.relay.com/")
private val moneyRelay = NormalizedRelayUrl("wss://wallet.relay.com/")
private val assumedRelay = NormalizedRelayUrl("wss://assumed.relay.com/")
private fun buildEvaluation(
torType: TorType = TorType.INTERNAL,
@@ -45,6 +46,7 @@ class TorRelayEvaluationTest {
dmRelays: Set<NormalizedRelayUrl> = setOf(dmRelay),
trustedRelays: Set<NormalizedRelayUrl> = setOf(trustedRelay),
moneyOpRelays: Set<NormalizedRelayUrl> = setOf(moneyRelay),
assumedRelays: Set<NormalizedRelayUrl> = setOf(assumedRelay),
) = TorRelayEvaluation(
torSettings =
TorRelaySettings(
@@ -55,11 +57,68 @@ class TorRelayEvaluationTest {
trustedRelaysViaTor = trustedViaTor,
moneyOperationsViaTor = moneyViaTor,
),
trustedRelayList = trustedRelays,
dmRelayList = dmRelays,
moneyOpRelayList = moneyOpRelays,
classification =
RelayClassification(
trusted = trustedRelays,
dm = dmRelays,
moneyOp = moneyOpRelays,
assumed = assumedRelays,
),
)
// --- assumed relays: the app's stand-in while the user's lists are unknown ---
/**
* The whole point: a guessed relay inherits the policy the user chose for their *own* lists, so
* the default configuration starts on clearnet and gets a fast first login.
*/
@Test
fun assumedRelay_followsTrustedPreference() {
assertFalse(buildEvaluation(trustedViaTor = false).useTor(assumedRelay))
assertTrue(buildEvaluation(trustedViaTor = true).useTor(assumedRelay))
}
/** Anyone who asked for Tor on their own relays keeps it here, with no separate opt-out. */
@Test
fun assumedRelay_hardenedUserStillUsesTor() {
assertTrue(buildEvaluation(trustedViaTor = true, newViaTor = true).useTor(assumedRelay))
}
/**
* The branch sits below every other classification, so being guessed can never downgrade a
* relay that already had a stricter policy.
*/
@Test
fun assumedRelay_neverOverridesOnionDmOrMoney() {
val eval =
buildEvaluation(
trustedViaTor = false,
onionViaTor = true,
dmViaTor = true,
moneyViaTor = true,
assumedRelays = setOf(assumedRelay, onionRelay, dmRelay, moneyRelay),
)
assertTrue(eval.useTor(onionRelay))
assertTrue(eval.useTor(dmRelay))
assertTrue(eval.useTor(moneyRelay))
}
/** A relay we are not guessing about is still a stranger. */
@Test
fun unknownRelay_stillFollowsNewPreference() {
assertTrue(buildEvaluation(newViaTor = true).useTor(clearnetRelay))
assertFalse(buildEvaluation(newViaTor = false).useTor(clearnetRelay))
}
/** With nothing guessed — every account that has any list — behaviour is exactly as before. */
@Test
fun emptyAssumedList_isTodaysBehaviour() {
val eval = buildEvaluation(assumedRelays = emptySet(), newViaTor = true, trustedViaTor = false)
assertTrue(eval.useTor(assumedRelay))
assertTrue(eval.useTor(clearnetRelay))
assertFalse(eval.useTor(trustedRelay))
}
// --- Tor OFF ---
@Test
fun torOff_alwaysFalse() {
@@ -48,8 +48,11 @@ class YggdrasilTorRoutingTest {
trustedRelaysViaTor = false,
moneyOperationsViaTor = false,
),
trustedRelayList = emptySet(),
dmRelayList = emptySet(),
classification =
RelayClassification(
trusted = emptySet(),
dm = emptySet(),
),
)
@Test
@@ -994,8 +994,10 @@ private fun AppInner(
newRelaysViaTor = torSettings.newRelaysViaTor,
trustedRelaysViaTor = torSettings.trustedRelaysViaTor,
),
trustedRelayList = emptySet(), // TODO: populate from account relay lists
dmRelayList = emptySet(), // TODO: populate from account relay lists
// TODO: populate from account relay lists
classification =
com.vitorpamplona.amethyst.commons.tor
.RelayClassification(),
)
}
@@ -64,7 +64,7 @@ class DesktopTorManager(
override val activePortOrNull: StateFlow<Int?> =
_status
.map { (it as? TorServiceStatus.Active)?.port }
.map { it.socksPort }
.stateIn(scope, SharingStarted.Eagerly, null)
private val runtime: TorRuntime by lazy {
+2
View File
@@ -224,6 +224,8 @@ verify_jni_symbols() {
"Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_initialize"
"Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_startSocksProxy"
"Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_stopSocksProxy"
"Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_isBootstrapped"
"Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_bootstrapProgressPermille"
"Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_destroy"
)
+114 -28
View File
@@ -3,7 +3,7 @@ use jni::objects::{JClass, JString, JObject, GlobalRef};
use jni::sys::{jint, jstring};
use jni::JavaVM;
use arti_client::TorClient;
use arti_client::{BootstrapBehavior, TorClient};
use arti_client::config::TorClientConfigBuilder;
// `kind()` is a trait method (tor_error::HasKind), not inherent, so the trait has to be in
// scope wherever we classify a connect failure. Both are re-exported by arti-client.
@@ -23,6 +23,10 @@ static TOKIO_RUNTIME: Mutex<Option<tokio::runtime::Runtime>> = Mutex::new(None);
static JAVA_VM: Mutex<Option<JavaVM>> = Mutex::new(None);
static LOG_CALLBACK: Mutex<Option<GlobalRef>> = Mutex::new(None);
static SOCKS_TASK: Mutex<Option<tokio::task::JoinHandle<()>>> = Mutex::new(None);
// The background directory download started by initialize(). It holds an Arc<TorClient>, so
// destroy() must abort it too — otherwise the client cannot drop, the state file lock is never
// released, and the next initialize() fails.
static BOOTSTRAP_TASK: Mutex<Option<tokio::task::JoinHandle<()>>> = Mutex::new(None);
// Per-connection handler tasks. Tracked so destroy() can abort in-flight handlers
// — otherwise their Arc<TorClient> clones keep the client alive and the
// state file lock would not be released for the next initialize().
@@ -170,14 +174,6 @@ pub extern "C" fn Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_initialize(
std::fs::create_dir_all(&cache_dir).ok();
std::fs::create_dir_all(&state_dir).ok();
// Bound the bootstrap so a hostile network (unreachable guards, wiped
// consensus) can't block this JNI call — and the Kotlin-side lifecycle lock
// it holds — indefinitely. On timeout the `create_bootstrapped` future is
// dropped, tearing down the partially built client, and -4 is returned so
// TorService can leave status Connecting and let the self-heal watchdog
// retry instead of wedging. The ABI is unchanged (still one String arg).
const BOOTSTRAP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
let outcome: jint = runtime.block_on(async {
log_info!("Creating Arti client...");
@@ -203,24 +199,62 @@ pub extern "C" fn Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_initialize(
}
};
match tokio::time::timeout(BOOTSTRAP_TIMEOUT, TorClient::create_bootstrapped(config)).await {
Ok(Ok(client)) => {
log_info!("Arti client created and bootstrapped");
*ARTI_CLIENT.lock().unwrap() = Some(Arc::new(client));
0
// Create the client WITHOUT waiting for the directory.
//
// `create_bootstrapped` used to block this JNI call — and the Kotlin lifecycle lock it
// holds — for the entire directory download: 12.6-33.8s measured on a cold install. The
// app treated Tor as absent for all of it, because `activePortOrNull` stays null until
// status flips to Active, so every relay dial fell back to 127.0.0.1:9050 (the Orbot
// default) where nothing listens, and failed instantly into backoff.
//
// With `BootstrapBehavior::OnDemand` the client is usable the moment it exists and each
// stream waits for the directory itself. The SOCKS proxy binds right away, so a dial
// issued mid-download queues on its own circuit instead of failing. Readiness stops being
// a global gate the app has to poll and becomes a property of individual connections.
//
// The `_async` variant is deliberate: it allows a short grace period for the state file
// lock, which a destroy()/initialize() cycle needs, where the sync one waits not at all.
let client = match TorClient::builder()
.config(config)
.bootstrap_behavior(BootstrapBehavior::OnDemand)
.create_unbootstrapped_async()
.await
{
Ok(c) => Arc::new(c),
Err(e) => {
log_error!("Failed to create Tor client: {:?}", e);
return -3;
}
Ok(Err(e)) => {
log_error!("Failed to bootstrap Tor client: {:?}", e);
-3
}
Err(_elapsed) => {
log_error!(
"Tor bootstrap timed out after {}s — aborting so the client can be retried",
BOOTSTRAP_TIMEOUT.as_secs()
);
-4
};
*ARTI_CLIENT.lock().unwrap() = Some(Arc::clone(&client));
// Start the download now rather than leaving it for the first stream to trigger, so it
// overlaps the login screen exactly as it used to, and publish the outcome so Kotlin can
// move the status from Bootstrapping to Active at the right moment.
let handle = tokio::spawn(async move {
let started = std::time::Instant::now();
match client.bootstrap().await {
Ok(()) => log_info!(
"Arti directory bootstrap complete after {}ms",
started.elapsed().as_millis()
),
// Not fatal and deliberately not latched anywhere: with OnDemand the next stream
// retries the bootstrap on its own, and isBootstrapped() reports live readiness, so
// a recovery after this point is picked up without us having to model it.
Err(e) => log_error!(
"Arti directory bootstrap failed after {}ms (streams will retry): {:?}",
started.elapsed().as_millis(),
e
),
}
});
if let Some(previous) = BOOTSTRAP_TASK.lock().unwrap().replace(handle) {
previous.abort();
}
log_info!("Arti client created (unbootstrapped; directory downloading in background)");
0
});
if outcome == 0 {
@@ -486,6 +520,51 @@ pub extern "C" fn Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_stopSocksPro
0
}
/// Directory-download progress in permille (0..1000), or -1 when there is no client.
///
/// Lets Kotlin tell a slow download from a stalled one. A timeout cannot: measured cold downloads
/// ran 12.6-34.4s on the same hardware and network, so any fixed patience is either short enough to
/// kill healthy ones or long enough to sit on a dead one. Forward progress separates them exactly.
///
/// Deliberately does NOT surface `BootstrapStatus::blocked()`. Arti documents it as best-effort and
/// warns it "may declare that Arti is stuck for reasons that are incorrect; or it may declare that
/// the client is not stuck when in fact no progress is being made" — acting on that would trade a
/// measurable signal for a guess.
#[no_mangle]
pub extern "C" fn Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_bootstrapProgressPermille(
_env: JNIEnv,
_class: JClass,
) -> jint {
match ARTI_CLIENT.lock().unwrap().as_ref() {
Some(client) => (client.bootstrap_status().as_frac() * 1000.0).clamp(0.0, 1000.0) as jint,
None => -1,
}
}
/// Live readiness: 1 = ready for traffic, 0 = not yet, -1 = no client at all.
///
/// Asks Arti itself (`bootstrap_status().ready_for_traffic()`, a cheap borrow-and-clone of a small
/// struct) rather than latching the outcome of the one background `bootstrap()` call. That call can
/// fail while the client stays perfectly usable — with OnDemand the next stream just retries — so a
/// latched failure would report "not bootstrapped" forever against a Tor that actually works,
/// leaving the UI wrong and the exit-rotation self-heal disabled.
#[no_mangle]
pub extern "C" fn Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_isBootstrapped(
_env: JNIEnv,
_class: JClass,
) -> jint {
match ARTI_CLIENT.lock().unwrap().as_ref() {
Some(client) => {
if client.bootstrap_status().ready_for_traffic() {
1
} else {
0
}
}
None => -1,
}
}
/// Destroy the TorClient — used by self-heal paths in Kotlin when Tor is
/// stuck and the in-memory state (guards, circuits) needs to be rebuilt
/// from scratch. Aborts the SOCKS listener and all in-flight per-connection
@@ -521,6 +600,13 @@ pub extern "C" fn Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_destroy(
});
}
// The background directory download holds an Arc<TorClient> too, for as long as it runs —
// which on a dead network is indefinitely. Abort it with the handlers or the state file lock
// outlives this destroy() and the next initialize() cannot take it.
if let Some(h) = BOOTSTRAP_TASK.lock().unwrap().take() {
h.abort();
}
// Abort all in-flight handlers — each holds an Arc<TorClient> clone, and
// the client cannot drop (state file lock cannot release) while any clone
// is alive.
@@ -538,10 +624,10 @@ pub extern "C" fn Java_com_vitorpamplona_amethyst_ui_tor_ArtiNative_destroy(
});
}
// Drop the static Arc. If any handler is still holding a clone, the
// TorClient stays alive until that handler finishes — in which case the
// next initialize() will fail and Kotlin's clearAllArtiData retry path
// will handle it.
// Drop the static Arc. If any handler is still holding a clone, the TorClient stays alive
// until that handler finishes — in which case the next initialize() fails and Kotlin leaves
// status Connecting for TorManager's self-heal watchdog to retry. (It no longer wipes all Arti
// data inline: that turned a transient failure into a lost guard sample and a lost consensus.)
let _ = ARTI_CLIENT.lock().unwrap().take();
log_info!("Arti client destroyed");