use http2 client for JPS cache (part 2)

GitOrigin-RevId: 09eac7e7f1114be98de5d686026cf454e1c69cc4
This commit is contained in:
Vladimir Krivosheev
2024-09-06 19:25:22 +00:00
committed by intellij-monorepo-bot
parent bfe323ec57
commit 29ef8fa3e3
13 changed files with 366 additions and 379 deletions
@@ -171,6 +171,9 @@ private fun writeAttributesAsHumanReadable(attributes: Attributes, sb: StringBui
if ((k.key == "modulesWithSearchableOptions" || v is List<*>) && (v as List<*>).size > 16) {
sb.append("…")
}
else if (v is List<*> && v.isEmpty()) {
sb.append("<empty list>")
}
else if (v is Iterable<*>) {
for (s in v) {
writeValueAsHumanReadable(s as String, sb)
@@ -16,18 +16,45 @@ import io.opentelemetry.sdk.OpenTelemetrySdk
import io.opentelemetry.sdk.resources.Resource
import io.opentelemetry.sdk.trace.SdkTracerProvider
import io.opentelemetry.sdk.trace.data.SpanData
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.Job
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.*
import org.jetbrains.intellij.build.dependencies.BuildDependenciesDownloader
import java.nio.file.Path
import java.util.concurrent.atomic.AtomicBoolean
import java.util.concurrent.atomic.AtomicReference
import kotlin.time.Duration.Companion.seconds
@Suppress("SSBasedInspection")
var traceManagerInitializer: () -> Pair<Tracer, BatchSpanProcessor> = {
// don't use JaegerJsonSpanExporter - not needed for clients, should be enabled only if needed to avoid writing a ~500KB JSON file
fun withTracer(block: suspend () -> Unit): Unit = runBlocking(Dispatchers.Default) {
val batchSpanProcessorScope = CoroutineScope(SupervisorJob(parent = coroutineContext.job)) + CoroutineName("BatchSpanProcessor")
@Suppress("ReplaceJavaStaticMethodWithKotlinAnalog")
val spanProcessor = BatchSpanProcessor(
coroutineScope = batchSpanProcessorScope,
spanExporters = java.util.List.of(ConsoleSpanExporter()),
scheduleDelay = 10.seconds,
)
try {
val tracerProvider = SdkTracerProvider.builder()
.addSpanProcessor(spanProcessor)
.setResource(Resource.create(Attributes.of(AttributeKey.stringKey("service.name"), "builder")))
.build()
traceManagerInitializer = {
val openTelemetry = OpenTelemetrySdk.builder()
.setTracerProvider(tracerProvider)
.build()
val tracer = openTelemetry.getTracer("build-script")
BuildDependenciesDownloader.TRACER = tracer
tracer to spanProcessor
}
block()
}
finally {
batchSpanProcessorScope.cancel()
traceManagerInitializer = { throw IllegalStateException("already built") }
}
}
private var traceManagerInitializer: () -> Pair<Tracer, BatchSpanProcessor> = {
val batchSpanProcessor = BatchSpanProcessor(
scheduleDelay = 10.seconds,
coroutineScope = CoroutineScope(Job()),
@@ -56,7 +83,7 @@ object TraceManager {
batchSpanProcessor = config.second
}
fun setTracer(tracer: Tracer){
fun setTracer(tracer: Tracer) {
this.tracer = tracer
}
@@ -4,24 +4,12 @@
package org.jetbrains.intellij.build.devServer
import com.intellij.openapi.application.PathManager
import com.intellij.platform.diagnostic.telemetry.exporters.BatchSpanProcessor
import com.intellij.platform.util.coroutines.childScope
import com.intellij.util.SystemProperties
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.sdk.OpenTelemetrySdk
import io.opentelemetry.sdk.resources.Resource
import io.opentelemetry.sdk.trace.SdkTracerProvider
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.cancel
import kotlinx.coroutines.runBlocking
import org.jetbrains.intellij.build.telemetry.ConsoleSpanExporter
import org.jetbrains.intellij.build.dependencies.BuildDependenciesDownloader
import org.jetbrains.intellij.build.dev.BuildRequest
import org.jetbrains.intellij.build.dev.buildProductInProcess
import org.jetbrains.intellij.build.dev.getAdditionalPluginMainModules
import org.jetbrains.intellij.build.dev.getIdeSystemProperties
import org.jetbrains.intellij.build.telemetry.traceManagerInitializer
import org.jetbrains.intellij.build.telemetry.withTracer
import java.io.File
import java.nio.file.Path
@@ -33,52 +21,29 @@ fun buildDevMain(): Collection<Path> {
var homePath: String? = null
var newClassPath: Collection<Path>? = null
runBlocking(Dispatchers.Default) {
val batchSpanProcessorScope = childScope("BatchSpanProcessor")
val spanProcessor = BatchSpanProcessor(coroutineScope = batchSpanProcessorScope, spanExporters = java.util.List.of(ConsoleSpanExporter()))
withTracer {
buildProductInProcess(
BuildRequest(
platformPrefix = System.getProperty("idea.platform.prefix", "idea"),
additionalModules = getAdditionalPluginMainModules(),
projectDir = ideaProjectRoot,
keepHttpClient = false,
platformClassPathConsumer = { classPath, runDir ->
newClassPath = classPath
homePath = runDir.toString().replace(File.separator, "/")
val tracerProvider = SdkTracerProvider.builder()
.addSpanProcessor(spanProcessor)
.setResource(Resource.create(Attributes.of(AttributeKey.stringKey("service.name"), "builder")))
.build()
try {
// don't use JaegerJsonSpanExporter - not needed for clients, should be enabled only if needed to avoid writing a ~500KB JSON file
traceManagerInitializer = {
val openTelemetry = OpenTelemetrySdk.builder()
.setTracerProvider(tracerProvider)
.build()
val tracer = openTelemetry.getTracer("build-script")
BuildDependenciesDownloader.TRACER = tracer
tracer to spanProcessor
}
buildProductInProcess(
BuildRequest(
platformPrefix = System.getProperty("idea.platform.prefix", "idea"),
additionalModules = getAdditionalPluginMainModules(),
projectDir = ideaProjectRoot,
keepHttpClient = false,
platformClassPathConsumer = { classPath, runDir ->
newClassPath = classPath
homePath = runDir.toString().replace(File.separator, "/")
@Suppress("SpellCheckingInspection")
val exceptions = setOf("jna.boot.library.path", "pty4j.preferred.native.folder", "jna.nosys", "jna.noclasspath", "jb.vmOptionsFile")
val systemProperties = System.getProperties()
for ((name, value) in getIdeSystemProperties(runDir).map) {
if (exceptions.contains(name) || !systemProperties.containsKey(name)) {
systemProperties.setProperty(name, value)
}
@Suppress("SpellCheckingInspection")
val exceptions = setOf("jna.boot.library.path", "pty4j.preferred.native.folder", "jna.nosys", "jna.noclasspath", "jb.vmOptionsFile")
val systemProperties = System.getProperties()
for ((name, value) in getIdeSystemProperties(runDir).map) {
if (exceptions.contains(name) || !systemProperties.containsKey(name)) {
systemProperties.setProperty(name, value)
}
},
generateRuntimeModuleRepository = SystemProperties.getBooleanProperty("intellij.build.generate.runtime.module.repository", false),
)
}
},
generateRuntimeModuleRepository = SystemProperties.getBooleanProperty("intellij.build.generate.runtime.module.repository", false),
)
}
finally {
batchSpanProcessorScope.cancel()
traceManagerInitializer = { throw IllegalStateException("already built") }
}
)
}
homePath?.let {
System.setProperty(PathManager.PROPERTY_HOME_PATH, it)
@@ -29,15 +29,6 @@ internal class Http2ClientConnection internal constructor(
connection.close()
}
fun withAuth(authHeader: CharSequence): Http2ClientConnection {
return Http2ClientConnection(
scheme = scheme,
authority = scheme,
commonHeaders = commonHeaders + arrayOf(HttpHeaderNames.AUTHORIZATION, AsciiString.of(authHeader)),
connection = connection,
)
}
suspend fun head(path: CharSequence): HttpResponseStatus {
return connection.stream { stream, result ->
stream.pipeline().addLast(object : InboundHandlerResultTracker<Http2HeadersFrame>(result) {
@@ -34,11 +34,15 @@ internal class Http2ClientConnectionFactory(
private val ioDispatcher: CoroutineDispatcher,
private val coroutineScope: CoroutineScope,
) {
fun connect(host: String, port: Int = 443): Http2ClientConnection {
return connect(InetSocketAddress.createUnresolved(host, port.let { if (it == -1) 443 else it }))
fun connect(host: String, port: Int = 443, authHeader: CharSequence? = null): Http2ClientConnection {
return connect(InetSocketAddress.createUnresolved(host, port.let { if (it == -1) 443 else it }), authHeader)
}
fun connect(server: InetSocketAddress): Http2ClientConnection {
fun connect(server: InetSocketAddress, auth: CharSequence? = null): Http2ClientConnection {
var commonHeaders = arrayOf(HttpHeaderNames.USER_AGENT, AsciiString.of("IJ Builder"))
if (auth != null) {
commonHeaders += arrayOf(HttpHeaderNames.AUTHORIZATION, AsciiString.of(auth))
}
return Http2ClientConnection(
connection = Http2ConnectionProvider(
server = server,
@@ -49,7 +53,7 @@ internal class Http2ClientConnectionFactory(
),
scheme = AsciiString.of(if (sslContext == null) "http" else "https"),
authority = AsciiString.of(server.hostString + ":" + server.port),
commonHeaders = arrayOf(HttpHeaderNames.USER_AGENT, AsciiString.of("IJ Builder")),
commonHeaders = commonHeaders,
)
}
@@ -268,7 +268,7 @@ class CompilationContextImpl private constructor(
overrideClassesOutputDirectory()
if (!this::compilationData.isInitialized) {
compilationData = JpsCompilationData(
dataStorageRoot = paths.buildOutputDir.resolve(".jps-build-data"),
dataStorageRoot = paths.buildOutputDir.resolve("jps-build-data"),
classesOutputDirectory = classesOutputDirectory,
buildLogFile = logDir.resolve("compilation.log"),
categoriesWithDebugLevelNullable = System.getProperty("intellij.build.debug.logging.categories", "")
@@ -337,8 +337,7 @@ suspend fun fetchAndUnpackCompiledClasses(
withHttp2ClientConnectionFactory(trustAll = metadata.serverUrl.contains("127.0.0.1")) { client ->
downloadCompilationCache(
client = client,
serverUrl = metadata.serverUrl,
prefix = metadata.prefix,
serverUrl = if (metadata.prefix.trim('/').isEmpty()) URI(metadata.serverUrl) else URI(metadata.serverUrl.trimEnd('/') + '/' + metadata.prefix),
toDownload = toDownload,
downloadedBytes = downloadedBytes,
skipUnpack = skipUnpack,
@@ -1,25 +1,25 @@
// Copyright 2000-2021 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license that can be found in the LICENSE file.
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package org.jetbrains.intellij.build.impl.compilation
import com.intellij.openapi.util.text.StringUtil
import com.intellij.util.io.awaitExit
import java.nio.file.Path
import java.util.concurrent.TimeUnit
class Git(private val dir: Path) {
fun log(commitsCount: Int): List<String> {
return execute("git", "log", "-$commitsCount", "--pretty=tformat:%H")
suspend fun log(commitCount: Int): List<String> {
@Suppress("SpellCheckingInspection")
return execute("git", "log", "-$commitCount", "--pretty=tformat:%H")
}
fun formatLatestCommit(format: String): String {
val lines = execute("git", "log", "--pretty=format:$format", "-n", "1")
return StringUtil.join(lines, "\n")
suspend fun formatLatestCommit(format: String): String {
return execute("git", "log", "--pretty=format:$format", "-n", "1").joinToString("\n")
}
fun listFilesUnderVersionControl(refSpec: String = "HEAD"): List<String> {
suspend fun listFilesUnderVersionControl(refSpec: String = "HEAD"): List<String> {
return execute("git", "ls-tree", "-r", refSpec, "--name-only")
}
fun currentCommitShortHash(): String {
suspend fun currentCommitShortHash(): String {
val lines = execute("git", "rev-parse", "--short", "HEAD")
if (lines.size != 1) {
throw IllegalStateException("Single line output is expected but got '$lines'")
@@ -31,24 +31,13 @@ class Git(private val dir: Path) {
return hash
}
fun lineBreaksConfig(): String {
val lines = maybeExecute("git", "config", "core.autocrlf").output.filter { !it.isBlank() }
if (lines.isEmpty()) {
return ""
}
if (lines.size != 1) {
throw IllegalStateException("Single line output is expected but got '$lines'")
}
return lines.first()
}
private fun maybeExecute(vararg command: String): ExecutionResult {
private suspend fun maybeExecute(vararg command: String): ExecutionResult {
val process = ProcessBuilder(*command).directory(dir.toFile()).start()
var output = process.inputStream.bufferedReader().use {
it.lines().map { line -> line.trim() }.toList()
}
if (!process.waitFor(1, TimeUnit.MINUTES)) {
process.destroyForcibly().waitFor()
process.destroyForcibly().awaitExit()
throw IllegalStateException("Cannot execute $command: 1 minute timeout")
}
if (process.exitValue() != 0) {
@@ -57,13 +46,13 @@ class Git(private val dir: Path) {
return ExecutionResult(process.exitValue(), output)
}
private fun execute(vararg command: String): List<String> {
private suspend fun execute(vararg command: String): List<String> {
val result = maybeExecute(*command)
if (result.exitCode != 0) {
throw IllegalStateException("git process failed with $result.exitCode:\n${result.output.joinToString("\n")}")
}
return result.output
}
private data class ExecutionResult(val exitCode: Int, val output: List<String>)
}
private data class ExecutionResult(@JvmField val exitCode: Int, @JvmField val output: List<String>)
@@ -3,22 +3,20 @@ package org.jetbrains.intellij.build.impl.compilation
import io.netty.util.AsciiString
import io.opentelemetry.api.trace.Span
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.runBlocking
import org.jetbrains.intellij.build.BuildOptions
import org.jetbrains.intellij.build.BuildPaths.Companion.ULTIMATE_HOME
import org.jetbrains.intellij.build.CompilationContext
import org.jetbrains.intellij.build.impl.cleanOutput
import org.jetbrains.intellij.build.impl.compilation.cache.CommitsHistory
import org.jetbrains.intellij.build.impl.createCompilationContext
import org.jetbrains.intellij.build.telemetry.TraceManager.spanBuilder
import org.jetbrains.intellij.build.telemetry.block
import org.jetbrains.intellij.build.telemetry.use
import org.jetbrains.intellij.build.telemetry.withTracer
import org.jetbrains.jps.incremental.storage.ProjectStamps
import java.net.URI
import java.nio.file.Path
import java.util.*
import java.util.concurrent.CancellationException
import kotlin.io.path.ExperimentalPathApi
import kotlin.io.path.deleteRecursively
internal val IS_PORTABLE_COMPILATION_CACHE_ENABLED: Boolean
get() = ProjectStamps.PORTABLE_CACHES && IS_JPS_CACHE_URL_CONFIGURED
@@ -26,30 +24,35 @@ internal val IS_PORTABLE_COMPILATION_CACHE_ENABLED: Boolean
private var isAlreadyUpdated = false
internal object TestJpsCompilationCacheDownload {
@ExperimentalPathApi
@JvmStatic
fun main(args: Array<String>) = runBlocking(Dispatchers.Default) {
fun main(args: Array<String>) = withTracer {
System.setProperty("jps.cache.test", "true")
System.setProperty("org.jetbrains.jps.portable.caches", "true")
if (System.getProperty(URL_PROPERTY) == null) {
System.setProperty(URL_PROPERTY, "https://127.0.0.1:1900/cache/jps")
}
val projectHome = ULTIMATE_HOME
val outputDir = projectHome.resolve("out/compilation")
val context = createCompilationContext(
val outputDir = projectHome.resolve("out/test-jps-cache-downloaded")
outputDir.deleteRecursively()
downloadJpsCache(
cacheUrl = URI(System.getProperty(URL_PROPERTY, "https://127.0.0.1:1900/cache/jps")),
gitUrl = computeRemoteGitUrl(),
authHeader = getAuthHeader(),
projectHome = projectHome,
defaultOutputRoot = outputDir,
options = BuildOptions(
incrementalCompilation = true,
useCompiledClassesFromProjectOutput = true,
),
classOutDir = outputDir.resolve("classes"),
cacheDestination = outputDir.resolve("jps-build-data"),
reportStatisticValue = { k, v ->
println("$k: $v")
}
)
downloadCacheAndCompileProject(forceDownload = false, gitUrl = computeRemoteGitUrl(), context = context)
}
}
internal suspend fun downloadJpsCacheAndCompileProject(context: CompilationContext) {
downloadCacheAndCompileProject(forceDownload = System.getProperty(FORCE_DOWNLOAD_PROPERTY).toBoolean(), gitUrl = computeRemoteGitUrl(), context = context)
downloadCacheAndCompileProject(
forceDownload = System.getProperty(FORCE_DOWNLOAD_PROPERTY).toBoolean(),
gitUrl = computeRemoteGitUrl(),
context = context,
)
}
/**
@@ -112,10 +115,11 @@ private suspend fun downloadCacheAndCompileProject(forceDownload: Boolean, gitUr
spanBuilder("download JPS cache and compile")
.setAttribute("forceRebuild", forceRebuild)
.setAttribute("forceDownload", forceDownload)
.block { span ->
.use { span ->
val cacheUrl = URI(require(URL_PROPERTY, "Remote Cache url"))
if (isAlreadyUpdated) {
span.addEvent("PortableCompilationCache is already updated")
return@block
return@use
}
check(IS_PORTABLE_COMPILATION_CACHE_ENABLED) {
@@ -126,11 +130,17 @@ private suspend fun downloadCacheAndCompileProject(forceDownload: Boolean, gitUr
cleanOutput(context = context, keepCompilationState = false)
}
val downloader = PortableCompilationCacheDownloader(context = context, git = Git(context.paths.projectHome))
val reportStatisticValue = context.messages::reportStatisticValue
val portableCompilationCache = PortableCompilationCache(forceDownload = forceDownload)
val availableCommitDepth = if (!forceRebuild && (forceDownload || !isIncrementalCompilationDataAvailable(context))) {
portableCompilationCache.downloadCache(downloader = downloader, gitUrl = computeRemoteGitUrl(), context = context)
portableCompilationCache.downloadCache(
cacheUrl = cacheUrl,
gitUrl = gitUrl,
reportStatisticValue = reportStatisticValue,
classOutDir = context.classesOutputDirectory,
projectHome = context.paths.projectHome,
context = context,
)
}
else {
-1
@@ -146,9 +156,10 @@ private suspend fun downloadCacheAndCompileProject(forceDownload: Boolean, gitUr
handleCompilationFailureBeforeRetry = { successMessage ->
portableCompilationCache.handleCompilationFailureBeforeRetry(
successMessage = successMessage,
portableCompilationCacheDownloader = downloader,
forceDownload = portableCompilationCache.forceDownload,
cacheUrl = cacheUrl,
gitUrl = gitUrl,
reportStatisticValue = reportStatisticValue,
context = context,
)
},
@@ -167,11 +178,12 @@ private class PortableCompilationCache(forceDownload: Boolean) {
* @return updated [successMessage]
*/
suspend fun handleCompilationFailureBeforeRetry(
cacheUrl: URI,
successMessage: String,
portableCompilationCacheDownloader: PortableCompilationCacheDownloader,
context: CompilationContext,
forceDownload: Boolean,
gitUrl: String,
reportStatisticValue: (key: String, value: String) -> Unit,
): String {
when {
forceDownload -> {
@@ -184,7 +196,14 @@ private class PortableCompilationCache(forceDownload: Boolean) {
// If download isn't forced, then locally available cache will be used which may suffer from those issues.
// Hence, compilation failure. Replacing local cache with remote one may help.
Span.current().addEvent("Incremental compilation using locally available caches failed. Re-trying using Remote Cache.")
val availableCommitDepth = downloadCache(portableCompilationCacheDownloader, gitUrl = gitUrl, context)
val availableCommitDepth = downloadCache(
cacheUrl = cacheUrl,
gitUrl = gitUrl,
reportStatisticValue = reportStatisticValue,
classOutDir = context.classesOutputDirectory,
projectHome = context.paths.projectHome,
context = context,
)
if (availableCommitDepth >= 0) {
return portableJpsCacheUsageStatus(availableCommitDepth)
}
@@ -193,10 +212,25 @@ private class PortableCompilationCache(forceDownload: Boolean) {
return successMessage
}
suspend fun downloadCache(downloader: PortableCompilationCacheDownloader, gitUrl: String, context: CompilationContext): Int {
return spanBuilder("downloading Portable Compilation Cache").use { span ->
suspend fun downloadCache(
cacheUrl: URI,
gitUrl: String,
reportStatisticValue: (key: String, value: String) -> Unit,
projectHome: Path,
classOutDir: Path,
context: CompilationContext,
): Int {
return spanBuilder("download Portable Compilation Cache").use { span ->
try {
downloader.download(cacheUrl = require(URL_PROPERTY, "Remote Cache url"), gitUrl = gitUrl, authHeader = getAuthHeader())
downloadJpsCache(
cacheUrl = cacheUrl,
gitUrl = gitUrl,
authHeader = getAuthHeader(),
projectHome = projectHome,
classOutDir = classOutDir,
cacheDestination = context.compilationData.dataStorageRoot,
reportStatisticValue = reportStatisticValue,
)
}
catch (e: CancellationException) {
throw e
@@ -277,11 +311,11 @@ internal class CompilationOutput(
private fun getJpsCacheUploadUrl(): URI = URI(require(UPLOAD_URL_PROPERTY, "Remote Cache upload url"))
private fun getAuthHeader(): CharSequence {
private fun getAuthHeader(): CharSequence? {
val username = System.getProperty("jps.auth.spaceUsername")
val password = System.getProperty("jps.auth.spacePassword")
return when {
password == null -> AsciiString.EMPTY_STRING
password == null -> null
username == null -> AsciiString.of("Bearer $password")
else -> AsciiString.of("Basic " + Base64.getEncoder().encodeToString("$username:$password".toByteArray()))
}
@@ -8,20 +8,17 @@ import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.launch
import org.jetbrains.intellij.build.CompilationContext
import kotlinx.coroutines.withContext
import org.jetbrains.intellij.build.forEachConcurrent
import org.jetbrains.intellij.build.http2Client.*
import org.jetbrains.intellij.build.impl.compilation.cache.CommitsHistory
import org.jetbrains.intellij.build.impl.compilation.cache.getAllCompilationOutputs
import org.jetbrains.intellij.build.retryWithExponentialBackOff
import org.jetbrains.intellij.build.telemetry.TraceManager.spanBuilder
import org.jetbrains.intellij.build.telemetry.use
import org.jetbrains.jps.incremental.storage.BuildTargetSourcesState
import java.net.URI
import java.nio.file.Files
import java.nio.file.NoSuchFileException
import java.nio.file.Path
import java.util.concurrent.CancellationException
import java.util.concurrent.TimeUnit
@@ -29,91 +26,128 @@ import java.util.concurrent.atomic.LongAdder
private const val COMMITS_COUNT = 1_000
internal class PortableCompilationCacheDownloader(private val context: CompilationContext, private val git: Git) {
private val lastCommits by lazy { git.log(COMMITS_COUNT) }
private suspend fun downloadToFile(urlPath: String, file: Path, spanName: String, notFound: LongAdder, connection: Http2ClientConnection?): Long {
return spanBuilder(spanName).setAttribute("urlPath", urlPath).setAttribute("path", file.toString()).use { span ->
if (connection == null) {
Files.createDirectories(file.parent)
require(urlPath.isS3())
retryWithExponentialBackOff {
awsS3Cli("cp", urlPath, file.toString())
}
return@use try {
Files.size(file)
}
catch (e: NoSuchFileException) {
0
}
}
val sizeOnDisk = connection.download(path = urlPath, file = file)
if (sizeOnDisk == -1L) {
span.addEvent("resource not found")
notFound.increment()
}
sizeOnDisk
}
}
private suspend fun getAvailableCachesKeysLazyTask(urlPathPrefix: String, gitUrl: String, connection: Http2ClientConnection?): Collection<String> {
val commitHistoryUrl = "$urlPathPrefix/${CommitsHistory.JSON_FILE}"
require(!commitHistoryUrl.isS3() && connection != null)
val json: Map<String, Set<String>> = connection.getJsonOrDefaultIfNotFound(path = commitHistoryUrl, defaultIfNotFound = emptyMap())
if (json.isEmpty()) {
return emptyList()
}
return CommitsHistory(json).commitsForRemote(gitUrl)
}
private suspend fun prepareDownload(urlPathPrefix: String, gitUrl: String, connection: Http2ClientConnection?): Pair<String, Int>? {
val availableCachesKeys = getAvailableCachesKeysLazyTask(urlPathPrefix = urlPathPrefix, gitUrl = gitUrl, connection = connection)
val availableCommitDepth = lastCommits.indexOfFirst {
availableCachesKeys.contains(it)
}
if (availableCommitDepth !in 0 until lastCommits.count()) {
Span.current().addEvent("unable to find cache for any of last ${lastCommits.count()} commits.")
return null
}
val lastCachedCommit = lastCommits.get(availableCommitDepth)
Span.current().addEvent(
"using cache for commit $lastCachedCommit ($availableCommitDepth behind last commit)",
Attributes.of(
AttributeKey.longKey("behind last commit"), availableCommitDepth.toLong(),
AttributeKey.stringArrayKey("available cache keys"), availableCachesKeys.toList(),
)
)
return lastCachedCommit to availableCommitDepth
}
suspend fun download(cacheUrl: String, authHeader: CharSequence, gitUrl: String): Int {
val start = System.nanoTime()
val totalDownloadedBytes = LongAdder()
val notFound = LongAdder()
var availableCommitDepth = -1
val total = if (cacheUrl.isS3()) {
val info = prepareDownload(urlPathPrefix = cacheUrl, gitUrl = gitUrl, connection = null) ?: return -1
internal suspend fun downloadJpsCache(
cacheUrl: URI,
authHeader: CharSequence?,
gitUrl: String,
projectHome: Path,
classOutDir: Path,
cacheDestination: Path,
reportStatisticValue: (key: String, value: String) -> Unit,
): Int {
val start = System.nanoTime()
val totalDownloadedBytes = LongAdder()
val notFound = LongAdder()
var availableCommitDepth = -1
val totalItemCount = withHttp2ClientConnectionFactory(trustAll = cacheUrl.host == "127.0.0.1") { client ->
checkMirrorAndConnect(initialServerUri = cacheUrl, client = client, authHeader = authHeader) { connection, urlPathPrefix ->
val info = spanBuilder("prepare downloading").use {
prepareDownload(urlPathPrefix = urlPathPrefix, gitUrl = gitUrl, connection = connection, lastCommits = Git(projectHome).log(COMMITS_COUNT))
} ?: return@checkMirrorAndConnect -1
availableCommitDepth = info.second
doDownload(
urlPathPrefix = cacheUrl,
lastCachedCommit = info.first,
spanBuilder("download jps cache").setAttribute("commit", info.first).use {
doDownload(
urlPathPrefix = urlPathPrefix,
lastCachedCommit = info.first,
notFound = notFound,
totalDownloadedBytes = totalDownloadedBytes,
connection = connection,
classOutDir = classOutDir,
cacheDestination = cacheDestination,
)
}
}
}
if (availableCommitDepth == -1) {
return -1
}
reportStatisticValue("jps-cache:download:time", TimeUnit.NANOSECONDS.toMillis((System.nanoTime() - start)).toString())
reportStatisticValue("jps-cache:downloaded:bytes", totalDownloadedBytes.sum().toString())
reportStatisticValue("jps-cache:downloaded:count", totalItemCount.toString())
reportStatisticValue("jps-cache:notFound:count", notFound.sum().toString())
return availableCommitDepth
}
private suspend fun downloadToFile(urlPath: String, file: Path, spanName: String, notFound: LongAdder, connection: Http2ClientConnection): Long {
return spanBuilder(spanName).setAttribute("urlPath", urlPath).setAttribute("path", file.toString()).use { span ->
val sizeOnDisk = connection.download(path = urlPath, file = file)
if (sizeOnDisk == -1L) {
span.addEvent("resource not found")
notFound.increment()
}
sizeOnDisk
}
}
private suspend fun prepareDownload(urlPathPrefix: String, gitUrl: String, connection: Http2ClientConnection?, lastCommits: List<String>): Pair<String, Int>? {
val availableCachesKeys = getAvailableCachesKeys(urlPathPrefix = urlPathPrefix, gitUrl = gitUrl, connection = connection)
val availableCommitDepth = lastCommits.indexOfFirst {
availableCachesKeys.contains(it)
}
if (availableCommitDepth !in 0 until lastCommits.count()) {
Span.current().addEvent(
"unable to find cache for any of last ${lastCommits.count()} commits",
Attributes.of(
AttributeKey.stringArrayKey("availableCacheKeys"), availableCachesKeys.toList()
),
)
return null
}
val lastCachedCommit = lastCommits.get(availableCommitDepth)
Span.current().addEvent(
"using cache for commit $lastCachedCommit ($availableCommitDepth behind last commit)",
Attributes.of(
AttributeKey.longKey("behind last commit"), availableCommitDepth.toLong(),
AttributeKey.stringArrayKey("available cache keys"), availableCachesKeys.toList(),
)
)
return lastCachedCommit to availableCommitDepth
}
private suspend fun getAvailableCachesKeys(urlPathPrefix: String, gitUrl: String, connection: Http2ClientConnection?): Collection<String> {
val commitHistoryUrl = "$urlPathPrefix/${CommitsHistory.JSON_FILE}"
require(!commitHistoryUrl.isS3() && connection != null)
val json: Map<String, Set<String>> = connection.getJsonOrDefaultIfNotFound(path = commitHistoryUrl, defaultIfNotFound = emptyMap())
if (json.isEmpty()) {
return emptyList()
}
return CommitsHistory(json).commitsForRemote(gitUrl)
}
private suspend fun doDownload(
urlPathPrefix: String,
lastCachedCommit: String,
notFound: LongAdder,
totalDownloadedBytes: LongAdder,
connection: Http2ClientConnection,
classOutDir: Path,
cacheDestination: Path,
): Int {
return withContext(Dispatchers.IO) {
launch {
downloadAndUnpackJpsCache(
urlPathPrefix = urlPathPrefix,
commitHash = lastCachedCommit,
notFound = notFound,
totalDownloadedBytes = totalDownloadedBytes,
connection = null,
totalBytes = totalDownloadedBytes,
connection = connection,
cacheDestination = cacheDestination,
)
}
else {
val serverUri = URI(cacheUrl)
withHttp2ClientConnectionFactory(trustAll = serverUri.host == "127.0.0.1") { client ->
client.connect(serverUri.host, serverUri.port).withAuth(authHeader).use { connection ->
val info = prepareDownload(urlPathPrefix = serverUri.path, gitUrl = gitUrl, connection = connection) ?: return@use -1
availableCommitDepth = info.second
doDownload(
urlPathPrefix = serverUri.path,
lastCachedCommit = info.first,
val json = connection.getString("$urlPathPrefix/metadata/$lastCachedCommit")
val outputs = getAllCompilationOutputs(sourceState = BuildTargetSourcesState.readJson(JsonReader(json.reader())), classOutDir = classOutDir)
spanBuilder("download compilation output parts").setAttribute(AttributeKey.longKey("count"), outputs.size.toLong()).use {
outputs.forEachConcurrent(downloadParallelism) { output ->
spanBuilder("get and unpack output").setAttribute("part", output.remotePath).use {
downloadAndUnpackCompilationOutput(
urlPathPrefix = urlPathPrefix,
compilationOutput = output,
notFound = notFound,
totalDownloadedBytes = totalDownloadedBytes,
connection = connection,
@@ -121,104 +155,52 @@ internal class PortableCompilationCacheDownloader(private val context: Compilati
}
}
}
if (availableCommitDepth == -1) {
return -1
}
reportStatisticValue("jps-cache:download:time", TimeUnit.NANOSECONDS.toMillis((System.nanoTime() - start)).toString())
reportStatisticValue("jps-cache:downloaded:bytes", totalDownloadedBytes.sum().toString())
reportStatisticValue("jps-cache:downloaded:count", total.toString())
reportStatisticValue("jps-cache:notFound:count", notFound.sum().toString())
return availableCommitDepth
outputs.size
}
}
private suspend fun PortableCompilationCacheDownloader.doDownload(
urlPathPrefix: String,
lastCachedCommit: String,
notFound: LongAdder,
totalDownloadedBytes: LongAdder,
connection: Http2ClientConnection?,
): Int {
return coroutineScope {
launch {
spanBuilder("get and unpack jps cache").setAttribute("commit", lastCachedCommit).use {
downloadAndUnpackJpsCache(
urlPathPrefix = urlPathPrefix,
commitHash = lastCachedCommit,
notFound = notFound,
totalBytes = totalDownloadedBytes,
connection = connection,
)
}
}
private suspend fun downloadAndUnpackJpsCache(
urlPathPrefix: String,
commitHash: String,
notFound: LongAdder,
totalBytes: LongAdder,
connection: Http2ClientConnection,
cacheDestination: Path,
) {
val cacheArchive = Files.createTempFile("cache", ".zip")
try {
val sizeOnDisk = downloadToFile(
urlPath = "$urlPathPrefix/caches/$commitHash",
file = cacheArchive,
spanName = "download jps cache",
notFound = notFound,
connection = connection,
)
val metadataUrlPath = "$urlPathPrefix/metadata/$lastCachedCommit"
val json = if (connection == null) {
require(metadataUrlPath.isS3())
retryWithExponentialBackOff {
awsS3Cli("cp", metadataUrlPath, "-")
}
require(sizeOnDisk > 0)
totalBytes.add(sizeOnDisk)
spanBuilder("unpack jps cache")
.setAttribute("archive", cacheArchive.toString())
.setAttribute("destination", cacheDestination.toString())
.use(Dispatchers.IO) {
unpackArchiveUsingNettyByteBufferPool(archiveFile = cacheArchive, outDir = cacheDestination, isCompressed = true)
}
else {
connection.getString(path = metadataUrlPath)
}
val outputs = getAllCompilationOutputs(sourceState = BuildTargetSourcesState.readJson(JsonReader(json.reader())), classOutDir = context.classesOutputDirectory)
spanBuilder("download compilation output parts").setAttribute(AttributeKey.longKey("count"), outputs.size.toLong()).use {
outputs.forEachConcurrent(downloadParallelism) { output ->
spanBuilder("get and unpack output").setAttribute("part", output.remotePath).use {
downloadAndUnpackCompilationOutput(
urlPathPrefix = urlPathPrefix,
compilationOutput = output,
notFound = notFound,
totalDownloadedBytes = totalDownloadedBytes,
connection = connection,
)
}
}
}
outputs.size
}
}
private suspend fun downloadAndUnpackJpsCache(urlPathPrefix: String, commitHash: String, notFound: LongAdder, totalBytes: LongAdder, connection: Http2ClientConnection?) {
val cacheArchive = Files.createTempFile("cache", ".zip")
try {
val sizeOnDisk = downloadToFile(
urlPath = "$urlPathPrefix/caches/$commitHash",
file = cacheArchive,
spanName = "download jps cache",
notFound = notFound,
connection = connection,
)
require(sizeOnDisk > 0)
totalBytes.add(sizeOnDisk)
val cacheDestination = context.compilationData.dataStorageRoot
spanBuilder("unpack jps cache")
.setAttribute("archive", cacheArchive.toString())
.setAttribute("destination", cacheDestination.toString())
.use(Dispatchers.IO) {
unpackArchiveUsingNettyByteBufferPool(archiveFile = cacheArchive, outDir = cacheDestination, isCompressed = true)
}
}
finally {
Files.deleteIfExists(cacheArchive)
}
finally {
Files.deleteIfExists(cacheArchive)
}
}
private suspend fun downloadAndUnpackCompilationOutput(
urlPathPrefix: String,
compilationOutput: CompilationOutput,
notFound: LongAdder,
totalDownloadedBytes: LongAdder,
connection: Http2ClientConnection?,
) {
val tempFile = compilationOutput.path.resolve("tmp-output.zip")
val urlPath = "$urlPathPrefix/${compilationOutput.remotePath.trimStart('/')}"
private suspend fun downloadAndUnpackCompilationOutput(
urlPathPrefix: String,
compilationOutput: CompilationOutput,
notFound: LongAdder,
totalDownloadedBytes: LongAdder,
connection: Http2ClientConnection,
) {
val urlPath = "$urlPathPrefix/${compilationOutput.remotePath}"
val tempFile = compilationOutput.path.resolve(java.lang.Long.toUnsignedString(System.nanoTime(), Character.MAX_RADIX) + ".tmp.zip")
try {
val sizeOnDisk = downloadToFile(
urlPath = urlPath,
file = tempFile,
@@ -236,26 +218,20 @@ internal class PortableCompilationCacheDownloader(private val context: Compilati
// For now, the JPS cache reports non-existent compilation outputs, and we have to use such a workaround.
totalDownloadedBytes.add(sizeOnDisk)
try {
spanBuilder("unpack output")
.setAttribute("archive", tempFile.toString())
.setAttribute("destination", compilationOutput.path.toString())
.use(Dispatchers.IO) {
unpackArchiveUsingNettyByteBufferPool(archiveFile = tempFile, outDir = compilationOutput.path, isCompressed = true)
}
}
catch (e: CancellationException) {
throw e
}
catch (e: Exception) {
throw Exception("Unable to unpack $urlPath to ${compilationOutput.path}", e)
}
finally {
Files.deleteIfExists(tempFile)
}
spanBuilder("unpack output")
.setAttribute("archive", tempFile.toString())
.setAttribute("destination", compilationOutput.path.toString())
.use {
unpackArchiveUsingNettyByteBufferPool(archiveFile = tempFile, outDir = compilationOutput.path, isCompressed = true)
}
}
private fun reportStatisticValue(key: String, value: String) {
context.messages.reportStatisticValue(key, value)
catch (e: CancellationException) {
throw e
}
catch (e: Exception) {
throw Exception("Unable to unpack $urlPath to ${compilationOutput.path}", e)
}
finally {
Files.deleteIfExists(tempFile)
}
}
@@ -40,7 +40,7 @@ private const val SOURCES_STATE_FILE_NAME = "target_sources_state.json"
internal suspend fun uploadJpsCache(
forcedUpload: Boolean,
commitHash: String,
authHeader: CharSequence,
authHeader: CharSequence?,
s3Dir: Path?,
uploadUrl: URI,
context: CompilationContext,
@@ -65,7 +65,7 @@ internal suspend fun uploadJpsCache(
val uploadedOutputCount = LongAdder()
val messages = context.messages
withHttp2ClientConnectionFactory(trustAll = uploadUrl.host == "127.0.0.1") { client ->
client.connect(uploadUrl.host, uploadUrl.port).withAuth(authHeader).use { connection ->
client.connect(host = uploadUrl.host, port = uploadUrl.port, authHeader = authHeader).use { connection ->
val urlPathPrefix = uploadUrl.path
withContext(Dispatchers.IO) {
launch {
@@ -204,7 +204,7 @@ internal suspend fun updateJpsCacheCommitHistory(
remoteGitUrl: String,
commitHash: String,
uploadUrl: URI,
authHeader: CharSequence,
authHeader: CharSequence?,
s3Dir: Path?,
context: CompilationContext,
) {
@@ -212,7 +212,7 @@ internal suspend fun updateJpsCacheCommitHistory(
val commitHistory = CommitsHistory(mapOf(remoteGitUrl to (overrideCommits ?: setOf(commitHash))))
withHttp2ClientConnectionFactory(trustAll = uploadUrl.host == "127.0.0.1") { client ->
val urlPathPrefix = uploadUrl.path
client.connect(uploadUrl.host, uploadUrl.port).withAuth(authHeader).use { connection ->
client.connect(uploadUrl.host, uploadUrl.port, authHeader = authHeader).use { connection ->
for (commitHashForRemote in commitHistory.commitsForRemote(remoteGitUrl)) {
val cacheUploaded = checkExists(connection, "$urlPathPrefix/caches/$commitHashForRemote")
val metadataUploaded = checkExists(connection, "$urlPathPrefix/metadata/$commitHashForRemote")
@@ -19,15 +19,14 @@ import java.util.concurrent.atomic.AtomicLong
import kotlin.io.path.name
internal suspend fun downloadCompilationCache(
serverUrl: String,
serverUrl: URI,
client: Http2ClientConnectionFactory,
prefix: String,
toDownload: Collection<FetchAndUnpackItem>,
downloadedBytes: AtomicLong,
skipUnpack: Boolean,
saveHash: Boolean,
) {
checkMirrorAndConnect(serverUrl = serverUrl, client = client, prefix = prefix) { connection, urlPathPrefix ->
checkMirrorAndConnect(initialServerUri = serverUrl, client = client) { connection, urlPathPrefix ->
val zstdDecompressContextPool = ZstdDecompressContextPool()
toDownload.forEachConcurrent(downloadParallelism) { item ->
val urlPath = "$urlPathPrefix/${item.name}/${item.file.fileName}"
@@ -59,47 +58,6 @@ internal suspend fun downloadCompilationCache(
}
}
internal suspend fun checkMirrorAndConnect(
serverUrl: String,
client: Http2ClientConnectionFactory,
prefix: String,
block: suspend (connection: Http2ClientConnection, urlPathPrefix: String) -> Unit,
) {
var urlPathWithSlash = "/$prefix/"
// first let's check for initial redirect (mirror selection)
val initialServerUri = URI(serverUrl)
var effectiveServerUri = initialServerUri
var connection: Http2ClientConnection? = client.connect(effectiveServerUri.host, effectiveServerUri.port)
var connectionToClose = connection
try {
spanBuilder("mirror selection").use { span ->
val newLocation = connection!!.getRedirectLocation(urlPathWithSlash)
if (newLocation == null) {
span.addEvent("origin server will be used", Attributes.of(AttributeKey.stringKey("url"), urlPathWithSlash))
connectionToClose = null
}
else {
effectiveServerUri = URI(newLocation.toString())
urlPathWithSlash = effectiveServerUri.path
span.addEvent("redirected to mirror", Attributes.of(AttributeKey.stringKey("url"), urlPathWithSlash))
connection = null
}
}
}
finally {
connectionToClose?.close()
connectionToClose = null
}
val effectiveConnection = connection ?: client.connect(effectiveServerUri.host, effectiveServerUri.port)
try {
block(effectiveConnection, urlPathWithSlash.trimEnd('/'))
}
finally {
effectiveConnection.close()
}
}
private suspend fun download(
item: FetchAndUnpackItem,
urlPath: String,
@@ -0,0 +1,41 @@
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package org.jetbrains.intellij.build.impl.compilation
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import org.jetbrains.intellij.build.http2Client.Http2ClientConnection
import org.jetbrains.intellij.build.http2Client.Http2ClientConnectionFactory
import org.jetbrains.intellij.build.telemetry.TraceManager.spanBuilder
import org.jetbrains.intellij.build.telemetry.use
import java.net.URI
internal suspend fun <T> checkMirrorAndConnect(
initialServerUri: URI,
client: Http2ClientConnectionFactory,
authHeader: CharSequence? = null,
block: suspend (connection: Http2ClientConnection, urlPathPrefix: String) -> T,
): T {
var urlPath = initialServerUri.path
// first let's check for initial redirect (mirror selection)
var connection = client.connect(host = initialServerUri.host, port = initialServerUri.port, authHeader = authHeader)
try {
spanBuilder("mirror selection").use { span ->
val newLocation = connection.getRedirectLocation("$urlPath/")?.toString()
if (newLocation == null) {
span.addEvent("origin server will be used", Attributes.of(AttributeKey.stringKey("url"), initialServerUri.toString()))
}
else {
connection.close()
val newServerUri = URI(newLocation)
urlPath = newServerUri.path.trimEnd('/')
span.addEvent("redirected to mirror", Attributes.of(AttributeKey.stringKey("url"), newLocation))
connection = client.connect(host = newServerUri.host, port = newServerUri.port, authHeader = authHeader)
}
}
return block(connection, urlPath)
}
finally {
connection.close()
}
}