From ed4769fae24a5840f9b686cf2ff19157e905ab94 Mon Sep 17 00:00:00 2001 From: sadiq1971 Date: Wed, 26 Aug 2026 00:37:25 +0600 Subject: [PATCH 1/3] feat: launch the sync operator app and test that it starts [ci] Wires the app into apps-app (config, environment, console references, metrics, config transforms) so it can be configured and started like the other apps. Adds an integration test that starts it against a base topology, checks it takes its synchronizer id from the sequencer it is configured with, then stops and restarts it. The operator runs against the splitwell synchronizer's sequencer, which stands in for a dedicated synchronizer until the test topologies include one. Signed-off-by: sadiq1971 --- .../splice/config/ConfigTransforms.scala | 24 +++++- .../splice/config/SpliceConfig.scala | 63 ++++++++++++++++ .../console/SyncOperatorReference.scala | 73 +++++++++++++++++++ .../SpliceConsoleEnvironment.scala | 54 ++++++++++++++ .../environment/SpliceEnvironment.scala | 35 ++++++++- .../splice/environment/SyncOperatorApps.scala | 35 +++++++++ .../splice/metrics/SpliceMetricsFactory.scala | 15 ++++ .../resources/sync-operator-topology.conf | 32 ++++++++ .../tests/SyncOperatorIntegrationTest.scala | 43 +++++++++++ .../util/CommonAppInstanceReferences.scala | 14 ++++ build.sbt | 1 + test-full-class-names.log | 1 + 12 files changed, 388 insertions(+), 2 deletions(-) create mode 100644 apps/app/src/main/scala/org/lfdecentralizedtrust/splice/console/SyncOperatorReference.scala create mode 100644 apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SyncOperatorApps.scala create mode 100644 apps/app/src/test/resources/sync-operator-topology.conf create mode 100644 apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala index 70795113ba..45781eab70 100644 --- a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala @@ -24,6 +24,7 @@ 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.validator.config.{ AnsAppExternalClientConfig, ValidatorAppBackendConfig, @@ -53,7 +54,8 @@ object ConfigTransforms { case object Scan extends ConfigurableApp case object Validator extends ConfigurableApp case object Splitwell extends ConfigurableApp - val All = Seq(Sv, Scan, Validator, Splitwell) + case object SyncOperator extends ConfigurableApp + val All = Seq(Sv, Scan, Validator, Splitwell, SyncOperator) } def makeAllTimeoutsBounded: ConfigTransform = { @@ -129,6 +131,7 @@ object ConfigTransforms { updateAllRemoteSplitwellAppConfigs_(c => c.copy(ledgerApiUser = s"${c.ledgerApiUser}-$suffix") ), + updateAllSyncOperatorAppConfigs_(c => c.copy(operatorUser = s"${c.operatorUser}-$suffix")), updateAllAnsAppExternalClientConfigs_(c => c.copy(ledgerApiUser = s"${c.ledgerApiUser}-$suffix") ), @@ -156,6 +159,8 @@ object ConfigTransforms { case Scan => updateAllScanAppConfigs_(c => c.focus(_.automation).modify(transform)) case Validator => updateAllValidatorConfigs_(c => c.focus(_.automation).modify(transform)) case Splitwell => updateAllSplitwellAppConfigs_(c => c.focus(_.automation).modify(transform)) + case SyncOperator => + updateAllSyncOperatorAppConfigs_(c => c.focus(_.automation).modify(transform)) } } @@ -166,6 +171,7 @@ object ConfigTransforms { updateAllScanAppConfigs_(c => c.focus(_.automation).modify(transform)), updateAllValidatorConfigs_(c => c.focus(_.automation).modify(transform)), updateAllSplitwellAppConfigs_(c => c.focus(_.automation).modify(transform)), + updateAllSyncOperatorAppConfigs_(c => c.focus(_.automation).modify(transform)), ) transforms.foldLeft(config)((c, tf) => tf(c)) } @@ -216,6 +222,7 @@ object ConfigTransforms { type ScanAppTransform = Endo[ScanAppBackendConfig] type SplitwellAppTransform = Endo[SplitwellAppBackendConfig] type RemoteSplitwellAppTransform = Endo[SplitwellAppClientConfig] + type SyncOperatorAppTransform = Endo[SyncOperatorAppBackendConfig] type AutomationConfigTransform = Endo[AutomationConfig] def withPausedSvDomainComponentsOffboardingTriggers(): ConfigTransform = @@ -413,6 +420,18 @@ object ConfigTransforms { ): ConfigTransform = updateAllRemoteSplitwellAppConfigs((_, config) => update(config)) + def updateAllSyncOperatorAppConfigs( + update: (String, SyncOperatorAppBackendConfig) => SyncOperatorAppBackendConfig + ): ConfigTransform = + _.focus(_.syncOperatorApps).modify(_.map { case (name, config) => + (name, update(name.unwrap, config)) + }) + + def updateAllSyncOperatorAppConfigs_( + update: SyncOperatorAppTransform + ): ConfigTransform = + updateAllSyncOperatorAppConfigs((_, config) => update(config)) + def bumpOptionalUrl(o: Option[String], bump: Int): Option[String] = { o.map(bumpUrl(bump, _)) } @@ -899,6 +918,9 @@ object ConfigTransforms { updateAllRemoteSplitwellAppConfigs_(c => { c.focus(_.participantClient.ledgerApi).modify(enableAuth(c.ledgerApiUser, _)) }), + updateAllSyncOperatorAppConfigs_(c => { + c.focus(_.participantClient.ledgerApi).modify(enableAuth(c.operatorUser, _)) + }), ) } diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/SpliceConfig.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/SpliceConfig.scala index fad206d8f9..9e48c0a4fe 100644 --- a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/SpliceConfig.scala +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/SpliceConfig.scala @@ -33,6 +33,11 @@ import org.lfdecentralizedtrust.splice.splitwell.config.{ import org.lfdecentralizedtrust.splice.sv.config.* import org.lfdecentralizedtrust.splice.sv.{SvAppClientConfig} import org.lfdecentralizedtrust.splice.sv.config.SvOnboardingConfig.FoundDso +import org.lfdecentralizedtrust.splice.syncoperator.config.{ + SyncOperatorAppBackendConfig, + SyncOperatorAppClientConfig, + SyncOperatorSequencerConfig, +} import org.lfdecentralizedtrust.splice.util.{Codec, SpliceRateLimitConfig} import org.lfdecentralizedtrust.splice.validator.config.* import org.lfdecentralizedtrust.splice.wallet.config.{ @@ -101,6 +106,8 @@ case class SpliceConfig( ansAppExternalClients: Map[InstanceName, AnsAppExternalClientConfig] = Map.empty, splitwellApps: Map[InstanceName, SplitwellAppBackendConfig] = Map.empty, splitwellAppClients: Map[InstanceName, SplitwellAppClientConfig] = Map.empty, + syncOperatorApps: Map[InstanceName, SyncOperatorAppBackendConfig] = Map.empty, + syncOperatorAppClients: Map[InstanceName, SyncOperatorAppClientConfig] = Map.empty, override val remoteParticipants: Map[InstanceName, RemoteParticipantConfig] = Map.empty, monitoring: MonitoringConfig = MonitoringConfig(), parameters: CantonParameters = CantonParameters( @@ -299,6 +306,50 @@ case class SpliceConfig( n.unwrap -> c } + private lazy val syncOperatorAppParameters_ : Map[InstanceName, SharedSpliceAppParameters] = + syncOperatorApps.fmap { syncOperatorConfig => + SharedSpliceAppParameters( + monitoring, + parameters.timeouts.processing, + parameters.timeouts.requestTimeout, + UpgradesConfig(), + syncOperatorConfig.parameters.circuitBreakers, + syncOperatorConfig.parameters.enabledFeatures, + syncOperatorConfig.parameters.caching, + parameters.enableAdditionalConsistencyChecks, + features.enablePreviewCommands, + parameters.nonStandardConfig, + syncOperatorConfig.sequencerClient, + dontWarnOnDeprecatedPV = false, + dbMigrateAndStart = true, + batchingConfig = new BatchingConfig(), + ) + } + + private[splice] def syncOperatorAppParameters( + appName: InstanceName + ): SharedSpliceAppParameters = + nodeParametersFor(syncOperatorAppParameters_, "sync-operator-app", appName) + + /** Use `syncOperatorAppParameters` instead! + */ + def trySyncOperatorAppParametersByString(name: String): SharedSpliceAppParameters = + syncOperatorAppParameters( + InstanceName.tryCreate(name) + ) + + /** Use `syncOperators` instead! + */ + def syncOperatorsByString: Map[String, SyncOperatorAppBackendConfig] = + syncOperatorApps.map { case (n, c) => + n.unwrap -> c + } + + def syncOperatorClientsByString: Map[String, SyncOperatorAppClientConfig] = + syncOperatorAppClients.map { case (n, c) => + n.unwrap -> c + } + override def dumpString: String = { val writers = new SpliceConfig.ConfigWriters(confidential = true) import writers.* @@ -872,6 +923,12 @@ object SpliceConfig { deriveReader[SplitwellAppBackendConfig] implicit val splitwellClientConfigReader: ConfigReader[SplitwellAppClientConfig] = deriveReader[SplitwellAppClientConfig] + implicit val syncOperatorSequencerConfigReader: ConfigReader[SyncOperatorSequencerConfig] = + deriveReader[SyncOperatorSequencerConfig] + implicit val syncOperatorConfigReader: ConfigReader[SyncOperatorAppBackendConfig] = + deriveReader[SyncOperatorAppBackendConfig] + implicit val syncOperatorClientConfigReader: ConfigReader[SyncOperatorAppClientConfig] = + deriveReader[SyncOperatorAppClientConfig] implicit val spliceConfigReader: ConfigReader[SpliceConfig] = deriveReader[SpliceConfig] } @@ -1180,6 +1237,12 @@ object SpliceConfig { deriveWriter[SplitwellAppBackendConfig] implicit val splitwellClientConfigWriter: ConfigWriter[SplitwellAppClientConfig] = deriveWriter[SplitwellAppClientConfig] + implicit val syncOperatorSequencerConfigWriter: ConfigWriter[SyncOperatorSequencerConfig] = + deriveWriter[SyncOperatorSequencerConfig] + implicit val syncOperatorConfigWriter: ConfigWriter[SyncOperatorAppBackendConfig] = + deriveWriter[SyncOperatorAppBackendConfig] + implicit val syncOperatorClientConfigWriter: ConfigWriter[SyncOperatorAppClientConfig] = + deriveWriter[SyncOperatorAppClientConfig] implicit val spliceConfigWriter: ConfigWriter[SpliceConfig] = deriveWriter[SpliceConfig] diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/console/SyncOperatorReference.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/console/SyncOperatorReference.scala new file mode 100644 index 0000000000..0e4ffccb64 --- /dev/null +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/console/SyncOperatorReference.scala @@ -0,0 +1,73 @@ +// Copyright (c) 2024 Digital Asset (Switzerland) GmbH and/or its affiliates. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package org.lfdecentralizedtrust.splice.console + +import com.digitalasset.canton.console.{BaseInspection, Help} +import org.lfdecentralizedtrust.splice.config.NetworkAppClientConfig +import org.lfdecentralizedtrust.splice.environment.SpliceConsoleEnvironment +import org.lfdecentralizedtrust.splice.syncoperator.{SyncOperatorApp, SyncOperatorAppBootstrap} +import org.lfdecentralizedtrust.splice.syncoperator.automation.SyncOperatorAutomationService +import org.lfdecentralizedtrust.splice.syncoperator.config.{ + SyncOperatorAppBackendConfig, + SyncOperatorAppClientConfig, +} + +/** Sync operator app reference. The app has no HTTP API of its own, so only the admin endpoints + * shared by every Splice app are available here. + */ +abstract class SyncOperatorAppReference( + override val spliceConsoleEnvironment: SpliceConsoleEnvironment, + override val name: String, +) extends HttpAppReference { + + override def basePath = "/api/syncoperator" +} + +final class SyncOperatorAppClientReference( + override val spliceConsoleEnvironment: SpliceConsoleEnvironment, + name: String, + val config: SyncOperatorAppClientConfig, +) extends SyncOperatorAppReference(spliceConsoleEnvironment, name) { + + override protected val instanceType = "Sync Operator Client" + + override def httpClientConfig = config.adminApi +} + +final class SyncOperatorAppBackendReference( + override val consoleEnvironment: SpliceConsoleEnvironment, + name: String, +) extends SyncOperatorAppReference(consoleEnvironment, name) + with AppBackendReference + with BaseInspection[SyncOperatorApp] { + + override def runningNode: Option[SyncOperatorAppBootstrap] = + consoleEnvironment.environment.syncOperators.getRunning(name) + + override def startingNode: Option[SyncOperatorAppBootstrap] = + consoleEnvironment.environment.syncOperators.getStarting(name) + + override protected val instanceType = "Sync Operator Backend" + + override def httpClientConfig = NetworkAppClientConfig( + s"http://127.0.0.1:${config.clientAdminApi.port}" + ) + + override val nodes: org.lfdecentralizedtrust.splice.environment.SyncOperatorApps = + consoleEnvironment.environment.syncOperators + + @Help.Summary( + "Returns the state of this app. May only be called while the app is running." + ) + def appState: SyncOperatorApp.State = _appState[SyncOperatorApp.State, SyncOperatorApp] + + @Help.Summary( + "Returns the automation service for the sync operator app. May only be called while the app is running." + ) + def syncOperatorAutomation: SyncOperatorAutomationService = appState.automation + + @Help.Summary("Return local sync operator app config") + def config: SyncOperatorAppBackendConfig = + consoleEnvironment.environment.config.syncOperatorsByString(name) +} diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceConsoleEnvironment.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceConsoleEnvironment.scala index 08fe1fc2cd..3dd2a59078 100644 --- a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceConsoleEnvironment.scala +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceConsoleEnvironment.scala @@ -86,11 +86,13 @@ class SpliceConsoleEnvironment( fullDsoApps.local, appsHostedByValidator.local, appsHostedByThirdParty.local, + syncOperators.local, ), mergeRemoteSpliceInstances( fullDsoApps.remote, appsHostedByValidator.remote, appsHostedByThirdParty.remote, + syncOperators.remote, ), ) } @@ -180,6 +182,18 @@ class SpliceConsoleEnvironment( environment.config.splitwellClientsByString.keys.map(createRemoteSplitwellReference).toSeq, ) + lazy val syncOperators: NodeReferences[ + SyncOperatorAppReference, + SyncOperatorAppClientReference, + SyncOperatorAppBackendReference, + ] = + NodeReferences( + environment.config.syncOperatorsByString.keys.map(createSyncOperatorReference).toSeq, + environment.config.syncOperatorClientsByString.keys + .map(createRemoteSyncOperatorReference) + .toSeq, + ) + private def createValidatorReference(name: String): ValidatorAppBackendReference = new ValidatorAppBackendReference(this, name) @@ -223,6 +237,16 @@ class SpliceConsoleEnvironment( private def createRemoteSplitwellReference(name: String): SplitwellAppClientReference = new SplitwellAppClientReference(this, name, environment.config.splitwellClientsByString(name)) + private def createSyncOperatorReference(name: String): SyncOperatorAppBackendReference = + new SyncOperatorAppBackendReference(this, name) + + private def createRemoteSyncOperatorReference(name: String): SyncOperatorAppClientReference = + new SyncOperatorAppClientReference( + this, + name, + environment.config.syncOperatorClientsByString(name), + ) + override protected def topLevelValues: Seq[TopLevelValue[?]] = { super.topLevelValues ++ @@ -309,6 +333,36 @@ class SpliceConsoleEnvironment( ), splitwells.remote, Seq("App References"), + ) :++ syncOperators.local.map(v => + TopLevelValue( + v.name, + helpText("local sync operator app", v.name), + v, + Seq("App References"), + ) + ) :++ syncOperators.remote.map(v => + TopLevelValue( + v.name, + helpText("sync operator app client", v.name), + v, + Seq("App References"), + ) + ) :+ TopLevelValue( + "syncOperators", + helpText( + "All local sync operator instances" + genericNodeReferencesDoc, + "SyncOperators", + ), + syncOperators.local, + Seq("App References"), + ) :+ TopLevelValue( + "syncOperatorClients", + helpText( + "All sync operator client instances" + genericNodeReferencesDoc, + "SyncOperators", + ), + syncOperators.remote, + Seq("App References"), ) :++ scans.local.map(scan => TopLevelValue(scan.name, helpText("Scan app", scan.name), scan, Seq("Scan")) ) :++ scans.remote.map(scan => diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceEnvironment.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceEnvironment.scala index 208426d487..444ae65597 100644 --- a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceEnvironment.scala +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SpliceEnvironment.scala @@ -19,6 +19,8 @@ import org.lfdecentralizedtrust.splice.splitwell.SplitwellAppBootstrap import org.lfdecentralizedtrust.splice.splitwell.config.SplitwellAppBackendConfig import org.lfdecentralizedtrust.splice.sv.SvAppBootstrap import org.lfdecentralizedtrust.splice.sv.config.SvAppBackendConfig +import org.lfdecentralizedtrust.splice.syncoperator.SyncOperatorAppBootstrap +import org.lfdecentralizedtrust.splice.syncoperator.config.SyncOperatorAppBackendConfig import org.lfdecentralizedtrust.splice.validator.ValidatorAppBootstrap import org.lfdecentralizedtrust.splice.validator.config.ValidatorAppBackendConfig @@ -169,9 +171,40 @@ class SpliceEnvironment( loggerFactory, ) + protected def createSyncOperator( + name: String, + syncOperatorConfig: SyncOperatorAppBackendConfig, + ): SyncOperatorAppBootstrap = { + val appLoggerFactory = loggerFactory.append(SyncOperatorAppBootstrap.LoggerFactoryKeyName, name) + SyncOperatorAppBootstrap( + name, + syncOperatorConfig, + config.trySyncOperatorAppParametersByString(name), + createClock(Some(SyncOperatorAppBootstrap.LoggerFactoryKeyName -> name)), + metrics.forSyncOperator(name), + testingConfig, + futureSupervisor, + appLoggerFactory, + configuredOpenTelemetry, + ) + .valueOr(err => + throw new RuntimeException( + s"Failed to create sync operator bootstrap: $err" + ) + ) + } + + lazy val syncOperators = new SyncOperatorApps( + createSyncOperator, + timeouts, + config.syncOperatorsByString, + config.trySyncOperatorAppParametersByString, + loggerFactory, + ) + // Ordering here matches SpliceConsoleEnvironment.startupOrderPrecedence def allSplices: List[Nodes[CantonNode, CantonNodeBootstrap[CantonNode]]] = - List(svs, scans, validators, splitwells) + List(svs, scans, validators, splitwells, syncOperators) override def allNodes: List[Nodes[CantonNode, CantonNodeBootstrap[CantonNode]]] = super.allNodes ::: allSplices diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SyncOperatorApps.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SyncOperatorApps.scala new file mode 100644 index 0000000000..cb37973130 --- /dev/null +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/environment/SyncOperatorApps.scala @@ -0,0 +1,35 @@ +// Copyright (c) 2024 Digital Asset (Switzerland) GmbH and/or its affiliates. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package org.lfdecentralizedtrust.splice.environment + +import org.lfdecentralizedtrust.splice.config.SharedSpliceAppParameters +import org.lfdecentralizedtrust.splice.syncoperator.{SyncOperatorApp, SyncOperatorAppBootstrap} +import org.lfdecentralizedtrust.splice.syncoperator.config.SyncOperatorAppBackendConfig +import com.digitalasset.canton.concurrent.ExecutionContextIdlenessExecutorService +import com.digitalasset.canton.config.ProcessingTimeout +import com.digitalasset.canton.environment.ManagedNodes +import com.digitalasset.canton.logging.NamedLoggerFactory + +/** Sync operator app instances. */ +class SyncOperatorApps( + create: (String, SyncOperatorAppBackendConfig) => SyncOperatorAppBootstrap, + _timeouts: ProcessingTimeout, + configs: Map[String, SyncOperatorAppBackendConfig], + parametersFor: String => SharedSpliceAppParameters, + _loggerFactory: NamedLoggerFactory, +)(implicit + protected val executionContext: ExecutionContextIdlenessExecutorService +) extends ManagedNodes[ + SyncOperatorApp, + SyncOperatorAppBackendConfig, + SharedSpliceAppParameters, + SyncOperatorAppBootstrap, + ]( + create, + _timeouts, + configs, + parametersFor, + startUpGroup = 0, + _loggerFactory, + ) {} diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/metrics/SpliceMetricsFactory.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/metrics/SpliceMetricsFactory.scala index 2b155ae273..b517d47a3a 100644 --- a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/metrics/SpliceMetricsFactory.scala +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/metrics/SpliceMetricsFactory.scala @@ -10,6 +10,7 @@ import com.digitalasset.canton.metrics.{DbStorageHistograms, MetricsFactoryProvi import org.lfdecentralizedtrust.splice.scan.metrics.ScanAppMetrics import org.lfdecentralizedtrust.splice.splitwell.metrics.SplitwellAppMetrics import org.lfdecentralizedtrust.splice.sv.metrics.SvAppMetrics +import org.lfdecentralizedtrust.splice.syncoperator.metrics.SyncOperatorAppMetrics import org.lfdecentralizedtrust.splice.validator.metrics.ValidatorAppMetrics import scala.collection.concurrent.TrieMap @@ -25,6 +26,7 @@ case class SpliceMetricsFactory( private val svs = TrieMap[String, SvAppMetrics]() private val scans = TrieMap[String, ScanAppMetrics]() private val splitwells = TrieMap[String, SplitwellAppMetrics]() + private val syncOperators = TrieMap[String, SyncOperatorAppMetrics]() def forValidator(name: String): ValidatorAppMetrics = { validators.getOrElseUpdate( @@ -78,4 +80,17 @@ case class SpliceMetricsFactory( }, ) } + + def forSyncOperator(name: String): SyncOperatorAppMetrics = { + syncOperators.getOrElseUpdate( + name, { + val metricsContext = MetricsContext("node_name" -> name, "node_type" -> "syncoperator") + new SyncOperatorAppMetrics( + metricsFactoryProvider.generateMetricsFactory(metricsContext), + storageHistograms, + loggerFactory, + ) + }, + ) + } } diff --git a/apps/app/src/test/resources/sync-operator-topology.conf b/apps/app/src/test/resources/sync-operator-topology.conf new file mode 100644 index 0000000000..528e01d5b9 --- /dev/null +++ b/apps/app/src/test/resources/sync-operator-topology.conf @@ -0,0 +1,32 @@ +# Adds a sync operator app to a base topology. +# +# It runs against the splitwell synchronizer's sequencer for now, until the test topologies +# include a dedicated synchronizer. +canton { + validator-apps { + splitwellValidator { + app-instances { + syncOperator { + service-user = "sync_operator_user" + dars = [] + } + } + } + } + + sync-operator-apps { + syncOperator { + include required("include/scan-client") + storage = ${_shared.storage} + storage.config.properties.databaseName = "splice_apps" + instance-lock-enabled = false + admin-api.address = 0.0.0.0 + admin-api.port = 5710 + participant-client = ${canton.validator-apps.splitwellValidator.participant-client} + operator-user = ${canton.validator-apps.splitwellValidator.app-instances.syncOperator.service-user} + # splitwellSequencer in simple-topology-canton.conf, which is a separate config so the + # port cannot be substituted here. + sequencer.admin-api.port = 5709 + } + } +} diff --git a/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala new file mode 100644 index 0000000000..6a71febe13 --- /dev/null +++ b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala @@ -0,0 +1,43 @@ +// 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 com.digitalasset.canton.SynchronizerAlias +import org.lfdecentralizedtrust.splice.integration.EnvironmentDefinition +import org.lfdecentralizedtrust.splice.integration.tests.SpliceTests.IntegrationTestWithIsolatedEnvironment + +class SyncOperatorIntegrationTest extends IntegrationTestWithIsolatedEnvironment { + + override def environmentDefinition: SpliceEnvironmentDefinition = + EnvironmentDefinition + .fromResources( + Seq("simple-topology-1sv.conf", "sync-operator-topology.conf"), + this.getClass.getSimpleName, + ) + .withStandardSetup + .withManualStart + + "sync operator app" should { + + "start and restart cleanly" in { implicit env => + initDsoWithSv1Only() + splitwellValidatorBackend.startSync() + + // startSync fails the test unless the app reports itself active within the timeout. + syncOperatorBackend.startSync() + + clue("it takes its synchronizer id from the sequencer it is configured with") { + val served = splitwellValidatorBackend.participantClientWithAdminToken.synchronizers + .id_of(SynchronizerAlias.tryCreate("splitwell")) + .logical + syncOperatorBackend.appState.store.key.synchronizerId shouldBe served + } + + syncOperatorBackend.stop() + syncOperatorBackend.is_running shouldBe false + + syncOperatorBackend.startSync() + } + } +} diff --git a/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/util/CommonAppInstanceReferences.scala b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/util/CommonAppInstanceReferences.scala index 33cb602bb0..f10d46512d 100644 --- a/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/util/CommonAppInstanceReferences.scala +++ b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/util/CommonAppInstanceReferences.scala @@ -10,6 +10,7 @@ import org.lfdecentralizedtrust.splice.console.{ SplitwellAppClientReference, SvAppBackendReference, SvAppClientReference, + SyncOperatorAppBackendReference, ValidatorAppBackendReference, ValidatorAppClientReference, WalletAppClientReference, @@ -297,6 +298,12 @@ trait CommonAppInstanceReferences { "providerSplitwellBackend" ) + def syncOperatorBackend(implicit + env: SpliceTestConsoleEnvironment + ): SyncOperatorAppBackendReference = syncop( + "syncOperator" + ) + def svb(name: String)(implicit env: SpliceTestConsoleEnvironment): SvAppBackendReference = env.svs.local .find(_.name == name) @@ -348,6 +355,13 @@ trait CommonAppInstanceReferences { .find(_.name == name) .getOrElse(sys.error(s"remote splitwell [$name] not configured")) + def syncop( + name: String + )(implicit env: SpliceTestConsoleEnvironment): SyncOperatorAppBackendReference = + env.syncOperators.local + .find(_.name == name) + .getOrElse(sys.error(s"local sync operator [$name] not configured")) + def scanb( name: String )(implicit env: SpliceTestConsoleEnvironment): ScanAppBackendReference = diff --git a/build.sbt b/build.sbt index 5804ad1e0e..3732c1275a 100644 --- a/build.sbt +++ b/build.sbt @@ -2389,6 +2389,7 @@ lazy val `apps-app`: Project = .dependsOn( `apps-common`, `apps-splitwell`, + `apps-syncoperator`, `apps-validator`, `apps-sv` % "compile->compile;test->test", `apps-scan`, diff --git a/test-full-class-names.log b/test-full-class-names.log index 759f3556ed..88aa72960e 100644 --- a/test-full-class-names.log +++ b/test-full-class-names.log @@ -53,6 +53,7 @@ org.lfdecentralizedtrust.splice.integration.tests.SvOnboardingViaNonFoundingSvIn org.lfdecentralizedtrust.splice.integration.tests.SvReconcileBftSequencingParametersIntegrationTest org.lfdecentralizedtrust.splice.integration.tests.SvReconcileSynchronizerConfigIntegrationTest org.lfdecentralizedtrust.splice.integration.tests.SvStateManagementIntegrationTest +org.lfdecentralizedtrust.splice.integration.tests.SyncOperatorIntegrationTest org.lfdecentralizedtrust.splice.integration.tests.TestTokenV2SettlementIntegrationTest org.lfdecentralizedtrust.splice.integration.tests.TokenStandardAllocationIntegrationTest org.lfdecentralizedtrust.splice.integration.tests.TokenStandardCliIntegrationTest From d7461ea87aed023ace51fc501ce8c3adb75abeb5 Mon Sep 17 00:00:00 2001 From: sadiq1971 Date: Wed, 26 Aug 2026 02:18:51 +0600 Subject: [PATCH 2/3] test: follow the splitwell shape for the sync operator smoke tests [ci] Splitwell and the validator both cover start and stop with a plain restart block and a separate liveness/readiness block, on an auto-started environment. Match that rather than hand-rolling a manual-start variant. Keeps one behaviour test of our own: that the operator takes its synchronizer id from the sequencer, which the log check cannot catch because a wrong id still starts cleanly. Signed-off-by: sadiq1971 --- .../tests/SyncOperatorIntegrationTest.scala | 35 ++++++++----------- 1 file changed, 15 insertions(+), 20 deletions(-) diff --git a/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala index 6a71febe13..6eef8a0d69 100644 --- a/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala +++ b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/tests/SyncOperatorIntegrationTest.scala @@ -5,9 +5,9 @@ package org.lfdecentralizedtrust.splice.integration.tests import com.digitalasset.canton.SynchronizerAlias import org.lfdecentralizedtrust.splice.integration.EnvironmentDefinition -import org.lfdecentralizedtrust.splice.integration.tests.SpliceTests.IntegrationTestWithIsolatedEnvironment +import org.lfdecentralizedtrust.splice.integration.tests.SpliceTests.IntegrationTest -class SyncOperatorIntegrationTest extends IntegrationTestWithIsolatedEnvironment { +class SyncOperatorIntegrationTest extends IntegrationTest { override def environmentDefinition: SpliceEnvironmentDefinition = EnvironmentDefinition @@ -16,28 +16,23 @@ class SyncOperatorIntegrationTest extends IntegrationTestWithIsolatedEnvironment this.getClass.getSimpleName, ) .withStandardSetup - .withManualStart - "sync operator app" should { - - "start and restart cleanly" in { implicit env => - initDsoWithSv1Only() - splitwellValidatorBackend.startSync() - - // startSync fails the test unless the app reports itself active within the timeout. + "sync operator" should { + "restart cleanly" in { implicit env => + syncOperatorBackend.stop() syncOperatorBackend.startSync() + } - clue("it takes its synchronizer id from the sequencer it is configured with") { - val served = splitwellValidatorBackend.participantClientWithAdminToken.synchronizers - .id_of(SynchronizerAlias.tryCreate("splitwell")) - .logical - syncOperatorBackend.appState.store.key.synchronizerId shouldBe served - } - - syncOperatorBackend.stop() - syncOperatorBackend.is_running shouldBe false + "report liveness and readiness" in { implicit env => + syncOperatorBackend.httpLive shouldBe true + syncOperatorBackend.httpReady shouldBe true + } - syncOperatorBackend.startSync() + "take its synchronizer id from the sequencer it is configured with" in { implicit env => + val served = splitwellValidatorBackend.participantClientWithAdminToken.synchronizers + .id_of(SynchronizerAlias.tryCreate("splitwell")) + .logical + syncOperatorBackend.appState.store.key.synchronizerId shouldBe served } } } From 5f1ff0ffecbd551f99375e1536c4c4e0f670e561 Mon Sep 17 00:00:00 2001 From: sadiq1971 Date: Wed, 26 Aug 2026 03:43:46 +0600 Subject: [PATCH 3/3] fix: port handling for the sync operator in test config [ci] Wait on the sync operator admin port in WaitForPorts, and bump both its participant and its sequencer admin API in bumpCantonPortsBy, so a test that composes the app with a port bump does not point at unbumped nodes. Move the admin port to 5115, alongside the other Splice app admin APIs, rather than the 57xx band that belongs to the splitwell Canton node. Note that the sequencer port is wall clock only. Signed-off-by: sadiq1971 --- .../splice/config/ConfigTransforms.scala | 9 +++++++++ apps/app/src/test/resources/sync-operator-topology.conf | 4 +--- .../splice/integration/plugins/WaitForPorts.scala | 1 + 3 files changed, 11 insertions(+), 3 deletions(-) diff --git a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala index 45781eab70..4932684a71 100644 --- a/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala +++ b/apps/app/src/main/scala/org/lfdecentralizedtrust/splice/config/ConfigTransforms.scala @@ -550,6 +550,15 @@ object ConfigTransforms { conf.focus(_.participantClient).modify(portTransform(bump, _)) else conf ), + updateAllSyncOperatorAppConfigs((name, conf) => + if (predicate(name)) + conf + .focus(_.participantClient) + .modify(portTransform(bump, _)) + .focus(_.sequencer.adminApi) + .modify(portTransform(bump, _)) + else conf + ), ) transforms.foldLeft((c: SpliceConfig) => c)((f, tf) => diff --git a/apps/app/src/test/resources/sync-operator-topology.conf b/apps/app/src/test/resources/sync-operator-topology.conf index 528e01d5b9..2fa277d35d 100644 --- a/apps/app/src/test/resources/sync-operator-topology.conf +++ b/apps/app/src/test/resources/sync-operator-topology.conf @@ -21,11 +21,9 @@ canton { storage.config.properties.databaseName = "splice_apps" instance-lock-enabled = false admin-api.address = 0.0.0.0 - admin-api.port = 5710 + admin-api.port = 5115 participant-client = ${canton.validator-apps.splitwellValidator.participant-client} operator-user = ${canton.validator-apps.splitwellValidator.app-instances.syncOperator.service-user} - # splitwellSequencer in simple-topology-canton.conf, which is a separate config so the - # port cannot be substituted here. sequencer.admin-api.port = 5709 } } diff --git a/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/plugins/WaitForPorts.scala b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/plugins/WaitForPorts.scala index cb7dc45ad1..97f556cf84 100644 --- a/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/plugins/WaitForPorts.scala +++ b/apps/app/src/test/scala/org/lfdecentralizedtrust/splice/integration/plugins/WaitForPorts.scala @@ -28,6 +28,7 @@ case class WaitForPorts(extraPortsToWaitFor: Seq[(String, Int)]) config.svApps.foreach(sv => waitForPort(sv._1, sv._2.adminApi.port.unwrap)) config.splitwellApps.foreach(sw => waitForPort(sw._1, sw._2.adminApi.port.unwrap)) config.scanApps.foreach(scan => waitForPort(scan._1, scan._2.adminApi.port.unwrap)) + config.syncOperatorApps.foreach(so => waitForPort(so._1, so._2.adminApi.port.unwrap)) extraPortsToWaitFor.foreach(p => waitForPort(InstanceName.tryCreate(p._1), p._2)) // Wait long enough that the prometheus server can start. config.monitoring.metrics.reporters.foreach {