From 050afe0e863d013ea16900cab070f76b0fb79fa8 Mon Sep 17 00:00:00 2001 From: Ruslan Cheremin Date: Thu, 25 Jul 2024 13:44:25 +0200 Subject: [PATCH] [indexes] IJP-1365: report DiskQueryRelay metrics to OTel + total time spent on task execution/waiting for result + total number of tasks requested/executed GitOrigin-RevId: 1d43b290f41ef0b97958e65a2db3d9e64d3ac79f --- .../intellij/openapi/vfs/DiskQueryRelay.java | 109 +++++++++++++++--- 1 file changed, 90 insertions(+), 19 deletions(-) diff --git a/platform/platform-impl/src/com/intellij/openapi/vfs/DiskQueryRelay.java b/platform/platform-impl/src/com/intellij/openapi/vfs/DiskQueryRelay.java index 2a5cc2e7ca3d..5a4e65d40ac8 100644 --- a/platform/platform-impl/src/com/intellij/openapi/vfs/DiskQueryRelay.java +++ b/platform/platform-impl/src/com/intellij/openapi/vfs/DiskQueryRelay.java @@ -5,18 +5,22 @@ import com.intellij.execution.process.ProcessIOExecutorService; import com.intellij.openapi.progress.ProcessCanceledException; import com.intellij.openapi.progress.util.ProgressIndicatorUtils; import com.intellij.openapi.util.ThrowableComputable; +import com.intellij.platform.diagnostic.telemetry.PlatformScopesKt; +import com.intellij.platform.diagnostic.telemetry.TelemetryManager; import com.intellij.util.ExceptionUtil; +import io.opentelemetry.api.metrics.Meter; import org.jetbrains.annotations.ApiStatus; import org.jetbrains.annotations.NotNull; import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Future; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.Function; import static com.intellij.openapi.progress.ContextKt.isInCancellableContext; +import static java.util.concurrent.TimeUnit.MICROSECONDS; +import static java.util.concurrent.TimeUnit.NANOSECONDS; /** * A utility to run a potentially long function on a pooled thread, wait for it in an interruptible way, @@ -46,27 +50,45 @@ public final class DiskQueryRelay { public DiskQueryRelay(@NotNull Function function, @NotNull ExecutorService executor) { - myFunction = function; + myFunction = arg -> { + long startedAtNs = System.nanoTime(); + try { + return function.apply(arg); + } + finally { + long elapsedNs = System.nanoTime() - startedAtNs; + taskExecutionTotalTimeNs.addAndGet(elapsedNs); + tasksExecutedCount.incrementAndGet(); + } + }; this.executor = executor; } public Result accessDiskWithCheckCanceled(@NotNull Param arg) { - if (!isInCancellableContext()) { - return myFunction.apply(arg); - } - Future future = myTasks.computeIfAbsent(arg, eachArg -> executor.submit(() -> { - try { - return myFunction.apply(eachArg); + long startedAtNs = System.nanoTime(); + try { + if (!isInCancellableContext()) { + return myFunction.apply(arg); } - finally { - myTasks.remove(eachArg); + Future future = myTasks.computeIfAbsent(arg, eachArg -> executor.submit(() -> { + try { + return myFunction.apply(eachArg); + } + finally { + myTasks.remove(eachArg); + } + })); + if (future.isDone()) { + // maybe it was very fast and completed before being put into a map + myTasks.remove(arg, future); } - })); - if (future.isDone()) { - // maybe it was very fast and completed before being put into a map - myTasks.remove(arg, future); + return ProgressIndicatorUtils.awaitWithCheckCanceled(future); + } + finally { + long elapsedNs = System.nanoTime() - startedAtNs; + taskWaitingTotalTimeNs.addAndGet(elapsedNs); + tasksRequestedCount.incrementAndGet(); } - return ProgressIndicatorUtils.awaitWithCheckCanceled(future); } /** @@ -75,7 +97,8 @@ public final class DiskQueryRelay { * To avoid deadlocks, please pay attention to locks held at the call time and try to abstain from taking locks * inside the {@code task} block. */ - public static Result compute(@NotNull ThrowableComputable task) throws E, ProcessCanceledException { + public static Result compute(@NotNull ThrowableComputable task) + throws E, ProcessCanceledException { return compute(task, ProcessIOExecutorService.INSTANCE); } @@ -117,4 +140,52 @@ public final class DiskQueryRelay { future.cancel(false); } } + + + // ==================================== monitoring: ====================================================== // + + /** total time (since app start) of actual task executions, ns */ + private static final AtomicLong taskExecutionTotalTimeNs = new AtomicLong(); + /** total time (since app start) spent waiting for the task result, ns */ + private static final AtomicLong taskWaitingTotalTimeNs = new AtomicLong(); + /** total (since app start) number of tasks actually executed */ + private static final AtomicInteger tasksExecutedCount = new AtomicInteger(); + /** total (since app start) number of tasks requested. Could be <= tasksExecuted because of task coalescing */ + private static final AtomicInteger tasksRequestedCount = new AtomicInteger(); + + public static long taskExecutionTotalTime(@NotNull TimeUnit unit) { + return unit.convert(taskExecutionTotalTimeNs.get(), NANOSECONDS); + } + + public static long taskWaitingTotalTime(@NotNull TimeUnit unit) { + return unit.convert(taskWaitingTotalTimeNs.get(), NANOSECONDS); + } + + public static int tasksExecuted() { + return tasksExecutedCount.get(); + } + + public static int tasksRequested() { + return tasksRequestedCount.get(); + } + + static { + Meter otelMeter = TelemetryManager.getInstance().getMeter(PlatformScopesKt.PlatformMetrics); + + var taskExecutionTimeUs = otelMeter.counterBuilder("DiskQueryRelay.taskExecutionTotalTimeUs").buildObserver(); + var taskWaitingTimeUs = otelMeter.counterBuilder("DiskQueryRelay.taskWaitingTotalTimeUs").buildObserver(); + var tasksExecuted = otelMeter.counterBuilder("DiskQueryRelay.tasksExecuted").buildObserver(); + var tasksRequested = otelMeter.counterBuilder("DiskQueryRelay.tasksRequested").buildObserver(); + + otelMeter.batchCallback( + () -> { + taskExecutionTimeUs.record(taskExecutionTotalTime(MICROSECONDS)); + taskWaitingTimeUs.record(taskWaitingTotalTime(MICROSECONDS)); + tasksExecuted.record(tasksExecuted()); + tasksRequested.record(tasksRequested()); + }, + taskExecutionTimeUs, taskWaitingTimeUs, + tasksExecuted, tasksRequested + ); + } }