///|
/// Online location and scale state for bounded-memory telemetry processing.
pub struct StreamingLocationScale {
window : Int
mut values : Array[Double]
mut count : Int
mut total : Double
mut squared_total : Double
}
///|
pub struct StreamingSnapshot {
count : Int
mean : Double
scale : Double
median : Double
minimum : Double
maximum : Double
last : Double
trend : Double
}
///|
pub struct StreamingRegression {
window : Int
x_values : Array[Double]
y_values : Array[Double]
mut count : Int
}
///|
pub struct StreamingAlert {
index : Int
value : Double
center : Double
scale : Double
score : Double
alert : Bool
}
///|
pub fn streaming_location_scale(window : Int) -> StreamingLocationScale {
{
window: if window < 2 {
2
} else {
window
},
values: [],
count: 0,
total: 0.0,
squared_total: 0.0,
}
}
///|
pub fn streaming_push(state : StreamingLocationScale, value : Double) -> Unit {
state.values.push(value)
state.count += 1
state.total += value
state.squared_total += value * value
if state.values.length() > state.window {
let removed = state.values.remove(0)
state.total -= removed
state.squared_total -= removed * removed
}
}
///|
pub fn streaming_reset(state : StreamingLocationScale) -> Unit {
state.values = []
state.count = 0
state.total = 0.0
state.squared_total = 0.0
}
///|
pub fn streaming_count(state : StreamingLocationScale) -> Int {
state.values.length()
}
///|
pub fn streaming_mean(state : StreamingLocationScale) -> Double {
if state.values.length() == 0 {
0.0
} else {
state.total / state.values.length().to_double()
}
}
///|
pub fn streaming_variance(state : StreamingLocationScale) -> Double {
let count = state.values.length()
if count <= 1 {
0.0
} else {
let value = (
state.squared_total - state.total * state.total / count.to_double()
) /
(count - 1).to_double()
if value < 0.0 {
0.0
} else {
value
}
}
}
///|
pub fn streaming_scale(state : StreamingLocationScale) -> Double {
streaming_variance(state).sqrt()
}
///|
pub fn streaming_median(state : StreamingLocationScale) -> Double {
median(state.values)
}
///|
pub fn streaming_snapshot(state : StreamingLocationScale) -> StreamingSnapshot {
{
count: state.values.length(),
mean: streaming_mean(state),
scale: streaming_scale(state),
median: streaming_median(state),
minimum: min_value(state.values),
maximum: max_value(state.values),
last: if state.values.length() == 0 {
0.0
} else {
state.values[state.values.length() - 1]
},
trend: if state.values.length() < 2 {
0.0
} else {
state.values[state.values.length() - 1] - state.values[0]
},
}
}
///|
pub fn streaming_snapshot_vector(snapshot : StreamingSnapshot) -> Array[Double] {
[
snapshot.count.to_double(),
snapshot.mean,
snapshot.scale,
snapshot.median,
snapshot.minimum,
snapshot.maximum,
snapshot.last,
snapshot.trend,
]
}
///|
pub fn streaming_snapshot_lines(snapshot : StreamingSnapshot) -> Array[String] {
[
"count=" + snapshot.count.to_string(),
"mean=" + snapshot.mean.to_string(),
"scale=" + snapshot.scale.to_string(),
"median=" + snapshot.median.to_string(),
"minimum=" + snapshot.minimum.to_string(),
"maximum=" + snapshot.maximum.to_string(),
"last=" + snapshot.last.to_string(),
"trend=" + snapshot.trend.to_string(),
]
}
///|
pub fn streaming_snapshot_string(snapshot : StreamingSnapshot) -> String {
streaming_snapshot_lines(snapshot).join("\n")
}
///|
pub fn streaming_series(
data : Array[Double],
window : Int,
) -> Array[StreamingSnapshot] {
let state = streaming_location_scale(window)
let result = []
for value in data {
streaming_push(state, value)
result.push(streaming_snapshot(state))
}
result
}
///|
pub fn streaming_alerts(
data : Array[Double],
window : Int,
threshold : Double,
) -> Array[StreamingAlert] {
let state = streaming_location_scale(window)
let result = []
let safe_threshold = if threshold < 0.0 { -threshold } else { threshold }
for index = 0; index < data.length(); index = index + 1 {
let center = streaming_mean(state)
let scale = streaming_scale(state)
let score = if scale == 0.0 {
0.0
} else {
abs_double(data[index] - center) / scale
}
result.push({
index,
value: data[index],
center,
scale,
score,
alert: score > safe_threshold && state.values.length() > 1,
})
streaming_push(state, data[index])
}
result
}
///|
pub fn streaming_alert_indices(alerts : Array[StreamingAlert]) -> Array[Int] {
let result = []
for alert in alerts {
if alert.alert {
result.push(alert.index)
}
}
result
}
///|
pub fn streaming_outlier_rate(alerts : Array[StreamingAlert]) -> Double {
if alerts.length() == 0 {
0.0
} else {
streaming_alert_indices(alerts).length().to_double() /
alerts.length().to_double()
}
}
///|
pub fn streaming_regression(window : Int) -> StreamingRegression {
{
window: if window < 2 {
2
} else {
window
},
x_values: [],
y_values: [],
count: 0,
}
}
///|
pub fn streaming_regression_push(
state : StreamingRegression,
x : Double,
y : Double,
) -> Unit {
state.x_values.push(x)
state.y_values.push(y)
state.count += 1
if state.x_values.length() > state.window {
let _ = state.x_values.remove(0)
let _ = state.y_values.remove(0)
}
}
///|
pub fn streaming_regression_slope(state : StreamingRegression) -> Double {
if state.x_values.length() < 2 {
return 0.0
}
let x_center = mean(state.x_values)
let y_center = mean(state.y_values)
let mut numerator = 0.0
let mut denominator = 0.0
for index = 0; index < state.x_values.length(); index = index + 1 {
let dx = state.x_values[index] - x_center
numerator += dx * (state.y_values[index] - y_center)
denominator += dx * dx
}
if denominator == 0.0 {
0.0
} else {
numerator / denominator
}
}
///|
pub fn streaming_regression_intercept(state : StreamingRegression) -> Double {
mean(state.y_values) -
streaming_regression_slope(state) * mean(state.x_values)
}
///|
pub fn streaming_regression_predict(
state : StreamingRegression,
x : Double,
) -> Double {
streaming_regression_intercept(state) + streaming_regression_slope(state) * x
}
///|
pub fn streaming_regression_residuals(
state : StreamingRegression,
) -> Array[Double] {
let result = []
for index = 0; index < state.x_values.length(); index = index + 1 {
result.push(
state.y_values[index] -
streaming_regression_predict(state, state.x_values[index]),
)
}
result
}
///|
pub fn streaming_regression_r_squared(state : StreamingRegression) -> Double {
if state.y_values.length() < 2 {
return 0.0
}
let residuals = streaming_regression_residuals(state)
let total = sum_squared(
transform_center(state.y_values, mean(state.y_values)),
)
if total == 0.0 {
1.0
} else {
1.0 - sum_squared(residuals) / total
}
}
///|
pub fn streaming_batch_snapshot(
data : Array[Double],
window : Int,
) -> StreamingSnapshot {
let state = streaming_location_scale(window)
for value in data {
streaming_push(state, value)
}
streaming_snapshot(state)
}
///|
pub fn streaming_batch_alert_score(
data : Array[Double],
window : Int,
threshold : Double,
) -> Double {
streaming_outlier_rate(streaming_alerts(data, window, threshold))
}
///|
pub fn streaming_change_score(
data : Array[Double],
window : Int,
) -> Array[Double] {
let snapshots = streaming_series(data, window)
let result = []
for snapshot in snapshots {
result.push(abs_double(snapshot.trend))
}
result
}
///|
pub fn streaming_online_mad(
data : Array[Double],
window : Int,
) -> Array[Double] {
let state = streaming_location_scale(window)
let result = []
for value in data {
streaming_push(state, value)
result.push(mad(state.values))
}
result
}
///|
pub fn streaming_quantile_path(
data : Array[Double],
window : Int,
probability : Double,
) -> Array[Double] {
let state = streaming_location_scale(window)
let result = []
for value in data {
streaming_push(state, value)
result.push(quantile(state.values, probability))
}
result
}
///|
pub fn streaming_ewma(data : Array[Double], alpha : Double) -> Array[Double] {
let result = []
if data.length() == 0 {
return result
}
let weight = if alpha < 0.0 { 0.0 } else if alpha > 1.0 { 1.0 } else { alpha }
let mut current = data[0]
for index = 0; index < data.length(); index = index + 1 {
current = weight * data[index] + (1.0 - weight) * current
result.push(current)
}
result
}
///|
pub fn streaming_ewma_residuals(
data : Array[Double],
alpha : Double,
) -> Array[Double] {
let fitted = streaming_ewma(data, alpha)
let result = []
for index = 0; index < data.length(); index = index + 1 {
result.push(data[index] - fitted[index])
}
result
}
///|
pub fn streaming_exponentially_weighted_scale(
data : Array[Double],
alpha : Double,
) -> Array[Double] {
let residuals = streaming_ewma_residuals(data, alpha)
let result = []
let mut variance = 0.0
let weight = if alpha < 0.0 { 0.0 } else if alpha > 1.0 { 1.0 } else { alpha }
for residual in residuals {
variance = weight * residual * residual + (1.0 - weight) * variance
result.push(variance.sqrt())
}
result
}
///|
pub fn streaming_forecast(
data : Array[Double],
horizon : Int,
alpha : Double,
) -> Array[Double] {
let result = []
let fitted = streaming_ewma(data, alpha)
let last = if fitted.length() == 0 {
0.0
} else {
fitted[fitted.length() - 1]
}
let count = if horizon < 0 { 0 } else { horizon }
for _ in 0.. Double {
let alerts = streaming_alerts(data, window, threshold)
if data.length() == 0 {
1.0
} else {
1.0 - streaming_outlier_rate(alerts)
}
}