package dev.dtrentin.chart.buffer internal class TieredBuffer { private val tier0 = CircularBuffer(TIER0_CAPACITY) private val tier1 = CircularBuffer(TIER1_CAPACITY) private val tier2 = CircularBuffer(TIER2_CAPACITY) // M4 per-bin accumulator: first/min/max/last (ts, value) pairs. // Preserves visual shape across bin boundaries vs (midTs, min)/(midTs, max) pairs // that collapse to a vertical spike at flush time (C7 artifact). private var tier1BinStartMs = -1L private var tier1BinFirstTs = 0L; private var tier1BinFirstV = 0f private var tier1BinLastTs = 0L; private var tier1BinLastV = 0f private var tier1BinMinTs = 0L; private var tier1BinMinV = 0f private var tier1BinMaxTs = 0L; private var tier1BinMaxV = 0f private var tier1BinHasData = false private var tier2BinStartMs = -1L private var tier2BinFirstTs = 0L; private var tier2BinFirstV = 0f private var tier2BinLastTs = 0L; private var tier2BinLastV = 0f private var tier2BinMinTs = 0L; private var tier2BinMinV = 0f private var tier2BinMaxTs = 0L; private var tier2BinMaxV = 0f private var tier2BinHasData = false private val t0Ts = LongArray(TIER0_CAPACITY) private val t0Vs = FloatArray(TIER0_CAPACITY) private val t1Ts = LongArray(TIER1_CAPACITY) private val t1Vs = FloatArray(TIER1_CAPACITY) private val t2Ts = LongArray(TIER2_CAPACITY) private val t2Vs = FloatArray(TIER2_CAPACITY) // Reusable scratch for M4 flush sort/dedup (4 records max per bin). Avoids per-call alloc. private val flushTs = LongArray(4) private val flushVs = FloatArray(4) fun push(timestampMs: Long, value: Float) { tier0.push(timestampMs, value) feedTier1(timestampMs, value) feedTier2(timestampMs, value) } private fun feedTier1(ts: Long, value: Float) { if (tier1BinStartMs < 0L) { tier1BinStartMs = (ts / TIER1_BIN_MS) * TIER1_BIN_MS } if (ts >= tier1BinStartMs + TIER1_BIN_MS) { if (tier1BinHasData) { flushM4ToTier1() } tier1BinStartMs = (ts / TIER1_BIN_MS) * TIER1_BIN_MS tier1BinFirstTs = ts; tier1BinFirstV = value tier1BinLastTs = ts; tier1BinLastV = value tier1BinMinTs = ts; tier1BinMinV = value tier1BinMaxTs = ts; tier1BinMaxV = value tier1BinHasData = true } else { if (!tier1BinHasData) { tier1BinFirstTs = ts; tier1BinFirstV = value tier1BinLastTs = ts; tier1BinLastV = value tier1BinMinTs = ts; tier1BinMinV = value tier1BinMaxTs = ts; tier1BinMaxV = value tier1BinHasData = true } else { // Samples arrive in chronological order → always update last. tier1BinLastTs = ts; tier1BinLastV = value if (value < tier1BinMinV) { tier1BinMinV = value; tier1BinMinTs = ts } if (value > tier1BinMaxV) { tier1BinMaxV = value; tier1BinMaxTs = ts } } } } private fun feedTier2(ts: Long, value: Float) { if (tier2BinStartMs < 0L) { tier2BinStartMs = (ts / TIER2_BIN_MS) * TIER2_BIN_MS } if (ts >= tier2BinStartMs + TIER2_BIN_MS) { if (tier2BinHasData) { flushM4ToTier2() } tier2BinStartMs = (ts / TIER2_BIN_MS) * TIER2_BIN_MS tier2BinFirstTs = ts; tier2BinFirstV = value tier2BinLastTs = ts; tier2BinLastV = value tier2BinMinTs = ts; tier2BinMinV = value tier2BinMaxTs = ts; tier2BinMaxV = value tier2BinHasData = true } else { if (!tier2BinHasData) { tier2BinFirstTs = ts; tier2BinFirstV = value tier2BinLastTs = ts; tier2BinLastV = value tier2BinMinTs = ts; tier2BinMinV = value tier2BinMaxTs = ts; tier2BinMaxV = value tier2BinHasData = true } else { tier2BinLastTs = ts; tier2BinLastV = value if (value < tier2BinMinV) { tier2BinMinV = value; tier2BinMinTs = ts } if (value > tier2BinMaxV) { tier2BinMaxV = value; tier2BinMaxTs = ts } } } } /** * Push M4 records (first/min/max/last) for current tier1 bin in chronological order, * deduplicated by (ts, v). 4-element insertion sort. Worst case 4 distinct pushes, * best case 1 (single-sample bin). */ private fun flushM4ToTier1() { flushTs[0] = tier1BinFirstTs; flushVs[0] = tier1BinFirstV flushTs[1] = tier1BinMinTs; flushVs[1] = tier1BinMinV flushTs[2] = tier1BinMaxTs; flushVs[2] = tier1BinMaxV flushTs[3] = tier1BinLastTs; flushVs[3] = tier1BinLastV sort4ByTs(flushTs, flushVs) pushDistinctRecords(tier1, flushTs, flushVs) } private fun flushM4ToTier2() { flushTs[0] = tier2BinFirstTs; flushVs[0] = tier2BinFirstV flushTs[1] = tier2BinMinTs; flushVs[1] = tier2BinMinV flushTs[2] = tier2BinMaxTs; flushVs[2] = tier2BinMaxV flushTs[3] = tier2BinLastTs; flushVs[3] = tier2BinLastV sort4ByTs(flushTs, flushVs) pushDistinctRecords(tier2, flushTs, flushVs) } private fun sort4ByTs(ts: LongArray, vs: FloatArray) { // Insertion sort on 4 elements. Co-sorts vs alongside ts. for (i in 1 until 4) { val tk = ts[i]; val vk = vs[i] var j = i - 1 while (j >= 0 && ts[j] > tk) { ts[j + 1] = ts[j]; vs[j + 1] = vs[j] j-- } ts[j + 1] = tk; vs[j + 1] = vk } } private fun pushDistinctRecords(target: CircularBuffer, ts: LongArray, vs: FloatArray) { // Push first; subsequent only if (ts, v) differs from previous pushed pair. target.push(ts[0], vs[0]) var prevTs = ts[0]; var prevV = vs[0] for (i in 1 until 4) { if (ts[i] != prevTs || vs[i] != prevV) { target.push(ts[i], vs[i]) prevTs = ts[i]; prevV = vs[i] } } } fun latestTimestampMs(): Long = tier0.latestTimestampMs() fun snapshot( windowStartMs: Long, windowMs: Long, outTimestamps: LongArray, outValues: FloatArray, ): Int { val nowMs = tier0.latestTimestampMs() if (nowMs < 0L) return 0 val windowEndMs = windowStartMs + windowMs val tier0BoundaryMs = nowMs - TIER0_DURATION_MS val tier1BoundaryMs = nowMs - TIER0_DURATION_MS - TIER1_DURATION_MS var outIdx = 0 if (windowStartMs < tier1BoundaryMs) { val n2 = tier2.snapshot(t2Ts, t2Vs) for (i in 0 until n2) { val ts = t2Ts[i] if (ts >= windowStartMs && ts < windowEndMs && ts < tier1BoundaryMs) { if (outIdx >= outTimestamps.size) break outTimestamps[outIdx] = ts; outValues[outIdx] = t2Vs[i]; outIdx++ } } } if (windowStartMs < tier0BoundaryMs && windowEndMs > tier1BoundaryMs) { val n1 = tier1.snapshot(t1Ts, t1Vs) for (i in 0 until n1) { val ts = t1Ts[i] if (ts >= windowStartMs && ts < windowEndMs && ts >= tier1BoundaryMs && ts < tier0BoundaryMs) { if (outIdx >= outTimestamps.size) break outTimestamps[outIdx] = ts; outValues[outIdx] = t1Vs[i]; outIdx++ } } } val tier0Start = maxOf(windowStartMs, tier0BoundaryMs) val n0 = tier0.snapshot(t0Ts, t0Vs) for (i in 0 until n0) { val ts = t0Ts[i] if (ts >= tier0Start && ts < windowEndMs) { if (outIdx >= outTimestamps.size) break outTimestamps[outIdx] = ts; outValues[outIdx] = t0Vs[i]; outIdx++ } } return outIdx } /** * Bisect-based variant of [snapshot]. Returns identical content + ordering for the * same (windowStartMs, windowMs) args, but locates the window-start index in each * tier's chronologically-sorted snapshot via O(log n) bisect instead of an * O(n) linear pre-scan. Linear walk runs only over the in-window subrange. * * Output ordering: tier2 oldest first, then tier1, then tier0 newest — same as [snapshot]. * Tier boundary clamps (tier0BoundaryMs / tier1BoundaryMs) preserved verbatim. * * Returns 0 on empty buffer or window entirely outside data. */ fun snapshotWindow( windowStartMs: Long, windowMs: Long, outTimestamps: LongArray, outValues: FloatArray, ): Int { val nowMs = tier0.latestTimestampMs() if (nowMs < 0L) return 0 val windowEndMs = windowStartMs + windowMs val tier0BoundaryMs = nowMs - TIER0_DURATION_MS val tier1BoundaryMs = nowMs - TIER0_DURATION_MS - TIER1_DURATION_MS var outIdx = 0 if (windowStartMs < tier1BoundaryMs) { val n2 = tier2.snapshot(t2Ts, t2Vs) // Tier2 records satisfy ts < tier1BoundaryMs (older than tier1 horizon), // so upper clamp is min(windowEndMs, tier1BoundaryMs). val tier2UpperExclusive = if (windowEndMs < tier1BoundaryMs) windowEndMs else tier1BoundaryMs val start = bisectStart(t2Ts, n2, windowStartMs) var i = start while (i < n2) { val ts = t2Ts[i] if (ts >= tier2UpperExclusive) break if (outIdx >= outTimestamps.size) return outIdx outTimestamps[outIdx] = ts; outValues[outIdx] = t2Vs[i]; outIdx++ i++ } } if (windowStartMs < tier0BoundaryMs && windowEndMs > tier1BoundaryMs) { val n1 = tier1.snapshot(t1Ts, t1Vs) // Tier1 records satisfy tier1BoundaryMs <= ts < tier0BoundaryMs. val tier1Lower = if (windowStartMs > tier1BoundaryMs) windowStartMs else tier1BoundaryMs val tier1UpperExclusive = if (windowEndMs < tier0BoundaryMs) windowEndMs else tier0BoundaryMs val start = bisectStart(t1Ts, n1, tier1Lower) var i = start while (i < n1) { val ts = t1Ts[i] if (ts >= tier1UpperExclusive) break if (outIdx >= outTimestamps.size) return outIdx outTimestamps[outIdx] = ts; outValues[outIdx] = t1Vs[i]; outIdx++ i++ } } val tier0Start = if (windowStartMs > tier0BoundaryMs) windowStartMs else tier0BoundaryMs val n0 = tier0.snapshot(t0Ts, t0Vs) val start0 = bisectStart(t0Ts, n0, tier0Start) var i = start0 while (i < n0) { val ts = t0Ts[i] if (ts >= windowEndMs) break if (outIdx >= outTimestamps.size) return outIdx outTimestamps[outIdx] = ts; outValues[outIdx] = t0Vs[i]; outIdx++ i++ } return outIdx } fun clear() { tier0.clear(); tier1.clear(); tier2.clear() tier1BinStartMs = -1L tier1BinFirstTs = 0L; tier1BinFirstV = 0f tier1BinLastTs = 0L; tier1BinLastV = 0f tier1BinMinTs = 0L; tier1BinMinV = 0f tier1BinMaxTs = 0L; tier1BinMaxV = 0f tier1BinHasData = false tier2BinStartMs = -1L tier2BinFirstTs = 0L; tier2BinFirstV = 0f tier2BinLastTs = 0L; tier2BinLastV = 0f tier2BinMinTs = 0L; tier2BinMinV = 0f tier2BinMaxTs = 0L; tier2BinMaxV = 0f tier2BinHasData = false } /** * Tier capacities. TIER1/TIER2 multiplied by 4 because M4 binning (T12) emits up to * 4 records per bin (first, min, max, last) instead of the previous 2 (min, max). * Steady-state average likely 2–3 records/bin after (ts, v) dedup; ×4 is safe upper bound. * * TIER0_CAPACITY: 5 min × 200 Hz = 60_000 raw samples. * TIER1_CAPACITY: 10 min × 10 Hz × 4 M4 records = 24_000 records. * TIER2_CAPACITY: 45 min × 1 Hz × 4 M4 records = 10_800 records. * TOTAL_CAPACITY: 94_800 samples per signal. */ internal companion object { const val TIER0_MAX_HZ = 200 const val TIER1_HZ = 10 const val TIER2_HZ = 1 const val TIER0_DURATION_MS = 5L * 60_000L const val TIER1_DURATION_MS = 10L * 60_000L const val TIER2_DURATION_MS = 45L * 60_000L const val TIER1_BIN_MS = 1000L / TIER1_HZ const val TIER2_BIN_MS = 1000L / TIER2_HZ val TIER0_CAPACITY = ((TIER0_DURATION_MS / 1000L) * TIER0_MAX_HZ).toInt() val TIER1_CAPACITY = ((TIER1_DURATION_MS / 1000L) * TIER1_HZ).toInt() * 4 val TIER2_CAPACITY = ((TIER2_DURATION_MS / 1000L) * TIER2_HZ).toInt() * 4 val TOTAL_CAPACITY = TIER0_CAPACITY + TIER1_CAPACITY + TIER2_CAPACITY } }