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
6 changes: 4 additions & 2 deletions src/main/kotlin/io/cloudshiftdev/mavensync/MavenHttpClient.kt
Original file line number Diff line number Diff line change
Expand Up @@ -30,13 +30,15 @@ internal class MavenHttpClient(
httpClient.close()
}

internal suspend fun upload(url: Url, file: File) {
internal suspend fun upload(url: Url, file: File): Long {
val resp =
httpClient.put(url) {
contentType(ContentType.Application.OctetStream)
setBody(LocalFileContent(file))
}
logger.info { "Uploaded $url: status=${resp.status} size=${file.length()}" }
val size = file.length()
logger.info { "Uploaded $url: status=${resp.status} size=$size" }
return size
}

internal suspend fun upload(url: Url, content: String) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,9 @@ internal interface MavenHttpRepository : AutoCloseable {
includeSignatures: Boolean,
): List<ArtifactVersionAsset>

suspend fun copyAsset(asset: ArtifactVersionAsset, targetRepository: MavenHttpRepository)
suspend fun copyAsset(asset: ArtifactVersionAsset, targetRepository: MavenHttpRepository): Long

suspend fun uploadAsset(asset: ArtifactVersionAsset, file: Path)
suspend fun uploadAsset(asset: ArtifactVersionAsset, file: Path): Long

suspend fun releaseVersion(coordinates: Coordinates)

Expand Down Expand Up @@ -100,14 +100,14 @@ internal class DefaultMavenHttpRepository(
override suspend fun copyAsset(
asset: ArtifactVersionAsset,
targetRepository: MavenHttpRepository,
) {
mavenHttpClient.download(url(asset.coordinates, asset.name)) { _, file ->
): Long {
return mavenHttpClient.download(url(asset.coordinates, asset.name)) { _, file ->
targetRepository.uploadAsset(asset, file.toPath())
}
}

override suspend fun uploadAsset(asset: ArtifactVersionAsset, file: Path) {
mavenHttpClient.upload(url(asset.coordinates, asset.name), file.toFile())
override suspend fun uploadAsset(asset: ArtifactVersionAsset, file: Path): Long {
return mavenHttpClient.upload(url(asset.coordinates, asset.name), file.toFile())
}

override suspend fun releaseVersion(coordinates: Coordinates) {
Expand Down
55 changes: 39 additions & 16 deletions src/main/kotlin/io/cloudshiftdev/mavensync/MavenSyncEngine.kt
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
package io.cloudshiftdev.mavensync

import io.github.oshai.kotlinlogging.KotlinLogging
import kotlin.coroutines.cancellation.CancellationException
import kotlin.time.TimeSource
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.channels.consumeEach
import kotlinx.coroutines.coroutineScope
Expand All @@ -15,6 +17,7 @@ internal class MavenSyncEngine(
private val source: MavenHttpRepository,
private val target: MavenHttpRepository,
private val options: SyncOptions,
private val metrics: SyncMetrics,
) {
@OptIn(ExperimentalCoroutinesApi::class)
suspend fun sync() = coroutineScope {
Expand All @@ -32,6 +35,9 @@ internal class MavenSyncEngine(
val targetMetadata = target.queryArtifactMetadata(metadata.group, metadata.artifact)
val sourceVersions = metadata.artifactVersions.toSet()
val targetVersions = targetMetadata.artifactVersions.toSet()
val inSyncVersions = sourceVersions intersect targetVersions
metrics.recordInSync(metadata.group, metadata.artifact, inSyncVersions)

val missingVersions = sourceVersions - targetVersions
if (missingVersions.isEmpty()) {
logger.debug { "No missing versions for ${metadata.group}:${metadata.artifact}" }
Expand All @@ -42,21 +48,38 @@ internal class MavenSyncEngine(
}
missingVersions
.map { Coordinates(metadata.group, metadata.artifact, it) }
.forEach { coordinates ->
val assets =
source.listArtifactVersionAssets(
coordinates,
options.transferChecksums,
options.transferSignatures,
)

if (assets.isNotEmpty()) {
assets.forEach { asset -> source.copyAsset(asset, target) }

target.releaseVersion(coordinates)

delay(options.downloadDelay)
}
}
.forEach { coordinates -> syncVersion(coordinates) }
}

private suspend fun syncVersion(coordinates: Coordinates) {
val mark = TimeSource.Monotonic.markNow()
try {
val assets =
source.listArtifactVersionAssets(
coordinates,
options.transferChecksums,
options.transferSignatures,
)

if (assets.isEmpty()) return

var bytes = 0L
assets.forEach { asset -> bytes += source.copyAsset(asset, target) }
target.releaseVersion(coordinates)
metrics.recordSynced(
coordinates = coordinates,
assetCount = assets.size,
bytes = bytes,
duration = mark.elapsedNow(),
)

delay(options.downloadDelay)
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
val msg = e.message ?: e.toString()
logger.error(e) { "Failed to sync $coordinates: $msg" }
metrics.recordFailure(coordinates, msg, mark.elapsedNow())
}
}
}
13 changes: 10 additions & 3 deletions src/main/kotlin/io/cloudshiftdev/mavensync/MavenSyncMain.kt
Original file line number Diff line number Diff line change
Expand Up @@ -27,10 +27,17 @@ public suspend fun main(args: Array<String>) {

logger.info { "Effective configuration: $config" }

config.source.toMavenHttpRepository().use { source ->
config.target.toMavenHttpRepository().use { target ->
MavenSyncEngine(source, target, config.toSyncOptions()).sync()
val metrics = SyncMetrics()
try {
config.source.toMavenHttpRepository().use { source ->
config.target.toMavenHttpRepository().use { target ->
MavenSyncEngine(source, target, config.toSyncOptions(), metrics).sync()
}
}
} finally {
val report = metrics.snapshot()
logger.info { "\n" + report.renderDetailed() }
logger.info { "\n" + report.renderSummary() }
}
}

Expand Down
154 changes: 154 additions & 0 deletions src/main/kotlin/io/cloudshiftdev/mavensync/SyncMetrics.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
package io.cloudshiftdev.mavensync

import java.util.concurrent.ConcurrentHashMap
import kotlin.time.Duration
import kotlin.time.TimeSource

internal class SyncMetrics {
private val mark = TimeSource.Monotonic.markNow()
private val artifacts = ConcurrentHashMap<Pair<Group, Artifact>, Entry>()

fun recordInSync(group: Group, artifact: Artifact, versions: Collection<ArtifactVersion>) {
if (versions.isEmpty()) return
entry(group, artifact).inSync.addAll(versions)
}

fun recordSynced(coordinates: Coordinates, assetCount: Int, bytes: Long, duration: Duration) {
entry(coordinates.group, coordinates.artifact)
.synced
.add(
VersionResult.Success(
version = coordinates.artifactVersion,
assetCount = assetCount,
bytes = bytes,
duration = duration,
)
)
}

fun recordFailure(coordinates: Coordinates, error: String, duration: Duration) {
entry(coordinates.group, coordinates.artifact)
.failed
.add(
VersionResult.Failure(
version = coordinates.artifactVersion,
error = error,
duration = duration,
)
)
}

fun snapshot(): SyncReport {
val artifactMetrics =
artifacts.values
.map { e ->
ArtifactMetrics(
group = e.group,
artifact = e.artifact,
inSync = e.inSync.toList(),
synced = e.synced.toList(),
failed = e.failed.toList(),
)
}
.sortedWith(compareBy({ it.group.value }, { it.artifact.value }))
return SyncReport(totalDuration = mark.elapsedNow(), artifacts = artifactMetrics)
}

private fun entry(group: Group, artifact: Artifact): Entry =
artifacts.computeIfAbsent(group to artifact) { Entry(group, artifact) }

private class Entry(val group: Group, val artifact: Artifact) {
val inSync: MutableList<ArtifactVersion> = mutableListOf()
val synced: MutableList<VersionResult.Success> = mutableListOf()
val failed: MutableList<VersionResult.Failure> = mutableListOf()
}
}

internal data class ArtifactMetrics(
val group: Group,
val artifact: Artifact,
val inSync: List<ArtifactVersion>,
val synced: List<VersionResult.Success>,
val failed: List<VersionResult.Failure>,
)

internal sealed interface VersionResult {
val version: ArtifactVersion
val duration: Duration

data class Success(
override val version: ArtifactVersion,
val assetCount: Int,
val bytes: Long,
override val duration: Duration,
) : VersionResult

data class Failure(
override val version: ArtifactVersion,
val error: String,
override val duration: Duration,
) : VersionResult
}

internal data class SyncReport(val totalDuration: Duration, val artifacts: List<ArtifactMetrics>) {
val inSyncTotal: Int = artifacts.sumOf { it.inSync.size }
val syncedTotal: Int = artifacts.sumOf { it.synced.size }
val failedTotal: Int = artifacts.sumOf { it.failed.size }
val assetsTotal: Int = artifacts.sumOf { a -> a.synced.sumOf { it.assetCount } }
val bytesTotal: Long = artifacts.sumOf { a -> a.synced.sumOf { it.bytes } }

fun renderDetailed(): String =
buildString {
appendLine("========== SYNC METRICS (detailed) ==========")
if (artifacts.isEmpty()) {
appendLine("(no artifacts processed)")
return@buildString
}
artifacts.forEach { a ->
appendLine("${a.group.value}:${a.artifact.value}")
if (a.inSync.isNotEmpty()) {
appendLine(
" in sync (${a.inSync.size}): ${a.inSync.joinToString(", ") { it.value }}"
)
}
if (a.synced.isNotEmpty()) {
appendLine(" synced (${a.synced.size}):")
a.synced.forEach { s ->
appendLine(
" ${s.version.value} — ${s.assetCount} assets, ${formatBytes(s.bytes)}, ${s.duration}"
)
}
}
if (a.failed.isNotEmpty()) {
appendLine(" failed (${a.failed.size}):")
a.failed.forEach { f ->
appendLine(" ${f.version.value} — ${f.error} (${f.duration})")
}
}
}
}
.trimEnd()

fun renderSummary(): String = buildString {
appendLine("========== SYNC METRICS (summary) ==========")
appendLine("artifacts: ${artifacts.size}")
appendLine("versions in sync: $inSyncTotal")
appendLine("versions synced: $syncedTotal")
appendLine("versions failed: $failedTotal")
appendLine("assets copied: $assetsTotal")
appendLine("bytes transferred: ${formatBytes(bytesTotal)}")
append("duration: $totalDuration")
}
}

private fun formatBytes(bytes: Long): String {
if (bytes < 1024) return "$bytes B"
val units = listOf("KB", "MB", "GB", "TB")
var value = bytes.toDouble() / 1024.0
var idx = 0
while (value >= 1024.0 && idx < units.lastIndex) {
value /= 1024.0
idx++
}
return "%.1f %s".format(value, units[idx])
}
Original file line number Diff line number Diff line change
Expand Up @@ -35,15 +35,19 @@ internal class FakeMavenHttpRepository(private val label: String) : MavenHttpRep
return assets[coordinates].orEmpty()
}

var copyAssetBehavior: suspend (ArtifactVersionAsset) -> Long = { 0L }

override suspend fun copyAsset(
asset: ArtifactVersionAsset,
targetRepository: MavenHttpRepository,
) {
): Long {
copyCalls += asset to targetRepository
return copyAssetBehavior(asset)
}

override suspend fun uploadAsset(asset: ArtifactVersionAsset, file: Path) {
override suspend fun uploadAsset(asset: ArtifactVersionAsset, file: Path): Long {
uploadCalls += asset to file
return 0L
}

override suspend fun releaseVersion(coordinates: Coordinates) {
Expand Down
Loading
Loading