diff --git a/app/src/main/java/com/bs/threadsimulator/common/MetricsExporter.kt b/app/src/main/java/com/bs/threadsimulator/common/MetricsExporter.kt index f0b501f..006d520 100644 --- a/app/src/main/java/com/bs/threadsimulator/common/MetricsExporter.kt +++ b/app/src/main/java/com/bs/threadsimulator/common/MetricsExporter.kt @@ -22,6 +22,11 @@ data class ExportedThreadMetric( val updateType: String, val updateCount: Long, val avgUpdateTimeMs: Long, + val peakUpdatesPerSec: Double = 0.0, + val stateTransitions: Int = 0, + val queueDepth: Int = 0, + val threadAllocatedBytes: Long = -1L, + val jitterMs: Double = 0.0, ) @Serializable @@ -110,6 +115,11 @@ class MetricsExporter "Update Type", "Update Count", "Avg Update Time (ms)", + "Peak Updates/Sec", + "State Transitions", + "Queue Depth", + "Thread Allocated Bytes", + "Jitter (ms)", ).joinToString(",") { escapeCsv(it) } val rows = metrics.joinToString("\n") { metric -> @@ -119,6 +129,11 @@ class MetricsExporter metric.updateType, metric.updateCount.toString(), metric.avgUpdateTimeMs.toString(), + "%.2f".format(metric.peakUpdatesPerSec), + metric.stateTransitions.toString(), + metric.queueDepth.toString(), + metric.threadAllocatedBytes.toString(), + "%.2f".format(metric.jitterMs), ).joinToString(",") { value -> escapeCsv(value) } } return header + "\n" + rows @@ -144,6 +159,11 @@ class MetricsExporter updateType = metric.updateType, updateCount = metric.updateCount, avgUpdateTimeMs = metric.avgUpdateTimeMs, + peakUpdatesPerSec = metric.peakUpdatesPerSec, + stateTransitions = metric.stateTransitions, + queueDepth = metric.queueDepth, + threadAllocatedBytes = metric.threadAllocatedBytes, + jitterMs = metric.jitterMs, ) } val metricsExport = diff --git a/app/src/main/java/com/bs/threadsimulator/common/ThreadMetrics.kt b/app/src/main/java/com/bs/threadsimulator/common/ThreadMetrics.kt index cb6ca81..dc88dd7 100644 --- a/app/src/main/java/com/bs/threadsimulator/common/ThreadMetrics.kt +++ b/app/src/main/java/com/bs/threadsimulator/common/ThreadMetrics.kt @@ -1,11 +1,14 @@ package com.bs.threadsimulator.common +import android.os.Debug import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.atomic.AtomicInteger import java.util.concurrent.atomic.AtomicLong import javax.inject.Inject +import kotlin.math.sqrt /** * Data class representing thread execution metrics. @@ -17,6 +20,11 @@ import javax.inject.Inject * @property updateType The type of update being tracked (e.g., "PE", "CurrentPrice", "HighLow") * @property updateCount Total number of updates performed by this thread * @property avgUpdateTimeMs Average time in milliseconds per update operation + * @property peakUpdatesPerSec Highest updates-per-second rate observed for this thread+type + * @property stateTransitions Count of thread state changes (e.g., RUNNABLE↔WAITING) + * @property queueDepth Current channel buffer occupancy (shared across threads) + * @property threadAllocatedBytes Cumulative bytes allocated by this thread (-1 if unavailable) + * @property jitterMs Standard deviation of update interval times in milliseconds */ data class ThreadMetrics( val threadId: Long, @@ -24,6 +32,11 @@ data class ThreadMetrics( val updateType: String, val updateCount: Long, val avgUpdateTimeMs: Long, + val peakUpdatesPerSec: Double = 0.0, + val stateTransitions: Int = 0, + val queueDepth: Int = 0, + val threadAllocatedBytes: Long = -1L, + val jitterMs: Double = 0.0, ) /** @@ -32,6 +45,9 @@ data class ThreadMetrics( * [ThreadMonitor] records update counts and timings for each thread and update type, * providing real-time visibility into multi-threaded performance. Useful for debugging * threading issues and monitoring performance bottlenecks in concurrent operations. + * + * Advanced metrics include peak updates/sec, thread state transitions, queue depth, + * per-thread memory allocation, and jitter (std-dev of update intervals). */ class ThreadMonitor @Inject @@ -46,20 +62,53 @@ class ThreadMonitor */ val metrics: StateFlow> = _metrics.asStateFlow() + // --- Basic tracking --- private val updateCounts = ConcurrentHashMap() private val updateTimes = ConcurrentHashMap() private val threadNames = ConcurrentHashMap() + // --- Advanced tracking --- + + /** Timestamps of each update for peak-UPS and jitter computation. */ + private val updateTimestamps = ConcurrentHashMap>() + + /** Per-thread state transition counter. */ + private val stateTransitionCounts = ConcurrentHashMap() + + /** Previous thread state for detecting transitions. */ + private val lastThreadStates = ConcurrentHashMap() + + /** Shared channel queue depth counter. */ + private val queueDepthCounter = AtomicInteger(0) + + /** + * Increments the queue depth counter. + * Should be called after successfully sending an element to the channel. + */ + fun incrementQueueDepth() { + queueDepthCounter.incrementAndGet() + } + + /** + * Decrements the queue depth counter. + * Should be called when an element is received from the channel. + */ + fun decrementQueueDepth() { + queueDepthCounter.decrementAndGet() + } + /** * Records an update operation with its execution time. * * Thread-safe. Should be called by worker threads to log their update operations. * Synchronized to prevent race conditions with [clearMetrics] operations. * + * Also tracks thread state transitions, timestamps for peak-UPS/jitter, and + * per-thread memory allocation. + * * @param updateType The type of update (e.g., "PE", "CurrentPrice", "HighLow") * @param updateTimeMs The time taken for this update operation in milliseconds */ - @Synchronized fun recordUpdate( updateType: String, updateTimeMs: Long, @@ -67,14 +116,35 @@ class ThreadMonitor val thread = Thread.currentThread() val key = "${thread.id}_$updateType" - threadNames.putIfAbsent(thread.id, thread.name) - updateCounts.getOrPut(key) { AtomicLong(0) }.incrementAndGet() - updateTimes.getOrPut(key) { AtomicLong(0) }.addAndGet(updateTimeMs) + synchronized(this) { + threadNames.putIfAbsent(thread.id, thread.name) + updateCounts.getOrPut(key) { AtomicLong(0) }.incrementAndGet() + updateTimes.getOrPut(key) { AtomicLong(0) }.addAndGet(updateTimeMs) + + // Record timestamp for peak-UPS and jitter with bounded 10-second retention window + val now = System.currentTimeMillis() + val timestamps = updateTimestamps.getOrPut(key) { mutableListOf() } + timestamps.add(now) + + val retentionWindowMs = 10_000L + val cutoff = now - retentionWindowMs + while (timestamps.isNotEmpty() && timestamps[0] < cutoff) { + timestamps.removeAt(0) + } + + // Track state transitions + val currentState = thread.state + val previousState = lastThreadStates.put(thread.id, currentState) + if (previousState != null && previousState != currentState) { + stateTransitionCounts.getOrPut(thread.id) { AtomicInteger(0) }.incrementAndGet() + } + } updateMetrics() } private fun updateMetrics() { + val currentQueueDepth = queueDepthCounter.get() val currentMetrics = updateCounts.keys.map { key -> val (threadId, updateType) = key.split("_", limit = 2) @@ -82,22 +152,89 @@ class ThreadMonitor val count = updateCounts[key]?.get() ?: 0 val totalTime = updateTimes[key]?.get() ?: 0 + val timestamps = + synchronized(this) { + updateTimestamps[key]?.toList() ?: emptyList() + } + ThreadMetrics( threadId = threadIdLong, threadName = threadNames[threadIdLong] ?: "Unknown", updateType = updateType, updateCount = count, avgUpdateTimeMs = if (count > 0) totalTime / count else 0, + peakUpdatesPerSec = computePeakUpdatesPerSec(timestamps), + stateTransitions = stateTransitionCounts[threadIdLong]?.get() ?: 0, + queueDepth = currentQueueDepth, + threadAllocatedBytes = getThreadAllocatedBytes(threadIdLong), + jitterMs = computeJitter(timestamps), ) } _metrics.value = currentMetrics } + /** + * Computes peak updates per second. + * Runs in O(n) time using a sliding window over the monotonically increasing timestamps. + */ + @Suppress("ReturnCount") + private fun computePeakUpdatesPerSec(timestamps: List): Double { + if (timestamps.isEmpty()) return 0.0 + if (timestamps.size == 1) return 1.0 + + val windowMillis = 1000L + var peak = 0 + var start = 0 + + for (end in timestamps.indices) { + val windowStart = timestamps[end] - windowMillis + + while (start < end && timestamps[start] <= windowStart) { + start++ + } + + val count = end - start + 1 + if (count > peak) { + peak = count + } + } + return peak.toDouble() + } + + /** + * Computes jitter (standard deviation) of inter-update intervals in milliseconds. + */ + private fun computeJitter(timestamps: List): Double { + if (timestamps.size < 2) return 0.0 + val intervals = + (1 until timestamps.size).map { i -> + (timestamps[i] - timestamps[i - 1]).toDouble() + } + val mean = intervals.average() + val variance = intervals.map { (it - mean) * (it - mean) }.average() + return sqrt(variance) + } + + /** + * Returns the process-wide native heap allocated size in bytes. + * Android does not support per-thread memory tracking, so this is + * a process-level metric shared across all threads. + * Returns -1 if unavailable. + */ + private fun getThreadAllocatedBytes( + @Suppress("UNUSED_PARAMETER") threadId: Long, + ): Long = + try { + Debug.getNativeHeapAllocatedSize() + } catch (_: Exception) { + -1L + } + /** * Clears all accumulated metrics. * * Useful for resetting metrics when starting a new simulation or test. - * Thread-safe. Synchronized to ensure all three maps are cleared atomically, + * Thread-safe. Synchronized to ensure all maps are cleared atomically, * preventing inconsistent state if [recordUpdate] is called concurrently. */ @Synchronized @@ -105,6 +242,10 @@ class ThreadMonitor updateCounts.clear() updateTimes.clear() threadNames.clear() + updateTimestamps.clear() + stateTransitionCounts.clear() + lastThreadStates.clear() + queueDepthCounter.set(0) _metrics.value = emptyList() } diff --git a/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeScreenRoute.kt b/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeScreenRoute.kt index b207a0e..401bf94 100644 --- a/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeScreenRoute.kt +++ b/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeScreenRoute.kt @@ -331,7 +331,7 @@ private fun ThreadMetricsDisplay(threadMetrics: List) { modifier = Modifier .padding(4.dp) - .width(200.dp), + .width(220.dp), ) { Column( modifier = Modifier.padding(8.dp), @@ -352,6 +352,28 @@ private fun ThreadMetricsDisplay(threadMetrics: List) { text = "Avg Time: ${metric.avgUpdateTimeMs}ms", style = MaterialTheme.typography.bodySmall, ) + Text( + text = "Peak UPS: ${"%.1f".format(metric.peakUpdatesPerSec)}", + style = MaterialTheme.typography.bodySmall, + ) + Text( + text = "State Transitions: ${metric.stateTransitions}", + style = MaterialTheme.typography.bodySmall, + ) + Text( + text = "Queue Depth: ${metric.queueDepth}", + style = MaterialTheme.typography.bodySmall, + ) + if (metric.threadAllocatedBytes >= 0) { + Text( + text = "Memory: ${"%.2f".format(metric.threadAllocatedBytes / 1024.0)} KB", + style = MaterialTheme.typography.bodySmall, + ) + } + Text( + text = "Jitter: ${"%.2f".format(metric.jitterMs)}ms", + style = MaterialTheme.typography.bodySmall, + ) } } } diff --git a/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeViewModel.kt b/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeViewModel.kt index 62f77d9..34bb3ed 100644 --- a/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeViewModel.kt +++ b/app/src/main/java/com/bs/threadsimulator/ui/screens/HomeViewModel.kt @@ -228,7 +228,13 @@ class HomeViewModel when (resource) { is Resource.Success -> { if (resource.data == null) return@collect - channel.send(resource.data) + threadMonitor.incrementQueueDepth() + try { + channel.send(resource.data) + } catch (e: Exception) { + threadMonitor.decrementQueueDepth() + throw e + } } is Resource.Error -> { @@ -254,7 +260,13 @@ class HomeViewModel when (resource) { is Resource.Success -> { if (resource.data == null) return@collect - channel.send(resource.data) + threadMonitor.incrementQueueDepth() + try { + channel.send(resource.data) + } catch (e: Exception) { + threadMonitor.decrementQueueDepth() + throw e + } } is Resource.Error -> { @@ -277,7 +289,13 @@ class HomeViewModel when (resource) { is Resource.Success -> { if (resource.data == null) return@collect - channel.send(resource.data) + threadMonitor.incrementQueueDepth() + try { + channel.send(resource.data) + } catch (e: Exception) { + threadMonitor.decrementQueueDepth() + throw e + } } is Resource.Error -> { @@ -296,6 +314,7 @@ class HomeViewModel viewModelScope.launch(appDispatchers.mainDispatcher) { try { for (companyData in channel) { + threadMonitor.decrementQueueDepth() ensureActive() val company = _companyList.getOrNull(companyData.id) ?: continue try { diff --git a/app/src/test/java/com/bs/threadsimulator/common/ThreadMonitorTest.kt b/app/src/test/java/com/bs/threadsimulator/common/ThreadMonitorTest.kt index a95f2c8..b941907 100644 --- a/app/src/test/java/com/bs/threadsimulator/common/ThreadMonitorTest.kt +++ b/app/src/test/java/com/bs/threadsimulator/common/ThreadMonitorTest.kt @@ -78,6 +78,11 @@ class ThreadMonitorTest { updateType = "PE", updateCount = 5, avgUpdateTimeMs = 20, + peakUpdatesPerSec = 10.0, + stateTransitions = 3, + queueDepth = 7, + threadAllocatedBytes = 1024L, + jitterMs = 2.5, ) assertEquals(1L, metrics.threadId) @@ -85,6 +90,11 @@ class ThreadMonitorTest { assertEquals("PE", metrics.updateType) assertEquals(5, metrics.updateCount) assertEquals(20, metrics.avgUpdateTimeMs) + assertEquals(10.0, metrics.peakUpdatesPerSec, 0.01) + assertEquals(3, metrics.stateTransitions) + assertEquals(7, metrics.queueDepth) + assertEquals(1024L, metrics.threadAllocatedBytes) + assertEquals(2.5, metrics.jitterMs, 0.01) } @Test @@ -152,4 +162,132 @@ class ThreadMonitorTest { cancelAndIgnoreRemainingEvents() } } + + // --- Advanced metrics tests --- + + @Test + fun testQueueDepthCanBeNegativeTemporarily() { + // Simulates consumer receiving before producer has recorded increment (edge case in channels) + threadMonitor.decrementQueueDepth() + threadMonitor.decrementQueueDepth() + + threadMonitor.recordUpdate("PE", 10) + var metrics = threadMonitor.metrics.value + var peMetric = metrics.find { it.updateType == "PE" } + assertEquals("Queue depth should be -2", -2, peMetric!!.queueDepth) + + threadMonitor.incrementQueueDepth() + threadMonitor.incrementQueueDepth() + threadMonitor.incrementQueueDepth() + + threadMonitor.recordUpdate("PE", 10) + metrics = threadMonitor.metrics.value + peMetric = metrics.find { it.updateType == "PE" } + assertEquals("Queue depth should be 1", 1, peMetric!!.queueDepth) + } + + @Test + fun testPeakUpdatesPerSecIsTracked() { + // Record several updates rapidly — they all happen in the same millisecond range + repeat(10) { + threadMonitor.recordUpdate("PE", 5) + } + + val metrics = threadMonitor.metrics.value + val peMetric = metrics.find { it.updateType == "PE" } + assertNotNull("PE metric should exist", peMetric) + assertTrue( + "Peak UPS should be >= 10 (all within 1s window)", + peMetric!!.peakUpdatesPerSec >= 10.0, + ) + } + + @Test + fun testQueueDepthIncrementAndDecrement() { + threadMonitor.incrementQueueDepth() + threadMonitor.incrementQueueDepth() + threadMonitor.incrementQueueDepth() + + // Record an update to trigger metrics refresh + threadMonitor.recordUpdate("PE", 10) + var metrics = threadMonitor.metrics.value + var peMetric = metrics.find { it.updateType == "PE" } + assertEquals("Queue depth should be 3", 3, peMetric!!.queueDepth) + + threadMonitor.decrementQueueDepth() + threadMonitor.recordUpdate("PE", 10) + metrics = threadMonitor.metrics.value + peMetric = metrics.find { it.updateType == "PE" } + assertEquals("Queue depth should be 2 after decrement", 2, peMetric!!.queueDepth) + } + + @Test + fun testJitterIsComputedForMultipleUpdates() { + // Rapid calls — jitter should be small or zero + repeat(5) { + threadMonitor.recordUpdate("PE", 10) + } + + val metrics = threadMonitor.metrics.value + val peMetric = metrics.find { it.updateType == "PE" } + assertNotNull("PE metric should exist", peMetric) + assertTrue( + "Jitter should be non-negative", + peMetric!!.jitterMs >= 0.0, + ) + } + + @Test + fun testThreadAllocatedBytesIsPopulated() { + threadMonitor.recordUpdate("PE", 10) + + val metrics = threadMonitor.metrics.value + val peMetric = metrics.find { it.updateType == "PE" } + assertNotNull("PE metric should exist", peMetric) + // threadAllocatedBytes is either >= 0 (supported) or -1 (unsupported) + assertTrue( + "Memory should be -1 or a positive value", + peMetric!!.threadAllocatedBytes == -1L || peMetric.threadAllocatedBytes > 0, + ) + } + + @Test + fun testClearMetricsResetsAdvancedFields() { + threadMonitor.incrementQueueDepth() + threadMonitor.incrementQueueDepth() + threadMonitor.recordUpdate("PE", 100) + + // Verify metrics exist + assertTrue(threadMonitor.metrics.value.isNotEmpty()) + + threadMonitor.clearMetrics() + + assertTrue("Metrics should be empty after clear", threadMonitor.metrics.value.isEmpty()) + + // Record again and verify fresh state + threadMonitor.recordUpdate("PE", 50) + val metrics = threadMonitor.metrics.value + val peMetric = metrics.find { it.updateType == "PE" } + assertNotNull(peMetric) + assertEquals("Count should be 1 after clear+record", 1L, peMetric!!.updateCount) + assertEquals("Queue depth should be 0 after clear", 0, peMetric.queueDepth) + } + + @Test + fun testThreadMetricsDefaultValues() { + val metrics = + ThreadMetrics( + threadId = 1L, + threadName = "Test", + updateType = "PE", + updateCount = 1, + avgUpdateTimeMs = 10, + ) + + assertEquals("Default peakUpdatesPerSec should be 0.0", 0.0, metrics.peakUpdatesPerSec, 0.01) + assertEquals("Default stateTransitions should be 0", 0, metrics.stateTransitions) + assertEquals("Default queueDepth should be 0", 0, metrics.queueDepth) + assertEquals("Default threadAllocatedBytes should be -1", -1L, metrics.threadAllocatedBytes) + assertEquals("Default jitterMs should be 0.0", 0.0, metrics.jitterMs, 0.01) + } }