diff --git a/build/jvm-rules/BUILD.bazel b/build/jvm-rules/BUILD.bazel index 5123cd93fea0..67d09144dc44 100644 --- a/build/jvm-rules/BUILD.bazel +++ b/build/jvm-rules/BUILD.bazel @@ -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 \ No newline at end of file diff --git a/build/jvm-rules/MODULE.bazel b/build/jvm-rules/MODULE.bazel index 59df6037d897..6aeee920ad14 100644 --- a/build/jvm-rules/MODULE.bazel +++ b/build/jvm-rules/MODULE.bazel @@ -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 diff --git a/build/jvm-rules/libs.lock.json b/build/jvm-rules/libs.lock.json index 7284569e4f5c..f8a6ec31b908 100644 --- a/build/jvm-rules/libs.lock.json +++ b/build/jvm-rules/libs.lock.json @@ -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" diff --git a/build/jvm-rules/libs.yaml b/build/jvm-rules/libs.yaml index 1027b24b6475..de4c2587d2e2 100644 --- a/build/jvm-rules/libs.yaml +++ b/build/jvm-rules/libs.yaml @@ -50,4 +50,6 @@ runtimeDeps: - ":arrow-memory-netty-buffer-patch" - id: org.apache.arrow:arrow-memory-netty-buffer-patch - version: 18.1.0 \ No newline at end of file + version: 18.1.0 + - id: io.opentelemetry:opentelemetry-exporter-logging-otlp + version: 1.45.0 \ No newline at end of file diff --git a/build/jvm-rules/src/jps-builder-test/TestJpsBuildWorker.kt b/build/jvm-rules/src/jps-builder-test/TestJpsBuildWorker.kt index 647056805663..4c4f2ce1a5d4 100644 --- a/build/jvm-rules/src/jps-builder-test/TestJpsBuildWorker.kt +++ b/build/jvm-rules/src/jps-builder-test/TestJpsBuildWorker.kt @@ -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, + ) + } } } } diff --git a/build/jvm-rules/src/jps-builder/BUILD.bazel b/build/jvm-rules/src/jps-builder/BUILD.bazel index f2ea6171ee77..2f75ca6d5b26 100644 --- a/build/jvm-rules/src/jps-builder/BUILD.bazel +++ b/build/jvm-rules/src/jps-builder/BUILD.bazel @@ -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", diff --git a/build/jvm-rules/src/jps-builder/JpsBuilder.kt b/build/jvm-rules/src/jps-builder/JpsBuilder.kt index ce158e066093..afaa86765024 100644 --- a/build/jvm-rules/src/jps-builder/JpsBuilder.kt +++ b/build/jvm-rules/src/jps-builder/JpsBuilder.kt @@ -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) { @@ -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() val sourceFileToDigest = hashMap(request.inputs.size) val sources = ArrayList() @@ -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, 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, ) } } diff --git a/build/jvm-rules/src/jps-builder/StorageInitializer.kt b/build/jvm-rules/src/jps-builder/StorageInitializer.kt index 9beb78d9c048..aea76ebf2322 100644 --- a/build/jvm-rules/src/jps-builder/StorageInitializer.kt +++ b/build/jvm-rules/src/jps-builder/StorageInitializer.kt @@ -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, ) } diff --git a/build/jvm-rules/src/jps-builder/impl/BazelKotlinBuilder.kt b/build/jvm-rules/src/jps-builder/impl/BazelKotlinBuilder.kt index c60449cdc678..6600b7de611c 100644 --- a/build/jvm-rules/src/jps-builder/impl/BazelKotlinBuilder.kt +++ b/build/jvm-rules/src/jps-builder/impl/BazelKotlinBuilder.kt @@ -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, 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()) } } diff --git a/build/jvm-rules/src/jps-builder/impl/BazelKotlinFsOperationsHelper.kt b/build/jvm-rules/src/jps-builder/impl/BazelKotlinFsOperationsHelper.kt index 8b8ef1584310..1d5cbe0e27e5 100644 --- a/build/jvm-rules/src/jps-builder/impl/BazelKotlinFsOperationsHelper.kt +++ b/build/jvm-rules/src/jps-builder/impl/BazelKotlinFsOperationsHelper.kt @@ -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) { - markFilesImpl(files, currentRound = false) { it.exists() } + markFilesImpl(files, currentRound = false, span = span) { it.exists() } } fun markInChunkOrDependents(files: Sequence, excludeFiles: Set) { - 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, 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, + )) } } \ No newline at end of file diff --git a/build/jvm-rules/src/jps-builder/impl/JpsTargetBuilder.kt b/build/jvm-rules/src/jps-builder/impl/JpsTargetBuilder.kt index e985fed4f847..28ce0f513a2d 100644 --- a/build/jvm-rules/src/jps-builder/impl/JpsTargetBuilder.kt +++ b/build/jvm-rules/src/jps-builder/impl/JpsTargetBuilder.kt @@ -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, - log: RequestLog, + span: Span, ): Boolean { val dirsToDelete = linkedSet() 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() })) } } diff --git a/build/jvm-rules/src/jps-builder/packager/RequestLog.kt b/build/jvm-rules/src/jps-builder/impl/RequestLog.kt similarity index 53% rename from build/jvm-rules/src/jps-builder/packager/RequestLog.kt rename to build/jvm-rules/src/jps-builder/impl/RequestLog.kt index fdd1f6e35741..e150ed3d4aed 100644 --- a/build/jvm-rules/src/jps-builder/packager/RequestLog.kt +++ b/build/jvm-rules/src/jps-builder/impl/RequestLog.kt @@ -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) } } diff --git a/build/jvm-rules/src/jps-builder/impl/fsState.kt b/build/jvm-rules/src/jps-builder/impl/fsState.kt index 5ce75da20a25..843c36c1242d 100644 --- a/build/jvm-rules/src/jps-builder/impl/fsState.kt +++ b/build/jvm-rules/src/jps-builder/impl/fsState.kt @@ -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() val deletedOutputFiles = ArrayList() @@ -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) diff --git a/build/jvm-rules/src/jps-builder/logger.kt b/build/jvm-rules/src/jps-builder/logger.kt index c0995f6b2568..dc00dfbc9a6a 100644 --- a/build/jvm-rules/src/jps-builder/logger.kt +++ b/build/jvm-rules/src/jps-builder/logger.kt @@ -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().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("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?) { diff --git a/build/jvm-rules/src/jps-builder/packager/BUILD.bazel b/build/jvm-rules/src/jps-builder/packager/BUILD.bazel index 4fe1e5f02134..abb3aa35ba42 100644 --- a/build/jvm-rules/src/jps-builder/packager/BUILD.bazel +++ b/build/jvm-rules/src/jps-builder/packager/BUILD.bazel @@ -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", diff --git a/build/jvm-rules/src/jps-builder/packager/packager.kt b/build/jvm-rules/src/jps-builder/packager/packager.kt index 40f5b18d194a..e61d1d427711 100644 --- a/build/jvm-rules/src/jps-builder/packager/packager.kt +++ b/build/jvm-rules/src/jps-builder/packager/packager.kt @@ -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, 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, classOutDir: Path, abiChannel: Channel?, - 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() + ), + ) } } } diff --git a/build/jvm-rules/src/jps-builder/state/BuildStateStorage.kt b/build/jvm-rules/src/jps-builder/state/BuildStateStorage.kt index 2387f3e08715..e4e7933c031f 100644 --- a/build/jvm-rules/src/jps-builder/state/BuildStateStorage.kt +++ b/build/jvm-rules/src/jps-builder/state/BuildStateStorage.kt @@ -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, 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 } diff --git a/build/jvm-rules/src/jps-builder/state/configuration.kt b/build/jvm-rules/src/jps-builder/state/configuration.kt index b76730ac93eb..008096f48153 100644 --- a/build/jvm-rules/src/jps-builder/state/configuration.kt +++ b/build/jvm-rules/src/jps-builder/state/configuration.kt @@ -19,4 +19,10 @@ internal value class TargetConfigurationDigestContainer( fun set(kind: TargetConfigurationDigestProperty, hash: Long) { list[kind.ordinal] = hash } + + fun asString(): List { + return TargetConfigurationDigestProperty.entries.map { kind -> + java.lang.Long.toUnsignedString(list[kind.ordinal], Character.MAX_RADIX) + } + } } \ No newline at end of file diff --git a/build/jvm-rules/src/kotlin-builder/KotlinBuilder.kt b/build/jvm-rules/src/kotlin-builder/KotlinBuilder.kt index 2d5d70b88a69..111d7c6e2e59 100644 --- a/build/jvm-rules/src/kotlin-builder/KotlinBuilder.kt +++ b/build/jvm-rules/src/kotlin-builder/KotlinBuilder.kt @@ -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) { 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() } diff --git a/build/jvm-rules/src/misc/BUILD.bazel b/build/jvm-rules/src/misc/BUILD.bazel index 5c0b51904efa..b32d825d32a2 100644 --- a/build/jvm-rules/src/misc/BUILD.bazel +++ b/build/jvm-rules/src/misc/BUILD.bazel @@ -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"], diff --git a/build/jvm-rules/src/misc/JvmWorker.kt b/build/jvm-rules/src/misc/JvmWorker.kt index c7ba12ae24bb..1bbd11c578eb 100644 --- a/build/jvm-rules/src/misc/JvmWorker.kt +++ b/build/jvm-rules/src/misc/JvmWorker.kt @@ -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) { - 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") diff --git a/build/jvm-rules/src/worker-framework/AsyncFileLogger.kt b/build/jvm-rules/src/worker-framework/AsyncFileLogger.kt deleted file mode 100644 index 8f7ed42cc154..000000000000 --- a/build/jvm-rules/src/worker-framework/AsyncFileLogger.kt +++ /dev/null @@ -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(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" diff --git a/build/jvm-rules/src/worker-framework/BUILD.bazel b/build/jvm-rules/src/worker-framework/BUILD.bazel index 8fdfecd33fba..d413e1216700 100644 --- a/build/jvm-rules/src/worker-framework/BUILD.bazel +++ b/build/jvm-rules/src/worker-framework/BUILD.bazel @@ -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"], ) diff --git a/build/jvm-rules/src/worker-framework/WorkRequestHandler.kt b/build/jvm-rules/src/worker-framework/WorkRequestHandler.kt index a8563e51b713..9bbcc89e507d 100644 --- a/build/jvm-rules/src/worker-framework/WorkRequestHandler.kt +++ b/build/jvm-rules/src/worker-framework/WorkRequestHandler.kt @@ -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, 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() @OptIn(DelicateCoroutinesApi::class) - internal suspend fun processRequests() { + internal suspend fun processRequests(tracingContext: Context) { val requestChannel = Channel(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) { + private suspend fun readRequests(requestChannel: Channel, span: Span) { val inputListToReuse = ArrayList() val argListToReuse = ArrayList() 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) { + private fun CoroutineScope.startTaskProcessing(requestChannel: Channel, 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) { + internal suspend fun handleRequest( + request: WorkRequest, + requestState: AtomicReference, + 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), + ) } } diff --git a/build/jvm-rules/src/worker-framework/WorkRequestHandlerTest.kt b/build/jvm-rules/src/worker-framework/WorkRequestHandlerTest.kt index a2812ac0169d..0b8541f77ced 100644 --- a/build/jvm-rules/src/worker-framework/WorkRequestHandlerTest.kt +++ b/build/jvm-rules/src/worker-framework/WorkRequestHandlerTest.kt @@ -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) { } diff --git a/build/jvm-rules/src/worker-framework/logging/json.kt b/build/jvm-rules/src/worker-framework/logging/json.kt deleted file mode 100644 index d505e3545408..000000000000 --- a/build/jvm-rules/src/worker-framework/logging/json.kt +++ /dev/null @@ -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() -} \ No newline at end of file diff --git a/build/jvm-rules/src/worker-framework/logging/logger.kt b/build/jvm-rules/src/worker-framework/logging/logger.kt deleted file mode 100644 index dd586021e331..000000000000 --- a/build/jvm-rules/src/worker-framework/logging/logger.kt +++ /dev/null @@ -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? = 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(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) { - 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(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() - } - } -} \ No newline at end of file diff --git a/build/jvm-rules/src/worker-framework/logging/yaml.kt b/build/jvm-rules/src/worker-framework/logging/yaml.kt deleted file mode 100644 index 240e31566eb1..000000000000 --- a/build/jvm-rules/src/worker-framework/logging/yaml.kt +++ /dev/null @@ -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) - } -} diff --git a/build/jvm-rules/src/worker-framework/tracing.kt b/build/jvm-rules/src/worker-framework/tracing.kt new file mode 100644 index 000000000000..ec58c48f5c53 --- /dev/null +++ b/build/jvm-rules/src/worker-framework/tracing.kt @@ -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 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() + } +} \ No newline at end of file