Skip to content
Merged
Show file tree
Hide file tree
Changes from 7 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,10 @@ import org.lfdecentralizedtrust.splice.sv.automation.singlesv.offboarding.{
}
import org.lfdecentralizedtrust.splice.sv.config.*
import org.lfdecentralizedtrust.splice.sv.SvAppClientConfig
import org.lfdecentralizedtrust.splice.syncoperator.config.SyncOperatorAppBackendConfig
import org.lfdecentralizedtrust.splice.syncoperator.config.{
SyncOperatorAppBackendConfig,
SyncOperatorSynchronizerNodeConfig,
}
import org.lfdecentralizedtrust.splice.validator.config.{
AnsAppExternalClientConfig,
ValidatorAppBackendConfig,
Expand Down Expand Up @@ -558,8 +561,10 @@ object ConfigTransforms {
conf
.focus(_.participantClient)
.modify(portTransform(bump, _))
.focus(_.sequencer.adminApi)
.focus(_.synchronizerNodes.current)
.modify(portTransform(bump, _))
.focus(_.synchronizerNodes.successor)
.modify(_.map(portTransform(bump, _)))
else conf
),
)
Expand Down Expand Up @@ -914,6 +919,19 @@ object ConfigTransforms {
private def portTransform(bump: Int, c: SvMediatorConfig): SvMediatorConfig =
c.focus(_.adminApi).modify(portTransform(bump, _))

private def portTransform(
bump: Int,
c: SyncOperatorSynchronizerNodeConfig,
): SyncOperatorSynchronizerNodeConfig =
c.focus(_.sequencer.adminApi)
.modify(portTransform(bump, _))
.focus(_.sequencer.internalApi)
.modify(_.map(portTransform(bump, _)))
.focus(_.sequencer.externalPublicApiUrl)
.modify(_.map(bumpUrl(bump, _)))
.focus(_.mediator)
.modify(_.map(_.focus(_.adminApi).modify(portTransform(bump, _))))

private def portTransform(
bump: Int,
c: SvSynchronizerNodeConfig,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,11 @@ import org.lfdecentralizedtrust.splice.sv.config.SvOnboardingConfig.FoundDso
import org.lfdecentralizedtrust.splice.syncoperator.config.{
SyncOperatorAppBackendConfig,
SyncOperatorAppClientConfig,
SyncOperatorLsuConfig,
SyncOperatorMediatorConfig,
SyncOperatorSequencerConfig,
SyncOperatorSynchronizerNodeConfig,
SyncOperatorSynchronizerNodesConfig,
}
import org.lfdecentralizedtrust.splice.util.{
Codec,
Expand Down Expand Up @@ -998,6 +1002,16 @@ object SpliceConfig {
deriveReader[SplitwellAppClientConfig]
implicit val syncOperatorSequencerConfigReader: ConfigReader[SyncOperatorSequencerConfig] =
deriveReader[SyncOperatorSequencerConfig]
implicit val syncOperatorMediatorConfigReader: ConfigReader[SyncOperatorMediatorConfig] =
deriveReader[SyncOperatorMediatorConfig]
implicit val syncOperatorSynchronizerNodeConfigReader
: ConfigReader[SyncOperatorSynchronizerNodeConfig] =
deriveReader[SyncOperatorSynchronizerNodeConfig]
implicit val syncOperatorSynchronizerNodesConfigReader
: ConfigReader[SyncOperatorSynchronizerNodesConfig] =
deriveReader[SyncOperatorSynchronizerNodesConfig]
implicit val syncOperatorLsuConfigReader: ConfigReader[SyncOperatorLsuConfig] =
deriveReader[SyncOperatorLsuConfig]
implicit val syncOperatorConfigReader: ConfigReader[SyncOperatorAppBackendConfig] =
deriveReader[SyncOperatorAppBackendConfig]
implicit val syncOperatorClientConfigReader: ConfigReader[SyncOperatorAppClientConfig] =
Expand Down Expand Up @@ -1330,6 +1344,16 @@ object SpliceConfig {
deriveWriter[SplitwellAppClientConfig]
implicit val syncOperatorSequencerConfigWriter: ConfigWriter[SyncOperatorSequencerConfig] =
deriveWriter[SyncOperatorSequencerConfig]
implicit val syncOperatorMediatorConfigWriter: ConfigWriter[SyncOperatorMediatorConfig] =
deriveWriter[SyncOperatorMediatorConfig]
implicit val syncOperatorSynchronizerNodeConfigWriter
: ConfigWriter[SyncOperatorSynchronizerNodeConfig] =
deriveWriter[SyncOperatorSynchronizerNodeConfig]
implicit val syncOperatorSynchronizerNodesConfigWriter
: ConfigWriter[SyncOperatorSynchronizerNodesConfig] =
deriveWriter[SyncOperatorSynchronizerNodesConfig]
implicit val syncOperatorLsuConfigWriter: ConfigWriter[SyncOperatorLsuConfig] =
deriveWriter[SyncOperatorLsuConfig]
implicit val syncOperatorConfigWriter: ConfigWriter[SyncOperatorAppBackendConfig] =
deriveWriter[SyncOperatorAppBackendConfig]
implicit val syncOperatorClientConfigWriter: ConfigWriter[SyncOperatorAppClientConfig] =
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
# The successor synchronizer node for a sync operator LSU. Left uninitialized on purpose: the
# operator's LSU trigger initializes both nodes from the synchronizer it is upgrading.
include required("include/canton-basic.conf")
include required("include/sequencers.conf")
include required("include/mediators.conf")

canton {
parameters {
non-standard-config = yes
}

sequencers {
syncOperatorSuccessorSequencer = ${_sequencer_reference_template} {
public-api.port = 27608
admin-api.port = 27609
storage.config.properties.databaseName = "sequencer_sync_operator_successor"
sequencer.config.storage.config.properties.databaseName = "sequencer_driver_sync_operator_successor"
}
syncOperatorSuccessorSequencer.storage.config.properties.databaseName = ${?SYNC_OPERATOR_SUCCESSOR_SEQUENCER_DB}
syncOperatorSuccessorSequencer.sequencer.config.storage.config.properties.databaseName = ${?SYNC_OPERATOR_SUCCESSOR_SEQUENCER_DRIVER_DB}
}

mediators {
syncOperatorSuccessorMediator = ${_mediator_template} {
admin-api.port = 27607
storage.config.properties.databaseName = "mediator_sync_operator_successor"
}
syncOperatorSuccessorMediator.storage.config.properties.databaseName = ${?SYNC_OPERATOR_SUCCESSOR_MEDIATOR_DB}
}
}
2 changes: 1 addition & 1 deletion apps/app/src/test/resources/sync-operator-topology.conf
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ canton {
admin-api.port = 5116
participant-client = ${canton.validator-apps.splitwellValidator.participant-client}
operator-user = ${canton.validator-apps.splitwellValidator.app-instances.syncOperator.service-user}
sequencer.admin-api.port = 5709
synchronizer-nodes.current.sequencer.admin-api.port = 5709
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,8 @@ import org.lfdecentralizedtrust.splice.sv.config.{
SvSynchronizerNodeConfig,
SvSynchronizerNodesConfig,
}
import org.lfdecentralizedtrust.splice.sv.lsu.{LsuTransferTrafficTrigger, LsuTrigger}
import org.lfdecentralizedtrust.splice.lsu.LsuTransferTrafficTrigger
import org.lfdecentralizedtrust.splice.sv.lsu.LsuTrigger
import org.lfdecentralizedtrust.splice.util.*
import org.lfdecentralizedtrust.splice.wallet.config.WalletAppClientConfig
import org.lfdecentralizedtrust.splice.wallet.store.TxLogEntry.Http.BuyTrafficRequestStatus
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,257 @@
// Copyright (c) 2024 Digital Asset (Switzerland) GmbH and/or its affiliates. All rights reserved.
// SPDX-License-Identifier: Apache-2.0

package org.lfdecentralizedtrust.splice.integration.tests

import better.files.*
import com.digitalasset.canton.SynchronizerAlias
import com.digitalasset.canton.config.{FullClientConfig, NonNegativeFiniteDuration}
import com.digitalasset.canton.config.RequireTypes.{NonNegativeInt, Port}
import com.digitalasset.canton.data.CantonTimestamp
import com.digitalasset.canton.version.ProtocolVersion
import org.lfdecentralizedtrust.splice.automation.Trigger
import org.lfdecentralizedtrust.splice.codegen.java.splice.decentralizedsynchronizer.GovernanceParameters
import org.lfdecentralizedtrust.splice.config.ConfigTransforms
import org.lfdecentralizedtrust.splice.integration.EnvironmentDefinition
import org.lfdecentralizedtrust.splice.integration.tests.SpliceTests.IntegrationTest
import org.lfdecentralizedtrust.splice.lsu.{LsuRollForwardTimestamp, LsuTransferTrafficTrigger}
import org.lfdecentralizedtrust.splice.syncoperator.automation.DedicatedLsuTrigger
import org.lfdecentralizedtrust.splice.syncoperator.config.{
SyncOperatorLsuConfig,
SyncOperatorMediatorConfig,
SyncOperatorSequencerConfig,
SyncOperatorSynchronizerNodeConfig,
}
import org.lfdecentralizedtrust.splice.util.{
StandaloneCanton,
SyncOperatorTestUtil,
TriggerTestUtil,
WalletTestUtil,
}

import java.time.Duration
import scala.concurrent.duration.*
import scala.jdk.CollectionConverters.*

/** Upgrades the synchronizer the operator runs, on a schedule the operator sets itself, and checks
* that members keep both what they purchased and what they have already spent.
*
* The successor's sequencer and mediator are started only for this test, uninitialized, so the
* operator's LSU trigger initializes them from the synchronizer it is upgrading.
*/
class SyncOperatorLsuIntegrationTest
extends IntegrationTest
with StandaloneCanton
with SyncOperatorTestUtil
with TriggerTestUtil
with WalletTestUtil {

override def dbsSuffix = "sync_operator_lsu"

override def usesDbs: Seq[String] = Seq(
s"sequencer_${dbsSuffix}_successor",
s"sequencer_driver_${dbsSuffix}_successor",
s"mediator_${dbsSuffix}_successor",
) ++ super.usesDbs

private val purchase = 2_000_000L
private val splitwellAlias = SynchronizerAlias.tryCreate("splitwell")

// The synchronizer is bootstrapped at serial 0 and protocol version 35, see bootstrap-canton.sc.
private val successorSerial = NonNegativeInt.one
private val successorPv = ProtocolVersion.v36

// Ports of the standalone successor nodes, see standalone-sync-operator-successor.conf.
private val successorSequencerAdminPort = Port.tryCreate(27609)
private val successorSequencerPublicPort = Port.tryCreate(27608)
private val successorMediatorAdminPort = Port.tryCreate(27607)

// The operator reads its schedule from files, so the test writes them once it is ready rather
// than picking a time while the environment is still starting.
private lazy val scheduleDir = File.newTemporaryDirectory("sync-operator-lsu")
private lazy val freezeTimeFile = scheduleDir / "topology-freeze-time"
private lazy val upgradeTimeFile = scheduleDir / "upgrade-time"

override def afterAll(): Unit = {
scheduleDir.delete(swallowIOExceptions = true)
super.afterAll()
}

override def environmentDefinition: SpliceEnvironmentDefinition =
EnvironmentDefinition
.fromResources(
Seq("simple-topology-1sv.conf", "sync-operator-topology.conf"),
this.getClass.getSimpleName,
)
.withOnlyAliceValidatorConnectingToSplitwell
.withStandardSetup
.addConfigTransform((_, config) =>
ConfigTransforms.updateAllSyncOperatorAppConfigs_ { c =>
c.copy(
synchronizerNodes = c.synchronizerNodes.copy(
// The mediator is only read during an upgrade, so the shared topology leaves it out.
current = c.synchronizerNodes.current.copy(
mediator = Some(
SyncOperatorMediatorConfig(FullClientConfig(port = Port.tryCreate(5707)))
)
),
successor = Some(
SyncOperatorSynchronizerNodeConfig(
sequencer = SyncOperatorSequencerConfig(
adminApi = FullClientConfig(port = successorSequencerAdminPort),
internalApi = Some(FullClientConfig(port = successorSequencerPublicPort)),
externalPublicApiUrl =
Some(s"http://localhost:${successorSequencerPublicPort.unwrap}"),
),
mediator = Some(
SyncOperatorMediatorConfig(FullClientConfig(port = successorMediatorAdminPort))
),
protocolVersion = successorPv,
serial = Some(successorSerial),
)
),
),
lsu = Some(
SyncOperatorLsuConfig(
topologyFreezeTime = LsuRollForwardTimestamp.TimestampFromFile(freezeTimeFile.path),
upgradeTime = LsuRollForwardTimestamp.TimestampFromFile(upgradeTimeFile.path),
newPhysicalSynchronizerSerial = successorSerial,
newPhysicalSynchronizerProtocolVersion = successorPv,
)
),
lsuDumpPath = Some((scheduleDir / "lsu-dump.json").path),
parameters = c.parameters.copy(
spliceCachingConfigs = c.parameters.spliceCachingConfigs.copy(
physicalSynchronizerExpiration = NonNegativeFiniteDuration.ofSeconds(1)
)
),
// Both triggers reach the successor's nodes and log if they cannot, so they stay
// paused until the test has started them.
automation = c.automation
.withPausedTrigger[DedicatedLsuTrigger]
.withPausedTrigger[LsuTransferTrafficTrigger],
)
}(config)
)

"sync operator" should {

"upgrade its synchronizer and carry the traffic state onto the successor" in { implicit env =>
val operatorParty = syncOperatorBackend.appState.store.key.operatorParty
val dsoParty = sv1Backend.getDsoInfo().dsoParty
val dsoRules = sv1Backend.getDsoInfo().dsoRules
val synchronizerId = aliceValidatorBackend.participantClientWithAdminToken.synchronizers
.id_of(splitwellAlias)
.logical
val member = aliceValidatorBackend.participantClient.id

clue("the DSO registers the synchronizer to this operator") {
sv1Backend.participantClientWithAdminToken.ledger_api_extensions.commands
.submitJava(
actAs = Seq(dsoParty),
readAs = Seq(dsoParty),
commands = dsoRules.contractId
.exerciseDsoRules_RegisterSynchronizer(
synchronizerId.toProtoPrimitive,
operatorParty.toProtoPrimitive,
new GovernanceParameters(java.math.BigDecimal.ONE.setScale(10)),
)
.commands
.asScala
.toSeq,
userId = sv1Backend.config.ledgerApiUser,
)
}

val registration = eventually() {
sv1ScanBackend.lookupSynchronizerRegistration(synchronizerId.toProtoPrimitive).value
}

val aliceParty = onboardWalletUser(aliceWalletClient, aliceValidatorBackend)
aliceWalletClient.tap(walletUsdToAmulet(200.0))

actAndCheck(
"alice buys traffic for the operator's synchronizer",
buyTraffic(aliceParty, member, synchronizerId, registration, dsoParty, purchase),
)(
"the operator grants it on the synchronizer it is about to upgrade",
_ => extraTrafficLimit(member) shouldBe purchase,
)

clue("some of it is spent, so there is consumption to carry across the upgrade") {
eventually() {
trafficState(member).map(_.extraTrafficConsumed.value).getOrElse(0L) should be > 0L
}
}

val consumedBefore =
trafficState(member).map(_.extraTrafficConsumed.value).getOrElse(0L)

val currentPsid = syncOperatorBackend.appState.sequencerAdminConnection
.getPhysicalSynchronizerId()
.futureValue
currentPsid.serial.value should be < successorSerial.value

withCanton(
Seq(testResourcesPath / "standalone-sync-operator-successor.conf"),
Seq.empty,
"sync-operator-lsu-successor",
"SYNC_OPERATOR_SUCCESSOR_SEQUENCER_DB" -> s"sequencer_${dbsSuffix}_successor",
"SYNC_OPERATOR_SUCCESSOR_SEQUENCER_DRIVER_DB" -> s"sequencer_driver_${dbsSuffix}_successor",
"SYNC_OPERATOR_SUCCESSOR_MEDIATOR_DB" -> s"mediator_${dbsSuffix}_successor",
) {
val lsuTrigger = syncOperatorBackend.appState.automation.trigger[DedicatedLsuTrigger]
val trafficTrigger =
syncOperatorBackend.appState.automation.trigger[LsuTransferTrafficTrigger]

setTriggersWithin(triggersToResumeAtStart = Seq[Trigger](lsuTrigger, trafficTrigger)) {
val upgradeTime = actAndCheck(
"the operator schedules the upgrade", {
val now = CantonTimestamp.now()
// Far enough out that a slow initialization cannot overshoot it.
val upgradeTime = now.plus(Duration.ofMinutes(5))
freezeTimeFile.overwrite(now.toInstant.toString)
upgradeTimeFile.overwrite(upgradeTime.toInstant.toString)
upgradeTime
},
)(
"the announcement is published on the synchronizer being upgraded",
_ =>
syncOperatorBackend.appState.sequencerAdminConnection
.listLsuAnnouncements(currentPsid.logical)
.futureValue
.map(_.mapping.successorSynchronizerId.serial) should contain(successorSerial),
)._1

val successorNode =
syncOperatorBackend.appState.synchronizerNodes.successor.value

clue(s"the successor's nodes are initialized from the predecessor before $upgradeTime") {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

let's run some transaction on the new physical synchronizer after upgrade time to check that it worked fully

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

added

eventuallySucceeds(5.minutes) {
val successorPsid =
successorNode.sequencerAdminConnection.getPhysicalSynchronizerId().futureValue
successorPsid.logical shouldBe currentPsid.logical
successorPsid.serial shouldBe successorSerial
successorPsid.protocolVersion shouldBe successorPv
}
eventually() {
successorNode.mediatorAdminConnection.isNodeInitialized().futureValue shouldBe true
}
}

clue("both halves of the traffic state are carried onto the successor") {
eventually(5.minutes) {
val state = successorNode.sequencerAdminConnection
.lookupSequencerTrafficControlState(member)
.futureValue
.value
state.extraTrafficLimit.value shouldBe purchase
// Without the transfer this resets to zero and members get back what they spent.
state.extraTrafficConsumed.value should be >= consumedBefore
}
}
}
}
}
}
}
Loading
Loading