jps bazel compiler - use OpenTelemetry

GitOrigin-RevId: 0f13a7d1d06f8c9f51c963be96156fa74e2452bd
This commit is contained in:
Vladimir Krivosheev
2025-01-21 04:29:09 +00:00
committed by intellij-monorepo-bot
parent 28ba7c357d
commit 2111b45002
29 changed files with 557 additions and 688 deletions
+7
View File
@@ -184,4 +184,11 @@ jvm_import(
visibility = ["//visibility:public"],
)
jvm_import(
name = "opentelemetry-exporter-logging-otlp",
jar = "@opentelemetry-exporter-logging-otlp//file",
source_jar = "@opentelemetry-exporter-logging-otlp-sources//file",
visibility = ["//visibility:public"],
)
### auto-generated section `libraries` end
+13
View File
@@ -313,6 +313,19 @@ http_file(
downloaded_file_path = "arrow-memory-netty-buffer-patch-18.1.0-sources.jar",
)
http_file(
name = "opentelemetry-exporter-logging-otlp",
url = "https://cache-redirector.jetbrains.com/repo1.maven.org/maven2/io/opentelemetry/opentelemetry-exporter-logging-otlp/1.45.0/opentelemetry-exporter-logging-otlp-1.45.0.jar",
sha256 = "7dcadc98f7f5bafa48d95c5d650d8247ad094812543072387ddee0adaafeff65",
downloaded_file_path = "opentelemetry-exporter-logging-otlp-1.45.0.jar",
)
http_file(
name = "opentelemetry-exporter-logging-otlp-sources",
url = "https://cache-redirector.jetbrains.com/repo1.maven.org/maven2/io/opentelemetry/opentelemetry-exporter-logging-otlp/1.45.0/opentelemetry-exporter-logging-otlp-1.45.0-sources.jar",
sha256 = "3fe5a2864d093cbb91a00661457493cc9bdef99b6e1d68f67e7c84f273b97c1d",
downloaded_file_path = "opentelemetry-exporter-logging-otlp-1.45.0-sources.jar",
)
### auto-generated section `libraries` end
# Test Libraries
+4
View File
@@ -11,6 +11,10 @@
"08edf341dfa4dd0a9e15b83a2a74850baad3a4a3ca2f93aed05ca6a481b1f394",
"02f7ac954f3a9574b37a54d7bdc79190325b9ee38a3244256a84d6d760eb6f1f"
],
"io.opentelemetry:opentelemetry-exporter-logging-otlp:1.45.0": [
"7dcadc98f7f5bafa48d95c5d650d8247ad094812543072387ddee0adaafeff65",
"3fe5a2864d093cbb91a00661457493cc9bdef99b6e1d68f67e7c84f273b97c1d"
],
"org.apache.arrow:arrow-algorithm:18.1.0": [
"8627d00c1fa3341ff49d894fb88f003e9a7ca4a3846527e31d881c162cf146ff",
"d5c6c4ed38c33a55e126f86bccfb4d7c87341d23b5f820392b332e011d22c7a0"
+3 -1
View File
@@ -50,4 +50,6 @@
runtimeDeps:
- ":arrow-memory-netty-buffer-patch"
- id: org.apache.arrow:arrow-memory-netty-buffer-patch
version: 18.1.0
version: 18.1.0
- id: io.opentelemetry:opentelemetry-exporter-logging-otlp
version: 1.45.0
@@ -3,17 +3,18 @@
package org.jetbrains.bazel.jvm.jps.test
import io.opentelemetry.context.Context
import org.apache.arrow.memory.RootAllocator
import org.jetbrains.bazel.jvm.TestModules
import org.jetbrains.bazel.jvm.collectSources
import org.jetbrains.bazel.jvm.configureOpenTelemetry
import org.jetbrains.bazel.jvm.getTestWorkerPaths
import org.jetbrains.bazel.jvm.jps.buildUsingJps
import org.jetbrains.bazel.jvm.jps.configureGlobalJps
import org.jetbrains.bazel.jvm.kotlin.JvmBuilderFlags
import org.jetbrains.bazel.jvm.kotlin.parseArgs
import org.jetbrains.bazel.jvm.logging.LogWriter
import org.jetbrains.bazel.jvm.performTestInvocation
import java.io.PrintStream
import org.jetbrains.bazel.jvm.use
import java.nio.file.Files
import java.security.MessageDigest
import kotlin.io.path.ExperimentalPathApi
@@ -32,33 +33,39 @@ internal object TestJpsBuildWorker {
performTestInvocation { out, coroutineScope ->
// IDEA console is bad and outdated, write to file and use modern tooling to view logs
// ${dateTimeFormatter.format(LocalDateTime.now())}
val logFile = testPaths.userHomeDir.resolve("kotlin-worker/log.ndjson")
val logFile = testPaths.userHomeDir.resolve("kotlin-worker/log.jsonl")
Files.createDirectories(logFile.parent)
configureGlobalJps(LogWriter(coroutineScope, PrintStream(Files.newOutputStream(logFile)), closeWriterOnShutdown = true))
val tracer = configureOpenTelemetry(Files.newOutputStream(logFile), "test-builder").getTracer("test-builder")
configureGlobalJps(tracer, coroutineScope)
val args = parseArgs(testParams.lines().toTypedArray())
val messageDigest = MessageDigest.getInstance("SHA-256")
RootAllocator(Long.MAX_VALUE).use { allocator ->
buildUsingJps(
baseDir = baseDir,
args = args,
out = out,
sources = sources,
dependencyFileToDigest = args.optionalList(JvmBuilderFlags.CLASSPATH).associate {
val file = baseDir.resolve(it).normalize()
val digest = messageDigest.digest(Files.readAllBytes(file))
messageDigest.reset()
file to digest
},
sourceFileToDigest = sources.associate {
val file = baseDir.resolve(it).normalize()
val digest = messageDigest.digest(Files.readAllBytes(file))
messageDigest.reset()
file to digest
},
isDebugEnabled = true,
allocator = allocator,
)
tracer.spanBuilder("build").use { span ->
buildUsingJps(
baseDir = baseDir,
args = args,
out = out,
sources = sources,
dependencyFileToDigest = args.optionalList(JvmBuilderFlags.CLASSPATH).associate {
val file = baseDir.resolve(it).normalize()
val digest = messageDigest.digest(Files.readAllBytes(file))
messageDigest.reset()
file to digest
},
sourceFileToDigest = sources.associate {
val file = baseDir.resolve(it).normalize()
val digest = messageDigest.digest(Files.readAllBytes(file))
messageDigest.reset()
file to digest
},
isDebugEnabled = true,
allocator = allocator,
tracingContext = Context.current(),
parentSpan = span,
tracer = tracer,
)
}
}
}
}
@@ -10,6 +10,7 @@ kt_jvm_library(
deps = [
"@lib//:kotlin-stdlib",
"@lib//:fastutil-min",
"@lib//:opentelemetry",
"//src/compiler-util",
"//src/worker-framework",
"@rules_java//java/runfiles",
+134 -84
View File
@@ -6,6 +6,11 @@ package org.jetbrains.bazel.jvm.jps
import com.intellij.openapi.diagnostic.IdeaLogRecordFormatter
import com.intellij.openapi.diagnostic.Logger
import com.intellij.openapi.util.io.FileUtilRt
import com.intellij.openapi.util.text.HtmlChunk.span
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.trace.Span
import io.opentelemetry.api.trace.Tracer
import io.opentelemetry.context.Context
import it.unimi.dsi.fastutil.objects.Object2ObjectArrayMap
import kotlinx.coroutines.*
import org.apache.arrow.memory.RootAllocator
@@ -24,8 +29,8 @@ import org.jetbrains.bazel.jvm.jps.state.saveBuildState
import org.jetbrains.bazel.jvm.kotlin.ArgMap
import org.jetbrains.bazel.jvm.kotlin.JvmBuilderFlags
import org.jetbrains.bazel.jvm.kotlin.parseArgs
import org.jetbrains.bazel.jvm.logging.LogWriter
import org.jetbrains.bazel.jvm.processRequests
import org.jetbrains.bazel.jvm.use
import org.jetbrains.jps.api.CanceledStatus
import org.jetbrains.jps.api.GlobalOptions
import org.jetbrains.jps.backwardRefs.JavaBackwardReferenceIndexBuilder
@@ -36,8 +41,6 @@ import org.jetbrains.jps.incremental.ModuleBuildTarget
import org.jetbrains.jps.incremental.RebuildRequestedException
import org.jetbrains.jps.incremental.java.JavaBuilder
import org.jetbrains.jps.incremental.relativizer.PathRelativizerService
import org.jetbrains.jps.incremental.storage.PathTypeAwareRelativizer
import org.jetbrains.jps.incremental.storage.RelativePathType
import org.jetbrains.jps.model.JpsModel
import org.jetbrains.kotlin.config.IncrementalCompilation
import org.jetbrains.kotlin.jps.incremental.KotlinCompilerReferenceIndexBuilder
@@ -61,8 +64,11 @@ private fun configureKotlincHome() {
System.setProperty("jps.kotlin.home", singleFile.parent.toString())
}
fun configureGlobalJps(logWriter: LogWriter) {
Logger.setFactory { BazelLogger(category = IdeaLogRecordFormatter.smartAbbreviate(it), writer = logWriter) }
fun configureGlobalJps(tracer: Tracer, scope: CoroutineScope) {
val globalSpanForIJLogger = tracer.spanBuilder("global").startSpan()
scope.coroutineContext.job.invokeOnCompletion { globalSpanForIJLogger.end() }
Logger.setFactory { BazelLogger(category = IdeaLogRecordFormatter.smartAbbreviate(it), span = globalSpanForIJLogger) }
System.setProperty("jps.service.manager.impl", BazelJpsServiceManager::class.java.name)
System.setProperty("jps.backward.ref.index.builder.fs.case.sensitive", "true")
System.setProperty(GlobalOptions.COMPILE_PARALLEL_MAX_THREADS_OPTION, Runtime.getRuntime().availableProcessors().toString())
@@ -74,7 +80,7 @@ fun configureGlobalJps(logWriter: LogWriter) {
configureKotlincHome()
}
class JpsBuildWorker private constructor(private val allocator: RootAllocator) : WorkRequestExecutor {
internal class JpsBuildWorker private constructor(private val allocator: RootAllocator) : WorkRequestExecutor {
companion object {
@JvmStatic
fun main(startupArgs: Array<String>) {
@@ -82,13 +88,14 @@ class JpsBuildWorker private constructor(private val allocator: RootAllocator) :
processRequests(
startupArgs = startupArgs,
executor = JpsBuildWorker(allocator),
setup = { configureGlobalJps(it) },
setup = { tracer, scope -> configureGlobalJps(tracer = tracer, scope = scope) },
serviceName = "jps-builder"
)
}
}
}
override suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path): Int {
override suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path, tracingContext: Context, tracer: Tracer): Int {
val dependencyFileToDigest = hashMap<Path, ByteArray>()
val sourceFileToDigest = hashMap<Path, ByteArray>(request.inputs.size)
val sources = ArrayList<Path>()
@@ -103,16 +110,30 @@ class JpsBuildWorker private constructor(private val allocator: RootAllocator) :
}
}
return buildUsingJps(
baseDir = baseDir,
args = parseArgs(request.arguments),
out = writer,
sources = sources,
dependencyFileToDigest = dependencyFileToDigest,
isDebugEnabled = request.verbosity > 0,
sourceFileToDigest = sourceFileToDigest,
allocator = allocator,
)
val isDebugEnabled = request.verbosity > 0
if (isDebugEnabled) {
tracer.spanBuilder("build")
.setParent(tracingContext)
.setAttribute(AttributeKey.stringArrayKey("sourceFileToDigest"), sourceFileToDigest.map { it.toString() })
.setAttribute(AttributeKey.stringArrayKey("dependencyFileToDigest"), dependencyFileToDigest.map { it.toString() })
}
else {
tracer.spanBuilder("build")
}.use { span ->
return buildUsingJps(
baseDir = baseDir,
args = parseArgs(request.arguments),
out = writer,
sources = sources,
dependencyFileToDigest = dependencyFileToDigest,
isDebugEnabled = isDebugEnabled,
sourceFileToDigest = sourceFileToDigest,
allocator = allocator,
parentSpan = span,
tracer = tracer,
tracingContext = tracingContext.with(span),
)
}
}
}
@@ -126,8 +147,11 @@ suspend fun buildUsingJps(
sourceFileToDigest: Map<Path, ByteArray>,
isDebugEnabled: Boolean,
allocator: RootAllocator,
parentSpan: Span,
tracer: Tracer,
tracingContext: Context,
): Int {
val log = RequestLog(out = out, isDebugEnabled = isDebugEnabled)
val log = RequestLog(out = out, parentSpan = parentSpan, tracer = tracer)
val abiJar = args.optionalSingle(JvmBuilderFlags.ABI_OUT)?.let { baseDir.resolve(it).normalize() }
val outJar = baseDir.resolve(args.mandatorySingle(JvmBuilderFlags.OUT)).normalize()
@@ -143,7 +167,7 @@ suspend fun buildUsingJps(
args = args,
classPathRootDir = baseDir,
classOutDir = classOutDir,
dependencyFileToDigest = dependencyFileToDigest
dependencyFileToDigest = dependencyFileToDigest,
)
val moduleTarget = BazelModuleBuildTarget(
outDir = classOutDir,
@@ -151,6 +175,14 @@ suspend fun buildUsingJps(
sources = sources,
)
if (isDebugEnabled) {
parentSpan.setAttribute("outJar", outJar.toString())
parentSpan.setAttribute("abiJar", abiJar?.toString() ?: "")
for (kind in TargetConfigurationDigestProperty.entries) {
parentSpan.setAttribute(kind.name, targetDigests.get(kind))
}
}
val relativizer = createPathRelativizer(baseDir = baseDir, classOutDir = classOutDir)
// if class output dir doesn't exist, make sure that we do not to use existing cache -
@@ -167,32 +199,36 @@ suspend fun buildUsingJps(
val buildStateFile = dataDir.resolve("$prefix-state-v1.arrow")
val typeAwareRelativizer = relativizer.typeAwareRelativizer!!
val buildState = if (isRebuild) {
null
}
else {
loadBuildState(
buildStateFile = buildStateFile,
relativizer = typeAwareRelativizer,
allocator = allocator,
val buildState = tracer.spanBuilder("load and check state").setParent(tracingContext).use { parentSpan ->
val buildState = if (isRebuild) {
null
}
else {
loadBuildState(
buildStateFile = buildStateFile,
relativizer = typeAwareRelativizer,
allocator = allocator,
actualDigestMap = sourceFileToDigest,
targetDigests = targetDigests,
parentSpan = parentSpan,
)
}
val forceFullRebuild = buildState != null && checkIsFullRebuildRequired(
buildState = buildState,
log = log,
actualDigestMap = sourceFileToDigest,
targetDigests = targetDigests,
sourceFileCount = sourceFileToDigest.size,
parentSpan = parentSpan,
)
}
val forceFullRebuild = buildState != null && checkIsFullRebuildRequired(
buildState = buildState,
log = log,
typeAwareRelativizer = typeAwareRelativizer,
sourceFileCount = sourceFileToDigest.size,
)
if (forceFullRebuild) {
FileUtilRt.deleteRecursively(dataDir)
FileUtilRt.deleteRecursively(classOutDir)
if (forceFullRebuild) {
FileUtilRt.deleteRecursively(dataDir)
FileUtilRt.deleteRecursively(classOutDir)
isRebuild = true
}
isRebuild = true
buildState
}
var exitCode = initAndBuild(
@@ -215,6 +251,8 @@ suspend fun buildUsingJps(
isCleanBuild = isRebuild,
),
buildState = buildState.takeIf { !isRebuild },
tracingContext = tracingContext,
parentSpan = parentSpan,
)
if (exitCode == -1) {
log.resetState()
@@ -238,6 +276,8 @@ suspend fun buildUsingJps(
isCleanBuild = true,
),
buildState = null,
tracingContext = tracingContext,
parentSpan = parentSpan,
)
}
@@ -247,26 +287,28 @@ suspend fun buildUsingJps(
private fun checkIsFullRebuildRequired(
buildState: LoadStateResult,
log: RequestLog,
typeAwareRelativizer: PathTypeAwareRelativizer,
sourceFileCount: Int
sourceFileCount: Int,
parentSpan: Span
): Boolean {
if (buildState.rebuildRequested != null) {
log.warn(buildState.rebuildRequested)
parentSpan.setAttribute("rebuildRequested", buildState.rebuildRequested)
return true
}
if (log.isDebugEnabled) {
log.info(
"changed ${buildState.changedFiles.size} files: ${buildState.changedFiles.joinToString(separator = "\n") { typeAwareRelativizer.toRelative(it, RelativePathType.SOURCE) }.prependIndent(" ")}" +
"\ndeleted ${buildState.deletedFiles.size} files: ${buildState.deletedFiles.joinToString(separator = "\n") { typeAwareRelativizer.toRelative(it.sourceFile, RelativePathType.SOURCE) }.prependIndent(" ")}"
)
}
val incrementalEffort = buildState.changedFiles.size + buildState.deletedFiles.size
val rebuildThreshold = sourceFileCount * thresholdPercentage
val forceFullRebuild = incrementalEffort >= rebuildThreshold
log.info("incrementalEffort=$incrementalEffort, rebuildThreshold=$rebuildThreshold, isFullRebuild=$forceFullRebuild")
log.out.appendLine("incrementalEffort=$incrementalEffort, rebuildThreshold=$rebuildThreshold, isFullRebuild=$forceFullRebuild")
if (parentSpan.isRecording) {
// do not use toRelative - print as is to show the actual path
parentSpan.setAttribute(AttributeKey.stringArrayKey("changedFiles"), buildState.changedFiles.map { it.toString() })
parentSpan.setAttribute(AttributeKey.stringArrayKey("deletedFiles"), buildState.deletedFiles.map { it.toString() })
parentSpan.setAttribute("incrementalEffort", incrementalEffort.toLong())
parentSpan.setAttribute("rebuildThreshold", rebuildThreshold)
parentSpan.setAttribute("forceFullRebuild", forceFullRebuild)
}
return forceFullRebuild
}
@@ -283,19 +325,19 @@ private suspend fun initAndBuild(
jpsModel: JpsModel,
buildDataProvider: BazelBuildDataProvider,
buildState: LoadStateResult?,
tracingContext: Context,
parentSpan: Span,
): Int {
if (messageHandler.isDebugEnabled) {
messageHandler.info("build (isRebuild=$isRebuild)")
}
val tracer = messageHandler.tracer
val storageInitializer = StorageInitializer(dataDir = dataDir, classOutDir = classOutDir)
val storageManager = if (isRebuild) {
storageInitializer.clearAndInit(messageHandler)
val storageManager = tracer.spanBuilder("init storage").setParent(tracingContext).use { span ->
if (isRebuild) {
storageInitializer.clearAndInit(span)
}
else {
storageInitializer.init(span)
}
}
else {
storageInitializer.init(messageHandler)
}
try {
val projectDescriptor = storageInitializer.createProjectDescriptor(
messageHandler = messageHandler,
@@ -303,6 +345,7 @@ private suspend fun initAndBuild(
moduleTarget = moduleTarget,
relativizer = relativizer,
buildDataProvider = buildDataProvider,
span = parentSpan,
)
try {
val compileScope = CompileScopeImpl(
@@ -323,23 +366,29 @@ private suspend fun initAndBuild(
},
)
val builders = arrayOf(
JavaBuilder(BazelSharedThreadPool),
//NotNullInstrumentingBuilder(),
JavaBackwardReferenceIndexBuilder(),
BazelKotlinBuilder(isKotlinBuilderInDumbMode = false, log = messageHandler, dataManager = buildDataProvider),
KotlinCompilerReferenceIndexBuilder(),
)
builders.sortBy { it.category.ordinal }
val exitCode = JpsTargetBuilder(
log = messageHandler,
isCleanBuild = storageInitializer.isCleanBuild,
dataManager = buildDataProvider,
).build(context = context, moduleTarget = moduleTarget, builders = builders, buildState = buildState)
val exitCode = tracer.spanBuilder("compile")
.setParent(tracingContext)
.setAttribute(AttributeKey.booleanKey("isRebuild"), isRebuild)
.use { span ->
val builders = arrayOf(
JavaBuilder(BazelSharedThreadPool),
//NotNullInstrumentingBuilder(),
JavaBackwardReferenceIndexBuilder(),
BazelKotlinBuilder(isKotlinBuilderInDumbMode = false, span = span, dataManager = buildDataProvider),
KotlinCompilerReferenceIndexBuilder(),
)
builders.sortBy { it.category.ordinal }
JpsTargetBuilder(
log = messageHandler,
isCleanBuild = storageInitializer.isCleanBuild,
dataManager = buildDataProvider,
span = span,
).build(context = context, moduleTarget = moduleTarget, builders = builders, buildState = buildState)
}
if (exitCode == 0) {
try {
postBuild(
messageHandler = messageHandler,
moduleTarget = moduleTarget,
outJar = outJar,
abiJar = abiJar,
@@ -347,6 +396,8 @@ private suspend fun initAndBuild(
context = context,
targetDigests = targetDigests,
buildDataProvider = buildDataProvider,
tracingContext = tracingContext,
tracer = tracer,
)
}
catch (e: Throwable) {
@@ -362,7 +413,7 @@ private suspend fun initAndBuild(
return exitCode
}
catch (e: RebuildRequestedException) {
messageHandler.info("RebuildRequestedException: ${e.cause?.message}: ${e.stackTraceToString()}")
parentSpan.recordException(e)
return -1
}
finally {
@@ -377,7 +428,6 @@ private suspend fun initAndBuild(
private val stateFileMetaNames = arrayOf(VERSION_META_NAME) + TargetConfigurationDigestProperty.entries.map { it.name }
private suspend fun postBuild(
messageHandler: RequestLog,
moduleTarget: ModuleBuildTarget,
outJar: Path,
abiJar: Path?,
@@ -385,6 +435,8 @@ private suspend fun postBuild(
context: CompileContextImpl,
targetDigests: TargetConfigurationDigestContainer,
buildDataProvider: BazelBuildDataProvider,
tracingContext: Context,
tracer: Tracer,
) {
coroutineScope {
val dataManager = context.projectDescriptor.dataManager
@@ -400,9 +452,7 @@ private suspend fun postBuild(
relativizer = buildDataProvider.relativizer,
metadata = Object2ObjectArrayMap(
stateFileMetaNames,
arrayOf(STATE_FILE_FORMAT_VERSION) + TargetConfigurationDigestProperty.entries.map {
java.lang.Long.toUnsignedString(targetDigests.get(it), Character.MAX_RADIX)
},
arrayOf(STATE_FILE_FORMAT_VERSION) + targetDigests.asString(),
),
allocator = buildDataProvider.allocator,
)
@@ -417,13 +467,13 @@ private suspend fun postBuild(
launch(CoroutineName("create output JAR and ABI JAR")) {
// pack to jar
messageHandler.measureTime("pack and abi") {
tracer.spanBuilder("create output JAR and ABI JAR").setParent(tracingContext).use { span ->
packageToJar(
outJar = outJar,
abiJar = abiJar,
sourceDescriptors = sourceDescriptors,
classOutDir = classOutDir,
log = messageHandler,
span = span,
)
}
}
@@ -4,10 +4,15 @@
package org.jetbrains.bazel.jvm.jps
import com.intellij.openapi.util.io.FileUtilRt
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import io.opentelemetry.api.trace.StatusCode
import kotlinx.coroutines.ensureActive
import org.h2.mvstore.MVStore
import org.jetbrains.bazel.jvm.jps.impl.BazelBuildDataProvider
import org.jetbrains.bazel.jvm.jps.impl.BazelModuleBuildTarget
import org.jetbrains.bazel.jvm.jps.impl.RequestLog
import org.jetbrains.bazel.jvm.jps.impl.loadJpsProject
import org.jetbrains.jps.cmdline.ProjectDescriptor
import org.jetbrains.jps.incremental.fs.BuildFSState
@@ -32,7 +37,7 @@ internal class StorageInitializer(private val dataDir: Path, private val classOu
var isCleanBuild: Boolean = false
private set
fun clearAndInit(messageHandler: RequestLog): StorageManager {
fun clearAndInit(span: Span): StorageManager {
isCheckRebuildRequired = false
wasCleared = true
isCleanBuild = true
@@ -40,20 +45,20 @@ internal class StorageInitializer(private val dataDir: Path, private val classOu
clearStorage()
Files.createDirectories(dataDir)
val logger = createLogger(messageHandler)
val logger = createLogger(span)
val store = tryOpenMvStore(file = cacheDbFile, readOnly = false, autoCommitDelay = 0, logger = logger)
return StorageManager(cacheDbFile, store)
.also { storageManager = it }
}
suspend fun init(messageHandler: RequestLog): StorageManager {
val logger = createLogger(messageHandler)
suspend fun init(span: Span): StorageManager {
val logger = createLogger(span)
coroutineContext.ensureActive()
isCleanBuild = Files.notExists(cacheDbFile)
if (isCleanBuild && Files.isDirectory(dataDir)) {
messageHandler.info("remove $dataDir and $classOutDir because no cache db file found: $cacheDbFile")
span.addEvent("remove $dataDir and $classOutDir because no cache db file found: $cacheDbFile")
// if no db file, make sure that data dir is also not reused
deleteDirs()
}
@@ -63,7 +68,7 @@ internal class StorageInitializer(private val dataDir: Path, private val classOu
tryOpenMvStore(file = cacheDbFile, readOnly = false, autoCommitDelay = 0, logger = logger)
}
catch (e: Throwable) {
messageHandler.info("rebuild due to internal error: ${e.stackTraceToString()}")
span.recordException(e, Attributes.of(AttributeKey.stringKey("message"), "rebuild due to internal error"))
clearStorage()
return StorageManager(cacheDbFile, createStoreAfterClear(logger))
@@ -83,14 +88,11 @@ internal class StorageInitializer(private val dataDir: Path, private val classOu
return tryOpenMvStore(file = cacheDbFile, readOnly = false, autoCommitDelay = 0, logger = logger)
}
private fun createLogger(messageHandler: RequestLog): StoreLogger {
private fun createLogger(span: Span): StoreLogger {
return { m: String, e: Throwable, isWarn: Boolean ->
val message = "$m: ${e.stackTraceToString()}"
if (isWarn) {
messageHandler.warn(message)
}
else {
messageHandler.error(message)
span.recordException(e, Attributes.of(AttributeKey.stringKey("message"), m))
if (!isWarn) {
span.setStatus(StatusCode.ERROR)
}
}
}
@@ -101,6 +103,7 @@ internal class StorageInitializer(private val dataDir: Path, private val classOu
moduleTarget: BazelModuleBuildTarget,
relativizer: PathRelativizerService,
buildDataProvider: BazelBuildDataProvider,
span: Span,
): ProjectDescriptor {
try {
return loadJpsProject(
@@ -121,7 +124,9 @@ internal class StorageInitializer(private val dataDir: Path, private val classOu
throw e
}
messageHandler.warn("Cannot open cache storage: ${e.stackTraceToString()}")
span.recordException(e, Attributes.of(
AttributeKey.stringKey("message"), "cannot open cache storage",
))
}
return createProjectDescriptor(
@@ -130,6 +135,7 @@ internal class StorageInitializer(private val dataDir: Path, private val classOu
moduleTarget = moduleTarget,
relativizer = relativizer,
buildDataProvider = buildDataProvider,
span = span,
)
}
@@ -1,7 +1,9 @@
@file:Suppress("HardCodedStringLiteral", "INVISIBLE_REFERENCE", "INVISIBLE_MEMBER", "DialogTitleCapitalization", "UnstableApiUsage", "ReplaceGetOrSet")
package org.jetbrains.bazel.jvm.jps.impl
import org.jetbrains.bazel.jvm.jps.RequestLog
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import org.jetbrains.jps.ModuleChunk
import org.jetbrains.jps.builders.DirtyFilesHolder
import org.jetbrains.jps.builders.FileProcessor
@@ -64,8 +66,8 @@ private val classesToLoadByParent = ClassCondition { className ->
internal class BazelKotlinBuilder(
private val isKotlinBuilderInDumbMode: Boolean,
private val enableLookupStorageFillingInDumbMode: Boolean = false,
private val log: RequestLog,
private val dataManager: BazelBuildDataProvider,
private val span: Span,
) : ModuleLevelBuilder(BuilderCategory.SOURCE_PROCESSOR) {
companion object {
const val JPS_KOTLIN_HOME_PROPERTY = "jps.kotlin.home"
@@ -110,7 +112,7 @@ internal class BazelKotlinBuilder(
kotlinContext.reportUnsupportedTargets()
}
log.info("Total Kotlin global compile context initialization time: $time ms")
span.addEvent("total Kotlin global compile context initialization time: $time ms")
return kotlinContext
}
@@ -172,7 +174,7 @@ internal class BazelKotlinBuilder(
}
}
)
val fsOperations = BazelKotlinFsOperationsHelper(context, chunk, dirtyFilesHolder, log, dataManager = dataManager)
val fsOperations = BazelKotlinFsOperationsHelper(context, chunk, dirtyFilesHolder, span, dataManager = dataManager)
val representativeTarget = kotlinContext.targetsBinding[chunk.representativeTarget()] ?: return
@@ -217,7 +219,7 @@ internal class BazelKotlinBuilder(
changeCollector = changesCollector,
caches = incrementalCaches.values,
lookupStorageManager = kotlinContext.lookupStorageManager,
reporter = BazelJpsICReporter(log),
reporter = BazelJpsICReporter(span),
)
fsOperations.markFilesForCurrentRound(affectedByRemovedClasses.dirtyFiles.asSequence() + affectedByRemovedClasses.forceRecompileTogether)
@@ -235,8 +237,8 @@ internal class BazelKotlinBuilder(
context = context,
chunk = chunk,
dirtyFilesHolder = kotlinDirtyFilesHolder,
log = log,
dataManager = dataManager,
span = span,
)
val proposedExitCode = doBuild(
chunk = chunk,
@@ -421,7 +423,7 @@ internal class BazelKotlinBuilder(
lookupStorageManager = kotlinContext.lookupStorageManager,
fsOperations = fsOperations,
caches = incrementalCaches.values,
reporter = BazelJpsICReporter(log),
reporter = BazelJpsICReporter(span),
)
}
}
@@ -471,8 +473,8 @@ internal class BazelKotlinBuilder(
context: CompileContext,
) {
val allDirtyFiles = dirtyFilesHolder.allDirtyFiles
if (log.isDebugEnabled) {
log.debug("compiling files: $allDirtyFiles")
if (span.isRecording) {
span.addEvent("compiling files", Attributes.of(AttributeKey.stringArrayKey("allDirtyFiles"), allDirtyFiles.map { it.path }))
}
JavaBuilderUtil.registerFilesToCompile(context, allDirtyFiles)
}
@@ -591,15 +593,12 @@ internal class BazelKotlinBuilder(
}
}
private class BazelJpsICReporter(private val log: RequestLog) : ICReporterBase() {
private class BazelJpsICReporter(private val span: Span) : ICReporterBase() {
override fun reportCompileIteration(incremental: Boolean, sourceFiles: Collection<File>, exitCode: ExitCode) {
}
override fun report(message: () -> String, severity: ReportSeverity) {
// Currently, all severity levels are mapped to debug
if (log.isDebugEnabled) {
log.debug(message())
}
span.addEvent(message())
}
}
@@ -1,8 +1,11 @@
@file:Suppress("HardCodedStringLiteral", "INVISIBLE_REFERENCE", "INVISIBLE_MEMBER", "DialogTitleCapitalization", "UnstableApiUsage", "ReplaceGetOrSet")
package org.jetbrains.bazel.jvm.jps.impl
import com.intellij.openapi.util.io.FileUtilRt
import org.jetbrains.bazel.jvm.jps.RequestLog
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import org.jetbrains.jps.ModuleChunk
import org.jetbrains.jps.builders.BuildRootDescriptor
import org.jetbrains.jps.builders.FileProcessor
@@ -21,7 +24,7 @@ internal class BazelKotlinFsOperationsHelper(
private val context: CompileContext,
private val chunk: ModuleChunk,
private val dirtyFilesHolder: KotlinDirtySourceFilesHolder,
private val log: RequestLog,
private val span: Span,
private val dataManager: BazelBuildDataProvider,
) {
internal var hasMarkedDirty = false
@@ -62,7 +65,7 @@ internal class BazelKotlinFsOperationsHelper(
}
}
markFilesImpl(files, currentRound = true) { it.exists() }
markFilesImpl(files, currentRound = true, span = span) { it.exists() }
}
/**
@@ -82,7 +85,7 @@ internal class BazelKotlinFsOperationsHelper(
dirtyFileToRoot[file] = root
}
markFilesImpl(files.asSequence(), currentRound = true) { it.exists() }
markFilesImpl(files.asSequence(), currentRound = true, span = span) { it.exists() }
cleanOutputsForNewDirtyFilesInCurrentRound(target, dirtyFileToRoot)
}
@@ -100,11 +103,11 @@ internal class BazelKotlinFsOperationsHelper(
}
fun markFiles(files: Sequence<File>) {
markFilesImpl(files, currentRound = false) { it.exists() }
markFilesImpl(files, currentRound = false, span = span) { it.exists() }
}
fun markInChunkOrDependents(files: Sequence<File>, excludeFiles: Set<File>) {
markFilesImpl(files, currentRound = false) {
markFilesImpl(files, currentRound = false, span = span) {
!excludeFiles.contains(it) && it.exists()
}
}
@@ -112,6 +115,7 @@ internal class BazelKotlinFsOperationsHelper(
private inline fun markFilesImpl(
files: Sequence<File>,
currentRound: Boolean,
span: Span,
shouldMark: (File) -> Boolean
) {
val filesToMark = files.filterTo(HashSet(), shouldMark)
@@ -130,6 +134,9 @@ internal class BazelKotlinFsOperationsHelper(
for (fileToMark in filesToMark) {
FSOperations.markDirty(context, compilationRound, fileToMark)
}
log.debug("Mark dirty: $filesToMark ($compilationRound)")
span.addEvent("mark dirty", Attributes.of(
AttributeKey.stringArrayKey("filesToMark"), filesToMark.map { it.toString() },
AttributeKey.stringKey("compilationRound"), compilationRound.name,
))
}
}
@@ -5,10 +5,12 @@ package org.jetbrains.bazel.jvm.jps.impl
import com.intellij.openapi.util.text.Formats.formatDuration
import com.intellij.tracing.Tracer.start
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import it.unimi.dsi.fastutil.objects.Object2ObjectArrayMap
import it.unimi.dsi.fastutil.objects.ObjectArraySet
import kotlinx.coroutines.ensureActive
import org.jetbrains.bazel.jvm.jps.RequestLog
import org.jetbrains.bazel.jvm.jps.hashMap
import org.jetbrains.bazel.jvm.jps.linkedSet
import org.jetbrains.bazel.jvm.jps.state.LoadStateResult
@@ -41,6 +43,7 @@ import kotlin.time.toJavaDuration
internal class JpsTargetBuilder(
private val log: RequestLog,
private val span: Span,
private val isCleanBuild: Boolean,
private val dataManager: BazelBuildDataProvider,
) {
@@ -81,7 +84,7 @@ internal class JpsTargetBuilder(
val message = "Build duration: ${builder.presentableName} took ${formatDuration(time.toJavaDuration())}; " +
processedSources + " sources processed" +
(if (processedSources == 0) "" else " (${time.inWholeMilliseconds / processedSources} ms per file)")
log.info(message)
span.addEvent(message)
}
}
catch (e: StopBuildException) {
@@ -92,13 +95,23 @@ internal class JpsTargetBuilder(
return if (log.hasErrors()) 1 else 0
}
catch (e: BuildDataCorruptedException) {
log.warn("Internal caches are corrupted or have outdated format, forcing project rebuild: $e")
span.recordException(
e,
Attributes.of(
AttributeKey.stringKey("message"), "internal caches are corrupted or have outdated format, forcing project rebuild"
)
)
throw RebuildRequestedException(e)
}
catch (e: ProjectBuildException) {
val cause = e.cause
if (cause is IOException || cause is BuildDataCorruptedException || (cause is RuntimeException && cause.cause is IOException)) {
log.warn("Internal caches are corrupted or have outdated format, forcing project rebuild: $e")
span.recordException(
e,
Attributes.of(
AttributeKey.stringKey("message"), "internal caches are corrupted or have outdated format, forcing project rebuild"
)
)
throw RebuildRequestedException(cause)
}
else {
@@ -143,9 +156,9 @@ internal class JpsTargetBuilder(
if (!isCleanBuild) {
cleanOutputsCorrespondingToChangedFiles(
context = context,
log = log,
target = target,
dataManager = dataManager,
span = span,
)
}
@@ -201,7 +214,9 @@ internal class JpsTargetBuilder(
break
}
else {
log.debug("Builder ${builder.presentableName} requested second chunk rebuild")
span.addEvent("builder requested second chunk rebuild", Attributes.of(
AttributeKey.stringKey("builder"), builder.presentableName,
))
}
}
@@ -270,7 +285,12 @@ internal class JpsTargetBuilder(
val deletedFiles = buildState.deletedFiles
if (!deletedFiles.isEmpty()) {
doneSomething = deleteOutputsAssociatedWithDeletedPaths(context = context, target = target, deletedFiles = deletedFiles, log = log)
doneSomething = deleteOutputsAssociatedWithDeletedPaths(
context = context,
target = target,
deletedFiles = deletedFiles,
span = span,
)
}
}
@@ -352,7 +372,7 @@ private fun deleteOutputsAssociatedWithDeletedPaths(
context: CompileContext,
target: ModuleBuildTarget,
deletedFiles: List<RemovedFileInfo>,
log: RequestLog,
span: Span,
): Boolean {
val dirsToDelete = linkedSet<Path>()
var doneSomething = false
@@ -370,7 +390,12 @@ private fun deleteOutputsAssociatedWithDeletedPaths(
}
if (!deletedOutputFiles.isEmpty()) {
doneSomething = true
log.info("Deleted files: $deletedOutputFiles")
if (span.isRecording) {
span.addEvent(
"deleted files",
Attributes.of(AttributeKey.stringArrayKey("deletedOutputFiles"), deletedOutputFiles.map { it.toString() }),
)
}
context.processMessage(FileDeletedEvent(deletedOutputFiles.map { it.toString() }))
}
}
@@ -1,16 +1,20 @@
// Copyright 2000-2025 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
@file:Suppress("HardCodedStringLiteral")
package org.jetbrains.bazel.jvm.jps
package org.jetbrains.bazel.jvm.jps.impl
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import io.opentelemetry.api.trace.Tracer
import org.jetbrains.jps.incremental.MessageHandler
import org.jetbrains.jps.incremental.messages.BuildMessage
import org.jetbrains.jps.incremental.messages.CompilerMessage
import org.jetbrains.jps.incremental.messages.ProgressMessage
class RequestLog(
@PublishedApi @JvmField internal val out: Appendable,
@JvmField val isDebugEnabled: Boolean,
internal class RequestLog(
@JvmField val out: Appendable,
@JvmField val parentSpan: Span,
@JvmField val tracer: Tracer,
) : MessageHandler {
@Volatile
private var hasErrors = false
@@ -19,31 +23,6 @@ class RequestLog(
hasErrors = false
}
fun warn(message: String) {
out.appendLine("WARN: $message")
}
fun error(message: String) {
out.appendLine("ERROR: $message")
}
fun error(message: String, error: Throwable) {
out.appendLine("ERROR: $message\n${error.stackTraceToString().prependIndent(" ")}")
}
fun info(message: String) {
out.appendLine("INFO: $message")
}
fun debug(message: String) {
out.appendLine("DEBUG: $message")
}
inline fun measureTime(label: String, block: () -> Unit) {
val duration = kotlin.time.measureTime(block)
out.appendLine("TIME: $label: $duration")
}
override fun processMessage(message: BuildMessage) {
val messageText = when (message) {
is CompilerMessage -> {
@@ -63,10 +42,11 @@ class RequestLog(
}
if (message.kind == BuildMessage.Kind.ERROR) {
parentSpan.addEvent("compilation error", Attributes.of(AttributeKey.stringKey("message"), messageText))
out.appendLine("Error: $messageText")
hasErrors = true
}
else if (message.kind !== BuildMessage.Kind.PROGRESS || !messageText.startsWith("Compiled") && !messageText.startsWith("Copying")) {
else if (!messageText.startsWith("Compiled") && !messageText.startsWith("Copying")) {
out.appendLine(messageText)
}
}
@@ -2,10 +2,12 @@
package org.jetbrains.bazel.jvm.jps.impl
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import it.unimi.dsi.fastutil.objects.Object2ObjectArrayMap
import it.unimi.dsi.fastutil.objects.ObjectArraySet
import kotlinx.coroutines.ensureActive
import org.jetbrains.bazel.jvm.jps.RequestLog
import org.jetbrains.jps.builders.BuildRootDescriptor
import org.jetbrains.jps.builders.BuildTarget
import org.jetbrains.jps.builders.FileProcessor
@@ -37,9 +39,9 @@ internal fun initFsStateForCleanBuild(context: CompileContext, target: BuildTarg
internal fun cleanOutputsCorrespondingToChangedFiles(
context: CompileContext,
log: RequestLog,
target: BazelModuleBuildTarget,
dataManager: BazelBuildDataProvider,
span: Span,
) {
val dirsToDelete = HashSet<Path>()
val deletedOutputFiles = ArrayList<Path>()
@@ -63,14 +65,20 @@ internal fun cleanOutputsCorrespondingToChangedFiles(
}
}
catch (e: IOException) {
log.warn("cannot delete output file (sourceFile=$sourceFile): $outputFile: $e")
span.recordException(e, Attributes.of(
AttributeKey.stringKey("message"), "cannot delete output file",
AttributeKey.stringKey("sourceFile"), sourceFile.toString(),
))
}
}
}
finally {
if (outputs.isNotEmpty()) {
sourceToOutputMapping.setOutputs(sourceFile, outputs)
log.warn("Some outputs were not removed for $sourceFile source file: $outputs")
span.addEvent("some outputs were not removed", Attributes.of(
AttributeKey.stringKey("sourceFile"), sourceFile.toString(),
AttributeKey.stringArrayKey("outputs"), outputs.map { it.toString() },
))
}
}
return true
@@ -78,8 +86,10 @@ internal fun cleanOutputsCorrespondingToChangedFiles(
})
if (!deletedOutputFiles.isEmpty()) {
if (JavaBuilderUtil.isCompileJavaIncrementally(context) && log.isDebugEnabled) {
log.info("allDeletedOutputPaths: $deletedOutputFiles")
if (JavaBuilderUtil.isCompileJavaIncrementally(context) && span.isRecording) {
span.addEvent("allDeletedOutputPaths", Attributes.of(
AttributeKey.stringArrayKey("deletedOutputFiles"), deletedOutputFiles.map { it.toString() },
))
}
context.processMessage(FileDeletedEvent(deletedOutputFiles.map { it.toString() }))
@@ -107,7 +117,7 @@ internal suspend fun markTargetUpToDate(
val delta = fsState.getDelta(target)
delta.lockData()
try {
val rootToRecompile = delta.getSourceMapToRecompile()
val rootToRecompile = delta.sourceMapToRecompile
if (!rootToRecompile.isEmpty()) {
for (entry in rootToRecompile) {
dataManager.stampStorage.markAsUpToDate(entry.value)
+9 -15
View File
@@ -5,17 +5,15 @@ package org.jetbrains.bazel.jvm.jps
import com.intellij.openapi.diagnostic.DefaultLogger
import com.intellij.openapi.diagnostic.LogLevel
import com.intellij.openapi.diagnostic.Logger
import org.jetbrains.bazel.jvm.logging.LogEvent
import org.jetbrains.bazel.jvm.logging.LogWriter
import java.time.Instant
private val levelToNameMap = enumValues<LogLevel>().associateWith { if ((it == LogLevel.INFO)) null else it.name.lowercase() }
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
// klogging uses a separate coroutine for log even instead of adding to the channel
// write to System.err as Bazel expects for worker log
internal class BazelLogger(category: String, private val writer: LogWriter) : Logger() {
internal class BazelLogger(category: String, private val span: Span) : Logger() {
private val maxLevel = LogLevel.INFO
private val categoryContext = arrayOf<Any>("c", category.trimStart('#'))
private val sharedAttributes = Attributes.of(AttributeKey.stringKey("category"), category)
override fun isDebugEnabled(): Boolean = maxLevel >= LogLevel.DEBUG
@@ -33,19 +31,15 @@ internal class BazelLogger(category: String, private val writer: LogWriter) : Lo
addEvent(LogLevel.DEBUG, message, t)
}
@Suppress("ReplaceGetOrSet")
private fun addEvent(level: LogLevel, message: String, t: Throwable?) {
if (level > maxLevel) {
return
}
writer.log(LogEvent(
timestamp = Instant.now(),
message = message,
context = categoryContext,
level = levelToNameMap.get(level),
exception = t,
))
span.addEvent(message, sharedAttributes)
t?.let {
span.recordException(it, sharedAttributes)
}
}
override fun info(message: String, t: Throwable?) {
@@ -13,6 +13,7 @@ kt_jvm_library(
"//:kotlinx-coroutines-core",
"@lib//:asm",
"@lib//:fastutil-min",
"@lib//:opentelemetry",
"//zip:build-zip",
"//src/jps-builder:jps-standalone",
"//:kotlin-metadata",
@@ -3,6 +3,9 @@
package org.jetbrains.bazel.jvm.jps
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.channels.Channel
import kotlinx.coroutines.launch
@@ -43,7 +46,7 @@ suspend fun packageToJar(
abiJar: Path?,
sourceDescriptors: Array<SourceDescriptor>,
classOutDir: Path,
log: RequestLog
span: Span,
) {
//var abiJar = Path.of(outJar.toString() + ".abi.jar")
if (abiJar == null) {
@@ -53,7 +56,7 @@ suspend fun packageToJar(
sourceDescriptors = sourceDescriptors,
classOutDir = classOutDir,
abiChannel = null,
messageHandler = log,
span = span,
)
}
return
@@ -67,7 +70,7 @@ suspend fun packageToJar(
sourceDescriptors = sourceDescriptors,
classOutDir = classOutDir,
abiChannel = classChannel,
messageHandler = log,
span = span,
)
classChannel.close()
}
@@ -84,7 +87,7 @@ private suspend fun createJar(
sourceDescriptors: Array<SourceDescriptor>,
classOutDir: Path,
abiChannel: Channel<JarContentToProcess>?,
messageHandler: RequestLog,
span: Span,
) {
val packageIndexBuilder = PackageIndexBuilder()
writeZipUsingTempFile(outJar, packageIndexBuilder.indexWriter) { stream ->
@@ -119,7 +122,13 @@ private suspend fun createJar(
stream.file(nameString = path, file = file)
}
catch (_: NoSuchFileException) {
messageHandler.warn("output file exists in src-to-output mapping, but not found on disk: $path (classOutDir=$classOutDir)")
span.addEvent(
"output file exists in src-to-output mapping, but not found on disk",
Attributes.of(
AttributeKey.stringKey("path"), path,
AttributeKey.stringKey("classOutDir"), classOutDir.toString()
),
)
}
}
}
@@ -3,6 +3,9 @@
package org.jetbrains.bazel.jvm.jps.state
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import org.apache.arrow.memory.RootAllocator
import org.apache.arrow.vector.FixedSizeBinaryVector
import org.apache.arrow.vector.VarCharVector
@@ -15,7 +18,6 @@ import org.apache.arrow.vector.types.pojo.Field
import org.apache.arrow.vector.types.pojo.FieldType
import org.apache.arrow.vector.types.pojo.Schema
import org.jetbrains.annotations.VisibleForTesting
import org.jetbrains.bazel.jvm.jps.RequestLog
import org.jetbrains.bazel.jvm.jps.SourceDescriptor
import org.jetbrains.bazel.jvm.jps.emptyList
import org.jetbrains.bazel.jvm.jps.hashMap
@@ -83,8 +85,8 @@ fun loadBuildState(
relativizer = relativizer,
allocator = allocator,
actualDigestMap = actualDigestMap,
log = null,
targetDigests = null,
parentSpan = null,
)
}
@@ -92,9 +94,9 @@ internal fun loadBuildState(
buildStateFile: Path,
relativizer: PathTypeAwareRelativizer,
allocator: RootAllocator,
log: RequestLog?,
actualDigestMap: Map<Path, ByteArray>,
targetDigests: TargetConfigurationDigestContainer?,
parentSpan: Span?,
): LoadStateResult? {
try {
FileChannel.open(buildStateFile, READ_FILE_OPTION).use { fileChannel ->
@@ -105,7 +107,7 @@ internal fun loadBuildState(
if (targetDigests != null) {
val rebuildRequested = checkConfiguration(metadata = fileReader.metaData, targetDigests = targetDigests)
if (rebuildRequested != null) {
if (log == null) {
if (parentSpan == null) {
throw IOException(rebuildRequested)
}
else {
@@ -131,11 +133,14 @@ internal fun loadBuildState(
return null
}
catch (e: Throwable) {
if (log == null) {
if (parentSpan == null) {
throw e
}
log.error("cannot load $buildStateFile", e)
parentSpan.recordException(e, Attributes.of(
AttributeKey.stringKey("message"), "cannot load build state file",
AttributeKey.stringKey("buildStateFile"), buildStateFile.toString(),
))
// will be deleted by caller
return null
}
@@ -19,4 +19,10 @@ internal value class TargetConfigurationDigestContainer(
fun set(kind: TargetConfigurationDigestProperty, hash: Long) {
list[kind.ordinal] = hash
}
fun asString(): List<String> {
return TargetConfigurationDigestProperty.entries.map { kind ->
java.lang.Long.toUnsignedString(list[kind.ordinal], Character.MAX_RADIX)
}
}
}
@@ -1,6 +1,8 @@
// 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.bazel.jvm.kotlin
import io.opentelemetry.api.trace.Tracer
import io.opentelemetry.context.Context
import org.jetbrains.bazel.jvm.WorkRequest
import org.jetbrains.bazel.jvm.WorkRequestExecutor
import org.jetbrains.bazel.jvm.processRequests
@@ -11,10 +13,10 @@ object KotlinBuildWorker : WorkRequestExecutor {
@JvmStatic
fun main(startupArgs: Array<String>) {
org.jetbrains.kotlin.cli.jvm.compiler.CompileEnvironmentUtil
processRequests(startupArgs, this)
processRequests(startupArgs = startupArgs, executor = this, serviceName = "kotlin-builder")
}
override suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path): Int {
override suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path, tracingContext: Context, tracer: Tracer): Int {
val sources = request.inputs.asSequence()
.filter { it.path.endsWith(".kt") || it.path.endsWith(".java") }
.map { baseDir.resolve(it.path).normalize() }
+1
View File
@@ -9,6 +9,7 @@ kt_jvm_library(
"//src/worker-framework",
"//zip:build-zip",
"//:protobuf-java",
"@lib//:opentelemetry",
"@bazel_tools//src/main/protobuf:deps_java_proto",
],
visibility = ["//visibility:public"],
+4 -2
View File
@@ -1,6 +1,8 @@
// 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.bazel.jvm
import io.opentelemetry.api.trace.Tracer
import io.opentelemetry.context.Context
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import org.jetbrains.intellij.build.io.*
@@ -12,10 +14,10 @@ import java.nio.file.Path
object JvmWorker : WorkRequestExecutor {
@JvmStatic
fun main(startupArgs: Array<String>) {
processRequests(startupArgs, this)
processRequests(startupArgs = startupArgs, executor = this, serviceName = "jvm-worker")
}
override suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path): Int {
override suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path, tracingContext: Context, tracer: Tracer): Int {
val args = request.arguments
if (args.isEmpty()) {
writer.appendLine("Command is not specified")
@@ -1,50 +0,0 @@
// 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.bazel.jvm
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.NonCancellable
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.launch
import java.nio.file.Files
import java.nio.file.Path
import java.nio.file.StandardOpenOption
import java.time.LocalDateTime
import java.time.format.DateTimeFormatter
internal class AsyncFileLogger(file: Path, coroutineScope: CoroutineScope) {
private val logChannel = Channel<String>(Channel.UNLIMITED)
private val writer = Files.newOutputStream(file, StandardOpenOption.APPEND, StandardOpenOption.CREATE).bufferedWriter()
// NonCancellable - make sure, that we process all messages. To stop processing, close the channel.
private val job = coroutineScope.launch(Dispatchers.IO + NonCancellable) { processQueue() }
fun log(message: String) {
val sendStatus = logChannel.trySend(formatMessage(message))
require(sendStatus.isSuccess) { "Cannot log: $sendStatus" }
}
suspend fun shutdown() {
try {
logChannel.close()
job.join()
writer.appendLine(formatMessage("logger shutdown"))
}
finally {
writer.close()
}
}
/**
* Processes the log queue and writes messages to the file.
*/
private suspend fun processQueue() {
for (message in logChannel) {
writer.appendLine(message)
writer.flush()
}
}
}
private fun formatMessage(message: String): String = "[${LocalDateTime.now().format(DateTimeFormatter.ISO_LOCAL_DATE_TIME)}] $message"
@@ -9,14 +9,21 @@ java_proto_library(
kt_jvm_library(
name = "worker-framework",
srcs = glob(["*.kt", "logging/*.kt"], exclude = ["*Test.kt"]),
srcs = glob(["*.kt"], exclude = ["*Test.kt"]),
kotlinc_opts = "//:rules_jvm_bootstrap_kotlinc_options",
deps = [
"@lib//:kotlin-stdlib",
"@lib//:opentelemetry",
"@lib//:opentelemetry-semconv",
"//:opentelemetry-exporter-logging-otlp",
"//:kotlinx-coroutines-core",
"//:protobuf-java",
"@lib//:jetbrains-annotations",
],
runtime_deps = [
"@lib//:opentelemetry-exporter-otlp-common",
"@lib//:jackson",
],
visibility = ["//visibility:public"],
)
@@ -4,14 +4,25 @@
package org.jetbrains.bazel.jvm
import com.google.protobuf.CodedOutputStream
import io.opentelemetry.api.OpenTelemetry
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import io.opentelemetry.api.trace.Tracer
import io.opentelemetry.context.Context
import io.opentelemetry.exporter.logging.otlp.internal.traces.OtlpStdoutSpanExporter
import io.opentelemetry.sdk.OpenTelemetrySdk
import io.opentelemetry.sdk.resources.Resource
import io.opentelemetry.sdk.trace.SdkTracerProvider
import io.opentelemetry.sdk.trace.export.BatchSpanProcessor
import io.opentelemetry.sdk.trace.samplers.Sampler
import io.opentelemetry.semconv.ServiceAttributes
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import org.jetbrains.annotations.VisibleForTesting
import org.jetbrains.bazel.jvm.WorkRequestState.*
import org.jetbrains.bazel.jvm.logging.LogEvent
import org.jetbrains.bazel.jvm.logging.LogWriter
import java.io.*
import java.nio.file.Path
import java.util.concurrent.ConcurrentHashMap
@@ -21,46 +32,69 @@ import kotlin.coroutines.coroutineContext
import kotlin.system.exitProcess
fun interface WorkRequestExecutor {
suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path): Int
suspend fun execute(request: WorkRequest, writer: Writer, baseDir: Path, tracingContext: Context, tracer: Tracer): Int
}
fun configureOpenTelemetry(out: OutputStream, serviceName: String): OpenTelemetrySdk {
// Set up a tracer provider with the desired configuration
val spanExporter = OtlpStdoutSpanExporter.builder()
.setOutput(out)
//.setMemoryMode(MemoryMode.REUSABLE_DATA)
.build()
val batchSpanProcessor = BatchSpanProcessor.builder(spanExporter)
.build()
val resource = Resource.create(
Attributes.of(ServiceAttributes.SERVICE_NAME, serviceName)
)
val tracerProvider = SdkTracerProvider.builder()
.setResource(resource)
.setSampler(Sampler.alwaysOn())
.addSpanProcessor(batchSpanProcessor) // For batch exporting (preferred in production)
.build()
// Build and set the OpenTelemetry SDK
val openTelemetrySdk = OpenTelemetrySdk.builder()
.setTracerProvider(tracerProvider)
.build()
return openTelemetrySdk
}
private val noopTracer = OpenTelemetry.noop().getTracer("noop")
fun processRequests(
startupArgs: Array<String>,
executor: WorkRequestExecutor,
setup: (LogWriter) -> Unit = {},
serviceName: String?,
setup: (Tracer, CoroutineScope) -> Unit = { _, _ -> },
) {
if (!startupArgs.contains("--persistent_worker")) {
System.err.println("Only persistent worker mode is supported")
exitProcess(1)
}
val tracer = if (serviceName == null) {
noopTracer
}
else {
configureOpenTelemetry(System.err, serviceName).getTracer(serviceName)
}
try {
runBlocking(Dispatchers.Default) {
val log = LogWriter(this, System.err)
try {
setup(log)
tracer.spanBuilder("process requests").use { span ->
setup(tracer, this@runBlocking)
WorkRequestHandler(requestExecutor = executor, input = System.`in`, out = System.out, tracer = tracer)
.processRequests(Context.current().with(span))
}
WorkRequestHandler(requestExecutor = executor, input = System.`in`, out = System.out, log = log)
.processRequests()
}
catch (e: CancellationException) {
log.log(LogEvent(message = "cancelled", exception = e))
throw e
}
catch (e: Throwable) {
log.log(LogEvent(message = "internal error", exception = e))
}
finally {
log.shutdown()
}
exitProcess(0)
}
}
catch (e: Throwable) {
e.printStackTrace(System.err)
exitProcess(1)
}
exitProcess(0)
}
private class RequestState(@JvmField val request: WorkRequest) {
@@ -89,7 +123,7 @@ internal class WorkRequestHandler internal constructor(
* Must be quick and safe - executed in a read thread
*/
private val cancelHandler: ((Int) -> Unit)? = null,
private val log: LogWriter? = null,
private val tracer: Tracer,
) {
private val workingDir = Path.of(".").toAbsolutePath().normalize()
@@ -99,13 +133,17 @@ internal class WorkRequestHandler internal constructor(
private val activeRequests = ConcurrentHashMap<Int, RequestState>()
@OptIn(DelicateCoroutinesApi::class)
internal suspend fun processRequests() {
internal suspend fun processRequests(tracingContext: Context) {
val requestChannel = Channel<RequestState>(Channel.UNLIMITED)
try {
coroutineScope {
startTaskProcessing(requestChannel)
tracer.spanBuilder("process requests").setParent(tracingContext).use { span ->
startTaskProcessing(requestChannel, tracingContext.with(span))
}
readRequests(requestChannel)
tracer.spanBuilder("read requests").setParent(tracingContext).use { span ->
readRequests(requestChannel, span)
}
}
}
finally {
@@ -127,7 +165,7 @@ internal class WorkRequestHandler internal constructor(
}
}
private suspend fun readRequests(requestChannel: Channel<RequestState>) {
private suspend fun readRequests(requestChannel: Channel<RequestState>, span: Span) {
val inputListToReuse = ArrayList<Input>()
val argListToReuse = ArrayList<String>()
while (coroutineContext.isActive) {
@@ -137,12 +175,12 @@ internal class WorkRequestHandler internal constructor(
}
}
catch (e: InterruptedIOException) {
log?.info("stop processing", e)
span.recordException(e)
null
}
if (request == null) {
log?.info("stop processing - no more requests")
span.addEvent("stop processing - no more requests")
requestChannel.close()
break
}
@@ -186,15 +224,44 @@ internal class WorkRequestHandler internal constructor(
}
}
private fun CoroutineScope.startTaskProcessing(requestChannel: Channel<RequestState>) {
private fun CoroutineScope.startTaskProcessing(requestChannel: Channel<RequestState>, tracingContext: Context) {
repeat(Runtime.getRuntime().availableProcessors().coerceAtLeast(2)) {
launch {
for (item in requestChannel) {
//logger?.log("request(id=${item.request.requestId}) started to execute")
val stateRef = item.state
when {
stateRef.compareAndSet(NOT_STARTED, STARTED) -> {
handleRequest(request = item.request, requestState = stateRef)
val request = item.request
val span = if (request.verbosity > 0) {
tracer.spanBuilder("execute request")
.setAllAttributes(Attributes.of(
AttributeKey.stringArrayKey("arguments"), request.arguments.toList(),
AttributeKey.stringArrayKey("inputs"), request.inputs.map { it.path },
AttributeKey.longKey("id"), request.requestId.toLong(),
AttributeKey.stringKey("sandboxDir"), request.sandboxDir?.toString() ?: "",
))
.setParent(tracingContext)
.startSpan()
}
else {
null
}
try {
handleRequest(
request = request,
requestState = stateRef,
tracingContext = span?.let { tracingContext.with(it) },
parentSpan = span,
)
}
catch (e: Throwable) {
span?.recordException(e)
throw e
}
finally {
span?.end()
}
}
else -> {
val state = stateRef.get()
@@ -262,30 +329,37 @@ internal class WorkRequestHandler internal constructor(
}
}
/**
* Handles and responds to the given [WorkRequest].
*
* @throws IOException if there is an error talking to the server. Errors from calling the [][.callback] are reported with exit code 1.
*/
// visible for tests
internal suspend fun handleRequest(request: WorkRequest, requestState: AtomicReference<WorkRequestState>) {
internal suspend fun handleRequest(
request: WorkRequest,
requestState: AtomicReference<WorkRequestState>,
tracingContext: Context?,
parentSpan: Span?,
) {
val baseDir = if (request.sandboxDir.isNullOrEmpty()) workingDir else workingDir.resolve(request.sandboxDir)
var exitCode = 1
val stringWriter = StringBuilderWriter()
var errorToThrow: Throwable? = null
val requestId = request.requestId
val span = if (tracingContext == null) {
null
}
else {
tracer.spanBuilder("requestExecutor.execute").setParent(tracingContext).startSpan()
}
try {
if (request.verbosity > 0) {
log?.log(LogEvent(
message = "execute request",
context = arrayOf(
"id", requestId,
"baseDir", baseDir,
"inputs", request.inputs.asSequence().map { it.path },
),
))
try {
val tracingContext = if (tracingContext == null) Context.root() else tracingContext.with(span!!)
exitCode = requestExecutor.execute(
request = request,
writer = stringWriter,
baseDir = baseDir,
tracingContext = tracingContext,
tracer = if (tracingContext == null) noopTracer else tracer,
)
}
finally {
span?.end()
}
exitCode = requestExecutor.execute(request = request, writer = stringWriter, baseDir = baseDir)
}
catch (e: CancellationException) {
errorToThrow = e
@@ -296,12 +370,16 @@ internal class WorkRequestHandler internal constructor(
errorToThrow = e
}
}
finally {
span?.end()
}
withContext(NonCancellable) {
if (!requestState.compareAndSet(STARTED, FINISHED)) {
if (request.verbosity > 0) {
log?.info("request state was modified during processing", arrayOf("id", requestId, "state", requestState.get()))
}
parentSpan?.addEvent(
"request state was modified during processing",
Attributes.of(AttributeKey.stringKey("state"), requestState.get().name),
)
return@withContext
}
@@ -313,10 +391,10 @@ internal class WorkRequestHandler internal constructor(
)
if (request.verbosity > 0) {
log?.log(LogEvent(
message = "request processed",
context = arrayOf("id", requestId, "exitCode", exitCode, "out", outString),
))
parentSpan?.addEvent(
"request processed",
Attributes.of(AttributeKey.longKey("exitCode"), exitCode.toLong(), AttributeKey.stringKey("out"), outString),
)
}
}
@@ -17,6 +17,8 @@ package org.jetbrains.bazel.jvm
import com.google.devtools.build.lib.worker.WorkerProtocol
import com.google.devtools.build.lib.worker.WorkerProtocol.WorkResponse
import io.opentelemetry.api.OpenTelemetry
import io.opentelemetry.context.Context
import kotlinx.coroutines.*
import org.assertj.core.api.Assertions.assertThat
import org.junit.jupiter.api.Test
@@ -39,9 +41,10 @@ class WorkRequestHandlerTest {
fun normalWorkRequest() {
val out = ByteArrayOutputStream()
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ -> 1 },
requestExecutor = { args, err, _, _, _ -> 1 },
input = ByteArrayInputStream(ByteArray(0)),
out = out,
tracer = OpenTelemetry.noop().getTracer("noop"),
)
val request = WorkRequest(
@@ -53,7 +56,7 @@ class WorkRequestHandlerTest {
sandboxDir = null,
)
runBlocking {
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED))
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED), tracingContext = null, parentSpan = null)
}
val response = WorkResponse.parseDelimitedFrom(out.toByteArray().inputStream())
@@ -67,14 +70,15 @@ class WorkRequestHandlerTest {
fun multiplexWorkRequest() {
val out = ByteArrayOutputStream()
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ -> 0 },
requestExecutor = { args, err, _, _, _ -> 0 },
input = ByteArray(0).inputStream(),
out,
tracer = OpenTelemetry.noop().getTracer("noop"),
)
val request = newWorkRequest(listOf("--sources", "A.java"))
runBlocking {
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED))
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED), tracingContext = null, parentSpan = null)
}
val response = WorkResponse.parseDelimitedFrom(out.toByteArray().inputStream())
@@ -92,7 +96,7 @@ class WorkRequestHandlerTest {
val started = Semaphore(0)
val workerThreads = AtomicInteger()
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ ->
requestExecutor = { args, err, _, _, _ ->
// each call to this, runs in its own thread
workerThreads.incrementAndGet()
started.release()
@@ -109,6 +113,7 @@ class WorkRequestHandlerTest {
},
input = PipedInputStream(src),
out = OutputStream.nullOutputStream(),
tracer = OpenTelemetry.noop().getTracer("noop"),
)
useHandler(handler = handler, waitForProcessing = true, errorFilter = { it.message != "Intentional death!" }) {
@@ -126,18 +131,19 @@ class WorkRequestHandlerTest {
fun testOutput() {
val out = ByteArrayOutputStream()
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ ->
requestExecutor = { args, err, _, _, _ ->
err.appendLine("Failed!")
1
},
input = ByteArray(0).inputStream(),
out,
tracer = OpenTelemetry.noop().getTracer("noop"),
)
val args = listOf("--sources", "A.java")
val request = newWorkRequest(args, requestId = 0)
runBlocking {
handler.handleRequest(request, AtomicReference(WorkRequestState.STARTED))
handler.handleRequest(request, AtomicReference(WorkRequestState.STARTED), tracingContext = null, parentSpan = null)
}
val response = WorkResponse.parseDelimitedFrom(out.toByteArray().inputStream())
@@ -151,17 +157,18 @@ class WorkRequestHandlerTest {
fun testException() {
val out = ByteArrayOutputStream()
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ ->
requestExecutor = { args, err, _, _, _ ->
throw RuntimeException("Exploded!")
},
input = ByteArray(0).inputStream(),
out,
tracer = OpenTelemetry.noop().getTracer("noop"),
)
val args = listOf("--sources", "A.java")
val request = newWorkRequest(args, 342)
runBlocking {
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED))
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED), tracingContext = null, parentSpan = null)
}
val response = WorkResponse.parseDelimitedFrom(ByteArrayInputStream(out.toByteArray()))
@@ -179,13 +186,14 @@ class WorkRequestHandlerTest {
val dest = PipedInputStream()
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ ->
requestExecutor = { args, err, _, _, _ ->
handlerCalled = true
err.appendLine("Such work! Much progress! Wow!")
1
},
input = PipedInputStream(src),
out = PipedOutputStream(dest),
tracer = OpenTelemetry.noop().getTracer("noop"),
cancelHandler = {
cancelCalled = true
},
@@ -227,7 +235,7 @@ class WorkRequestHandlerTest {
// we force the regular handling to not finish until after we have read the cancel response, to avoid flakiness
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ ->
requestExecutor = { args, err, _, _, _ ->
// this handler waits until the main thread has sent a cancel request
handlerCalled.release()
try {
@@ -241,6 +249,7 @@ class WorkRequestHandlerTest {
},
input = PipedInputStream(src),
out = PipedOutputStream(dest),
tracer = OpenTelemetry.noop().getTracer("noop"),
cancelHandler = { i -> cancelCalled.incrementAndGet() }
)
@@ -276,7 +285,7 @@ class WorkRequestHandlerTest {
// we force the regular handling to not finish until after we have read the cancel response, to avoid flakiness
val inputStream = PipedInputStream(src)
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ ->
requestExecutor = { args, err, _, _, _ ->
try {
waitForCancel.acquire()
}
@@ -288,6 +297,7 @@ class WorkRequestHandlerTest {
},
input = inputStream,
out = PipedOutputStream(dest),
tracer = OpenTelemetry.noop().getTracer("noop"),
cancelHandler = {
cancelCalled.incrementAndGet()
},
@@ -322,13 +332,14 @@ class WorkRequestHandlerTest {
// we force the cancel request to not happen until after we have read the normal response, to avoid flakiness
val handler = WorkRequestHandler(
requestExecutor = { args, err, _ ->
requestExecutor = { args, err, _, _, _ ->
handlerCalled.release()
err.appendLine("Such work! Much progress! Wow!")
2
},
input = PipedInputStream(src),
out = PipedOutputStream(dest),
tracer = OpenTelemetry.noop().getTracer("noop"),
)
var r: WorkResponse? = null
@@ -356,15 +367,16 @@ class WorkRequestHandlerTest {
fun workRequestHandlerWithWorkRequestCallback() {
val out = ByteArrayOutputStream()
val handler = WorkRequestHandler(
requestExecutor = { request, err, _ -> request.arguments.size },
requestExecutor = { request, err, _, _, _ -> request.arguments.size },
ByteArrayInputStream(ByteArray(0)),
out,
tracer = OpenTelemetry.noop().getTracer("noop"),
)
val args = listOf("--sources", "B.java")
val request = newWorkRequest(args, requestId = 0)
runBlocking {
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED))
handler.handleRequest(request = request, requestState = AtomicReference(WorkRequestState.STARTED), tracingContext = null, parentSpan = null)
}
val response = WorkResponse.parseDelimitedFrom(ByteArrayInputStream(out.toByteArray()))
@@ -388,7 +400,7 @@ private inline fun useHandler(
) {
val processor = GlobalScope.async(Dispatchers.Default) {
try {
handler.processRequests()
handler.processRequests(Context.current())
}
catch (_: CancellationException) {
}
@@ -1,174 +0,0 @@
package org.jetbrains.bazel.jvm.logging
import java.io.PrintWriter
import java.io.Writer
import java.time.format.DateTimeFormatter
object JsonLogRenderer : LogRenderer {
override fun createStringBuilder(): StringBuilder {
val sb = StringBuilder()
sb.append("{\"@t\":\"")
return sb
}
override fun render(sb: StringBuilder, event: LogEvent) {
renderToJson(sb, event)
}
}
private fun renderToJson(sb: StringBuilder, event: LogEvent) {
sb.setLength(7)
DateTimeFormatter.ISO_INSTANT.formatTo(event.timestamp, sb)
sb.append('"')
if (event.level != null) {
sb.append(",\"@l\":\"")
sb.append(event.level)
sb.append('"')
}
if (event.message != null) {
sb.append(",\"@m\":\"")
escapeToJsonStringValue(event.message, sb)
sb.append('"')
}
if (event.messageTemplate != null) {
sb.append(",\"@m\":\"")
escapeToJsonStringValue(event.messageTemplate, sb)
sb.append('"')
}
writeCustomFields(event, sb)
if (event.exception != null) {
sb.append(",\"@x\":\"")
event.exception.printStackTrace(object : PrintWriter(JsonStringStringBuilderWriter(sb)) {
override fun println() {
sb.append("\\n")
}
override fun print(value: String) {
escapeToJsonStringValue(value, sb)
}
override fun print(c: Char) {
escapeChar(c, sb)
}
})
sb.append('"')
}
sb.append("}\n")
}
private fun writeCustomFields(event: LogEvent, sb: StringBuilder) {
val extraFields = event.context ?: return
for (i in extraFields.indices step 2) {
sb.append(",\"").append(extraFields[i] as String).append("\":")
val v = extraFields[i + 1]
when (v) {
is Int -> {
sb.append(v)
}
is List<*> -> serializeList(sb, v.asSequence())
is Array<*> -> serializeList(sb, v.asSequence())
is Sequence<*> -> serializeList(sb, v)
else -> {
sb.append('"')
escapeToJsonStringValue(v.toString(), sb)
sb.append('"')
}
}
}
}
private fun serializeList(sb: StringBuilder, v: Sequence<*>) {
sb.append('[')
for (any in v) {
sb.append('"')
escapeToJsonStringValue(any.toString(), sb)
sb.append('"')
sb.append(',')
}
sb.setLength(sb.length - 1)
sb.append(']')
}
private fun escapeToJsonStringValue(input: CharSequence, sb: StringBuilder) {
for (char in input) {
escapeChar(char, sb)
}
}
private fun escapeChar(char: Char, sb: StringBuilder) {
when (char) {
'"' -> sb.append("\\\"") // Escape double quotes
'\\' -> sb.append("\\\\") // Escape backslashes
'\b' -> sb.append("\\b") // Escape backspace
'\u000C' -> sb.append("\\f") // Escape form feed
'\n' -> sb.append("\\n") // Escape newline
'\r' -> sb.append("\\r") // Escape carriage return
'\t' -> sb.append("\\t") // Escape tab
in '\u0000'..'\u001F' -> { // Escape other control characters as Unicode
sb.append(String.format("\\u%04x", char.code))
}
else -> sb.append(char)
}
}
private class JsonStringStringBuilderWriter(private val sb: java.lang.StringBuilder) : Writer() {
override fun write(value: Int) {
escapeChar(value.toChar(), sb)
}
override fun write(cbuf: CharArray) {
for (char in cbuf) {
escapeChar(char, sb)
}
}
override fun append(value: Char): Writer {
escapeChar(value, sb)
return this
}
override fun append(value: CharSequence): Writer {
escapeToJsonStringValue(value, sb)
return this
}
override fun append(value: CharSequence, start: Int, end: Int): Writer {
for (i in start until end) {
escapeChar(value[i], sb)
}
return this
}
override fun write(value: String) {
sb.append(value)
}
override fun write(value: String, offset: Int, length: Int) {
for (i in offset until offset + length) {
escapeChar(value[i], sb)
}
}
override fun write(value: CharArray, offset: Int, length: Int) {
for (i in offset until offset + length) {
escapeChar(value[i], sb)
}
}
override fun flush() {
}
override fun close() {
}
override fun toString(): String = sb.toString()
}
@@ -1,88 +0,0 @@
package org.jetbrains.bazel.jvm.logging
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.channels.Channel
import java.io.PrintStream
import java.time.Instant
import kotlin.time.Duration.Companion.seconds
class LogEvent(
@JvmField val timestamp: Instant = Instant.now(),
@JvmField val message: String? = null,
@JvmField val messageTemplate: String? = null,
@JvmField val level: String? = null,
@JvmField val exception: Throwable? = null,
@JvmField val context: Array<Any>? = null,
)
interface LogRenderer {
fun createStringBuilder(): StringBuilder
fun render(sb: StringBuilder, event: LogEvent)
}
class LogWriter(
coroutineScope: CoroutineScope,
private val writer: PrintStream,
private val closeWriterOnShutdown: Boolean = false,
private val renderer: LogRenderer = JsonLogRenderer,
) {
private val logChannel = Channel<LogEvent>(Channel.UNLIMITED)
// NonCancellable - make sure, that we process all messages. To stop processing, close the channel.
private val job = coroutineScope.launch(NonCancellable) { processQueue() }
fun log(event: LogEvent) {
val sendStatus = logChannel.trySend(event)
require(sendStatus.isSuccess) { "Cannot log: $sendStatus" }
}
fun info(message: String, context: Array<Any>) {
log(LogEvent(message = message, context = context))
}
fun info(message: String, exception: Throwable? = null) {
log(LogEvent(message = message, exception = exception))
}
suspend fun shutdown() {
try {
logChannel.close()
job.join()
}
finally {
if (closeWriterOnShutdown) {
writer.close()
}
else {
writer.flush()
}
}
}
private suspend fun processQueue() {
coroutineScope {
val flushRequestChannel = Channel<Unit>(capacity = 1, onBufferOverflow = BufferOverflow.DROP_OLDEST)
val flushJob = launch {
for (ignored in flushRequestChannel) {
delay(10.seconds)
withContext(Dispatchers.IO) {
writer.flush()
}
}
}
val sb = renderer.createStringBuilder()
for (event in logChannel) {
renderer.render(sb, event)
withContext(Dispatchers.IO) {
writer.append(sb)
}
flushRequestChannel.trySend(Unit)
}
flushRequestChannel.close()
flushJob.join()
}
}
}
@@ -1,82 +0,0 @@
package org.jetbrains.bazel.jvm.logging
import java.io.PrintWriter
import java.io.StringWriter
import java.time.format.DateTimeFormatter
object YamlLogRenderer : LogRenderer {
override fun createStringBuilder(): StringBuilder {
val sb = StringBuilder()
sb.append("---\n")
return sb
}
override fun render(sb: StringBuilder, event: LogEvent) {
sb.setLength(4)
sb.append("t: ")
DateTimeFormatter.ISO_INSTANT.formatTo(event.timestamp, sb)
sb.append('\n')
if (event.level != null) {
sb.append("l: ").append(event.level).append('\n')
}
val extraFields = event.context
if (!extraFields.isNullOrEmpty()) {
for (i in extraFields.indices step 2) {
sb.append(",\"").append(extraFields[i] as String).append("\":")
val v = extraFields[i + 1]
if (v is Int) {
sb.append(v)
}
else {
appendMessage(v.toString(), sb)
}
}
}
val message = event.message
if (message != null) {
sb.append("m: ")
appendMessage(message, sb)
}
if (event.messageTemplate != null) {
sb.append("mt: ")
appendMessage(event.messageTemplate, sb)
}
if (event.exception != null) {
sb.append("x: |")
val sw = StringWriter()
val pw = PrintWriter(sw)
pw.flush()
event.exception.printStackTrace(pw)
for (string in sw.buffer.lineSequence()) {
if (string.isNotEmpty()) {
sb.append("\n ").append(string)
}
}
sb.append('\n')
}
}
}
private fun appendMessage(message: String, sb: StringBuilder) {
if (message.contains('\n')) {
sb.append('|')
appendMultiline(message, sb)
}
else {
sb.append('"').append(message).append('"')
}
sb.append('\n')
}
private fun appendMultiline(message: CharSequence, sb: StringBuilder) {
for (string in message.lineSequence()) {
sb.append("\n ").append(string)
}
}
@@ -0,0 +1,35 @@
// 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.bazel.jvm
import io.opentelemetry.api.common.Attributes
import io.opentelemetry.api.trace.Span
import io.opentelemetry.api.trace.SpanBuilder
import io.opentelemetry.api.trace.StatusCode
import io.opentelemetry.semconv.ExceptionAttributes
import java.util.concurrent.CancellationException
//fun getExceptionAttributes(e: Throwable): Attributes {
// return Attributes.of(
// ExceptionAttributes.EXCEPTION_MESSAGE, e.message ?: "",
// ExceptionAttributes.EXCEPTION_STACKTRACE, e.stackTraceToString()
// )
//}
inline fun <T> SpanBuilder.use(block: (Span) -> T): T {
val span = startSpan()
try {
return block(span)
}
catch (e: CancellationException) {
span.recordException(e, Attributes.of(ExceptionAttributes.EXCEPTION_ESCAPED, true))
throw e
}
catch (e: Throwable) {
span.recordException(e, Attributes.of(ExceptionAttributes.EXCEPTION_ESCAPED, true))
span.setStatus(StatusCode.ERROR)
throw e
}
finally {
span.end()
}
}