Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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 @@ -15,8 +15,8 @@ import scala.util.control.NonFatal
private[ml] object FabricArtifactCleanup {
val Owner = "SynapseML OSS Fabric E2E"
private val RetentionSeconds = TimeUnit.HOURS.toSeconds(24)
private val ConfirmationAttempts = 31
private val ConfirmationDelaySeconds = 2L
private val ConfirmationAttempts = 11
private val ConfirmationDelayMillis = TimeUnit.SECONDS.toMillis(30)
private val Guid = "[0-9a-fA-F]{8}(?:-[0-9a-fA-F]{4}){3}-[0-9a-fA-F]{12}"
private val UniqueStore = "(Lakehouse|Warehouse)[0-9]{14}[0-9a-fA-F]{32}".r
private val RelationFields = Seq("artifactRelations", "datasetRelations", "dataflowRelations", "datamartRelations")
Expand Down Expand Up @@ -44,12 +44,11 @@ private[ml] object FabricArtifactCleanup {
Try(OffsetDateTime.parse(s).toInstant).getOrElse(LocalDateTime.parse(s).toInstant(ZoneOffset.UTC))
}

private def references(value: JsValue): Set[String] = value match {
private def references(value: JsValue, field: String): Set[String] = value match {
case JsString(s) if s.matches(Guid) => Set(java.util.UUID.fromString(s).toString)
case JsObject(fields) => fields.values.flatMap(references).toSet
case JsArray(values) => values.flatMap(references).toSet
case JsNull => Set.empty
case _ => Set.empty
case JsObject(fields) if fields.nonEmpty => fields.values.flatMap(references(_, field)).toSet
case JsArray(values) if values.nonEmpty => values.flatMap(references(_, field)).toSet
case _ => throw new IllegalArgumentException(s"Invalid or unknown $field relation metadata")
}

private def reference(value: JsValue, field: String): Set[String] = value.asJsObject.fields.get(field) match {
Expand All @@ -66,8 +65,7 @@ private[ml] object FabricArtifactCleanup {
fields.get(field) match {
case Some(JsNull) => Set.empty[String]
case Some(JsArray(values)) =>
require(values.forall(v => references(v).nonEmpty), s"Unknown $field metadata for $id")
values.flatMap(references).toSet
values.flatMap(references(_, field)).toSet
case _ => throw new IllegalArgumentException(s"Missing or invalid $field metadata for $id")
}
}.toSet
Expand Down Expand Up @@ -161,13 +159,13 @@ private[ml] object FabricArtifactCleanup {
store.kind == "Lakehouse" && i.kind == "SQLEndpoint" &&
i.references == Set(store.id) && i.expired(cutoff)

private def confirmAbsent(client: Client, id: String, pause: () => Unit): Unit = {
private def confirmAbsent(client: Client, id: String, pause: Long => Unit): Unit = {
@tailrec
def check(remaining: Int): Unit = {
if (index(client.inventory()).contains(id)) {
require(remaining > 1,
s"Fabric cleanup could not confirm deletion of $id after $ConfirmationAttempts reads")
pause()
pause(ConfirmationDelayMillis)
check(remaining - 1)
}
}
Expand All @@ -191,51 +189,59 @@ private[ml] object FabricArtifactCleanup {
managedEndpoint(i, candidate, cutoff) && neighbors(i.id, current) == Set(candidate.id)))
}

private def deleteAndConfirm(client: Client, id: String, pause: () => Unit, log: String => Unit): Unit = {
private def tryDeleteItem(client: Client, id: String, log: String => Unit): Option[Throwable] = {
try {
client.delete(id)
None
} catch {
case e: RuntimeException if Option(e.getMessage).exists(_.contains("PowerBIEntityNotFound")) =>
log(s"Fabric cleanup item $id was concurrently deleted; confirming absence")
None
case NonFatal(e) => Some(e)
}
confirmAbsent(client, id, pause)
}

def run(client: Client, now: Instant, dryRun: Boolean = false,
pause: () => Unit = () => TimeUnit.SECONDS.sleep(ConfirmationDelaySeconds),
pause: Long => Unit = millis => Thread.sleep(millis),
log: String => Unit = println): Vector[String] = {
val cutoff = now.minusSeconds(RetentionSeconds)
val initial = index(client.inventory())
val jobs = initial.values.filter(ownedJob).toVector.sortBy(_.id)
val stores = initial.values.filter(i => ownedStore(i, initial)).toVector.sortBy(_.id)
var deleted = Vector.empty[String]
var failures = Vector.empty[Throwable]
(jobs ++ stores).filter(_.expired(cutoff)).foreach { candidate =>
val current = index(client.inventory())
val expected = candidate.copy(references = candidate.references -- deleted)
val unchanged = current.get(candidate.id).contains(expected)
val safe = unchanged && (if (ownedJob(candidate)) {
safeJob(candidate, current, initial, client, cutoff)
} else {
failures.isEmpty && safeStore(candidate, current, cutoff)
})
if (safe) {
log(s"Fabric cleanup ${if (dryRun) "would delete" else "deleting"} ${candidate.kind} " +
s"${candidate.id} (${candidate.name}); created=${candidate.created}, updated=${candidate.updated}")
if (!dryRun) {
try {
deleteAndConfirm(client, candidate.id, pause, log)
deleted :+= candidate.id
log(s"Fabric cleanup confirmed deletion of ${candidate.id}")
} catch {
case NonFatal(e) =>
failures :+= e
log(s"Fabric cleanup failed for ${candidate.id}: ${e.getClass.getSimpleName}; retaining stores")
try {
(jobs ++ stores).filter(_.expired(cutoff)).foreach { candidate =>
val current = index(client.inventory())
val expected = candidate.copy(references = candidate.references -- deleted)
val unchanged = current.get(candidate.id).contains(expected)
val safe = unchanged && (if (ownedJob(candidate)) {
safeJob(candidate, current, initial, client, cutoff)
} else {
failures.isEmpty && safeStore(candidate, current, cutoff)
})
if (safe) {
log(s"Fabric cleanup ${if (dryRun) "would delete" else "deleting"} ${candidate.kind} " +
s"${candidate.id} (${candidate.name}); created=${candidate.created}, updated=${candidate.updated}")
if (!dryRun) {
tryDeleteItem(client, candidate.id, log) match {
case Some(e) =>
failures :+= e
log(s"Fabric cleanup failed for ${candidate.id}: ${e.getClass.getSimpleName}; retaining stores")
case None =>
confirmAbsent(client, candidate.id, pause)
deleted :+= candidate.id
log(s"Fabric cleanup confirmed deletion of ${candidate.id}")
}
}
} else {
log(s"Fabric cleanup retains ${candidate.id}: changed metadata, active jobs, schedules, or dependencies")
}
} else {
log(s"Fabric cleanup retains ${candidate.id}: changed metadata, active jobs, schedules, or dependencies")
}
} catch {
case NonFatal(e) =>
failures.filterNot(_ eq e).foreach(e.addSuppressed)
throw e
}
failures.headOption.foreach { first =>
failures.tail.filterNot(_ eq first).foreach(first.addSuppressed)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,18 @@ trait HasFabricNotebookTestConnection extends HasFabricOperationsConnection {
protected def cleanupTrackedArtifacts(): Unit = artifactTracker.cleanup()

protected def createTrackedStore(): String = trackArtifact(fabric.createStoreArtifact())

protected final def withFabricJobFailure[T](notebookName: String)(job: => T): T = {
try {
job
} catch {
case error: InterruptedException =>
Thread.currentThread().interrupt()
throw error
case NonFatal(t) =>
throw new RuntimeException(s"Job failed for $notebookName", t)
}
}
}

class FabricTestCleanup extends TestBase with HasFabricNotebookTestConnection {
Expand Down Expand Up @@ -117,14 +129,11 @@ class FabricSmokeTests extends TestBase with HasFabricNotebookTestConnection {
blocking {
Thread.sleep(10000) //scalastyle:ignore
}
try {
withFabricJobFailure(notebookName) {
val result = Await.ready(
fabric.monitorJob(artifactId, jobInstanceId),
Duration(fabric.timeoutInMillis.toLong, TimeUnit.MILLISECONDS)).value.get
assert(result.isSuccess)
} catch {
case t: Throwable =>
throw new RuntimeException(s"Job failed for $notebookName", t)
}
}
}
Expand Down Expand Up @@ -194,14 +203,8 @@ class FabricNotebookTests extends TestBase with HasFabricNotebookTestConnection
test(notebookName) {
ensureFabricPreflight()
val (future, submittedNotebookName) = futures(index)
try {
withFabricJobFailure(submittedNotebookName) {
Await.result(future, notebookTimeout)
} catch {
case error: InterruptedException =>
Thread.currentThread().interrupt()
throw error
case NonFatal(t) =>
throw new RuntimeException(s"Job failed for $submittedNotebookName", t)
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,12 +29,15 @@ private[nbtest] final class FabricTestArtifactTracker(deleteArtifact: String =>
deleteTrackedArtifact(artifactId)
artifactIds.remove(artifactId)
} catch {
case cleanupError: Throwable =>
case NonFatal(cleanupError) =>
failure match {
case Some(original) =>
if (original ne cleanupError) original.addSuppressed(cleanupError)
case None => throw cleanupError
}
case cleanupError: Throwable =>
failure.filterNot(_ eq cleanupError).foreach(cleanupError.addSuppressed)
throw cleanupError
}
}
}
Expand Down Expand Up @@ -62,7 +65,7 @@ private[nbtest] final class FabricTestArtifactTracker(deleteArtifact: String =>
}

failures.headOption.foreach { failure =>
failures.tail.foreach(failure.addSuppressed)
failures.tail.filterNot(_ eq failure).foreach(failure.addSuppressed)
throw failure
}
}
Expand Down
Loading
Loading