Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions app/src/main/java/com/bs/threadsimulator/common/MetricsExporter.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 ->
Expand All @@ -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
Expand All @@ -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 =
Expand Down
151 changes: 146 additions & 5 deletions app/src/main/java/com/bs/threadsimulator/common/ThreadMetrics.kt
Original file line number Diff line number Diff line change
@@ -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.
Expand All @@ -17,13 +20,23 @@ 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)
Comment thread
buddhasaikia marked this conversation as resolved.
* @property jitterMs Standard deviation of update interval times in milliseconds
*/
data class ThreadMetrics(
val threadId: Long,
val threadName: String,
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,
)

/**
Expand All @@ -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
Expand All @@ -46,65 +62,190 @@ class ThreadMonitor
*/
val metrics: StateFlow<List<ThreadMetrics>> = _metrics.asStateFlow()

// --- Basic tracking ---
private val updateCounts = ConcurrentHashMap<String, AtomicLong>()
private val updateTimes = ConcurrentHashMap<String, AtomicLong>()
private val threadNames = ConcurrentHashMap<Long, String>()

// --- Advanced tracking ---

/** Timestamps of each update for peak-UPS and jitter computation. */
private val updateTimestamps = ConcurrentHashMap<String, MutableList<Long>>()
Comment thread
buddhasaikia marked this conversation as resolved.

/** Per-thread state transition counter. */
private val stateTransitionCounts = ConcurrentHashMap<Long, AtomicInteger>()

/** Previous thread state for detecting transitions. */
private val lastThreadStates = ConcurrentHashMap<Long, Thread.State>()

/** 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,
) {
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()

Copilot AI Feb 27, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The @Synchronized annotation on recordUpdate() protects against concurrent modifications to the data structures, but updateMetrics() is called at the end of recordUpdate() while still holding the lock. This means every update operation blocks all other threads trying to record updates while metrics are being computed (including the O(n²) peak UPS calculation).

For high-frequency concurrent updates, this can create a significant bottleneck. Consider either: (1) moving updateMetrics() outside the synchronized block and accepting eventual consistency, (2) using more granular locking, or (3) deferring metrics computation to a separate background task that runs periodically rather than on every update.

Copilot uses AI. Check for mistakes.
}

private fun updateMetrics() {
val currentQueueDepth = queueDepthCounter.get()
val currentMetrics =
updateCounts.keys.map { key ->
val (threadId, updateType) = key.split("_", limit = 2)
val threadIdLong = threadId.toLong()
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<Long>): 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
}
}
Comment on lines +179 to +200

Copilot AI Feb 27, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The peak UPS computation has O(n²) time complexity in the worst case, where n is the number of timestamps. For high-frequency updates over long periods, this can become a performance bottleneck since it runs on every recordUpdate() call through updateMetrics().

Consider optimizing this by: (1) maintaining a rolling window of timestamps from the last N seconds instead of all timestamps, (2) using a more efficient algorithm like a sliding window with a deque, or (3) computing this metric less frequently (e.g., only when metrics are explicitly requested for display/export rather than on every update).

Suggested change
*/
private fun computePeakUpdatesPerSec(timestamps: List<Long>): Double {
if (timestamps.size < 2) return timestamps.size.toDouble()
var peak = 0
for (i in timestamps.indices) {
val windowEnd = timestamps[i] + 1000L
var count = 0
for (j in i until timestamps.size) {
if (timestamps[j] <= windowEnd) count++ else break
}
if (count > peak) peak = count
}
*
* This implementation runs in O(n) time by using a two-pointer sliding window
* over the (assumed) monotonically increasing list of timestamps.
*/
private fun computePeakUpdatesPerSec(timestamps: List<Long>): 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
// Advance start index while the timestamp is older than the 1-second window
while (start < end && timestamps[start] < windowStart) {
start++
}
val count = end - start + 1
if (count > peak) {
peak = count
}
}

Copilot uses AI. Check for mistakes.
return peak.toDouble()
}

/**
* Computes jitter (standard deviation) of inter-update intervals in milliseconds.
*/
private fun computeJitter(timestamps: List<Long>): 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
fun clearMetrics() {
updateCounts.clear()
updateTimes.clear()
threadNames.clear()
updateTimestamps.clear()
stateTransitionCounts.clear()
lastThreadStates.clear()
queueDepthCounter.set(0)
_metrics.value = emptyList()
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -331,7 +331,7 @@ private fun ThreadMetricsDisplay(threadMetrics: List<ThreadMetrics>) {
modifier =
Modifier
.padding(4.dp)
.width(200.dp),
.width(220.dp),
) {
Column(
modifier = Modifier.padding(8.dp),
Expand All @@ -352,6 +352,28 @@ private fun ThreadMetricsDisplay(threadMetrics: List<ThreadMetrics>) {
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,
Comment on lines +367 to +370

Copilot AI Feb 27, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The memory display shows bytes converted to KB using integer division. For very small allocations (less than 1024 bytes), this will display as "0KB" which might be confusing. Additionally, for large allocations, KB might not be the most readable unit.

Consider using a more sophisticated formatting approach that selects appropriate units (B, KB, MB, GB) based on the size, or at least format with decimal places for KB (e.g., "%.2f KB").

Copilot uses AI. Check for mistakes.
)
}
Text(
text = "Jitter: ${"%.2f".format(metric.jitterMs)}ms",
style = MaterialTheme.typography.bodySmall,
)
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 -> {
Expand All @@ -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 -> {
Expand All @@ -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 -> {
Expand All @@ -296,6 +314,7 @@ class HomeViewModel
viewModelScope.launch(appDispatchers.mainDispatcher) {
try {
for (companyData in channel) {
threadMonitor.decrementQueueDepth()
Comment thread
buddhasaikia marked this conversation as resolved.
ensureActive()
val company = _companyList.getOrNull(companyData.id) ?: continue
try {
Expand Down
Loading
Loading