From 47c46d5ae07b64a4ce97a52dfee678de24c65851 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Tue, 14 Jan 2025 13:30:43 -0500 Subject: [PATCH 01/10] Explore parallelization --- build.sbt | 11 ++++++++++- .../sirius/uberstore/segmented/ParallelHelpers.scala | 3 +++ .../sirius/uberstore/segmented/ParallelHelpers.scala | 8 ++++++++ .../sirius/uberstore/segmented/ParallelHelpers.scala | 8 ++++++++ .../uberstore/segmented/SegmentedUberStore.scala | 10 +++++++++- .../xfinity/sirius/writeaheadlog/SiriusLog.scala | 10 +++++++++- 6 files changed, 47 insertions(+), 3 deletions(-) create mode 100644 src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala create mode 100644 src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala create mode 100644 src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala diff --git a/build.sbt b/build.sbt index 6c1b21b4..ce3590d2 100644 --- a/build.sbt +++ b/build.sbt @@ -18,7 +18,7 @@ name := "sirius" version := "2.4.0" -scalaVersion := "2.13.6" +scalaVersion := "2.12.14" crossScalaVersions := Seq("2.11.12", "2.12.14", "2.13.6") // NOTE: keep sync'd with .travis.yml organization := "com.comcast" @@ -50,6 +50,15 @@ libraryDependencies ++= { ) } +libraryDependencies ++= { + CrossVersion.partialVersion(scalaVersion.value) match { + case Some((2, major)) if major <= 12 => + Seq("org.scala-lang.modules" %% "scala-collection-compat" % "2.12.0") + case _ => + Seq("org.scala-lang.modules" %% "scala-parallel-collections" % "1.0.4") + } +} + // Set the artifact names. artifactName := { (scalaVersion: ScalaVersion, module: ModuleID, artifact: Artifact) => artifact.`type` match { diff --git a/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala new file mode 100644 index 00000000..bca816d1 --- /dev/null +++ b/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -0,0 +1,3 @@ +package com.comcast.xfinity.sirius.uberstore.segmented + +object ParallelHelpers { } diff --git a/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala new file mode 100644 index 00000000..ea446040 --- /dev/null +++ b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -0,0 +1,8 @@ +package com.comcast.xfinity.sirius.uberstore.segmented + +import scala.collection.parallel.immutable.ParSeq + +object ParallelHelpers { + def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par +} + diff --git a/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala new file mode 100644 index 00000000..12adda29 --- /dev/null +++ b/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -0,0 +1,8 @@ +package com.comcast.xfinity.sirius.uberstore.segmented + +import scala.collection.parallel.immutable.ParSeq +import scala.collection.parallel.CollectionConverters._ + +object ParallelHelpers { + def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par +} diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala index 5be76948..d80d78b9 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala @@ -16,7 +16,6 @@ package com.comcast.xfinity.sirius.uberstore.segmented import java.io.{File => JFile} - import better.files.File import com.comcast.xfinity.sirius.api.SiriusConfiguration import com.comcast.xfinity.sirius.api.impl.OrderedEvent @@ -136,6 +135,15 @@ class SegmentedUberStore private[segmented] (base: JFile, */ def getNextSeq = nextSeq + override def foreach[T](parallel: Boolean, fun: OrderedEvent => T): Unit = + if (parallel) parallelForeach(fun) else foldLeft(())((_, e) => fun(e)) + + private def parallelForeach[T](fun: OrderedEvent => T): Unit = { + ParallelHelpers.parallelize(liveDir :: readOnlyDirs).foldLeft(())( + (_, dir) => dir.foldLeftRange(0, Long.MaxValue)(())((_, e) => fun(e)) + ) + } + /** * @inheritdoc */ diff --git a/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala b/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala index b8b12070..c831a14b 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala @@ -34,7 +34,15 @@ trait SiriusLog { * * @param fun function to apply */ - def foreach[T](fun: OrderedEvent => T): Unit = foldLeft(())((_, e) => fun(e)) + def foreach[T](fun: OrderedEvent => T): Unit = foreach[T](parallel = false, fun) + + /** + * Apply fun to each entry in the log, optionally parallel and out of order + * + * @param parallel scan in parallel and out of order + * @param fun function to apply + */ + def foreach[T](parallel: Boolean, fun: OrderedEvent => T): Unit = foldLeft(())((_, e) => fun(e)) /** * Fold left across the log entries From 611366bc96d7d4ccb993e76db43c79ec545ddf0a Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Tue, 14 Jan 2025 13:37:06 -0500 Subject: [PATCH 02/10] Change package --- .../xfinity/sirius/uberstore/segmented/ParallelHelpers.scala | 2 +- .../xfinity/sirius/uberstore/segmented/ParallelHelpers.scala | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala index ea446040..c3d7b4c1 100644 --- a/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala +++ b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -1,6 +1,6 @@ package com.comcast.xfinity.sirius.uberstore.segmented -import scala.collection.parallel.immutable.ParSeq +import scala.collection.parallel.ParSeq object ParallelHelpers { def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par diff --git a/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala index 12adda29..523f15bb 100644 --- a/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala +++ b/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -1,7 +1,7 @@ package com.comcast.xfinity.sirius.uberstore.segmented -import scala.collection.parallel.immutable.ParSeq import scala.collection.parallel.CollectionConverters._ +import scala.collection.parallel.ParSeq object ParallelHelpers { def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par From ebead4b497c573fadab58c3e818de0b65b4d6914 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Tue, 14 Jan 2025 16:03:00 -0500 Subject: [PATCH 03/10] Refactor and add configuration options --- .../sirius/api/SiriusConfiguration.scala | 5 +++ .../sirius/api/impl/state/StateSup.scala | 31 +++++++++++-------- .../segmented/SegmentedUberStore.scala | 7 ++--- .../sirius/writeaheadlog/SiriusLog.scala | 7 ++--- 4 files changed, 28 insertions(+), 22 deletions(-) diff --git a/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala b/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala index 38af1b78..24ff4aa1 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala @@ -308,6 +308,11 @@ object SiriusConfiguration { * Maximum akka message size in KB. Default is 1024. Type is Integer. */ final val MAX_AKKA_MESSAGE_SIZE_KB = "sirius.akka.maximum-frame-size-kb" + + /** + * Whether or not to bootstrap the log in parallel, only applies to segmented uberstores + */ + final val LOG_PARALLEL_ENABLED = "sirius.log.parallel-enabled" } /** diff --git a/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala b/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala index cc9ad3c1..c9635279 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala @@ -119,22 +119,15 @@ class StateSup(requestHandler: RequestHandler, bootstrapTime = Some(0L) case _ => + val parallel = config.getProp(SiriusConfiguration.LOG_PARALLEL_ENABLED, default = false) val start = System.currentTimeMillis logger.info("Beginning SiriusLog replay at {}", start) requestHandler.onBootstrapStarting() - siriusLog.foreach( - orderedEvent => - try { - orderedEvent.request match { - case Put(key, body) => requestHandler.handlePut(orderedEvent.sequence, key, body) - case Delete(key) => requestHandler.handleDelete(orderedEvent.sequence, key) - } - } catch { - case rte: RuntimeException => - eventReplayFailureCount += 1 - logger.error("Exception replaying {}: {}", orderedEvent, rte) - } - ) + if (parallel) { + siriusLog.parallelForeach(bootstrapEvent) + } else { + siriusLog.foreach(bootstrapEvent) + } requestHandler.onBootstrapComplete() val totalBootstrapTime = System.currentTimeMillis - start bootstrapTime = Some(totalBootstrapTime) @@ -142,6 +135,18 @@ class StateSup(requestHandler: RequestHandler, } } + private def bootstrapEvent(orderedEvent : OrderedEvent): Unit = + try { + orderedEvent.request match { + case Put(key, body) => requestHandler.handlePut(orderedEvent.sequence, key, body) + case Delete(key) => requestHandler.handleDelete(orderedEvent.sequence, key) + } + } catch { + case rte: RuntimeException => + eventReplayFailureCount += 1 + logger.error("Exception replaying {}: {}", orderedEvent, rte) + } + trait StateInfoMBean { def getEventReplayFailureCount: Long def getBootstrapTime: String diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala index d80d78b9..510f1a53 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala @@ -135,11 +135,8 @@ class SegmentedUberStore private[segmented] (base: JFile, */ def getNextSeq = nextSeq - override def foreach[T](parallel: Boolean, fun: OrderedEvent => T): Unit = - if (parallel) parallelForeach(fun) else foldLeft(())((_, e) => fun(e)) - - private def parallelForeach[T](fun: OrderedEvent => T): Unit = { - ParallelHelpers.parallelize(liveDir :: readOnlyDirs).foldLeft(())( + override def parallelForeach[T](fun: OrderedEvent => T): Unit = { + ParallelHelpers.parallelize(readOnlyDirs ::: liveDir :: Nil).foldLeft(())( (_, dir) => dir.foldLeftRange(0, Long.MaxValue)(())((_, e) => fun(e)) ) } diff --git a/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala b/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala index c831a14b..53254349 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala @@ -34,15 +34,14 @@ trait SiriusLog { * * @param fun function to apply */ - def foreach[T](fun: OrderedEvent => T): Unit = foreach[T](parallel = false, fun) + def foreach[T](fun: OrderedEvent => T): Unit = foldLeft(())((_, e) => fun(e)) /** - * Apply fun to each entry in the log, optionally parallel and out of order + * Apply fun to each entry in the log in parallel and potentially out of order * - * @param parallel scan in parallel and out of order * @param fun function to apply */ - def foreach[T](parallel: Boolean, fun: OrderedEvent => T): Unit = foldLeft(())((_, e) => fun(e)) + def parallelForeach[T](fun: OrderedEvent => T): Unit = foreach[T](fun) /** * Fold left across the log entries From 657015f63bbb5b49867046f6d5c6e86f6b8d5398 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Tue, 14 Jan 2025 16:04:20 -0500 Subject: [PATCH 04/10] Fix Scala 2.11/2.12 parallel helpers --- .../sirius/uberstore/segmented/ParallelHelpers.scala | 6 +++++- .../sirius/uberstore/segmented/ParallelHelpers.scala | 1 - 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala index bca816d1..2af3dc89 100644 --- a/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala +++ b/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -1,3 +1,7 @@ package com.comcast.xfinity.sirius.uberstore.segmented -object ParallelHelpers { } +import scala.collection.parallel.ParSeq + +object ParallelHelpers { + def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par +} diff --git a/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala index c3d7b4c1..2af3dc89 100644 --- a/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala +++ b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -5,4 +5,3 @@ import scala.collection.parallel.ParSeq object ParallelHelpers { def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par } - From 30a72817fb68289c6c97347c0608f4b12d1d4eb8 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Tue, 14 Jan 2025 16:08:39 -0500 Subject: [PATCH 05/10] Remove hopefully unnecessary dependency and reset Scala version to 2.13 --- build.sbt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/build.sbt b/build.sbt index ce3590d2..30950a29 100644 --- a/build.sbt +++ b/build.sbt @@ -18,7 +18,7 @@ name := "sirius" version := "2.4.0" -scalaVersion := "2.12.14" +scalaVersion := "2.13.6" crossScalaVersions := Seq("2.11.12", "2.12.14", "2.13.6") // NOTE: keep sync'd with .travis.yml organization := "com.comcast" @@ -53,7 +53,7 @@ libraryDependencies ++= { libraryDependencies ++= { CrossVersion.partialVersion(scalaVersion.value) match { case Some((2, major)) if major <= 12 => - Seq("org.scala-lang.modules" %% "scala-collection-compat" % "2.12.0") + Seq() case _ => Seq("org.scala-lang.modules" %% "scala-parallel-collections" % "1.0.4") } From 3b70ec18883c05314b51f8750e71a7055a1bb935 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Tue, 14 Jan 2025 17:31:58 -0500 Subject: [PATCH 06/10] Quick sanity check unit test --- .../segmented/SegmentedUberStoreTest.scala | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala index acee7ee5..94b19f1c 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala @@ -16,6 +16,8 @@ package com.comcast.xfinity.sirius.uberstore.segmented +import scala.collection.concurrent._ + import com.comcast.xfinity.sirius.NiceTest import java.io.{File => JFile} @@ -195,6 +197,21 @@ class SegmentedUberStoreTest extends NiceTest { } } + describe("parallelForeach") { + it("should bootstrap the uberstore in parallel") { + createPopulatedSegment(dir, "1", Range.inclusive(1, 3).toList, isApplied = true) + createPopulatedSegment(dir, "2", Range.inclusive(4, 6).toList, isApplied = true) + createPopulatedSegment(dir, "3", Range.inclusive(7, 9).toList, isApplied = true) + val config = new SiriusConfiguration + config.setProp(SiriusConfiguration.LOG_PARALLEL_ENABLED, true) + uberstore = SegmentedUberStore(dir.getAbsolutePath, config) + val map = new TrieMap[Long, SiriusRequest]() + uberstore.parallelForeach(event => map.put(event.sequence, event.request)) + + assert(map.size == 9) + } + } + describe("close") { it("should close all of the associated uberdirs") { uberstore.close() From 05c1f754c9ff8e6672eaa96a53824ca45bb82a26 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Tue, 14 Jan 2025 22:00:21 -0500 Subject: [PATCH 07/10] Fix parallelization --- .../sirius/uberstore/segmented/SegmentedUberStore.scala | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala index 510f1a53..9fbbc2f9 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala @@ -136,9 +136,8 @@ class SegmentedUberStore private[segmented] (base: JFile, def getNextSeq = nextSeq override def parallelForeach[T](fun: OrderedEvent => T): Unit = { - ParallelHelpers.parallelize(readOnlyDirs ::: liveDir :: Nil).foldLeft(())( - (_, dir) => dir.foldLeftRange(0, Long.MaxValue)(())((_, e) => fun(e)) - ) + ParallelHelpers.parallelize(readOnlyDirs :+ liveDir) + .foreach(_.foldLeftRange(0, Long.MaxValue)(())((_, e) => fun(e))) } /** From da90fdad31ff177597f2ecedff1bfa1d294e8163 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Wed, 15 Jan 2025 09:54:26 -0500 Subject: [PATCH 08/10] Add SiriusConfiguration option to disable checksum validation on reads --- .../xfinity/sirius/api/SiriusConfiguration.scala | 5 +++++ .../comcast/xfinity/sirius/uberstore/UberPair.scala | 4 ++-- .../comcast/xfinity/sirius/uberstore/UberStore.scala | 4 +++- .../xfinity/sirius/uberstore/common/Checksummer.scala | 10 ++++++++++ .../uberstore/common/SkipValidationChecksummer.scala | 8 ++++++++ .../xfinity/sirius/uberstore/data/UberDataFile.scala | 9 ++++++--- .../sirius/uberstore/data/UberStoreBinaryFileOps.scala | 2 +- .../xfinity/sirius/uberstore/segmented/Segment.scala | 8 ++++---- .../uberstore/segmented/SegmentedUberStore.scala | 4 +++- 9 files changed, 42 insertions(+), 12 deletions(-) create mode 100644 src/main/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummer.scala diff --git a/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala b/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala index 24ff4aa1..e3f5b857 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala @@ -313,6 +313,11 @@ object SiriusConfiguration { * Whether or not to bootstrap the log in parallel, only applies to segmented uberstores */ final val LOG_PARALLEL_ENABLED = "sirius.log.parallel-enabled" + + /** + * Whether to skip checksum validation when reading events from the log + */ + final val LOG_SKIP_CHECKSUM_VALIDATION = "sirius.log.skip-checksum-validation" } /** diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala index 439549d1..b8e9d691 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala @@ -30,9 +30,9 @@ object UberPair { * * @return an instantiated UberStoreFilePair */ - def apply(baseDir: String, startingSeq: Long, fileHandleFactory: UberDataFileHandleFactory): UberPair = { + def apply(baseDir: String, startingSeq: Long, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean): UberPair = { val baseName = "%s/%s".format(baseDir, startingSeq) - val dataFile = UberDataFile("%s.data".format(baseName), fileHandleFactory) + val dataFile = UberDataFile("%s.data".format(baseName), fileHandleFactory, validateChecksum) val index = DiskOnlySeqIndex("%s.index".format(baseName)) repairIndex(index, dataFile) new UberPair(dataFile, index) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala index 10e026cf..5b9c0b23 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala @@ -39,8 +39,10 @@ object UberStore { } val fileHandleFactory = UberDataFileHandleFactory(siriusConfig) + val skipChecksumValidation = siriusConfig.getProp(SiriusConfiguration.LOG_SKIP_CHECKSUM_VALIDATION, false) + val validateChecksum = !skipChecksumValidation - new UberStore(baseDir, UberPair(baseDir, 1L, fileHandleFactory)) + new UberStore(baseDir, UberPair(baseDir, 1L, fileHandleFactory, validateChecksum)) } /** diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala index e464623a..d6e63392 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala @@ -19,6 +19,16 @@ package com.comcast.xfinity.sirius.uberstore.common * Trait supplying checksumming capabilities */ trait Checksummer { + /** + * Determines if the array of bytes has a calculated checksum + * that matches the provided checksum + * + * @param chksum the checksum + * @param bytes Array[Byte] to checksum + * @return if the checksum of the bytes matches the provided checksum + */ + def validate(bytes: Array[Byte], chksum: Long): Boolean = + checksum(bytes) == chksum /** * Given an array of bytes will calculate a Long checksum diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummer.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummer.scala new file mode 100644 index 00000000..bcb8974d --- /dev/null +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummer.scala @@ -0,0 +1,8 @@ +package com.comcast.xfinity.sirius.uberstore.common + +trait SkipValidationChecksummer extends Checksummer { + /** + * Skips calculating and validating the checksum on the Array[Byte] + */ + override def validate(bytes: Array[Byte], chksum: Long): Boolean = true +} diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala index f77bfaa6..42526b80 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala @@ -16,7 +16,7 @@ package com.comcast.xfinity.sirius.uberstore.data import com.comcast.xfinity.sirius.api.impl.OrderedEvent -import com.comcast.xfinity.sirius.uberstore.common.Fnv1aChecksummer +import com.comcast.xfinity.sirius.uberstore.common.{Fnv1aChecksummer, SkipValidationChecksummer} import scala.annotation.tailrec @@ -35,8 +35,11 @@ object UberDataFile { * * @return fully constructed UberDataFile */ - def apply(dataFileName: String, fileHandleFactory: UberDataFileHandleFactory): UberDataFile = { - val fileOps = new UberStoreBinaryFileOps with Fnv1aChecksummer + def apply(dataFileName: String, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean): UberDataFile = { + val fileOps = if (validateChecksum) + new UberStoreBinaryFileOps with Fnv1aChecksummer + else + new UberStoreBinaryFileOps with Fnv1aChecksummer with SkipValidationChecksummer val codec = new BinaryEventCodec new UberDataFile(dataFileName, fileHandleFactory, fileOps, codec) } diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala index b8819c73..575fa41e 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala @@ -53,7 +53,7 @@ class UberStoreBinaryFileOps extends UberStoreFileOps { } else { val (bodyLen, chksum) = readHeader(readHandle) val body = readBody(readHandle, bodyLen) - if (chksum == checksum(body)) { + if (validate(body, chksum)) { Some(body) // [that i used to know | to love] } else { throw new IllegalStateException("File corrupted at offset " + readHandle.offset()) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala index d7497e33..2966e991 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala @@ -30,8 +30,8 @@ object Segment { * * @return an Segment instance, fully repaired and usable */ - def apply(base: File, name: String, fileHandleFactory: UberDataFileHandleFactory): Segment = { - apply(new File(base, name), fileHandleFactory) + def apply(base: File, name: String, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean = true): Segment = { + apply(new File(base, name), fileHandleFactory, validateChecksum) } /** @@ -42,7 +42,7 @@ object Segment { * * @return an Segment instance, fully repaired and usable */ - def apply(location: File, fileHandleFactory: UberDataFileHandleFactory): Segment = { + def apply(location: File, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean = true): Segment = { location.mkdirs() val dataFile = new File(location, "data") @@ -50,7 +50,7 @@ object Segment { val compactionFlagFile = new File(location, "keys-collected") val internalCompactionFlagFile = new File(location, "internally-compacted") - val data = UberDataFile(dataFile.getAbsolutePath, fileHandleFactory) + val data = UberDataFile(dataFile.getAbsolutePath, fileHandleFactory, validateChecksum) val index = DiskOnlySeqIndex(indexFile.getAbsolutePath) val compactionFlag = FlagFile(compactionFlagFile.getAbsolutePath) val internalCompactionFlag = FlagFile(internalCompactionFlagFile.getAbsolutePath) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala index 9fbbc2f9..4b8b6fe7 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala @@ -86,8 +86,10 @@ object SegmentedUberStore { val MAX_EVENTS_PER_SEGMENT = siriusConfig.getProp(SiriusConfiguration.LOG_EVENTS_PER_SEGMENT, 1000000L) val fileHandleFactory = UberDataFileHandleFactory(siriusConfig) + val skipChecksumValidation = siriusConfig.getProp(SiriusConfiguration.LOG_SKIP_CHECKSUM_VALIDATION, false) + val validateChecksum = !skipChecksumValidation - def buildSegment(location: JFile) = Segment(location, fileHandleFactory) + def buildSegment(location: JFile) = Segment(location, fileHandleFactory, validateChecksum) val segmentedCompactor = SegmentedCompactor(siriusConfig, buildSegment) new SegmentedUberStore(new JFile(base), MAX_EVENTS_PER_SEGMENT, segmentedCompactor, buildSegment) From 2900513e9e2167206d583d10ee6bbf937b46de51 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Wed, 15 Jan 2025 10:12:08 -0500 Subject: [PATCH 09/10] Refactor and remove unused overload --- .../xfinity/sirius/uberstore/segmented/Segment.scala | 12 ------------ .../sirius/uberstore/segmented/SegmentTest.scala | 2 +- .../uberstore/segmented/SegmentedCompactorTest.scala | 2 +- .../uberstore/segmented/SegmentedUberStoreTest.scala | 2 +- 4 files changed, 3 insertions(+), 15 deletions(-) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala index 2966e991..a5d8f6f0 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala @@ -22,18 +22,6 @@ import java.io.File object Segment { - /** - * Create an Segment based in baseDir named "name". - * - * @param base directory containing the Segment - * @param name the name of this dir - * - * @return an Segment instance, fully repaired and usable - */ - def apply(base: File, name: String, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean = true): Segment = { - apply(new File(base, name), fileHandleFactory, validateChecksum) - } - /** * Create an Segment at the specified location. * diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala index 9b9c92bd..5001e9d8 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala @@ -58,7 +58,7 @@ class SegmentTest extends NiceTest with BeforeAndAfterAll { val fileHandleFactory: UberDataFileHandleFactory = RandomAccessFileHandleFactory def buildSegment(base: JFile, name: String) = - Segment(base, name, fileHandleFactory) + Segment(new JFile(base, name), fileHandleFactory) override def afterAll(): Unit = { File(tempDir.getPath).delete() diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala index 78177724..807a741f 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala @@ -51,7 +51,7 @@ class SegmentedCompactorTest extends NiceTest with BeforeAndAfterAll { } def buildSegment(base: JFile, name: String): Segment = { - Segment(base, name, fileHandleFactory) + Segment(new JFile(base, name), fileHandleFactory) } def buildSegment(fullPath: String): Segment = { diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala index 94b19f1c..9fe9bb5f 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala @@ -79,7 +79,7 @@ class SegmentedUberStoreTest extends NiceTest { } def buildSegment(base: JFile, name: String): Segment = - Segment(base, name, fileHandleFactory) + Segment(new JFile(base, name), fileHandleFactory) def buildSegment(location: JFile): Segment = Segment(location, fileHandleFactory) From 820c1aae1995d7a962180653a8f0b606d51756e6 Mon Sep 17 00:00:00 2001 From: "Spindler, Justin" Date: Wed, 15 Jan 2025 10:28:23 -0500 Subject: [PATCH 10/10] Tests --- .../common/Fnv1aChecksummerTest.scala | 5 +++++ .../SkipValidationChecksummerTest.scala | 19 +++++++++++++++++++ 2 files changed, 24 insertions(+) create mode 100644 src/test/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummerTest.scala diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala index 34fbfc86..193beb72 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala @@ -32,4 +32,9 @@ class Fnv1aChecksummerTest extends NiceTest { val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes assert(-2758076559093427003L === underTest.checksum(referenceBytes)) } + + it ("validates the checksum") { + val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes + assert(true === underTest.validate(referenceBytes, -2758076559093427003L)) + } } diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummerTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummerTest.scala new file mode 100644 index 00000000..e6341e84 --- /dev/null +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummerTest.scala @@ -0,0 +1,19 @@ +package com.comcast.xfinity.sirius.uberstore.common + +import com.comcast.xfinity.sirius.NiceTest + +class SkipValidationChecksummerTest extends NiceTest { + val underTest = new Object with Fnv1aChecksummer with SkipValidationChecksummer + + it ("must match the reference impl") { + // from: + // http://trac.tools.ietf.org/wg/tls/draft-ietf-tls-cached-info/draft-ietf-tls-cached-info-06-from-05.diff.txt + val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes + assert(-2758076559093427003L === underTest.checksum(referenceBytes)) + } + + it ("skips validating the checksum") { + val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes + assert(true === underTest.validate(referenceBytes, 0L)) + } +}