[fleet] Move fleet.codecache to Community

(cherry picked from commit a8ecc20e305c2592b26bf1740f2228bae741755a)

FLEET-MR-6894

GitOrigin-RevId: facb856bea2f4a5a5578b6997e31b10a3b34bba3
This commit is contained in:
Titouan Bion
2025-10-09 10:46:21 +00:00
committed by intellij-monorepo-bot
parent e609b67cfd
commit accd769736
15 changed files with 933 additions and 0 deletions
+2
View File
@@ -3,6 +3,8 @@
<component name="ProjectModuleManager">
<modules>
<module fileurl="file://$PROJECT_DIR$/fleet/andel/fleet.andel.iml" filepath="$PROJECT_DIR$/fleet/andel/fleet.andel.iml" />
<module fileurl="file://$PROJECT_DIR$/fleet/codecache/fleet.codecache.iml" filepath="$PROJECT_DIR$/community/fleet/codecache/fleet.codecache.iml" />
<module fileurl="file://$PROJECT_DIR$/fleet/codecache/test/fleet.codecache.test.iml" filepath="$PROJECT_DIR$/community/fleet/codecache/test/fleet.codecache.test.iml" />
<module fileurl="file://$PROJECT_DIR$/fleet/bifurcan/fleet.bifurcan.iml" filepath="$PROJECT_DIR$/fleet/bifurcan/fleet.bifurcan.iml" />
<module fileurl="file://$PROJECT_DIR$/fleet/bundles/fleet.bundles.iml" filepath="$PROJECT_DIR$/community/fleet/bundles/fleet.bundles.iml" />
<module fileurl="file://$PROJECT_DIR$/fleet/compiler-plugins/fleet.compiler.plugins.iml" filepath="$PROJECT_DIR$/fleet/compiler-plugins/fleet.compiler.plugins.iml" />
+2
View File
@@ -222,6 +222,8 @@ community-resources
fleet/andel
fleet/bifurcan
fleet/bundles
fleet/codecache
fleet/codecache/test
fleet/compiler-plugins
fleet/fastutil
fleet/junit4
+36
View File
@@ -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
+45
View File
@@ -0,0 +1,45 @@
<?xml version="1.0" encoding="UTF-8"?>
<module type="JAVA_MODULE" version="4">
<component name="FacetManager">
<facet type="kotlin-language" name="Kotlin">
<configuration version="5" platform="JVM 21" allPlatforms="JVM [21]" useProjectSettings="false">
<compilerSettings>
<option name="additionalArguments" value="-opt-in=kotlinx.coroutines.ExperimentalCoroutinesApi -opt-in=kotlin.ExperimentalStdlibApi -Xlambdas=class -Xconsistent-data-class-copy-visibility -Xcontext-parameters -XXLanguage:+AllowEagerSupertypeAccessibilityChecks" />
</compilerSettings>
<compilerArguments>
<stringArguments>
<stringArg name="jvmTarget" arg="21" />
<stringArg name="apiVersion" arg="2.2" />
<stringArg name="languageVersion" arg="2.2" />
</stringArguments>
<arrayArguments>
<arrayArg name="pluginClasspaths">
<args>$KOTLIN_BUNDLED$/lib/kotlinx-serialization-compiler-plugin.jar</args>
</arrayArg>
<arrayArg name="pluginOptions" />
</arrayArguments>
</compilerArguments>
</configuration>
</facet>
</component>
<component name="NewModuleRootManager" inherit-compiler-output="true">
<exclude-output />
<content url="file://$MODULE_DIR$">
<sourceFolder url="file://$MODULE_DIR$/srcCommonMain" isTestSource="false" />
<sourceFolder url="file://$MODULE_DIR$/srcJvmMain" isTestSource="false" />
<excludeFolder url="file://$MODULE_DIR$/gradlebuild/build" />
</content>
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" name="kotlinx-coroutines-core" level="project" />
<orderEntry type="library" name="kotlin-stdlib" level="project" />
<orderEntry type="module" module-name="fleet.bundles" />
<orderEntry type="library" name="jetbrains-annotations" level="project" />
<orderEntry type="module" module-name="fleet.util.logging.api" />
<orderEntry type="library" name="kotlinx-serialization-core" level="project" />
<orderEntry type="library" name="kotlinx-serialization-json" level="project" />
<orderEntry type="library" name="kotlinx-collections-immutable" level="project" />
<orderEntry type="module" module-name="fleet.ktor.client.core" />
<orderEntry type="module" module-name="fleet.modules.api" />
</component>
</module>
@@ -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
}
@@ -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")
}
}
@@ -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<String>, 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<VersionInfo> = emptyList(),
)
//@fleet.kernel.plugins.InternalInPluginModules(where = ["fleet.common"])
@Serializable
data class VersionInfo(
val xmlId: String,
val version: String,
)
@@ -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<Marketplace>()
}
//@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<PluginName>, shipVersion: PluginVersion): Map<PluginName, PluginVersion> {
Marketplace.logger.debug("Querying latest versions of '$names' compatible with '$shipVersion'...")
tailrec suspend fun loop(
out: HashMap<PluginName, PluginVersion> = HashMap(),
offset: Int = 0,
): Map<PluginName, PluginVersion> {
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"
}
}
@@ -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<LayerSelector, ResolvedPluginLayer>,
sortedSelectors: List<LayerSelector>,
externalDependencies: Map<LayerSelector, List<FleetModuleLayer>>,
baseLayers: List<FleetModuleLayer>,
loadDockModuleLayer: (modulePath: Set<FleetModuleInfo>) -> FleetModuleLayer): Map<LayerSelector, PluginModulesAndResources> {
val result = HashMap<LayerSelector, PluginModulesAndResources>(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<FleetModuleInfo> {
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<FleetModule>, val layer: FleetModuleLayer, val resources: Set<ResourceBundle>)
@@ -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<CodeCachePath>,
private val lockDelay: Long = LOCK_DELAY_MS,
private val maxLockWaitTime: Long = MAX_LOCK_WAIT_MS,
private val queryParams: (suspend () -> Map<String, String>)? = 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<String, String>): 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 <T> 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<String, String>): 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<String>?, build: String?): Map<String, String> {
return buildMap {
machineId?.await()?.let { put("uuid", it) }
build?.let { put("build", it) }
}
}
+44
View File
@@ -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
@@ -0,0 +1,45 @@
<?xml version="1.0" encoding="UTF-8"?>
<module type="JAVA_MODULE" version="4">
<component name="FacetManager">
<facet type="kotlin-language" name="Kotlin">
<configuration version="5" platform="JVM 21" allPlatforms="JVM [21]" useProjectSettings="false">
<compilerSettings>
<option name="additionalArguments" value="-opt-in=kotlinx.coroutines.ExperimentalCoroutinesApi -opt-in=kotlin.ExperimentalStdlibApi -Xlambdas=class -Xconsistent-data-class-copy-visibility -Xcontext-parameters -XXLanguage:+AllowEagerSupertypeAccessibilityChecks" />
</compilerSettings>
<compilerArguments>
<stringArguments>
<stringArg name="jvmTarget" arg="21" />
<stringArg name="apiVersion" arg="2.2" />
<stringArg name="languageVersion" arg="2.2" />
</stringArguments>
<arrayArguments>
<arrayArg name="pluginClasspaths">
<args>
<arg>$KOTLIN_BUNDLED$/lib/kotlinx-serialization-compiler-plugin.jar</arg>
</args>
</arrayArg>
<arrayArg name="pluginOptions" />
</arrayArguments>
</compilerArguments>
</configuration>
</facet>
</component>
<component name="NewModuleRootManager" inherit-compiler-output="true">
<exclude-output />
<content url="file://$MODULE_DIR$">
<sourceFolder url="file://$MODULE_DIR$/src" isTestSource="true" />
<excludeFolder url="file://$MODULE_DIR$/gradlebuild/build" />
</content>
<orderEntry type="inheritedJdk" />
<orderEntry type="sourceFolder" forTests="false" />
<orderEntry type="library" scope="TEST" name="kotlin-stdlib" level="project" />
<orderEntry type="module" module-name="fleet.junit4" scope="TEST" />
<orderEntry type="module" module-name="fleet.codecache" scope="TEST" />
<orderEntry type="library" scope="TEST" name="kotlinx-coroutines-core" level="project" />
<orderEntry type="module" module-name="fleet.bundles" scope="TEST" />
<orderEntry type="library" scope="TEST" name="kotlin-test" level="project" />
<orderEntry type="module" module-name="fleet.util.network" scope="TEST" />
<orderEntry type="module" module-name="fleet.ktor.client.core" scope="TEST" />
<orderEntry type="library" scope="TEST" name="kotlinx-serialization-core" level="project" />
</component>
</module>
@@ -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
}
@@ -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)
}
}
}
}