From accd76973639aa5262c58c921094608bf6dfd1e6 Mon Sep 17 00:00:00 2001 From: Titouan Bion Date: Wed, 8 Oct 2025 18:53:33 +0200 Subject: [PATCH] [fleet] Move `fleet.codecache` to Community (cherry picked from commit a8ecc20e305c2592b26bf1740f2228bae741755a) FLEET-MR-6894 GitOrigin-RevId: facb856bea2f4a5a5578b6997e31b10a3b34bba3 --- .idea/modules.xml | 2 + build/bazel-generated-file-list.txt | 2 + fleet/codecache/BUILD.bazel | 36 +++ fleet/codecache/fleet.codecache.iml | 45 ++++ fleet/codecache/gradlebuild/build.gradle.kts | 78 ++++++ .../.pseudoCommonKotlinSourceSet | 0 .../fleet/codecache/CodeCache.kt | 21 ++ .../fleet/codecache/Marketplace.kt | 82 ++++++ .../fleet/codecache/MarketplaceRepository.kt | 92 +++++++ .../fleet/codecache/ModuleLayerLoader.kt | 61 +++++ .../fleet/codecache/CodeCache.jvm.kt | 242 ++++++++++++++++++ fleet/codecache/test/BUILD.bazel | 44 ++++ fleet/codecache/test/fleet.codecache.test.iml | 45 ++++ .../test/gradlebuild/build.gradle.kts | 70 +++++ .../codecache/test/CodeCacheFileLockTest.kt | 113 ++++++++ 15 files changed, 933 insertions(+) create mode 100644 fleet/codecache/BUILD.bazel create mode 100644 fleet/codecache/fleet.codecache.iml create mode 100644 fleet/codecache/gradlebuild/build.gradle.kts create mode 100644 fleet/codecache/srcCommonMain/.pseudoCommonKotlinSourceSet create mode 100644 fleet/codecache/srcCommonMain/fleet/codecache/CodeCache.kt create mode 100644 fleet/codecache/srcCommonMain/fleet/codecache/Marketplace.kt create mode 100644 fleet/codecache/srcCommonMain/fleet/codecache/MarketplaceRepository.kt create mode 100644 fleet/codecache/srcCommonMain/fleet/codecache/ModuleLayerLoader.kt create mode 100644 fleet/codecache/srcJvmMain/fleet/codecache/CodeCache.jvm.kt create mode 100644 fleet/codecache/test/BUILD.bazel create mode 100644 fleet/codecache/test/fleet.codecache.test.iml create mode 100644 fleet/codecache/test/gradlebuild/build.gradle.kts create mode 100644 fleet/codecache/test/src/fleet/codecache/test/CodeCacheFileLockTest.kt diff --git a/.idea/modules.xml b/.idea/modules.xml index bdb749bc865e..22b5cd32c1dc 100644 --- a/.idea/modules.xml +++ b/.idea/modules.xml @@ -3,6 +3,8 @@ + + diff --git a/build/bazel-generated-file-list.txt b/build/bazel-generated-file-list.txt index 32b9ef8830d5..e9d4d6ab580c 100644 --- a/build/bazel-generated-file-list.txt +++ b/build/bazel-generated-file-list.txt @@ -222,6 +222,8 @@ community-resources fleet/andel fleet/bifurcan fleet/bundles +fleet/codecache +fleet/codecache/test fleet/compiler-plugins fleet/fastutil fleet/junit4 diff --git a/fleet/codecache/BUILD.bazel b/fleet/codecache/BUILD.bazel new file mode 100644 index 000000000000..567aa2ed0050 --- /dev/null +++ b/fleet/codecache/BUILD.bazel @@ -0,0 +1,36 @@ +### auto-generated section `build fleet.codecache` start +load("//build:compiler-options.bzl", "create_kotlinc_options") +load("@rules_jvm//:jvm.bzl", "jvm_library") + +create_kotlinc_options( + name = "custom_codecache", + opt_in = [ + "kotlinx.coroutines.ExperimentalCoroutinesApi", + "kotlin.ExperimentalStdlibApi", + ], + x_consistent_data_class_copy_visibility = True, + x_context_parameters = True, + x_jvm_default = "all-compatibility", + x_lambdas = "class" +) + +jvm_library( + name = "codecache", + module_name = "fleet.codecache", + visibility = ["//visibility:public"], + srcs = glob(["srcCommonMain/**/*.kt", "srcCommonMain/**/*.java", "srcCommonMain/**/*.form", "srcJvmMain/**/*.kt", "srcJvmMain/**/*.java", "srcJvmMain/**/*.form"], allow_empty = True, exclude = ["**/module-info.java"]), + kotlinc_opts = ":custom_codecache", + deps = [ + "@lib//:kotlinx-coroutines-core", + "@lib//:kotlin-stdlib", + "//fleet/bundles", + "@lib//:jetbrains-annotations", + "//fleet/util/logging/api", + "@lib//:kotlinx-serialization-core", + "@lib//:kotlinx-serialization-json", + "@lib//:kotlinx-collections-immutable", + "//fleet/ktor/client/core", + "//fleet/modules/api", + ] +) +### auto-generated section `build fleet.codecache` end \ No newline at end of file diff --git a/fleet/codecache/fleet.codecache.iml b/fleet/codecache/fleet.codecache.iml new file mode 100644 index 000000000000..70280cf62bc3 --- /dev/null +++ b/fleet/codecache/fleet.codecache.iml @@ -0,0 +1,45 @@ + + + + + + + + + + + + + + + + $KOTLIN_BUNDLED$/lib/kotlinx-serialization-compiler-plugin.jar + + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/fleet/codecache/gradlebuild/build.gradle.kts b/fleet/codecache/gradlebuild/build.gradle.kts new file mode 100644 index 000000000000..5f9bcb142c9d --- /dev/null +++ b/fleet/codecache/gradlebuild/build.gradle.kts @@ -0,0 +1,78 @@ +// IMPORT__MARKER_START +import fleet.buildtool.conventions.configureAtMostOneJvmTargetOrThrow +import fleet.buildtool.conventions.withJavaSourceSet +// IMPORT__MARKER_END +plugins { + alias(libs.plugins.kotlin.multiplatform) + id("fleet.project-module-conventions") + id("fleet.toolchain-conventions") + id("fleet.module-publishing-conventions") + id("fleet.sdk-repositories-publishing-conventions") + id("fleet.open-source-module-conventions") + alias(libs.plugins.dokka) + // GRADLE_PLUGINS__MARKER_START + id("fleet-module") + alias(jps.plugins.kotlin.serialization) + // GRADLE_PLUGINS__MARKER_END +} + +fleetModule { + module { + name = "fleet.codecache" + importedFromJps {} + } +} + +@OptIn(org.jetbrains.kotlin.gradle.ExperimentalWasmDsl::class) +kotlin { + // KOTLIN__MARKER_START + compilerOptions.freeCompilerArgs = listOf( + "-opt-in=kotlinx.coroutines.ExperimentalCoroutinesApi", + "-opt-in=kotlin.ExperimentalStdlibApi", + "-Xlambdas=class", + "-Xconsistent-data-class-copy-visibility", + "-Xcontext-parameters", + "-XXLanguage:+AllowEagerSupertypeAccessibilityChecks", + ) + jvm {} + wasmJs { + browser {} + } + sourceSets.commonMain.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcCommonMain")) } + sourceSets.commonMain.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesCommonMain")) } + sourceSets.commonTest.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcCommonTest")) } + sourceSets.commonTest.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesCommonTest")) } + sourceSets.jvmMain.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcJvmMain")) } + configureAtMostOneJvmTargetOrThrow { compilations.named("main") { withJavaSourceSet { javaSourceSet -> javaSourceSet.java.srcDir(layout.projectDirectory.dir("../srcJvmMain")) } } } + sourceSets.jvmMain.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesJvmMain")) } + sourceSets.jvmTest.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcJvmTest")) } + configureAtMostOneJvmTargetOrThrow { compilations.named("test") { withJavaSourceSet { javaSourceSet -> javaSourceSet.java.srcDir(layout.projectDirectory.dir("../srcJvmTest")) } } } + sourceSets.jvmTest.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesJvmTest")) } + sourceSets.wasmJsMain.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcWasmJsMain")) } + sourceSets.wasmJsMain.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesWasmJsMain")) } + sourceSets.wasmJsTest.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcWasmJsTest")) } + sourceSets.wasmJsTest.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesWasmJsTest")) } + sourceSets.commonMain.dependencies { + implementation(jps.com.intellij.platform.kotlinx.coroutines.core.jvm134738847.get().let { "${it.group}:kotlinx-coroutines-core:${it.version}" }) { + isTransitive = false + } + implementation(jps.org.jetbrains.kotlin.kotlin.stdlib1993400674.get().let { "${it.group}:${it.name}:${it.version}" }) { + exclude(group = "org.jetbrains", module = "annotations") + } + implementation(jps.org.jetbrains.annotations1504825916.get()) + implementation(jps.org.jetbrains.kotlinx.kotlinx.serialization.core.jvm1739247612.get().let { "${it.group}:kotlinx-serialization-core:${it.version}" }) { + isTransitive = false + } + implementation(jps.org.jetbrains.kotlinx.kotlinx.serialization.json.jvm231489733.get().let { "${it.group}:kotlinx-serialization-json:${it.version}" }) { + isTransitive = false + } + implementation(jps.org.jetbrains.kotlinx.kotlinx.collections.immutable.jvm717536558.get().let { "${it.group}:kotlinx-collections-immutable:${it.version}" }) { + isTransitive = false + } + implementation(project(":fleet.bundles")) + implementation(project(":fleet.util.logging.api")) + implementation(project(":fleet.ktor.client.core")) + implementation(project(":fleet.modules.api")) + } + // KOTLIN__MARKER_END +} \ No newline at end of file diff --git a/fleet/codecache/srcCommonMain/.pseudoCommonKotlinSourceSet b/fleet/codecache/srcCommonMain/.pseudoCommonKotlinSourceSet new file mode 100644 index 000000000000..e69de29bb2d1 diff --git a/fleet/codecache/srcCommonMain/fleet/codecache/CodeCache.kt b/fleet/codecache/srcCommonMain/fleet/codecache/CodeCache.kt new file mode 100644 index 000000000000..24bc7cd331a9 --- /dev/null +++ b/fleet/codecache/srcCommonMain/fleet/codecache/CodeCache.kt @@ -0,0 +1,21 @@ +package fleet.codecache + +import fleet.bundles.Coordinates +import io.ktor.http.decodeURLQueryComponent + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.commander.workspace"]) +fun Coordinates.relativePathToCodeCache(): String { + return when (this) { + is Coordinates.Remote -> { + // todo hacky, to be reconsidered + val moduleDividedParts = url.split("&module=") + val fileName = when (moduleDividedParts.size) { + 1 -> url.substringAfterLast('/') // for "$host/fleet-parts/*" we don't use the new endpoint? + else -> moduleDividedParts.last() + }.decodeURLQueryComponent(plusIsSpace = false) + "$hash/$fileName" + } + //is Coordinates.Workspace -> return "$filename#$sha2" + else -> throw RuntimeException("Unsupported coordinates: $this") + } +} diff --git a/fleet/codecache/srcCommonMain/fleet/codecache/Marketplace.kt b/fleet/codecache/srcCommonMain/fleet/codecache/Marketplace.kt new file mode 100644 index 000000000000..47156e928f9c --- /dev/null +++ b/fleet/codecache/srcCommonMain/fleet/codecache/Marketplace.kt @@ -0,0 +1,82 @@ +package fleet.codecache + +import fleet.bundles.PluginName +import fleet.bundles.PluginVersion +import io.ktor.http.encodeURLQueryComponent +import kotlinx.serialization.Serializable + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +@Serializable +data class VersionsQuery(val query: String) + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +fun versionsQuery(names: Set, offset: Int, version: PluginVersion?): VersionsQuery { + val joinedName = names.joinToString(separator = ",") { name -> "\"${name}\"" } + + val compatRange = version?.toLong()?.let { long -> + "compatibility: {range: {gte: $long, lte: $long}}" + }.orEmpty() + + return VersionsQuery(""" + query { + updates( + search: { + filters: [{ field: "family", value: "fleet" }, { field: "xmlId", value: [$joinedName] }] + max: 10 + $compatRange + offset: $offset + collapseField: PLUGIN_ID + } + ) { + total + updates { + xmlId + version + } + } + } + """.trimIndent()) +} + +/** + * Marketplace URL to resolve a resource of an uploaded Fleet plugin. + * + * Resources referenced by [resourceFilename] could be jars, icons, Fleet parts file, etc. + */ +fun resourceUrl(host: String, pluginName: PluginName, pluginVersion: PluginVersion, resourceFilename: String): String = + encodedResourceUrl(host, pluginName.name, pluginVersion.versionString, resourceFilename) + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +fun bundleUri(host: String, bundleName: String, bundleVersion: String): String = + encodedResourceUrl(host, bundleName, bundleVersion, "extension.json") + +private fun encodedResourceUrl(host: String, xmlId: String, version: String, resourceFilename: String): String = + "$host/api/fleet/download?xmlId=${xmlId.encodeURLQueryComponentForMarketplace()}&version=${version.encodeURLQueryComponentForMarketplace()}&module=${resourceFilename.encodeURLQueryComponentForMarketplace()}" + +private fun String.encodeURLQueryComponentForMarketplace() = encodeURLQueryComponent(spaceToPlus = false, encodeFull = true) + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +fun searchUri(host: String): String = + "$host/api/search/graphql" + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +@Serializable +data class VersionsResponse(val data: VersionsData) + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +@Serializable +data class VersionsData(val updates: VersionsList) + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +@Serializable +data class VersionsList( + val total: Int = 0, + val updates: List = emptyList(), +) + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"]) +@Serializable +data class VersionInfo( + val xmlId: String, + val version: String, +) diff --git a/fleet/codecache/srcCommonMain/fleet/codecache/MarketplaceRepository.kt b/fleet/codecache/srcCommonMain/fleet/codecache/MarketplaceRepository.kt new file mode 100644 index 000000000000..103b7f716dab --- /dev/null +++ b/fleet/codecache/srcCommonMain/fleet/codecache/MarketplaceRepository.kt @@ -0,0 +1,92 @@ +package fleet.codecache + +import fleet.bundles.PluginDescriptor +import fleet.bundles.PluginName +import fleet.bundles.PluginRepository +import fleet.bundles.PluginVersion +import fleet.util.logging.logger +import io.ktor.client.HttpClient +import io.ktor.client.request.get +import io.ktor.client.request.post +import io.ktor.client.request.setBody +import io.ktor.client.statement.bodyAsText +import io.ktor.http.ContentType +import io.ktor.http.content.TextContent +import kotlinx.coroutines.CancellationException +import kotlinx.serialization.json.Json +import kotlin.uuid.ExperimentalUuidApi + +private object Marketplace { + val logger = logger() +} + +//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common", "fleet.plugins.pluginManagement.test"]) +fun marketPlaceRepository(httpClient: suspend () -> HttpClient, host: String, defaultRepository: PluginRepository? = null): PluginRepository { + return object : PluginRepository { + override suspend fun getLatestVersions(names: Set, shipVersion: PluginVersion): Map { + Marketplace.logger.debug("Querying latest versions of '$names' compatible with '$shipVersion'...") + tailrec suspend fun loop( + out: HashMap = HashMap(), + offset: Int = 0, + ): Map { + val query = versionsQuery(names.mapTo(HashSet()) { it.name }, offset, shipVersion) + val queryStr = Json.encodeToString(VersionsQuery.serializer(), query) + val versions = kotlin.runCatching { + val url = searchUri(host) + val body = httpClient().post(url) { + setBody(TextContent(queryStr, ContentType.Application.Json)) + }.bodyAsText() + val response = Json.decodeFromString(VersionsResponse.serializer(), body) + Marketplace.logger.debug("URL:\n$url\nQuery:\n$queryStr\nResponse:\n$response") + response.data.updates.updates.map { versionInfo -> + PluginName(versionInfo.xmlId) to PluginVersion.fromString(versionInfo.version) + } + }.onFailure { x -> + if (x is CancellationException) { + throw x + } + Marketplace.logger.warn(x, "got error from versions query for $names, query = $query") + }.getOrNull() ?: emptyList() + + out.putAll(versions) + return when { + versions.size < 10 -> out + else -> loop(out, offset + versions.size) + } + } + + val versionMap = loop() + Marketplace.logger.debug("Found latest versions '$versionMap'") + return versionMap + } + + override suspend fun getPlugin(pluginName: PluginName, pluginVersion: PluginVersion): PluginDescriptor? = run { + // we should not download the descriptor from the internet if we already did it once, we have it in the trusted repo + // this optimization should speed up the plugin management significantly + defaultRepository?.getPlugin(pluginName, pluginVersion) ?: run { + Marketplace.logger.debug("Querying marketplace $host for a plugin $pluginName") + val uri = bundleUri(host, pluginName.name, pluginVersion.versionString) + try { + val bundleStr = httpClient().get(uri).bodyAsText() + Marketplace.logger.debug("getPlugin($pluginName, $pluginVersion) => ${bundleStr.length} bytes") + Json.decodeFromString(PluginDescriptor.serializer(), bundleStr) + } + catch (x: CancellationException) { + throw x + } + catch (x: Throwable) { + Marketplace.logger.debug(x, "got error from $uri") + throw x + } + } + } + + override fun presentableName(): String = "Marketplace PluginRepository `$host`" + + @OptIn(ExperimentalUuidApi::class) + /** + * Marketplace repository is always considered identical in context of cache invalidation, we want reproducibility (after the first resolution) instead of freshness. + */ + override fun cacheKey(): String = "never changing Marketplace" + } +} diff --git a/fleet/codecache/srcCommonMain/fleet/codecache/ModuleLayerLoader.kt b/fleet/codecache/srcCommonMain/fleet/codecache/ModuleLayerLoader.kt new file mode 100644 index 000000000000..231a9109f522 --- /dev/null +++ b/fleet/codecache/srcCommonMain/fleet/codecache/ModuleLayerLoader.kt @@ -0,0 +1,61 @@ +@file:Suppress("ReplacePutWithAssignment", "ReplaceGetOrSet") + +package fleet.codecache + +import fleet.bundles.LayerSelector +import fleet.bundles.ResolvedPluginLayer +import fleet.bundles.ResourceBundle +import fleet.bundles.dockLayer +import fleet.bundles.internalReadability +import fleet.modules.api.FleetModule +import fleet.modules.api.FleetModuleInfo +import fleet.modules.api.FleetModuleLayer +import fleet.modules.api.FleetModuleLayerLoader +import fleet.util.logging.KLoggers.logger + +private val logger by lazy { logger(ResolvedPluginLayer::class) } + +fun loadPluginModulesAndResources(moduleLayerLoader: FleetModuleLayerLoader, + resolvedLayers: Map, + sortedSelectors: List, + externalDependencies: Map>, + baseLayers: List, + loadDockModuleLayer: (modulePath: Set) -> FleetModuleLayer): Map { + val result = HashMap(sortedSelectors.size) + for (selector in sortedSelectors) { + val resolvedLayer = resolvedLayers.get(selector) ?: continue + val moduleLayer = when (selector) { + dockLayer -> loadDockModuleLayer(moduleInfos(resolvedLayer)) + else -> { + val internal = internalReadability[selector]?.mapNotNull { dependency -> + result[dependency]?.layer + } ?: emptyList() + val external = externalDependencies[selector] ?: emptyList() + moduleLayerLoader.moduleLayer(modulePath = moduleInfos(resolvedLayer), parentLayers = baseLayers + external + internal) + } + } + val modules = resolvedLayer.modules.mapNotNull { moduleName -> + moduleLayer.findModule(moduleName) + } + result.put(selector, PluginModulesAndResources(layer = moduleLayer, modules = modules, resources = resolvedLayer.resources)) + } + return result +} + +fun moduleInfos(resolvedPluginLayer: ResolvedPluginLayer): Set { + return resolvedPluginLayer.modulePath.mapTo(HashSet(resolvedPluginLayer.modulePath.size)) { path -> + val filePath = path.path + when (val second = path.serializedModuleDescriptor) { + null -> FleetModuleInfo.Path(filePath) + else -> + resolvedPluginLayer.runCatching { + FleetModuleInfo.WithDescriptor(serializedModuleDescriptor = second, path = filePath) + }.getOrElse { message -> + logger.warn("$message: Cannot deserialize descriptor for ${resolvedPluginLayer}") + FleetModuleInfo.Path(filePath) + } + } + } +} + +data class PluginModulesAndResources(val modules: List, val layer: FleetModuleLayer, val resources: Set) \ No newline at end of file diff --git a/fleet/codecache/srcJvmMain/fleet/codecache/CodeCache.jvm.kt b/fleet/codecache/srcJvmMain/fleet/codecache/CodeCache.jvm.kt new file mode 100644 index 000000000000..59bed41730a0 --- /dev/null +++ b/fleet/codecache/srcJvmMain/fleet/codecache/CodeCache.jvm.kt @@ -0,0 +1,242 @@ +package fleet.codecache + +import fleet.bundles.Coordinates +import fleet.bundles.CoordinatesResolution +import fleet.bundles.ResolutionException +import fleet.util.logging.KLoggers.logger +import io.ktor.client.HttpClient +import io.ktor.client.call.body +import io.ktor.client.plugins.expectSuccess +import io.ktor.client.request.prepareGet +import io.ktor.client.statement.bodyAsText +import io.ktor.http.isSuccess +import io.ktor.utils.io.ByteReadChannel +import io.ktor.utils.io.jvm.javaio.copyTo +import kotlinx.coroutines.* +import org.jetbrains.annotations.VisibleForTesting +import java.nio.charset.StandardCharsets +import java.nio.file.FileAlreadyExistsException +import java.nio.file.Files +import java.nio.file.Path +import java.nio.file.StandardOpenOption +import java.security.MessageDigest +import java.text.DecimalFormat +import kotlin.io.path.absolutePathString +import kotlin.io.path.isDirectory +import kotlin.io.path.outputStream +import kotlin.math.log10 +import kotlin.math.pow + +private const val MAX_LOCK_WAIT_MS = 60000L +private const val LOCK_DELAY_MS = 100L + +data class CodeCachePath(val path: Path, val writable: Boolean) + +class CodeCache( + private val httpClientFn: suspend () -> HttpClient, + private val paths: List, + private val lockDelay: Long = LOCK_DELAY_MS, + private val maxLockWaitTime: Long = MAX_LOCK_WAIT_MS, + private val queryParams: (suspend () -> Map)? = null +) { + private val hasher: CodeCacheHasher = CodeCacheHasher() + + companion object { + private val logger = logger(CodeCache::class) + } + + suspend fun resolve(coord: Coordinates): String { + return when (coord) { + is Coordinates.Local -> { + Path.of(coord.path).takeIf { it.exists() } ?: throw ResolutionException(coord) + } + is Coordinates.Remote -> { + val relativePath = coord.relativePathToCodeCache() + val resolved = paths.map { it.path.resolve(relativePath) }.firstOrNull { it.exists() } + if (resolved != null) { + resolved + } + else { + val writableCachePath = paths.firstOrNull { it.writable }?.path ?: error("No writable code cache is provided") + val targetFile = writableCachePath.resolve(relativePath) + val tmpFile = writableCachePath.resolve("tmp").resolve(relativePath) + if (downloadWithLock(targetFile, tmpFile, coord, queryParams?.invoke() ?: emptyMap())) { + targetFile + } + else { + throw ResolutionException(coord) + } + } + } + }.absolutePathString() + } + + private suspend fun downloadWithLock(targetFile: Path, + tmpFile: Path, + coord: Coordinates.Remote, + queryParams: Map): Boolean { + logger.debug("Downloading $coord to $targetFile") + + ensureDirExists(targetFile.parent) + ensureDirExists(tmpFile.parent) + + return withFileLock(tmpFile.parent, tmpFile.name, coord) { + if (!targetFile.exists()) { + httpClientFn().downloadFile(coord.url, tmpFile, queryParams) + val actualHash = hash(tmpFile) + val hashesMatch = actualHash == coord.hash + if (hashesMatch) { + Files.move(tmpFile, targetFile) // due to compatibility with the Gradle plugin and kotlin 1.6 + } + else { + logger.error(Throwable("Hash mismatch")) { + val fileSize = tmpFile.fileSize() + val buffer = CharArray(1024 * 4) + runInterruptible { + kotlin.runCatching { + Files.newBufferedReader(tmpFile, StandardCharsets.ISO_8859_1) + .use { reader -> reader.read(buffer) } + }.onFailure { + logger.error(it) { + "empty message" + } + } + } + "Hash mismatch (${coord.url}). Expected: ${coord.hash}, actual: $actualHash, " + + "downloaded file size: ${formatFileSize(fileSize)}, target: ${targetFile.absolute()}, tmp: ${tmpFile.absolute()}, " + + "content: `${buffer.concatToString()}`" + } + Files.deleteIfExists(tmpFile) // due to compatibility with the Gradle plugin and kotlin 1.6 + } + hashesMatch + } + else { + true + } + } + } + + private fun ensureDirExists(dir: Path) { + try { + Files.createDirectories(dir) + } + catch (e: FileAlreadyExistsException) { + if (dir.isDirectory()) { + // it was a symlink + } + else { + throw e + } + } + catch (e: Exception) { + throw RuntimeException("Couldn't create cache directory: $dir", e) + } + } + + @VisibleForTesting + suspend fun withFileLock(folder: Path, filename: String, coord: Coordinates, body: suspend CoroutineScope.() -> T): T { + val lockFile = folder.resolve("$filename.lock") + fun tryLock(lockFile: Path): Boolean = try { + Files.createFile(lockFile) // due to compatibility with the Gradle plugin and kotlin 1.6 + true + } + catch (_: FileAlreadyExistsException) { + false + } + + try { + val success = withTimeoutOrNull(maxLockWaitTime) { + while (!tryLock(lockFile)) { + delay(lockDelay) + } + true + } + + if (success == null) { + throw ResolutionException(coord, Throwable( + "Waited on a lock for more than $maxLockWaitTime ms. Please remove the lock file by hand: $lockFile")) + } + + return coroutineScope(body) + } + finally { + Files.deleteIfExists(lockFile) // due to compatibility with the Gradle plugin and kotlin 1.6 + } + } + + //@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.commander.workspace"]) + fun hash(path: Path): String = hasher.hash(path) +} + +class CodeCacheHasher(hashAlgorithm: String = Coordinates.Remote.HASH_ALGORITHM) { + private val digestToClone = MessageDigest.getInstance(hashAlgorithm) + + fun hash(path: Path): String { + val buffer = ByteArray(1024 * 1024) + val digest = digestToClone.clone() as MessageDigest + digest.reset() + path.inputStream().buffered().use { + while (true) { + val read = it.read(buffer) + if (read <= 0) break + digest.update(buffer, 0, read) + } + } + + val shaBytes = digest.digest() + return buildString { + shaBytes.forEach { byte -> append(String.format("%1$02x", byte)) } + } + } +} + +// due to compatibility with the Gradle plugin and kotlin 1.6 +private fun Path.exists() = toFile().exists() +private fun Path.inputStream() = toFile().inputStream() +private fun Path.absolute() = toFile().isAbsolute +private val Path.name: String + get() { + return toFile().name + } + +private fun Path.fileSize() = Files.size(this) + +private suspend fun HttpClient.downloadFile(url: String, destination: Path, queryParams: Map): Path { + prepareGet(url) { + url { + queryParams.forEach { + parameters.append(it.key, it.value) + } + } + expectSuccess = false + }.execute { response -> + if (!response.status.isSuccess()) { + withContext(Dispatchers.IO) { Files.deleteIfExists(destination) } + error("Couldn't download $url. Status code: ${response.status}. Body: ${response.bodyAsText()}") + } + + val channel: ByteReadChannel = response.body() + withContext(Dispatchers.IO) { + destination.outputStream(StandardOpenOption.CREATE).buffered().use { fos -> + channel.copyTo(fos) + } + } + } + + return destination +} + +private fun formatFileSize(fileSize: Long): String { + if (fileSize == 0L) return "0 B" + val rank = ((log10(fileSize.toDouble()) + 0.0000021714778384307465) / 3).toInt() + val value = fileSize / 1000.0.pow(rank.toDouble()) + val units = arrayOf("B", "kB", "MB", "GB", "TB", "PB", "EB") + return DecimalFormat("0.##").format(value) + units[rank] +} + +suspend fun codeCacheParams(machineId: Deferred?, build: String?): Map { + return buildMap { + machineId?.await()?.let { put("uuid", it) } + build?.let { put("build", it) } + } +} \ No newline at end of file diff --git a/fleet/codecache/test/BUILD.bazel b/fleet/codecache/test/BUILD.bazel new file mode 100644 index 000000000000..f67b15c496b5 --- /dev/null +++ b/fleet/codecache/test/BUILD.bazel @@ -0,0 +1,44 @@ +### auto-generated section `build fleet.codecache.test` start +load("//build:compiler-options.bzl", "create_kotlinc_options") +load("@rules_jvm//:jvm.bzl", "jvm_library") + +create_kotlinc_options( + name = "custom_test", + opt_in = [ + "kotlinx.coroutines.ExperimentalCoroutinesApi", + "kotlin.ExperimentalStdlibApi", + ], + x_consistent_data_class_copy_visibility = True, + x_context_parameters = True, + x_jvm_default = "all-compatibility", + x_lambdas = "class" +) + +jvm_library( + name = "test_test_lib", + module_name = "fleet.codecache.test", + visibility = ["//visibility:public"], + srcs = glob(["src/**/*.kt", "src/**/*.java", "src/**/*.form"], allow_empty = True, exclude = ["**/module-info.java"]), + kotlinc_opts = ":custom_test", + deps = [ + "@lib//:kotlin-stdlib", + "//fleet/junit4", + "//fleet/codecache", + "@lib//:kotlinx-coroutines-core", + "//fleet/bundles", + "@lib//:kotlin-test", + "//fleet/util/network", + "//fleet/ktor/client/core", + "@lib//:kotlinx-serialization-core", + ] +) +### auto-generated section `build fleet.codecache.test` end + +### auto-generated section `test fleet.codecache.test` start +load("@community//build:tests-options.bzl", "jps_test") + +jps_test( + name = "test_test", + runtime_deps = [":test_test_lib"] +) +### auto-generated section `test fleet.codecache.test` end \ No newline at end of file diff --git a/fleet/codecache/test/fleet.codecache.test.iml b/fleet/codecache/test/fleet.codecache.test.iml new file mode 100644 index 000000000000..e381268b4298 --- /dev/null +++ b/fleet/codecache/test/fleet.codecache.test.iml @@ -0,0 +1,45 @@ + + + + + + + + + + + + + + + + + $KOTLIN_BUNDLED$/lib/kotlinx-serialization-compiler-plugin.jar + + + + + + + + + + + + + + + + + + + + + + + + + + + \ No newline at end of file diff --git a/fleet/codecache/test/gradlebuild/build.gradle.kts b/fleet/codecache/test/gradlebuild/build.gradle.kts new file mode 100644 index 000000000000..e76d130393e0 --- /dev/null +++ b/fleet/codecache/test/gradlebuild/build.gradle.kts @@ -0,0 +1,70 @@ +// IMPORT__MARKER_START +import fleet.buildtool.conventions.configureAtMostOneJvmTargetOrThrow +import fleet.buildtool.conventions.withJavaSourceSet +// IMPORT__MARKER_END +plugins { + alias(libs.plugins.kotlin.multiplatform) + id("fleet.project-module-conventions") + id("fleet.toolchain-conventions") + alias(libs.plugins.dokka) + // GRADLE_PLUGINS__MARKER_START + id("fleet-module") + alias(jps.plugins.kotlin.serialization) + // GRADLE_PLUGINS__MARKER_END +} + +fleetModule { + module { + name = "fleet.codecache.test" + importedFromJps {} + test {} + } +} + +@OptIn(org.jetbrains.kotlin.gradle.ExperimentalWasmDsl::class) +kotlin { + // KOTLIN__MARKER_START + compilerOptions.freeCompilerArgs = listOf( + "-opt-in=kotlinx.coroutines.ExperimentalCoroutinesApi", + "-opt-in=kotlin.ExperimentalStdlibApi", + "-Xlambdas=class", + "-Xconsistent-data-class-copy-visibility", + "-Xcontext-parameters", + "-XXLanguage:+AllowEagerSupertypeAccessibilityChecks", + ) + jvm {} + sourceSets.jvmTest.configure { kotlin.srcDir(layout.projectDirectory.dir("../src")) } + configureAtMostOneJvmTargetOrThrow { compilations.named("test") { withJavaSourceSet { javaSourceSet -> javaSourceSet.java.srcDir(layout.projectDirectory.dir("../src")) } } } + sourceSets.commonMain.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcCommonMain")) } + sourceSets.commonMain.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesCommonMain")) } + sourceSets.commonTest.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcCommonTest")) } + sourceSets.commonTest.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesCommonTest")) } + sourceSets.jvmMain.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcJvmMain")) } + configureAtMostOneJvmTargetOrThrow { compilations.named("main") { withJavaSourceSet { javaSourceSet -> javaSourceSet.java.srcDir(layout.projectDirectory.dir("../srcJvmMain")) } } } + sourceSets.jvmMain.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesJvmMain")) } + sourceSets.jvmTest.configure { kotlin.srcDir(layout.projectDirectory.dir("../srcJvmTest")) } + configureAtMostOneJvmTargetOrThrow { compilations.named("test") { withJavaSourceSet { javaSourceSet -> javaSourceSet.java.srcDir(layout.projectDirectory.dir("../srcJvmTest")) } } } + sourceSets.jvmTest.configure { resources.srcDir(layout.projectDirectory.dir("../resourcesJvmTest")) } + sourceSets.commonTest.dependencies { + implementation(jps.org.jetbrains.kotlin.kotlin.stdlib1993400674.get().let { "${it.group}:${it.name}:${it.version}" }) { + exclude(group = "org.jetbrains", module = "annotations") + } + implementation(jps.com.intellij.platform.kotlinx.coroutines.core.jvm134738847.get().let { "${it.group}:kotlinx-coroutines-core:${it.version}" }) { + isTransitive = false + } + implementation(jps.org.jetbrains.kotlin.kotlin.test542871666.get().let { "${it.group}:${it.name}:${it.version}" }) { + exclude(group = "org.jetbrains.kotlin", module = "kotlin-stdlib") + } + implementation(jps.org.jetbrains.kotlinx.kotlinx.serialization.core.jvm1739247612.get().let { "${it.group}:kotlinx-serialization-core:${it.version}" }) { + isTransitive = false + } + implementation(project(":fleet.codecache")) + implementation(project(":fleet.bundles")) + implementation(project(":fleet.util.network")) + implementation(project(":fleet.ktor.client.core")) + } + sourceSets.jvmTest.dependencies { + implementation(project(":fleet.junit4")) + } + // KOTLIN__MARKER_END +} \ No newline at end of file diff --git a/fleet/codecache/test/src/fleet/codecache/test/CodeCacheFileLockTest.kt b/fleet/codecache/test/src/fleet/codecache/test/CodeCacheFileLockTest.kt new file mode 100644 index 000000000000..12f087892b2b --- /dev/null +++ b/fleet/codecache/test/src/fleet/codecache/test/CodeCacheFileLockTest.kt @@ -0,0 +1,113 @@ +package fleet.codecache.test + +import fleet.bundles.Coordinates +import fleet.bundles.ResolutionException +import fleet.codecache.CodeCache +import fleet.codecache.CodeCachePath +import io.ktor.client.HttpClient +import kotlinx.coroutines.* +import kotlinx.coroutines.sync.Semaphore +import kotlin.test.Test +import java.nio.file.Files +import java.nio.file.Path +import java.util.concurrent.atomic.AtomicBoolean +import kotlin.test.fail + +class CodeCacheFileLockTest { + @Test + fun `consequential requests do not hang`() { + val testTimeout = 1000L + codeCacheTest(testTimeout, testTimeout / 2) { cc, cacheDir -> + val filename = "file.txt" + val mockCoord = Coordinates.Local(filename) + + cc.withFileLock(cacheDir, filename, mockCoord) {} + cc.withFileLock(cacheDir, filename, mockCoord) {} + } + } + + @Test + fun `simultaneous requests stress test`() { + val testTimeout = 10000L + codeCacheTest(testTimeout, testTimeout) { cc, cacheDir -> + val acquired = AtomicBoolean(false) + + repeat(50) { + val filename = "file.txt" + val mockCoord = Coordinates.Local(filename) + + launch { + cc.withFileLock(cacheDir, filename, mockCoord) { + try { + if (acquired.getAndSet(true)) fail("More than one client") + delay(10) + } + finally { + acquired.set(false) + } + } + } + } + } + } + + @Test + fun `forgotten lock file ends up with exception`() { + val testTimeout = 3000L + codeCacheTest(testTimeout, testTimeout / 3) { cc, cacheDir -> + val filename = "file.txt" + val mockCoord = Coordinates.Local(filename) + + val semaphore = Semaphore(1, 1) + // lock the file "forever" + val lockJob = launch { + cc.withFileLock(cacheDir, filename, mockCoord) { + semaphore.release() + delay(testTimeout * 2) + } + } + + semaphore.acquire() + //file here is already locked + try { + cc.withFileLock(cacheDir, filename, mockCoord) { + fail("This code should not be reached, the file must be still locked") + } + } + catch (e: ResolutionException) { + //this is ok, cancel the "forever" job, otherwise it will be locked + lockJob.cancel() + } + } + } + + @Test + fun `exception does not break locking mechanism`() { + val testTimeout = 1000L + codeCacheTest(testTimeout, testTimeout / 2) { cc, cacheDir -> + val filename = "file.txt" + val mockCoord = Coordinates.Local(filename) + + try { + cc.withFileLock(cacheDir, filename, mockCoord) { + throw RuntimeException() + } + } + catch (_: RuntimeException) { + } + + cc.withFileLock(cacheDir, filename, mockCoord) {} + } + } + + private fun codeCacheTest(testTimeout: Long, maxLockWait: Long, body: suspend CoroutineScope.(CodeCache, Path) -> Unit) { + runBlocking { + withTimeout(testTimeout) { + val cacheDir = Files.createTempDirectory("codeCacheTest") + val codeCachePaths = listOf(CodeCachePath(cacheDir, writable = true)) + val cc = CodeCache({ HttpClient { } }, codeCachePaths, maxLockWaitTime = maxLockWait) + body(cc, cacheDir) + } + } + } +}