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
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ public final class sk/ainet/context/DirectCpuExecutionContext : sk/ainet/context
public fun getPhase ()Lsk/ainet/context/Phase;
public fun getScratch ()Lsk/ainet/lang/tensor/scratch/ScratchPool;
public fun getTensorDataFactory ()Lsk/ainet/lang/tensor/data/TensorDataFactory;
public fun getTraceSink ()Lsk/ainet/lang/memory/trace/TraceSink;
public fun isRecording ()Z
public fun ones (Lsk/ainet/lang/tensor/Shape;Lkotlin/reflect/KClass;)Lsk/ainet/lang/tensor/Tensor;
public fun placeholder (Lsk/ainet/lang/tensor/Shape;Lkotlin/reflect/KClass;)Lsk/ainet/lang/tensor/Tensor;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -182,6 +182,7 @@ public final class sk/ainet/lang/graph/DefaultGraphExecutionContext : sk/ainet/l
public final fun getSession ()Lsk/ainet/lang/trace/TraceSession;
public fun getTapeStack ()Lsk/ainet/tape/TapeStack;
public fun getTensorDataFactory ()Lsk/ainet/lang/tensor/data/TensorDataFactory;
public fun getTraceSink ()Lsk/ainet/lang/memory/trace/TraceSink;
public fun isRecording ()Z
public fun ones (Lsk/ainet/lang/tensor/Shape;Lkotlin/reflect/KClass;)Lsk/ainet/lang/tensor/Tensor;
public fun placeholder (Lsk/ainet/lang/tensor/Shape;Lkotlin/reflect/KClass;)Lsk/ainet/lang/tensor/Tensor;
Expand Down Expand Up @@ -450,6 +451,7 @@ public final class sk/ainet/lang/graph/exec/GraphExecutionContext$DefaultImpls {
public static fun getMemoryPlanner (Lsk/ainet/lang/graph/exec/GraphExecutionContext;)Lsk/ainet/lang/tensor/storage/MemoryPlanner;
public static fun getMemoryTracker (Lsk/ainet/lang/graph/exec/GraphExecutionContext;)Lsk/ainet/lang/tensor/storage/MemoryTracker;
public static fun getScratch (Lsk/ainet/lang/graph/exec/GraphExecutionContext;)Lsk/ainet/lang/tensor/scratch/ScratchPool;
public static fun getTraceSink (Lsk/ainet/lang/graph/exec/GraphExecutionContext;)Lsk/ainet/lang/memory/trace/TraceSink;
public static fun isRecording (Lsk/ainet/lang/graph/exec/GraphExecutionContext;)Z
public static fun ones (Lsk/ainet/lang/graph/exec/GraphExecutionContext;Lsk/ainet/lang/tensor/Shape;Lkotlin/reflect/KClass;)Lsk/ainet/lang/tensor/Tensor;
public static fun placeholder (Lsk/ainet/lang/graph/exec/GraphExecutionContext;Lsk/ainet/lang/tensor/Shape;Lkotlin/reflect/KClass;)Lsk/ainet/lang/tensor/Tensor;
Expand Down
384 changes: 384 additions & 0 deletions skainet-lang/skainet-lang-core/api/jvm/skainet-lang-core.api

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,14 @@ import sk.ainet.lang.types.DType
import kotlin.reflect.KClass

public interface ExecutionContext {
/**
* Where this context's trace events go (SKEEP-003 §4.9): phases, kernel runs, adapter
* insertions, allocations. Default [sk.ainet.lang.memory.trace.NoopTraceSink] — nothing is
* recorded until a context opts in with a recording or exporting sink.
*/
@sk.ainet.lang.memory.ExperimentalMemoryApi
public val traceSink: sk.ainet.lang.memory.trace.TraceSink get() = sk.ainet.lang.memory.trace.NoopTraceSink

public val ops: TensorOps

// Optional forward hooks for recording or diagnostics (null → disabled)
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,188 @@
@file:OptIn(kotlin.concurrent.atomics.ExperimentalAtomicApi::class)

package sk.ainet.lang.memory

import sk.ainet.lang.memory.trace.NoopTraceSink
import sk.ainet.lang.memory.trace.TraceEvent
import sk.ainet.lang.memory.trace.TraceSink
import sk.ainet.lang.tensor.TensorId
import sk.ainet.lang.tensor.storage.MemoryDomain
import kotlin.concurrent.atomics.AtomicLong
import kotlin.concurrent.atomics.fetchAndIncrement
import kotlin.jvm.JvmInline

/**
* Monotonic per-process identity of one allocation (SKEEP-003 §0 *StorageId*): what the memory
* debugger keys on. One `TensorId` maps to many storage ids over time (a `Forward` scope is
* recycled every step); one storage may back many `TensorId`s (views, KV ring).
*/
@ExperimentalMemoryApi
@JvmInline
public value class StorageId(public val value: Long) {
override fun toString(): String = "#$value"

public companion object {
private val counter = AtomicLong(1L)
/** The next id; thread-safe. */
public fun next(): StorageId = StorageId(counter.fetchAndIncrement())
}
}

/**
* How a [Storage] came to hold its bytes (SKEEP-003 §0 *Owner*, §4.4). Ownership is a
* constructor argument, then it is enforced: [Borrowed] storage cannot be freed through us,
* [Alias] keeps its parent alive and cannot free or resize, [Owned] storage is freed exactly once,
* by its scope.
*/
@ExperimentalMemoryApi
public sealed interface Owner {
/** We allocated the bytes; the scope of [scope] kind frees them. */
public data class Owned(val scope: ScopeKind) : Owner
/** The caller's array / buffer / segment / mmap; we never free it. [external] identifies the lender for debugging. */
public data class Borrowed(val external: Any? = null) : Owner
/** A view's storage reference: a strong reference to [parent]; mutability delegated. */
public class Alias(public val parent: Storage) : Owner {
override fun toString(): String = "Alias(parent=${parent.id})"
}
}

/** Thrown on any access to a storage after its scope or the storage itself was closed (SKEEP-003 rule 2). */
@ExperimentalMemoryApi
public class StorageClosedException(
public val storageId: StorageId,
public val origin: TensorId?,
message: String = "Storage ${storageId}${origin?.let { " (" + it.canonical + ")" } ?: ""} is closed",
) : IllegalStateException(message)

/**
* The one and only owner of bytes (SKEEP-003 §0, §4.2). `TensorView` interprets a storage, `Tensor`
* is the DSL handle over a view — neither owns bytes. Sealed over the four kinds; [Heap] is final
* and common, [OffHeap] / [Mapped] / [Device] are abstract here and bound per platform
* (`MemorySegment` / `FileChannel.map` on the JVM, `malloc` / `mmap` on Native, heap fallbacks on
* JS/Wasm — slices #1019, #1020).
*
* Rules enforced here: exactly one byte owner; closing invalidates every alias; a borrowed storage
* is released (forgotten) but never freed; every access after close throws
* [StorageClosedException] carrying the id and origin — not a JVM crash, not silent corruption.
*/
@ExperimentalMemoryApi
public sealed class Storage : AutoCloseable {
public abstract val id: StorageId
public abstract val sizeBytes: Long
public abstract val owner: Owner
public abstract val domain: MemoryDomain
/** The `TensorId` these bytes back, for diagnostics; `null` for anonymous storage. */
public abstract val debugOrigin: TensorId?
/** Where trace events about this storage go (allocation, close). */
protected abstract val sink: TraceSink

/** The lifetime class: from [Owner.Owned], else the parent's, else `AMBIENT`. */
public val scope: ScopeKind
get() = when (val o = owner) {
is Owner.Owned -> o.scope
is Owner.Alias -> o.parent.scope
is Owner.Borrowed -> ScopeKind.AMBIENT
}

private var closed: Boolean = false

/** `true` until this storage — or, for an alias, its parent — is closed. */
public val isAlive: Boolean
get() = !closed && ((owner as? Owner.Alias)?.parent?.isAlive ?: true)

/** Whether writes are allowed: owned and borrowed-mutable storage yes; an alias delegates to its parent. */
public abstract val isMutable: Boolean

/** Throws [StorageClosedException] if this storage is no longer alive. Called by every accessor. */
public fun checkAlive() { if (!isAlive) throw StorageClosedException(id, debugOrigin) }

/**
* Close: an [Owner.Owned] storage releases its bytes (exactly once); an [Owner.Borrowed] storage
* is forgotten (the lender's bytes are untouched); an [Owner.Alias] is detached (its parent is
* unaffected). Idempotent.
*/
final override fun close() {
if (closed) return
closed = true
onClose()
if (sink.isEnabled && owner !is Owner.Alias) sink.emit(TraceEvent.Free(id.value, scope, sizeBytes))
}

/** Release platform resources (owned storage only); default nothing. */
protected open fun onClose() {}

/** A zero-copy alias over `[offsetBytes, offsetBytes + lengthBytes)` of this storage. */
public abstract fun slice(offsetBytes: Long, lengthBytes: Long): Storage

override fun toString(): String = "${this::class.simpleName}(${id}, ${sizeBytes} B, $owner, $domain${debugOrigin?.let { ", $it" } ?: ""}${if (isAlive) "" else ", closed"})"

/**
* Heap storage: a Kotlin array on the managed heap — the JIT-friendliest kind, the only kind on
* JS/Wasm, the default for `Ambient` scope. Exactly one of [floats], [ints], [bytes] is non-null;
* [arrayOffset] (in elements of that array) and [sizeBytes] delimit the region.
*
* Kernels unwrap once per call (`floats` / `ints` / `bytes` + [arrayOffset]) — the Phase-2 spike
* showed per-element access through a view is the slow path by design.
*/
public class Heap private constructor(
override val id: StorageId,
public val floats: FloatArray?,
public val ints: IntArray?,
public val bytes: ByteArray?,
public val arrayOffset: Int,
override val sizeBytes: Long,
override val owner: Owner,
override val debugOrigin: TensorId?,
override val sink: TraceSink,
private val mutable: Boolean,
) : Storage() {
override val domain: MemoryDomain get() = MemoryDomain.HOST_HEAP
override val isMutable: Boolean get() = (owner as? Owner.Alias)?.parent?.isMutable ?: mutable

/** Bytes per element of the backing array (4 for floats/ints, 1 for bytes). */
public val elementBytes: Int get() = if (bytes != null) 1 else 4
/** Number of array elements this storage spans. */
public val elementCount: Int get() = (sizeBytes / elementBytes).toInt()

override fun slice(offsetBytes: Long, lengthBytes: Long): Heap {
checkAlive()
require(offsetBytes >= 0 && lengthBytes >= 0 && offsetBytes + lengthBytes <= sizeBytes) { "slice [$offsetBytes, ${offsetBytes + lengthBytes}) outside $sizeBytes bytes" }
require(offsetBytes % elementBytes == 0L && lengthBytes % elementBytes == 0L) { "slice must align to $elementBytes-byte elements" }
return Heap(StorageId.next(), floats, ints, bytes, arrayOffset + (offsetBytes / elementBytes).toInt(), lengthBytes, Owner.Alias(this), debugOrigin, sink, mutable)
}

public companion object {
private fun create(floats: FloatArray?, ints: IntArray?, bytes: ByteArray?, offset: Int, count: Int, owner: Owner, origin: TensorId?, sink: TraceSink, mutable: Boolean): Heap {
val eb = if (bytes != null) 1 else 4
val s = Heap(StorageId.next(), floats, ints, bytes, offset, count.toLong() * eb, owner, origin, sink, mutable)
if (sink.isEnabled && owner is Owner.Owned) sink.emit(TraceEvent.Allocation(s.id.value, owner.scope, s.sizeBytes, origin))
return s
}

/** Allocate [count] zeroed floats on the heap, owned by a scope of kind [scope]. */
public fun floats(count: Int, scope: ScopeKind = ScopeKind.AMBIENT, origin: TensorId? = null, sink: TraceSink = NoopTraceSink): Heap =
create(FloatArray(count), null, null, 0, count, Owner.Owned(scope), origin, sink, true)
public fun ints(count: Int, scope: ScopeKind = ScopeKind.AMBIENT, origin: TensorId? = null, sink: TraceSink = NoopTraceSink): Heap =
create(null, IntArray(count), null, 0, count, Owner.Owned(scope), origin, sink, true)
public fun bytes(count: Int, scope: ScopeKind = ScopeKind.AMBIENT, origin: TensorId? = null, sink: TraceSink = NoopTraceSink): Heap =
create(null, null, ByteArray(count), 0, count, Owner.Owned(scope), origin, sink, true)

/** Wrap the caller's array without copying — never freed by us (the #782 `copyOf` replacement). */
public fun wrap(array: FloatArray, offset: Int = 0, count: Int = array.size - offset, mutable: Boolean = true, origin: TensorId? = null, sink: TraceSink = NoopTraceSink): Heap =
create(array, null, null, offset, count, Owner.Borrowed(array), origin, sink, mutable)
public fun wrap(array: IntArray, offset: Int = 0, count: Int = array.size - offset, mutable: Boolean = true, origin: TensorId? = null, sink: TraceSink = NoopTraceSink): Heap =
create(null, array, null, offset, count, Owner.Borrowed(array), origin, sink, mutable)
public fun wrap(array: ByteArray, offset: Int = 0, count: Int = array.size - offset, mutable: Boolean = true, origin: TensorId? = null, sink: TraceSink = NoopTraceSink): Heap =
create(null, null, array, offset, count, Owner.Borrowed(array), origin, sink, mutable)
}
}

/** Off-heap storage (`MemorySegment` / direct buffer / `malloc`): bound per platform in #1019/#1020. */
public abstract class OffHeap : Storage() { override val domain: MemoryDomain get() = MemoryDomain.HOST_OFFHEAP }

/** A mapped file region (`FileChannel.map` / `mmap`): bound per platform in #1019/#1020. */
public abstract class Mapped : Storage() { override val domain: MemoryDomain get() = MemoryDomain.MMAP_FILE }

/** An accelerator buffer — placeholder until a device backend is scheduled (PRD non-goal). */
public abstract class Device : Storage() { override val domain: MemoryDomain get() = MemoryDomain.DEVICE_LOCAL }
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
package sk.ainet.lang.memory.trace

import sk.ainet.lang.memory.ExperimentalMemoryApi
import sk.ainet.lang.memory.plan.MemoryPlan

/** Emit this plan as a [TraceEvent.Plan] so the plan-vs-actual check (#1030) can find it in the stream. */
@ExperimentalMemoryApi
public fun MemoryPlan.emit(sink: TraceSink) {
if (!sink.isEnabled) return
sink.emit(
TraceEvent.Plan(
model = input.modelName, ctx = input.ctx,
weightsBytes = weightsBytes, kvBytes = kvBytes, forwardBytes = forwardBytes, headroomBytes = headroomBytes,
budgetBytes = budget?.bytes, fits = fits,
),
)
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
package sk.ainet.lang.memory.trace

import sk.ainet.lang.memory.ExperimentalMemoryApi
import sk.ainet.lang.memory.Format
import sk.ainet.lang.memory.ScopeKind
import sk.ainet.lang.tensor.TensorId

/**
* One event of the SKaiNET observability stream (SKEEP-003 §4.9): phases, kernel runs, adapter
* insertions, allocations and platform counters share one event model, keyed by [TensorId],
* storage id and [ScopeKind]. Exporters (Perfetto / JFR / `android.os.Trace`, M1 slice #1025),
* the memory debugger and the benchmark report are consumers of the same stream.
*
* Events are small value objects; [timeNanos] is a monotonic timestamp in nanoseconds
* ([TraceClock.nowNanos]), comparable only within one process.
*/
@ExperimentalMemoryApi
public sealed interface TraceEvent {
public val timeNanos: Long

/** A phase opened: `load`, `compile`, `prefill`, `decode` (with [step]), `sample`, or a module span like `layers[3].attn`. */
public data class PhaseBegin(val phase: String, val step: Int? = null, val attributes: Map<String, String> = emptyMap(), override val timeNanos: Long = TraceClock.nowNanos()) : TraceEvent

/** The matching phase closed; [durationNanos] is filled by [TraceSink.phase]. */
public data class PhaseEnd(val phase: String, val step: Int? = null, val durationNanos: Long = 0L, override val timeNanos: Long = TraceClock.nowNanos()) : TraceEvent

/** A kernel ran: which op, which registered kernel (the `KernelKey` string once M1 has it), on which tensors, how many bytes it touched. */
public data class KernelRun(
val op: String,
val kernel: String,
val inputs: List<TensorId?> = emptyList(),
val output: TensorId? = null,
val bytesRead: Long = 0L,
val bytesWritten: Long = 0L,
val durationNanos: Long = 0L,
override val timeNanos: Long = TraceClock.nowNanos(),
) : TraceEvent

/** The dispatcher inserted a conversion (dequantize, requantize, gather) — always visible (§5.1). */
public data class AdapterInserted(
val kind: String,
val from: Format,
val to: Format,
val bytes: Long,
val target: TensorId? = null,
val scope: ScopeKind = ScopeKind.FORWARD,
override val timeNanos: Long = TraceClock.nowNanos(),
) : TraceEvent

/** A storage was allocated. [site] is the allocation site in debug mode, [origin] the TensorId it backs. */
public data class Allocation(
val storageId: Long,
val scope: ScopeKind,
val bytes: Long,
val origin: TensorId? = null,
val site: String? = null,
override val timeNanos: Long = TraceClock.nowNanos(),
) : TraceEvent

/** A storage was freed / closed. */
public data class Free(val storageId: Long, val scope: ScopeKind, val bytes: Long, override val timeNanos: Long = TraceClock.nowNanos()) : TraceEvent

/** A `Forward` (or other) scope was reset: how many bytes were live before and after. */
public data class ScopeReset(val scope: ScopeKind, val liveBytesBefore: Long, val liveBytesAfter: Long, override val timeNanos: Long = TraceClock.nowNanos()) : TraceEvent

/** A platform counter sample (RSS, page faults, heap, direct memory …). */
public data class Counter(val name: String, val value: Long, val unit: String = "bytes", override val timeNanos: Long = TraceClock.nowNanos()) : TraceEvent

/** A memory plan was computed (M0 `MemoryPlan`): the plan-vs-actual check (#1030) compares this with the allocation events. */
public data class Plan(
val model: String,
val ctx: Int,
val weightsBytes: Long,
val kvBytes: Long,
val forwardBytes: Long,
val headroomBytes: Long,
val budgetBytes: Long? = null,
val fits: Boolean? = null,
override val timeNanos: Long = TraceClock.nowNanos(),
) : TraceEvent {
val totalBytes: Long get() = weightsBytes + kvBytes + forwardBytes + headroomBytes
}
}

/** Monotonic clock for trace timestamps (`kotlin.time.TimeSource.Monotonic`). */
@ExperimentalMemoryApi
public object TraceClock {
private val start = kotlin.time.TimeSource.Monotonic.markNow()
public fun nowNanos(): Long = start.elapsedNow().inWholeNanoseconds
}
Loading
Loading