From 596955f6b4ba5227068b89e803f6d75746c3c64b Mon Sep 17 00:00:00 2001 From: Michal Harakal Date: Mon, 10 Aug 2026 10:53:05 +0200 Subject: [PATCH] fix(lang): per-source copy attribution in MemoryTracker; volatile ActiveMemoryTracker MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit MemoryTracker.recordCopy(sourceName, bytes) discarded sourceName — it only bumped copyCount/copyBytes, while every instrumented call site passes a meaningful label (CopyMaterializationStrategy, DenseTensorDataFactory.createFloatTensorData, ...). The API promised per-source attribution and threw it away. Aggregate reports now carry copiesBySource: Map (count + bytes per code path), included in the report's text form sorted by volume, and reset by clear(). ActiveMemoryTracker.current becomes @Volatile so installing/clearing a tracker is visible across threads, and its doc now states the honest contract: the tracker itself is not synchronized, concurrent sessions should install their own around their critical section. Replacing the process-wide hook with a per-execution-context tracker is part of the SKEEP-003 storage-model discussion (#932) and intentionally out of scope here. Tests: per-source aggregation across repeated and distinct sources, clear() resetting attribution, and the text report containing the breakdown. Closes #931 --- CHANGELOG.md | 9 +++++ .../tensor/storage/ActiveMemoryTracker.kt | 12 +++++- .../lang/tensor/storage/MemoryTracker.kt | 32 +++++++++++++-- .../tensor/storage/ActiveMemoryTrackerTest.kt | 39 +++++++++++++++++++ 4 files changed, 87 insertions(+), 5 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index fca59a3c2..54c844888 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,15 @@ ## [Unreleased] +### Fixed + +- **Memory-copy diagnostics attribute copies to their source.** `MemoryTracker.recordCopy` + discarded the `sourceName` every instrumented call site passes; reports now carry a + per-source breakdown (`copiesBySource: Map`, included in the + report's text form, sorted by volume). `ActiveMemoryTracker.current` is now `@Volatile` + with an honest thread-safety contract in its docs; making the tracker per-execution-context + instead of a process-wide hook is part of the SKEEP-003 storage-model discussion. (#931) + ## [0.38.0] - 2026-07-30 ### Added diff --git a/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTracker.kt b/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTracker.kt index 8c79c90dc..10c957b76 100644 --- a/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTracker.kt +++ b/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTracker.kt @@ -1,5 +1,7 @@ package sk.ainet.lang.tensor.storage +import kotlin.concurrent.Volatile + /** * Global hook for the active [MemoryTracker]. * @@ -7,10 +9,16 @@ package sk.ainet.lang.tensor.storage * from instrumented copy paths (e.g. CopyMaterializationStrategy, * DenseTensorDataFactory.from*Array). Set to `null` to disable tracking. * - * Thread-safety note: on JVM this should ideally be a ThreadLocal. - * For now, a simple global works for single-threaded inference. + * Thread-safety: [current] is `@Volatile`, so installing or clearing a + * tracker is immediately visible to other threads. [MemoryTracker] itself + * is not synchronized — concurrent loads/inference sessions that need + * isolated attribution should each install their own tracker around their + * critical section, or serialize access. Making the tracker installable + * per execution context (instead of a process-wide hook) is part of the + * storage-model discussion in SKEEP-003. */ public object ActiveMemoryTracker { + @Volatile public var current: MemoryTracker? = null /** Record a copy event on the active tracker, if any. */ diff --git a/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/MemoryTracker.kt b/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/MemoryTracker.kt index e723748fb..2fba95032 100644 --- a/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/MemoryTracker.kt +++ b/skainet-lang/skainet-lang-core/src/commonMain/kotlin/sk/ainet/lang/tensor/storage/MemoryTracker.kt @@ -13,16 +13,27 @@ public class MemoryTracker { private val entries = mutableListOf() private var copyCount: Long = 0 private var copyBytes: Long = 0 + private val copiesBySource = mutableMapOf() /** Record a tensor storage allocation. */ public fun record(name: String, storage: TensorStorage) { entries.add(TrackedEntry(name, storage.memoryReport())) } - /** Record an explicit copy event (for copy-tracing). */ + /** + * Record an explicit copy event (for copy-tracing). + * + * [sourceName] is aggregated per source, so reports can attribute copy + * volume to the code path that produced it (e.g. which factory or + * materialization strategy) — previously the label was discarded. + */ public fun recordCopy(sourceName: String, bytes: Long) { copyCount++ copyBytes += bytes + val prev = copiesBySource[sourceName] + copiesBySource[sourceName] = + if (prev == null) CopySourceStat(count = 1, bytes = bytes) + else CopySourceStat(count = prev.count + 1, bytes = prev.bytes + bytes) } /** Reset all tracked entries. */ @@ -30,6 +41,7 @@ public class MemoryTracker { entries.clear() copyCount = 0 copyBytes = 0 + copiesBySource.clear() } /** Generate an aggregate memory report. */ @@ -69,11 +81,18 @@ public class MemoryTracker { fileBackedCount = fileBackedCount, copyCount = copyCount, copyBytes = copyBytes, - entries = entries.toList() + entries = entries.toList(), + copiesBySource = copiesBySource.toMap() ) } } +/** Per-source copy statistics: how many copies a code path produced and their total volume. */ +public data class CopySourceStat( + val count: Long, + val bytes: Long +) + public data class TrackedEntry( val name: String, val report: StorageMemoryReport @@ -90,7 +109,8 @@ public data class AggregateMemoryReport( val fileBackedCount: Int, val copyCount: Long, val copyBytes: Long, - val entries: List + val entries: List, + val copiesBySource: Map = emptyMap() ) { val overallCompressionRatio: Double get() = if (totalPhysicalBytes > 0) totalLogicalBytes.toDouble() / totalPhysicalBytes else 1.0 @@ -103,6 +123,12 @@ public data class AggregateMemoryReport( appendLine("File-backed: $fileBackedCount ($fileBackedBytes bytes)") appendLine("Owned: $ownedCount, Borrowed: $borrowedCount, Aliased: $aliasedCount") appendLine("Copies: $copyCount ($copyBytes bytes)") + if (copiesBySource.isNotEmpty()) { + appendLine("--- Copies by source ---") + for ((source, stat) in copiesBySource.entries.sortedByDescending { it.value.bytes }) { + appendLine(" $source: ${stat.count} (${stat.bytes} bytes)") + } + } if (entries.isNotEmpty()) { appendLine("--- Per-tensor ---") for (e in entries) { diff --git a/skainet-lang/skainet-lang-core/src/commonTest/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTrackerTest.kt b/skainet-lang/skainet-lang-core/src/commonTest/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTrackerTest.kt index f8da35dff..d50e104b4 100644 --- a/skainet-lang/skainet-lang-core/src/commonTest/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTrackerTest.kt +++ b/skainet-lang/skainet-lang-core/src/commonTest/kotlin/sk/ainet/lang/tensor/storage/ActiveMemoryTrackerTest.kt @@ -27,6 +27,45 @@ class ActiveMemoryTrackerTest { assertEquals(100L, report.copyBytes) } + @Test + fun recordCopy_attributesPerSource() { + // The source label used to be discarded (#931) — every call site + // passes a meaningful one, and reports must break copies down by it. + val tracker = MemoryTracker() + ActiveMemoryTracker.current = tracker + + ActiveMemoryTracker.recordCopy("factory", 100) + ActiveMemoryTracker.recordCopy("factory", 50) + ActiveMemoryTracker.recordCopy("materialize", 200) + + val report = tracker.report() + assertEquals(3L, report.copyCount) + assertEquals(350L, report.copyBytes) + assertEquals(CopySourceStat(count = 2, bytes = 150), report.copiesBySource["factory"]) + assertEquals(CopySourceStat(count = 1, bytes = 200), report.copiesBySource["materialize"]) + } + + @Test + fun clear_resetsPerSourceAttribution() { + val tracker = MemoryTracker() + tracker.recordCopy("a", 10) + tracker.clear() + tracker.recordCopy("b", 20) + + val report = tracker.report() + assertEquals(1L, report.copyCount) + assertEquals(mapOf("b" to CopySourceStat(1, 20)), report.copiesBySource) + } + + @Test + fun report_toString_includesPerSourceBreakdown() { + val tracker = MemoryTracker() + tracker.recordCopy("DenseTensorDataFactory.createFloatTensorData", 4096) + val text = tracker.report().toString() + kotlin.test.assertTrue("DenseTensorDataFactory.createFloatTensorData" in text, text) + kotlin.test.assertTrue("4096" in text, text) + } + @Test fun recordCopy_withNullTracker_noOp() { ActiveMemoryTracker.current = null