mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-05 19:28:25 +00:00
feat(marmot): add the lifecycle state machine and convergence branch selection
Stages 5 and 6, protocol cores. The lifecycle model: six canonical states with their legal-transition table, four derived convergence statuses with the legal-combination table, and the durable local gates (Leaving, Disbanding, realized removal) that restrict outbound work without being lifecycle states. Two table entries are load-bearing rather than bookkeeping, and both have tests. There is no Merging -> Recovering edge: a competing branch observed while applying our own confirmed commit is retained, the merge completes to Stable, and admission into a bounded pass then triggers Stable -> Recovering. Diverting mid-merge would leave a half-applied epoch. And Disbanded has no outgoing edge at all — no later branch supersedes a terminalized disband. Branch selection replaces the superseded MIP-03 rule, which broke a same-epoch tie on the outer Nostr created_at and then the event id. Both are transport evidence: timestamps are chosen by senders, and each transport copy of one MLS message carries a different event id. The replacement reads only authenticated values. Three details that decide whether two clients agree: Byte ordering is unsigned. Account keys and SHA-256 digests are uniformly distributed, so a signed comparison inverts roughly half of all final ties, and two implementations would disagree that often. raw_commit_depth gets no comparison step of its own — it is already inside effective_commit_depth, so once effective depth and quorum status tie, a further raw-depth comparison is necessarily tied too. A widely circulated write-up of this algorithm lists raw depth as a step; the spec does not, and there is a test that fails if it is added. Witnesses count distinct sender ACCOUNTS per branch epoch, capped at the quorum size, and epochs at or before fork_epoch do not count. Counting by account stops a multi-device member counting twice; counting distinct senders stops one member inflating a branch by sending a lot; the per-epoch cap stops one busy epoch outweighing several quiet ones. The policy constructor enforces max_witness_override_depth <= max_rewind_commits, because without that bound app-payload traffic could push a branch past the rollback horizon and beat an arbitrarily longer valid commit branch. Twenty-five tests, including the worked example: a three-commit branch with witness quorum ties a four-commit branch without one at effective depth four and then wins on quorum, while a five-commit branch beats both because the boost is capped at one. Selection is asserted invariant under input order, reversal, shuffling and every rotation. Full quartz jvmTest: 4,582 tests, 0 failures. Still open in these stages: the bounded pass scheduler and the candidate-graph builder that replays MLS bytes against retained states, plus wiring the lifecycle states into MlsGroup so they gate anything. CommitOrdering's transport-metadata tiebreak therefore still stands — deleting it is only safe once something replaces it end to end, and selection alone does not. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_016kCuA6tc4JQzHPCDd39GHq
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
# Marmot: resync against the adopted spec and current MDK
|
||||
|
||||
Status: Stages 0-4 done, plus the Stage 3 image-crypto follow-up. Stages 5-7 open.
|
||||
Status: Stages 0-4 done. Stages 5-6 have their protocol cores landed; the pass scheduler,
|
||||
candidate-graph replay and Stage 7 remain.
|
||||
|
||||
Sources checked on 2026-09-08:
|
||||
|
||||
@@ -404,13 +405,43 @@ canonical epoch, retained epochs inside the rollback horizon, and any staged-but
|
||||
commit — and no more") is defined in terms of the retained-state set that convergence owns, so
|
||||
it lands with Stage 6 rather than ahead of it.
|
||||
|
||||
**Stage 5 — lifecycle state machine + publish-before-apply.**
|
||||
The six canonical states, the `Leaving` / `Disbanding` gates, and the publish-obligation
|
||||
record (bytes + recipient scope + prior state + pending state) surviving restart.
|
||||
**Stage 5 — lifecycle state machine. CORE DONE.**
|
||||
|
||||
**Stage 6 — convergence engine.**
|
||||
Bounded passes, candidate graph, eligibility, witnesses, six-step selection, dispositions and
|
||||
withdrawal. Delete `CommitOrdering`'s transport-metadata tiebreak at this point, not before.
|
||||
`GroupLifecycleState` (the six canonical states with their legal-transition table),
|
||||
`ConvergenceStatus` (the four derived statuses with the legal-combination table), and
|
||||
`LocalOutboundGate` for `Leaving` / `Disbanding` / realized-removal.
|
||||
|
||||
Two table entries are load-bearing and tested as such: there is NO `Merging -> Recovering`
|
||||
edge — a competing branch seen mid-merge is retained, the merge finishes to `Stable`, and the
|
||||
bounded pass then triggers `Stable -> Recovering`, because diverting mid-merge would leave a
|
||||
half-applied epoch — and `Disbanded` has no outgoing edge at all.
|
||||
|
||||
**Still open:** the publish-obligation record (bytes + recipient scope + prior state + pending
|
||||
state) surviving restart, and wiring these states into `MlsGroup`/`MarmotManager` so they
|
||||
actually gate anything. Today they are a correct model with no enforcement behind them.
|
||||
|
||||
**Stage 6 — convergence engine. SELECTION DONE.**
|
||||
|
||||
`ConvergencePolicy` (the v1 constants, with the `max_witness_override_depth <=
|
||||
max_rewind_commits` bound enforced in the constructor), `CandidateBranch` (fork/tip epochs,
|
||||
raw depth, tip priority/committer/digest, per-epoch witnesses), `BranchSelector` (the six-step
|
||||
comparison), and the `ConvergenceDisposition` / `ConvergenceCategory` vocabularies.
|
||||
|
||||
Things worth knowing about the implementation:
|
||||
|
||||
- Byte ordering is UNSIGNED. Account keys and digests are uniformly distributed, so a signed
|
||||
comparison would invert about half of all final ties — and two clients would then disagree
|
||||
that often.
|
||||
- `raw_commit_depth` has no comparison step of its own. It is already inside
|
||||
`effective_commit_depth`. The widely circulated write-up of this algorithm lists it as a
|
||||
step; the spec explicitly does not, and there is a test that fails if it is added.
|
||||
- Witnesses count DISTINCT sender accounts per branch epoch, capped at the quorum size, and
|
||||
epochs at or before `fork_epoch` do not count at all.
|
||||
|
||||
**Still open:** the bounded pass scheduler (quiescence/deadline timers, the frozen batch,
|
||||
`pass_base_epoch`) and the candidate-graph builder that replays MLS bytes against retained
|
||||
states. `CommitOrdering`'s transport-metadata tiebreak therefore still stands — it is only safe
|
||||
to delete once something replaces it end to end, and selection alone does not.
|
||||
|
||||
**Stage 7 — durability/restart conformance, app payload kinds (1009/1210), encrypted-media
|
||||
v2, push owner proof.**
|
||||
|
||||
@@ -10,7 +10,7 @@ _Audited 2026-09-08. 12 plans: 7 shipped (archived), 0 in-progress, 4 queued, 1
|
||||
| [2026-07-03-incremental-negentropy-storage.md](2026-07-03-incremental-negentropy-storage.md) | Always-current (created_at, id) index so cold NEG-OPENs stop paying a full scan + seal (~340 ms at 50k vs strfry's ~21 ms). |
|
||||
| [2026-07-04-small-req-floor.md](2026-07-04-small-req-floor.md) | Small-REQ dispatch floor: decomposed, inline fast path tried and reverted (no wire-level win); floor is transport-side. |
|
||||
| [2026-08-13-gpu-pow-mining.md](2026-08-13-gpu-pow-mining.md) | GPU NIP-13 mining declined (ARMv8 has SHA-256 in silicon, mobile GPUs do not). Midstate is ~3x on JVM targets; Android hinges on Conscrypt per-digest JNI cost, still unmeasured. created_at refresh while mining shipped. |
|
||||
| [2026-09-08-marmot-spec-resync.md](2026-09-08-marmot-spec-resync.md) | Marmot moved off the MIP-era spec (2026-07-02): group state split into `app_data_dictionary` components, account identity proof v2, and a convergence engine. Current MDK rejects our groups outright. Gap analysis + 8-stage plan; Stages 0-3 done (mdk interop reference, app_data_dictionary + AppDataUpdate, account-identity-proof v2, the six group components + current-profile authorization). |
|
||||
| [2026-09-08-marmot-spec-resync.md](2026-09-08-marmot-spec-resync.md) | Marmot moved off the MIP-era spec (2026-07-02): group state split into `app_data_dictionary` components, account identity proof v2, and a convergence engine. Current MDK rejects our groups outright. Gap analysis + 8-stage plan; Stages 0-4 done (mdk interop reference, app_data_dictionary, identity proof v2, the six group components, transport corrections); lifecycle + branch selection landed. |
|
||||
|
||||
## Archived (shipped)
|
||||
| Plan | Summary |
|
||||
|
||||
+126
@@ -0,0 +1,126 @@
|
||||
/*
|
||||
* 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.quartz.marmot.protocolCore
|
||||
|
||||
/**
|
||||
* Branch selection (`protocol-core/convergence.md`, "Branch selection").
|
||||
*
|
||||
* Two clients that start from the same retained anchor, use the same policy,
|
||||
* and resolve the same frozen input batch MUST select the same branch —
|
||||
* regardless of transport arrival order, local scheduling, or device speed.
|
||||
* This object is where that determinism lives, so everything it reads is
|
||||
* authenticated and nothing it reads is transport metadata.
|
||||
*
|
||||
* This REPLACES the superseded MIP-03 rule, which broke a same-epoch tie on the
|
||||
* outer Nostr `created_at` and then the Nostr event id. Both are transport
|
||||
* evidence: timestamps are chosen by senders, and each transport copy of one
|
||||
* MLS message carries a different event id.
|
||||
*/
|
||||
object BranchSelector {
|
||||
/**
|
||||
* Compare two eligible branches, most-preferred first.
|
||||
*
|
||||
* The order is exactly:
|
||||
*
|
||||
* 1. higher `effective_commit_depth`
|
||||
* 2. witness quorum beats no quorum
|
||||
* 3. higher `app_witness_score`
|
||||
* 4. lower `tip_priority` (privileged before ordinary)
|
||||
* 5. lower `tip_committer`
|
||||
* 6. lower `tip_digest`
|
||||
*
|
||||
* `raw_commit_depth` has NO step of its own: it is already inside
|
||||
* `effective_commit_depth`, so once effective depth and quorum status are
|
||||
* both tied a further raw-depth comparison is necessarily tied too. (A
|
||||
* widely circulated write-up of this algorithm lists raw depth as step 2 —
|
||||
* it is not in the spec, and adding it would change nothing except to make
|
||||
* two implementations disagree about where the comparison ended.)
|
||||
*/
|
||||
fun comparator(policy: ConvergencePolicy): Comparator<CandidateBranch> =
|
||||
Comparator { a, b ->
|
||||
var result = b.effectiveCommitDepth(policy).compareTo(a.effectiveCommitDepth(policy))
|
||||
if (result != 0) return@Comparator result
|
||||
|
||||
result = quorumRank(b, policy).compareTo(quorumRank(a, policy))
|
||||
if (result != 0) return@Comparator result
|
||||
|
||||
result = b.appWitnessScore(policy).compareTo(a.appWitnessScore(policy))
|
||||
if (result != 0) return@Comparator result
|
||||
|
||||
result = a.tipPriority.order.compareTo(b.tipPriority.order)
|
||||
if (result != 0) return@Comparator result
|
||||
|
||||
result = compareUnsigned(a.tipCommitter, b.tipCommitter)
|
||||
if (result != 0) return@Comparator result
|
||||
|
||||
compareUnsigned(a.tipDigest, b.tipDigest)
|
||||
}
|
||||
|
||||
/**
|
||||
* The canonical branch among [branches], or null when none is eligible.
|
||||
*
|
||||
* [passBaseEpoch] is the pass's frozen base; branches whose fork lies
|
||||
* outside the rollback horizon MUST NOT be selected.
|
||||
*/
|
||||
fun select(
|
||||
branches: Collection<CandidateBranch>,
|
||||
passBaseEpoch: Long,
|
||||
policy: ConvergencePolicy = ConvergencePolicy.V1,
|
||||
): CandidateBranch? =
|
||||
branches
|
||||
.filter { it.isEligible(passBaseEpoch, policy) }
|
||||
.minWithOrNull(comparator(policy))
|
||||
|
||||
/** [branches] ordered most-preferred first, eligible ones only. */
|
||||
fun rank(
|
||||
branches: Collection<CandidateBranch>,
|
||||
passBaseEpoch: Long,
|
||||
policy: ConvergencePolicy = ConvergencePolicy.V1,
|
||||
): List<CandidateBranch> =
|
||||
branches
|
||||
.filter { it.isEligible(passBaseEpoch, policy) }
|
||||
.sortedWith(comparator(policy))
|
||||
|
||||
private fun quorumRank(
|
||||
branch: CandidateBranch,
|
||||
policy: ConvergencePolicy,
|
||||
): Int = if (branch.witnessQuorumMet(policy)) 1 else 0
|
||||
|
||||
/**
|
||||
* Lexicographic order over raw bytes, compared UNSIGNED.
|
||||
*
|
||||
* Both a 32-byte x-only account key and a SHA-256 digest are uniformly
|
||||
* distributed, so about half of all comparisons involve a byte above 0x7f.
|
||||
* A signed comparison would invert those and two implementations would
|
||||
* disagree about the winner roughly half the time a tie reached this far.
|
||||
*/
|
||||
private fun compareUnsigned(
|
||||
a: ByteArray,
|
||||
b: ByteArray,
|
||||
): Int {
|
||||
val common = minOf(a.size, b.size)
|
||||
for (i in 0 until common) {
|
||||
val diff = (a[i].toInt() and 0xFF) - (b[i].toInt() and 0xFF)
|
||||
if (diff != 0) return diff
|
||||
}
|
||||
return a.size - b.size
|
||||
}
|
||||
}
|
||||
+154
@@ -0,0 +1,154 @@
|
||||
/*
|
||||
* 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.quartz.marmot.protocolCore
|
||||
|
||||
import com.vitorpamplona.quartz.nip01Core.core.HexKey
|
||||
import com.vitorpamplona.quartz.nip01Core.core.toHexKey
|
||||
|
||||
/**
|
||||
* The authorization class of a branch's tip Commit
|
||||
* (`protocol-core/convergence.md`, "Candidate branches").
|
||||
*
|
||||
* A Commit is [PRIVILEGED] exactly when its applicable Marmot authorization
|
||||
* rule REQUIRES an active admin in the candidate parent state, and [ORDINARY]
|
||||
* when the rule permits a non-admin committer. This follows the authorization
|
||||
* rule, not the operation's apparent importance: a component change whose
|
||||
* owning document explicitly permits a non-admin committer is `ordinary`.
|
||||
*/
|
||||
enum class TipPriority {
|
||||
PRIVILEGED,
|
||||
ORDINARY,
|
||||
;
|
||||
|
||||
/** Lower sorts first, and privileged wins a tie. */
|
||||
val order: Int get() = if (this == PRIVILEGED) 0 else 1
|
||||
}
|
||||
|
||||
/**
|
||||
* One candidate branch in a convergence pass.
|
||||
*
|
||||
* Every value here MUST come from MLS-valid bytes, retained state, decrypted
|
||||
* app payloads, or the pinned policy. Transport arrival order, transport
|
||||
* timestamps, outer event ids and local receive order MUST NOT appear — which
|
||||
* is exactly what the superseded MIP-03 rule got wrong by ranking on
|
||||
* `created_at` and the Nostr event id.
|
||||
*/
|
||||
data class CandidateBranch(
|
||||
/** Epoch where this branch diverged from retained canonical state. */
|
||||
val forkEpoch: Long,
|
||||
/** Epoch reached after replaying this branch's valid commits. */
|
||||
val tipEpoch: Long,
|
||||
/** Number of valid commits from [forkEpoch] to [tipEpoch]. */
|
||||
val rawCommitDepth: Long,
|
||||
val tipPriority: TipPriority,
|
||||
/** Authenticated Marmot account identity of the tip committer (32 raw bytes). */
|
||||
val tipCommitter: ByteArray,
|
||||
/** `SHA-256` of the tip Commit's serialized MLS message bytes (32 bytes). */
|
||||
val tipDigest: ByteArray,
|
||||
/**
|
||||
* Distinct valid app-payload sender identities per branch epoch.
|
||||
*
|
||||
* Keyed by epoch; the value is the set of ACCOUNT identities (hex) that
|
||||
* sent a fully validated app payload decrypting on this branch at that
|
||||
* epoch. Counting by account rather than by leaf is what stops a
|
||||
* multi-device member from counting several times, and counting distinct
|
||||
* senders is what stops one member inflating a branch by sending a lot.
|
||||
*/
|
||||
val witnessesByEpoch: Map<Long, Set<HexKey>>,
|
||||
) {
|
||||
init {
|
||||
require(tipCommitter.size == 32) { "tip_committer must be a 32-byte account identity" }
|
||||
require(tipDigest.size == 32) { "tip_digest must be a 32-byte SHA-256" }
|
||||
require(rawCommitDepth >= 0) { "raw_commit_depth must not be negative" }
|
||||
}
|
||||
|
||||
/** Epochs strictly greater than [forkEpoch] and at most [tipEpoch]. */
|
||||
val branchEpochs: LongRange get() = (forkEpoch + 1)..tipEpoch
|
||||
|
||||
/**
|
||||
* `sum over branch epochs of min(distinct senders, quorum size)`.
|
||||
*
|
||||
* The per-epoch `min` is why one very chatty epoch cannot outweigh several
|
||||
* quiet ones.
|
||||
*/
|
||||
fun appWitnessScore(policy: ConvergencePolicy): Long =
|
||||
branchEpochs.sumOf { epoch ->
|
||||
val senders = witnessesByEpoch[epoch]?.size ?: 0
|
||||
minOf(senders, policy.witnessQuorumSendersPerEpoch).toLong()
|
||||
}
|
||||
|
||||
/**
|
||||
* True when at least [ConvergencePolicy.witnessQuorumEpochs] branch epochs
|
||||
* each had at least [ConvergencePolicy.witnessQuorumSendersPerEpoch]
|
||||
* distinct senders.
|
||||
*
|
||||
* Each epoch is evaluated independently: the qualifying sender set MAY
|
||||
* differ from epoch to epoch, and no cohort has to span them all.
|
||||
*/
|
||||
fun witnessQuorumMet(policy: ConvergencePolicy): Boolean =
|
||||
branchEpochs.count { epoch ->
|
||||
(witnessesByEpoch[epoch]?.size ?: 0) >= policy.witnessQuorumSendersPerEpoch
|
||||
} >= policy.witnessQuorumEpochs
|
||||
|
||||
/** `raw_commit_depth` plus the bounded witness boost. */
|
||||
fun effectiveCommitDepth(policy: ConvergencePolicy): Long = rawCommitDepth + if (witnessQuorumMet(policy)) policy.maxWitnessOverrideDepth else 0L
|
||||
|
||||
/**
|
||||
* Eligible only when the fork lies inside the rollback horizon measured
|
||||
* from the pass's FROZEN base epoch.
|
||||
*
|
||||
* The base is frozen for the whole pass so candidate replay cannot move
|
||||
* the horizon while candidates are being compared. Deferred-commit expiry
|
||||
* uses the LIVE canonical tip instead, so obsolete input still ages out as
|
||||
* canonical state advances across passes — the two use different reference
|
||||
* points on purpose.
|
||||
*/
|
||||
fun isEligible(
|
||||
passBaseEpoch: Long,
|
||||
policy: ConvergencePolicy,
|
||||
): Boolean = passBaseEpoch - forkEpoch <= policy.maxRewindCommits
|
||||
|
||||
val tipCommitterHex: HexKey get() = tipCommitter.toHexKey()
|
||||
val tipDigestHex: HexKey get() = tipDigest.toHexKey()
|
||||
|
||||
override fun equals(other: Any?): Boolean {
|
||||
if (this === other) return true
|
||||
if (other !is CandidateBranch) return false
|
||||
return forkEpoch == other.forkEpoch &&
|
||||
tipEpoch == other.tipEpoch &&
|
||||
rawCommitDepth == other.rawCommitDepth &&
|
||||
tipPriority == other.tipPriority &&
|
||||
tipCommitter.contentEquals(other.tipCommitter) &&
|
||||
tipDigest.contentEquals(other.tipDigest) &&
|
||||
witnessesByEpoch == other.witnessesByEpoch
|
||||
}
|
||||
|
||||
override fun hashCode(): Int {
|
||||
var result = forkEpoch.hashCode()
|
||||
result = 31 * result + tipEpoch.hashCode()
|
||||
result = 31 * result + rawCommitDepth.hashCode()
|
||||
result = 31 * result + tipPriority.hashCode()
|
||||
result = 31 * result + tipCommitter.contentHashCode()
|
||||
result = 31 * result + tipDigest.contentHashCode()
|
||||
result = 31 * result + witnessesByEpoch.hashCode()
|
||||
return result
|
||||
}
|
||||
}
|
||||
+66
@@ -0,0 +1,66 @@
|
||||
/*
|
||||
* 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.quartz.marmot.protocolCore
|
||||
|
||||
/**
|
||||
* What convergence decided about one retained input
|
||||
* (`foundation/errors.md`, "Convergence dispositions").
|
||||
*
|
||||
* A disposition says WHAT happened; [ConvergenceCategory] says why.
|
||||
*/
|
||||
enum class ConvergenceDisposition {
|
||||
/** On, or consumed by, the selected canonical branch. */
|
||||
ACCEPTED,
|
||||
|
||||
/**
|
||||
* No terminal classification yet; reconsidered when more input arrives.
|
||||
*
|
||||
* This is where valid commits on a NON-selected but still-eligible branch
|
||||
* sit. Losing one pass is not permanent ineligibility — the branch stays
|
||||
* eligible while its fork is inside the rollback horizon measured from the
|
||||
* live tip.
|
||||
*/
|
||||
DEFERRED,
|
||||
|
||||
/** Can no longer affect the group. */
|
||||
STALE,
|
||||
|
||||
/**
|
||||
* An MLS application message that decrypts only on a losing branch. Its
|
||||
* app payload is WITHDRAWN from application output — the counterpart to
|
||||
* withdrawing state notifications from a superseded commit.
|
||||
*/
|
||||
INVALIDATED,
|
||||
}
|
||||
|
||||
/**
|
||||
* Why an input received its disposition. A `stale` or `deferred` input SHOULD
|
||||
* carry one.
|
||||
*/
|
||||
enum class ConvergenceCategory {
|
||||
DUPLICATE,
|
||||
UNKNOWN_GROUP,
|
||||
STALE_EPOCH,
|
||||
AUTHORIZATION_FAILED,
|
||||
MISSING_HISTORY,
|
||||
TRANSPORT_DEFERRED,
|
||||
RESOURCE_REFUSED,
|
||||
}
|
||||
+77
@@ -0,0 +1,77 @@
|
||||
/*
|
||||
* 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.quartz.marmot.protocolCore
|
||||
|
||||
/**
|
||||
* The Marmot convergence policy (`protocol-core/convergence.md`).
|
||||
*
|
||||
* Version 1 is a set of PROTOCOL CONSTANTS, not a preference. It is not carried
|
||||
* in group state and cannot be negotiated: every client uses exactly these
|
||||
* values, and every branch scored within a pass uses the same ones.
|
||||
*
|
||||
* That rigidity is the point. Convergence is deliberately not group-tunable
|
||||
* because a bad policy choice forks a group — two clients scoring the same
|
||||
* candidates under different constants can each be internally consistent and
|
||||
* still disagree about which branch is canonical. A future change ships as a
|
||||
* NEW app component behind a required capability; clients MUST NOT infer the
|
||||
* active policy from a software version, and until such a component exists
|
||||
* there is no mechanism to change these at all.
|
||||
*/
|
||||
data class ConvergencePolicy(
|
||||
/** How far back from the tip a branch MAY fork and still be eligible. */
|
||||
val maxRewindCommits: Long,
|
||||
/** How many past epochs may still produce delivered payloads or witnesses. */
|
||||
val appPayloadPastEpochLimit: Long,
|
||||
/** Minimum quiet time before a pass MAY be treated as settled. */
|
||||
val settlementQuiescenceMs: Long,
|
||||
/** Maximum duration of one input-collection window; never extended by later input. */
|
||||
val maxConvergencePassMs: Long,
|
||||
/** Distinct senders needed for one branch epoch to count toward quorum. */
|
||||
val witnessQuorumSendersPerEpoch: Int,
|
||||
/** How many branch epochs must meet sender quorum. */
|
||||
val witnessQuorumEpochs: Int,
|
||||
/** Maximum commit-depth boost a branch may receive from witness quorum. */
|
||||
val maxWitnessOverrideDepth: Long,
|
||||
) {
|
||||
init {
|
||||
// The bound exists so app-payload traffic can never push a branch past
|
||||
// the rollback horizon: without it, message volume could beat an
|
||||
// arbitrarily longer valid commit branch.
|
||||
require(maxWitnessOverrideDepth <= maxRewindCommits) {
|
||||
"max_witness_override_depth ($maxWitnessOverrideDepth) must not exceed " +
|
||||
"max_rewind_commits ($maxRewindCommits)"
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
/** Marmot convergence policy, version 1. */
|
||||
val V1 =
|
||||
ConvergencePolicy(
|
||||
maxRewindCommits = 5,
|
||||
appPayloadPastEpochLimit = 5,
|
||||
settlementQuiescenceMs = 1_000,
|
||||
maxConvergencePassMs = 5_000,
|
||||
witnessQuorumSendersPerEpoch = 2,
|
||||
witnessQuorumEpochs = 1,
|
||||
maxWitnessOverrideDepth = 1,
|
||||
)
|
||||
}
|
||||
}
|
||||
+213
@@ -0,0 +1,213 @@
|
||||
/*
|
||||
* 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.quartz.marmot.protocolCore
|
||||
|
||||
/**
|
||||
* A Marmot group's canonical lifecycle state (`protocol-core/group-state.md`).
|
||||
*
|
||||
* Each group has exactly one canonical MLS state at a time. A client may hold
|
||||
* candidate or pending state alongside it, but only one state is visible as
|
||||
* canonical, and this says which phase the group is in.
|
||||
*/
|
||||
enum class GroupLifecycleState {
|
||||
/** Has a canonical epoch. The ONLY state where a new local commit may be prepared. */
|
||||
STABLE,
|
||||
|
||||
/** A local commit is prepared but its publish obligation is unconfirmed. */
|
||||
PENDING_PUBLISH,
|
||||
|
||||
/** Publication confirmed; the staged commit is being applied. */
|
||||
MERGING,
|
||||
|
||||
/**
|
||||
* Selecting a branch after a fork-shaped conflict, or after admitting a
|
||||
* valid disband candidate — which forces this state even with no fork, so
|
||||
* terminalization can only happen after selection.
|
||||
*/
|
||||
RECOVERING,
|
||||
|
||||
/**
|
||||
* Required retained material is permanently missing or corrupt and no
|
||||
* verified repair path exists.
|
||||
*
|
||||
* Local to one client: it does NOT mean the group is dead. It means this
|
||||
* client must repair, restore, rejoin, or discard its copy before it can
|
||||
* safely apply more group traffic. A client here MUST NOT settle for its
|
||||
* current local state just because it is the only one available.
|
||||
*/
|
||||
UNRECOVERABLE,
|
||||
|
||||
/** An authenticated disband Commit was selected. Absorbing; no rejoin. */
|
||||
DISBANDED,
|
||||
;
|
||||
|
||||
val isTerminal: Boolean get() = this == DISBANDED
|
||||
|
||||
/** Only `Stable` may prepare a new local group-state commit. */
|
||||
val canPrepareLocalCommit: Boolean get() = this == STABLE
|
||||
|
||||
/**
|
||||
* Whether retained inbound may change canonical group state here.
|
||||
*
|
||||
* False during `PendingPublish` and `Merging` (a local transition is in
|
||||
* flight), during `Unrecoverable` (nothing is safe to apply), and in
|
||||
* `Disbanded` (inbound is not even retained). `Recovering` is false too:
|
||||
* canonical state changes only when a SELECTED branch is applied, and
|
||||
* replaying candidates before then is not application.
|
||||
*/
|
||||
val mayApplyInboundToCanonicalState: Boolean get() = this == STABLE
|
||||
|
||||
fun canTransitionTo(next: GroupLifecycleState): Boolean = next in LEGAL_TRANSITIONS.getValue(this)
|
||||
|
||||
companion object {
|
||||
/**
|
||||
* The legal transition table.
|
||||
*
|
||||
* Two absences are deliberate rather than oversights:
|
||||
*
|
||||
* - There is no `Merging -> Recovering` edge. A competing branch seen
|
||||
* while applying our own confirmed commit is retained, the merge
|
||||
* completes to `Stable`, and admission into a bounded pass then
|
||||
* triggers `Stable -> Recovering`. Diverting mid-merge would leave a
|
||||
* half-applied epoch.
|
||||
* - `Disbanded` has no outgoing edge at all. It is absorbing: no later
|
||||
* branch supersedes a terminalized disband, and a replacement
|
||||
* conversation is a new MLS group.
|
||||
*
|
||||
* `Recovering` re-entry is implicit rather than a self-edge: input
|
||||
* arriving during recovery, before the pass cutoff, folds into the
|
||||
* pass already running.
|
||||
*/
|
||||
private val LEGAL_TRANSITIONS: Map<GroupLifecycleState, Set<GroupLifecycleState>> =
|
||||
mapOf(
|
||||
STABLE to setOf(PENDING_PUBLISH, RECOVERING, UNRECOVERABLE),
|
||||
PENDING_PUBLISH to setOf(MERGING, STABLE, UNRECOVERABLE),
|
||||
MERGING to setOf(STABLE, UNRECOVERABLE),
|
||||
RECOVERING to setOf(STABLE, DISBANDED, UNRECOVERABLE),
|
||||
UNRECOVERABLE to setOf(STABLE),
|
||||
DISBANDED to emptySet(),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Convergence's derived status (`protocol-core/group-state.md`, "Convergence
|
||||
* status").
|
||||
*
|
||||
* Derived from stored input and policy — never a claim made by the transport.
|
||||
* The lifecycle state is authoritative; this is a view of how convergence is
|
||||
* progressing within it.
|
||||
*/
|
||||
enum class ConvergenceStatus {
|
||||
/** A bounded pass is collecting selection-relevant input. */
|
||||
SYNCING,
|
||||
|
||||
/**
|
||||
* The batch is frozen and a deterministic fixed point is being computed
|
||||
* over already-retained state. Does NOT wait for fetches or admit later
|
||||
* input.
|
||||
*/
|
||||
RESOLVING,
|
||||
|
||||
/**
|
||||
* A fixed point was reached and any selected branch applied.
|
||||
*
|
||||
* Local to the input this client retained and admitted. It is NOT global
|
||||
* finality, not proof that transports finished synchronizing, and not a
|
||||
* promise that later valid input cannot open another pass.
|
||||
*/
|
||||
SETTLED,
|
||||
|
||||
/** Cannot continue without a repair path, missing material, or resources. */
|
||||
BLOCKED,
|
||||
;
|
||||
|
||||
/**
|
||||
* Whether this status permits preparing a group-state change or encrypting
|
||||
* an app payload.
|
||||
*
|
||||
* Only `Settled`. Outbound work is held while convergence is unresolved
|
||||
* because a payload MUST be encrypted against the SELECTED canonical state
|
||||
* — encrypting against a state that later loses branch selection produces
|
||||
* a message the group will invalidate.
|
||||
*/
|
||||
val allowsOutboundWork: Boolean get() = this == SETTLED
|
||||
|
||||
fun isLegalIn(lifecycle: GroupLifecycleState): Boolean = lifecycle in LEGAL_COMBINATIONS.getValue(this)
|
||||
|
||||
companion object {
|
||||
/**
|
||||
* Which lifecycle states each status may appear in.
|
||||
*
|
||||
* `PendingPublish` and `Merging` appear nowhere: they are local-publish
|
||||
* states, not convergence passes, so convergence status is not
|
||||
* meaningful in them.
|
||||
*/
|
||||
private val LEGAL_COMBINATIONS: Map<ConvergenceStatus, Set<GroupLifecycleState>> =
|
||||
mapOf(
|
||||
SYNCING to setOf(GroupLifecycleState.STABLE, GroupLifecycleState.RECOVERING),
|
||||
RESOLVING to setOf(GroupLifecycleState.STABLE, GroupLifecycleState.RECOVERING),
|
||||
SETTLED to setOf(GroupLifecycleState.STABLE, GroupLifecycleState.DISBANDED),
|
||||
BLOCKED to
|
||||
setOf(
|
||||
GroupLifecycleState.STABLE,
|
||||
GroupLifecycleState.RECOVERING,
|
||||
GroupLifecycleState.UNRECOVERABLE,
|
||||
),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Durable local gates that restrict outbound work without being canonical
|
||||
* lifecycle states.
|
||||
*
|
||||
* Each survives restart and each is one-way until the protocol event that
|
||||
* clears it. They are separate from [GroupLifecycleState] because the MLS group
|
||||
* state does not change when they are set — the member is still in the tree.
|
||||
*/
|
||||
enum class LocalOutboundGate {
|
||||
/**
|
||||
* `Leaving`: a SelfRemove proposal was sent. The member is still in the MLS
|
||||
* group until a commit removes it, and the gate may span several
|
||||
* epoch-bound SelfRemove proposals.
|
||||
*/
|
||||
LEAVING,
|
||||
|
||||
/**
|
||||
* `Disbanding`: an admin's irreversible disband request is unresolved. It
|
||||
* blocks all new outbound work while the request is prepared, published,
|
||||
* retried, or evaluated by convergence, and survives publication failure
|
||||
* and restart.
|
||||
*/
|
||||
DISBANDING,
|
||||
|
||||
/**
|
||||
* The local member's own removal has been realized. The group is held as a
|
||||
* removed, inactive copy: history may be kept, but the group must not be
|
||||
* presented as active and nothing may be sent to it. Cleared only by an
|
||||
* authenticated re-join.
|
||||
*/
|
||||
REMOVED,
|
||||
;
|
||||
|
||||
val blocksOutbound: Boolean get() = true
|
||||
}
|
||||
+281
@@ -0,0 +1,281 @@
|
||||
/*
|
||||
* 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.quartz.marmot.protocolCore
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFailsWith
|
||||
import kotlin.test.assertNull
|
||||
import kotlin.test.assertSame
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/**
|
||||
* Branch selection under convergence policy v1.
|
||||
*
|
||||
* These are the rules that decide which history a group keeps, so a
|
||||
* disagreement here is a group split rather than a wrong answer. Everything is
|
||||
* pure — no timing, no transport — which is exactly why it can be pinned this
|
||||
* precisely.
|
||||
*/
|
||||
class BranchSelectorTest {
|
||||
private val policy = ConvergencePolicy.V1
|
||||
|
||||
private fun key(b: Int) = ByteArray(32) { b.toByte() }
|
||||
|
||||
private fun branch(
|
||||
forkEpoch: Long = 8,
|
||||
depth: Long,
|
||||
priority: TipPriority = TipPriority.ORDINARY,
|
||||
committer: Int = 0x10,
|
||||
digest: Int = 0x20,
|
||||
witnesses: Map<Long, Set<String>> = emptyMap(),
|
||||
) = CandidateBranch(
|
||||
forkEpoch = forkEpoch,
|
||||
tipEpoch = forkEpoch + depth,
|
||||
rawCommitDepth = depth,
|
||||
tipPriority = priority,
|
||||
tipCommitter = key(committer),
|
||||
tipDigest = key(digest),
|
||||
witnessesByEpoch = witnesses,
|
||||
)
|
||||
|
||||
/** Two distinct senders on one branch epoch — the quorum shape. */
|
||||
private fun quorumAt(epoch: Long) = mapOf(epoch to setOf("a".repeat(64), "b".repeat(64)))
|
||||
|
||||
@Test
|
||||
fun policyV1MatchesTheAdoptedConstants() {
|
||||
assertEquals(5, policy.maxRewindCommits)
|
||||
assertEquals(5, policy.appPayloadPastEpochLimit)
|
||||
assertEquals(1_000, policy.settlementQuiescenceMs)
|
||||
assertEquals(5_000, policy.maxConvergencePassMs)
|
||||
assertEquals(2, policy.witnessQuorumSendersPerEpoch)
|
||||
assertEquals(1, policy.witnessQuorumEpochs)
|
||||
assertEquals(1, policy.maxWitnessOverrideDepth)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun theWitnessBoostCannotExceedTheRewindHorizon() {
|
||||
// Without this bound, app-payload traffic could push a branch past the
|
||||
// rollback horizon and beat an arbitrarily longer valid commit branch.
|
||||
assertFailsWith<IllegalArgumentException> {
|
||||
policy.copy(maxWitnessOverrideDepth = policy.maxRewindCommits + 1)
|
||||
}
|
||||
}
|
||||
|
||||
// --- the worked example ---------------------------------------------------
|
||||
|
||||
@Test
|
||||
fun aWitnessedThreeCommitBranchBeatsAnUnwitnessedFour() {
|
||||
// Effective depth ties at 4 (3 + the 1-commit boost), so selection
|
||||
// falls to step 2, where quorum beats no quorum.
|
||||
val witnessed = branch(depth = 3, witnesses = quorumAt(10), digest = 0xff)
|
||||
val plain = branch(depth = 4, digest = 0x01)
|
||||
|
||||
assertEquals(4, witnessed.effectiveCommitDepth(policy))
|
||||
assertEquals(4, plain.effectiveCommitDepth(policy))
|
||||
assertTrue(witnessed.witnessQuorumMet(policy))
|
||||
assertTrue(!plain.witnessQuorumMet(policy))
|
||||
|
||||
assertSame(witnessed, BranchSelector.select(listOf(plain, witnessed), passBaseEpoch = 8, policy = policy))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun aFiveCommitBranchStillBeatsBoth() {
|
||||
// The boost is capped at one commit, so raw depth wins outright here.
|
||||
val witnessed = branch(depth = 3, witnesses = quorumAt(10))
|
||||
val plain = branch(depth = 4)
|
||||
val longest = branch(depth = 5, digest = 0xfe)
|
||||
|
||||
assertSame(
|
||||
longest,
|
||||
BranchSelector.select(listOf(witnessed, plain, longest), passBaseEpoch = 8, policy = policy),
|
||||
)
|
||||
}
|
||||
|
||||
// --- the comparison chain -------------------------------------------------
|
||||
|
||||
@Test
|
||||
fun rawDepthIsNotASeparateComparisonStep() {
|
||||
// Both reach effective depth 4 and both have quorum, so step 3 is the
|
||||
// witness SCORE — not raw depth. The shorter branch has the higher
|
||||
// score here, and it must win: an implementation that inserted a
|
||||
// raw-depth step would pick the other one.
|
||||
val shortHighScore =
|
||||
branch(depth = 3, witnesses = mapOf(10L to setOf("a".repeat(64), "b".repeat(64)), 11L to setOf("c".repeat(64), "d".repeat(64))))
|
||||
val longLowScore = branch(depth = 3, witnesses = quorumAt(9), digest = 0x01)
|
||||
|
||||
assertEquals(shortHighScore.effectiveCommitDepth(policy), longLowScore.effectiveCommitDepth(policy))
|
||||
assertTrue(shortHighScore.appWitnessScore(policy) > longLowScore.appWitnessScore(policy))
|
||||
assertSame(
|
||||
shortHighScore,
|
||||
BranchSelector.select(listOf(longLowScore, shortHighScore), passBaseEpoch = 8, policy = policy),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun privilegedTipsBeatOrdinaryOnesOnceEverythingElseTies() {
|
||||
// Stops commit-byte choice alone from letting an ordinary self-update
|
||||
// beat a tied admin removal.
|
||||
val admin = branch(depth = 1, priority = TipPriority.PRIVILEGED, committer = 0xff, digest = 0xff)
|
||||
val ordinary = branch(depth = 1, priority = TipPriority.ORDINARY, committer = 0x01, digest = 0x01)
|
||||
|
||||
assertSame(admin, BranchSelector.select(listOf(ordinary, admin), passBaseEpoch = 8, policy = policy))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun committerThenDigestBreakTheFinalTie() {
|
||||
val lowCommitter = branch(depth = 1, committer = 0x01, digest = 0xff)
|
||||
val highCommitter = branch(depth = 1, committer = 0x02, digest = 0x00)
|
||||
assertSame(
|
||||
lowCommitter,
|
||||
BranchSelector.select(listOf(highCommitter, lowCommitter), passBaseEpoch = 8, policy = policy),
|
||||
)
|
||||
|
||||
// Digest decides only when the SAME committer produced both.
|
||||
val lowDigest = branch(depth = 1, committer = 0x01, digest = 0x01)
|
||||
val highDigest = branch(depth = 1, committer = 0x01, digest = 0x02)
|
||||
assertSame(
|
||||
lowDigest,
|
||||
BranchSelector.select(listOf(highDigest, lowDigest), passBaseEpoch = 8, policy = policy),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun byteOrderingIsUnsigned() {
|
||||
// 0x01 must sort before 0x80. Account keys and digests are uniformly
|
||||
// distributed, so a signed comparison would invert roughly half of all
|
||||
// final ties — and two clients would disagree that often.
|
||||
val low = branch(depth = 1, committer = 0x01, digest = 0x01)
|
||||
val high = branch(depth = 1, committer = 0x80, digest = 0x01)
|
||||
assertSame(low, BranchSelector.select(listOf(high, low), passBaseEpoch = 8, policy = policy))
|
||||
|
||||
val lowDigest = branch(depth = 1, committer = 0x01, digest = 0x01)
|
||||
val highDigest = branch(depth = 1, committer = 0x01, digest = 0x80)
|
||||
assertSame(
|
||||
lowDigest,
|
||||
BranchSelector.select(listOf(highDigest, lowDigest), passBaseEpoch = 8, policy = policy),
|
||||
)
|
||||
}
|
||||
|
||||
// --- witness scoring ------------------------------------------------------
|
||||
|
||||
@Test
|
||||
fun oneSenderCannotInflateABranchByShouting() {
|
||||
// Witnesses are counted by DISTINCT sender per epoch, so a hundred
|
||||
// messages from one account score exactly one.
|
||||
val loner = branch(depth = 1, witnesses = mapOf(9L to setOf("a".repeat(64))))
|
||||
assertEquals(1, loner.appWitnessScore(policy))
|
||||
assertTrue(!loner.witnessQuorumMet(policy), "one sender is not a quorum")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun perEpochScoreIsCappedAtTheQuorumSize() {
|
||||
// Five senders in one epoch score 2, not 5 — so one very busy epoch
|
||||
// cannot outweigh several quiet ones.
|
||||
val busy = branch(depth = 1, witnesses = mapOf(9L to (1..5).map { "$it".repeat(64) }.toSet()))
|
||||
assertEquals(2, busy.appWitnessScore(policy))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun witnessesAtOrBeforeTheForkEpochDoNotCount() {
|
||||
// Branch epochs are strictly greater than fork_epoch: traffic from
|
||||
// before the divergence is shared history, not evidence for a branch.
|
||||
val branch = branch(forkEpoch = 8, depth = 2, witnesses = mapOf(8L to setOf("a".repeat(64), "b".repeat(64))))
|
||||
assertEquals(0, branch.appWitnessScore(policy))
|
||||
assertTrue(!branch.witnessQuorumMet(policy))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun quorumEpochsAreEvaluatedIndependently() {
|
||||
// The qualifying sender set may differ per epoch; no cohort has to
|
||||
// span them all.
|
||||
val branch =
|
||||
branch(
|
||||
depth = 2,
|
||||
witnesses =
|
||||
mapOf(
|
||||
9L to setOf("a".repeat(64), "b".repeat(64)),
|
||||
10L to setOf("c".repeat(64), "d".repeat(64)),
|
||||
),
|
||||
)
|
||||
assertTrue(branch.witnessQuorumMet(policy))
|
||||
assertEquals(4, branch.appWitnessScore(policy))
|
||||
}
|
||||
|
||||
// --- eligibility ----------------------------------------------------------
|
||||
|
||||
@Test
|
||||
fun branchesForkedOutsideTheRewindHorizonAreNeverSelected() {
|
||||
val inHorizon = branch(forkEpoch = 5, depth = 1)
|
||||
val outside = branch(forkEpoch = 4, depth = 50, digest = 0x01)
|
||||
|
||||
assertTrue(inHorizon.isEligible(passBaseEpoch = 10, policy = policy))
|
||||
assertTrue(!outside.isEligible(passBaseEpoch = 10, policy = policy))
|
||||
|
||||
// Even though it is far longer, the out-of-horizon branch is not a
|
||||
// candidate at all.
|
||||
assertSame(
|
||||
inHorizon,
|
||||
BranchSelector.select(listOf(outside, inHorizon), passBaseEpoch = 10, policy = policy),
|
||||
)
|
||||
assertNull(BranchSelector.select(listOf(outside), passBaseEpoch = 10, policy = policy))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun selectionIsIndependentOfInputOrder() {
|
||||
// The whole point of the algorithm: two clients that received the same
|
||||
// candidates in different orders must choose the same branch.
|
||||
val branches =
|
||||
listOf(
|
||||
branch(depth = 2, committer = 0x30, digest = 0x40),
|
||||
branch(depth = 3, witnesses = quorumAt(10), committer = 0x10, digest = 0x90),
|
||||
branch(depth = 4, committer = 0x20, digest = 0x10),
|
||||
branch(depth = 1, priority = TipPriority.PRIVILEGED, committer = 0x05, digest = 0x05),
|
||||
)
|
||||
val expected = BranchSelector.select(branches, passBaseEpoch = 8, policy = policy)
|
||||
|
||||
assertEquals(expected, BranchSelector.select(branches.reversed(), passBaseEpoch = 8, policy = policy))
|
||||
assertEquals(expected, BranchSelector.select(branches.shuffled(), passBaseEpoch = 8, policy = policy))
|
||||
for (rotation in branches.indices) {
|
||||
val rotated = branches.drop(rotation) + branches.take(rotation)
|
||||
assertEquals(expected, BranchSelector.select(rotated, passBaseEpoch = 8, policy = policy))
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun rankIsATotalOrderOverEligibleBranches() {
|
||||
val branches =
|
||||
listOf(
|
||||
branch(depth = 1, committer = 0x03),
|
||||
branch(depth = 1, committer = 0x01),
|
||||
branch(depth = 1, committer = 0x02),
|
||||
branch(forkEpoch = 0, depth = 9),
|
||||
)
|
||||
val ranked = BranchSelector.rank(branches, passBaseEpoch = 8, policy = policy)
|
||||
assertEquals(3, ranked.size, "the out-of-horizon branch is filtered out")
|
||||
assertEquals(listOf(0x01, 0x02, 0x03), ranked.map { it.tipCommitter[0].toInt() })
|
||||
}
|
||||
|
||||
@Test
|
||||
fun noEligibleBranchesSelectsNothing() {
|
||||
assertNull(BranchSelector.select(emptyList(), passBaseEpoch = 8, policy = policy))
|
||||
}
|
||||
}
|
||||
+143
@@ -0,0 +1,143 @@
|
||||
/*
|
||||
* 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.quartz.marmot.protocolCore
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
/** The canonical lifecycle transition table and its couplings. */
|
||||
class GroupLifecycleStateTest {
|
||||
@Test
|
||||
fun onlyStableMayPrepareALocalCommit() {
|
||||
for (state in GroupLifecycleState.entries) {
|
||||
assertEquals(state == GroupLifecycleState.STABLE, state.canPrepareLocalCommit, "$state")
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun mergingNeverDivertsIntoRecovering() {
|
||||
// A competing branch observed mid-merge is retained; the merge finishes
|
||||
// to Stable and the bounded pass then triggers Stable -> Recovering.
|
||||
// Diverting here would leave a half-applied epoch.
|
||||
assertFalse(GroupLifecycleState.MERGING.canTransitionTo(GroupLifecycleState.RECOVERING))
|
||||
assertTrue(GroupLifecycleState.MERGING.canTransitionTo(GroupLifecycleState.STABLE))
|
||||
assertTrue(GroupLifecycleState.STABLE.canTransitionTo(GroupLifecycleState.RECOVERING))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun disbandedIsAbsorbing() {
|
||||
for (target in GroupLifecycleState.entries) {
|
||||
assertFalse(
|
||||
GroupLifecycleState.DISBANDED.canTransitionTo(target),
|
||||
"no later branch may supersede a terminalized disband ($target)",
|
||||
)
|
||||
}
|
||||
assertTrue(GroupLifecycleState.DISBANDED.isTerminal)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun onlyRecoveringReachesDisbanded() {
|
||||
for (state in GroupLifecycleState.entries) {
|
||||
assertEquals(
|
||||
state == GroupLifecycleState.RECOVERING,
|
||||
state.canTransitionTo(GroupLifecycleState.DISBANDED),
|
||||
"$state",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun unrecoverableIsLocalAndRepairable() {
|
||||
// Local to one client, and it does have a way out — a repair, a
|
||||
// restore, or a verified re-join.
|
||||
assertTrue(GroupLifecycleState.UNRECOVERABLE.canTransitionTo(GroupLifecycleState.STABLE))
|
||||
assertFalse(GroupLifecycleState.UNRECOVERABLE.isTerminal)
|
||||
assertFalse(GroupLifecycleState.UNRECOVERABLE.canTransitionTo(GroupLifecycleState.RECOVERING))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun inboundChangesCanonicalStateOnlyInStable() {
|
||||
for (state in GroupLifecycleState.entries) {
|
||||
assertEquals(
|
||||
state == GroupLifecycleState.STABLE,
|
||||
state.mayApplyInboundToCanonicalState,
|
||||
"$state",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun everyStateReachableFromStableEventuallyReturnsToIt() {
|
||||
// Sanity on the table: nothing except Disbanded is a dead end.
|
||||
for (state in GroupLifecycleState.entries) {
|
||||
if (state == GroupLifecycleState.DISBANDED) continue
|
||||
assertTrue(
|
||||
reaches(state, GroupLifecycleState.STABLE),
|
||||
"$state must be able to reach Stable again",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
private fun reaches(
|
||||
from: GroupLifecycleState,
|
||||
target: GroupLifecycleState,
|
||||
seen: MutableSet<GroupLifecycleState> = mutableSetOf(),
|
||||
): Boolean {
|
||||
if (from == target) return true
|
||||
if (!seen.add(from)) return false
|
||||
return GroupLifecycleState.entries.any { from.canTransitionTo(it) && reaches(it, target, seen) }
|
||||
}
|
||||
|
||||
// --- convergence status ---------------------------------------------------
|
||||
|
||||
@Test
|
||||
fun onlySettledReleasesOutboundWork() {
|
||||
for (status in ConvergenceStatus.entries) {
|
||||
assertEquals(status == ConvergenceStatus.SETTLED, status.allowsOutboundWork, "$status")
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun theStatusLifecycleCombinationTableHolds() {
|
||||
assertTrue(ConvergenceStatus.SYNCING.isLegalIn(GroupLifecycleState.RECOVERING))
|
||||
assertTrue(ConvergenceStatus.SETTLED.isLegalIn(GroupLifecycleState.DISBANDED))
|
||||
// A group leaves Recovering for Stable only after Settled — so Settled
|
||||
// is never legal *in* Recovering.
|
||||
assertFalse(ConvergenceStatus.SETTLED.isLegalIn(GroupLifecycleState.RECOVERING))
|
||||
// Blocked on missing state is the Unrecoverable condition.
|
||||
assertTrue(ConvergenceStatus.BLOCKED.isLegalIn(GroupLifecycleState.UNRECOVERABLE))
|
||||
assertFalse(ConvergenceStatus.SYNCING.isLegalIn(GroupLifecycleState.UNRECOVERABLE))
|
||||
|
||||
// PendingPublish and Merging are local-publish states, not convergence
|
||||
// passes, so no status is meaningful in them.
|
||||
for (status in ConvergenceStatus.entries) {
|
||||
assertFalse(status.isLegalIn(GroupLifecycleState.PENDING_PUBLISH), "$status")
|
||||
assertFalse(status.isLegalIn(GroupLifecycleState.MERGING), "$status")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun assertEquals(
|
||||
expected: Boolean,
|
||||
actual: Boolean,
|
||||
message: String,
|
||||
) = kotlin.test.assertEquals(expected, actual, message)
|
||||
Reference in New Issue
Block a user