Skip to main content

Time-Series Native

Core Concepts

Time-Series Native

Time-Series Native is a core design philosophy of EnerOS: time-series data is not an auxiliary feature of an external database, but a native storage engine of the operating system kernel. All operational states of the power system (voltage, current, power, frequency, temperature) are essentially time-series data. EnerOS builds the time-series storage engine as a kernel component, providing nanosecond precision and million-level throughput for writes and queries, without relying on external time-series databases such as InfluxDB, TimescaleDB, or Prometheus.

Design Motivation

Pain Points of Traditional Time-Series Solutions

Using external time-series databases in power systems has the following issues:

Pain PointDescriptionImpact
Network latencyApplication layer accesses database over networkIncreases decision latency by 1-50ms
Deployment complexityRequires independent deployment and operationsHigh operations cost
Data inconsistencyKernel state and database state are asynchronousDecisions based on stale data
Double-write problemKernel events write to both log and databaseResource waste
Multi-tenant limitationsExternal databases have weak multi-tenant supportHard to isolate different Agents
Failure propagationDatabase failure affects all applicationsReduced system availability

Time-Series Native Solution

EnerOS builds the time-series storage engine as a kernel component, accessible by all Agents through system calls:

┌─────────────────────────────────────────┐
│            Agent / Application Layer     │
├─────────────────────────────────────────┤
│  Time-Series System Call Interface       │
│  (TimeSeriesSyscall)                     │
├─────────────────────────────────────────┤
│         Time-Series Engine               │
│  ├─ Write Path (WAL + LSM-Tree)          │
│  ├─ Query Path (Index + Vectorized Exec) │
│  ├─ Aggregation Engine (Stream + Precompute) │
│  └─ Downsampling Engine (LTTB / M4 / MinMax) │
├─────────────────────────────────────────┤
│         Storage Backend                  │
│  ├─ SQLite + WAL (default)              │
│  ├─ In-Memory Cache (LRU)               │
│  └─ Columnar Compression (Gorilla / Delta) │
└─────────────────────────────────────────┘

Time-Series Engine Architecture

Write Path

Measurement Data → Memory Buffer → WAL Log → LSM-Tree → Compaction
   ↓          ↓          ↓
  Validate   Index Update  Persist

Query Path

Query Request → Time Index → Data Load → Vectorized Filter → Aggregation → Result Return
              ↓          ↓           ↓            ↓
           B+tree lookup  Columnar read  SIMD acceleration  Stream aggregation

Data Model

Point (Data Point)

The smallest unit of time-series data, containing measurement name, timestamp, value, and optional tags:

use eneros_timeseries::{Point, Timestamp, Value};

let point = Point::new(
    "bus_1_voltage",                    // Measurement name
    Timestamp::now(),                    // Nanosecond timestamp
    Value::Float(1.024),                 // Measured value
)
.with_tag("region", "east")             // Tag
.with_tag("voltage_level", "110kV")
.with_tag("substation", "ss-001");

Series (Time Series)

A unique combination of a measurement name and a set of tags constitutes a time series:

use eneros_timeseries::{Series, SeriesKey};

let series_key = SeriesKey::new("bus_1_voltage")
    .with_tag("region", "east")
    .with_tag("voltage_level", "110kV");

let series = ts.get_series(&series_key)?;
println!("Series ID: {}", series.id());
println!("Data point count: {}", series.count());
println!("Time range: {} - {}", series.start(), series.end());

Tag (Label)

Tags are used to identify and filter time series, supporting multi-dimensional queries:

Tag KeyExample ValuesPurpose
regioneast / west / northDispatch region
voltage_level500kV / 220kV / 110kVVoltage level
substationss-001 / ss-002Substation
device_typebreaker / transformer / lineDevice type
phaseA / B / C / NPhase
measurement_typeinstantaneous / accumulatedMeasurement type

Write API

Single Point Write

use eneros_timeseries::{TimeSeriesEngine, Point, Value};

let ts = TimeSeriesEngine::open("data/eneros.db")?;

// Single point write
ts.write(Point::new(
    "bus_1_voltage",
    Timestamp::now(),
    Value::Float(1.024),
))?;

Batch Write

// Batch write (recommended, 10x higher throughput)
let batch: Vec<Point> = vec![
    Point::new("bus_1_voltage",   Timestamp::now(), Value::Float(1.024)),
    Point::new("bus_1_frequency", Timestamp::now(), Value::Float(50.001)),
    Point::new("bus_1_p",         Timestamp::now(), Value::Float(45.5)),
    Point::new("bus_1_q",         Timestamp::now(), Value::Float(12.3)),
    Point::new("bus_2_voltage",   Timestamp::now(), Value::Float(0.998)),
    Point::new("bus_2_frequency", Timestamp::now(), Value::Float(49.998)),
];

ts.batch_write(&batch)?;

Stream Write

use eneros_timeseries::StreamWriter;

// Stream write (suitable for SCADA real-time data)
let mut writer = ts.stream_writer("bus_1_voltage")?;

for measurement in scada_stream {
    writer.write(measurement.timestamp, measurement.value)?;
}

writer.flush()?; // Explicit flush

Async Write

use eneros_timeseries::AsyncWriter;

// Async writer (background batch commit)
let writer = AsyncWriter::new(ts.clone())
    .batch_size(1000)                // 1000 points per batch
    .flush_interval(Duration::milliseconds(100)); // Or flush every 100ms

for point in measurement_stream {
    writer.send(point).await?; // Returns immediately, writes in background
}

writer.flush().await?; // Wait for all data to be written

Query API

Basic Query

use eneros_timeseries::{Query, Duration, Aggregation, Downsample};

// 1. Time range query
let data: Vec<Point> = ts.query()
    .point("bus_1_voltage")
    .range(now!() - Duration::hours(24), now!())
    .execute()?;

// 2. Multi-measurement query
let multi: Vec<SeriesResult> = ts.query()
    .points(vec!["bus_1_voltage", "bus_2_voltage", "bus_3_voltage"])
    .range(now!() - Duration::hours(1), now!())
    .execute()?;

// 3. Tag filter query
let filtered: Vec<Point> = ts.query()
    .point("voltage")
    .tag("region", "east")
    .tag("voltage_level", "110kV")
    .range(now!() - Duration::hours(1), now!())
    .execute()?;

Aggregation Query

// 4. Aggregation query (15-minute average)
let avg: Vec<Point> = ts.query()
    .point("bus_1_voltage")
    .range(now!() - Duration::days(7), now!())
    .aggregation(Aggregation::Avg, Duration::minutes(15))
    .execute()?;

// 5. Multi-aggregation query
let multi_agg: Vec<Point> = ts.query()
    .point("bus_1_voltage")
    .range(now!() - Duration::days(1), now!())
    .aggregations(vec![
        Aggregation::Avg,
        Aggregation::Min,
        Aggregation::Max,
        Aggregation::StdDev,
    ], Duration::minutes(15))
    .execute()?;

Downsampling Query

// 6. Downsampling (take the latest 1000 points)
let downsampled: Vec<Point> = ts.query()
    .point("bus_1_voltage")
    .range(now!() - Duration::days(30), now!())
    .downsample(Downsample::LTTB, 1000)
    .execute()?;

// 7. Aggregation + downsampling combination
let combined: Vec<Point> = ts.query()
    .point("bus_1_voltage")
    .range(now!() - Duration::days(30), now!())
    .aggregation(Aggregation::Avg, Duration::minutes(5))
    .downsample(Downsample::M4, 2000)
    .execute()?;

Aggregation Functions

EnerOS supports a rich set of aggregation functions:

Aggregation FunctionDescriptionTypical Use
AvgAverageVoltage/frequency mean
MinMinimumMinimum voltage analysis
MaxMaximumMaximum load analysis
SumSumCumulative energy
CountCountData completeness rate
StdDevStandard deviationFluctuation analysis
Percentile(p)PercentileP99 latency analysis
FirstFirst valueInitial state value
LastLast valueCurrent state
MedianMedianOutlier-resistant mean
RangeRangeFluctuation range
// Custom aggregation window
let result = ts.query()
    .point("total_load")
    .range(now!() - Duration::days(1), now!())
    .aggregation(Aggregation::Percentile(99.0), Duration::minutes(15))
    .execute()?;

Downsampling Algorithms

AlgorithmDescriptionApplicable Scenarios
LTTBLargest-Triangle-Three-BucketsVisualization (preserves trends)
M4Min-Min-Max-MaxCandlestick charts
MinMaxExtrema samplingProtection logic (preserves extrema)
AverageAverage samplingSmooth curves
FirstFirst point samplingState sequences
EveryEqual interval samplingSimple downsampling
// Visualization scenario: LTTB preserves trends
let viz_data = ts.query()
    .point("bus_1_voltage")
    .range(now!() - Duration::days(30), now!())
    .downsample(Downsample::LTTB, 1000)
    .execute()?;

// Protection scenario: MinMax preserves extrema
let protection_data = ts.query()
    .point("bus_1_voltage")
    .range(now!() - Duration::days(1), now!())
    .downsample(Downsample::MinMax, 5000)
    .execute()?;

Data Retention Policy

EnerOS supports multi-tier data retention policies to automatically manage data lifecycle:

use eneros_timeseries::{RetentionPolicy, TieredStorage};

let policy = RetentionPolicy::new()
    // Raw data: retain 7 days
    .tier("raw", Duration::days(7), Resolution::Raw)
    // 1-minute aggregation: retain 90 days
    .tier("1min", Duration::days(90), Resolution::Minutes(1))
    // 15-minute aggregation: retain 1 year
    .tier("15min", Duration::days(365), Resolution::Minutes(15))
    // 1-hour aggregation: retain permanently
    .tier("1hour", Duration::infinity(), Resolution::Hours(1));

ts.set_retention_policy("bus_.*_voltage", policy)?;

Retention Policy Configuration

# config/timeseries.toml

[[retention]]
pattern = "bus_.*_voltage"          # Measurement name pattern
raw_days = 7                        # Raw data retention days
tiers = [
    { name = "1min", days = 90, resolution = "1min" },
    { name = "15min", days = 365, resolution = "15min" },
    { name = "1hour", days = -1, resolution = "1hour" },  # -1 = permanent
]

[[retention]]
pattern = ".*_frequency"
raw_days = 30
tiers = [
    { name = "1sec", days = 7, resolution = "1sec" },
    { name = "1min", days = 365, resolution = "1min" },
]

[[retention]]
pattern = ".*_power"
raw_days = 14
tiers = [
    { name = "1min", days = 180, resolution = "1min" },
    { name = "15min", days = 1825, resolution = "15min" },  # 5 years
]

Integration with Agents

Agent Directly Queries Time-Series

use eneros_timeseries::{Duration, Aggregation};

impl Agent for ForecastAgent {
    async fn run(&mut self, ctx: &mut AgentContext) -> AgentResult<()> {
        // 1. Query historical load data (via system call)
        let history = ctx.query_timeseries()
            .point("total_load")
            .range(now!() - Duration::days(7), now!())
            .aggregation(Aggregation::Avg, Duration::minutes(15))
            .execute()?;

        // 2. Query weather data
        let weather = ctx.query_timeseries()
            .point("temperature")
            .range(now!() - Duration::days(7), now!())
            .aggregation(Aggregation::Avg, Duration::hours(1))
            .execute()?;

        // 3. Forecast based on historical data
        let forecast = self.forecast(&history, &weather)?;

        // 4. Write forecast results
        for point in forecast {
            ctx.write_timeseries(
                "forecast_load",
                point.timestamp,
                point.value,
            ).await?;
        }

        Ok(())
    }
}

Real-Time Subscription

Agents can subscribe to time-series data streams to achieve event-driven behavior:

use eneros_timeseries::Subscription;

impl Agent for ProtectionAgent {
    async fn run(&mut self, ctx: &mut AgentContext) -> AgentContext<()> {
        // Subscribe to voltage violation events
        let mut sub = ctx.subscribe_timeseries()
            .point("bus_1_voltage")
            .condition(|v| v < 0.95 || v > 1.05)
            .build()?;

        while let Some(event) = sub.recv().await {
            log::warn!("Voltage violation: {:?}", event);
            self.handle_voltage_violation(event).await?;
        }

        Ok(())
    }
}

Storage Backend

Default Backend: SQLite + WAL

use eneros_timeseries::StorageConfig;

let config = StorageConfig::sqlite("data/eneros.db")
    .wal_mode(true)                    // Enable WAL
    .cache_size_mb(256)                // 256MB cache
    .page_size(4096)                   // Page size
    .checkpoint_interval(Duration::minutes(5));

let ts = TimeSeriesEngine::open_with_config(config)?;

In-Memory Backend

// All-in-memory storage (suitable for testing or temporary data)
let config = StorageConfig::memory()
    .max_size_mb(512);

let ts = TimeSeriesEngine::open_with_config(config)?;

Compression Algorithms

AlgorithmApplicable Data TypesCompression RatioDecompression Speed
GorillaFloating point (voltage/current)8-12xExtremely fast
Delta-of-DeltaTimestamps10-20xExtremely fast
RLEState quantities (switch positions)100x+Extremely fast
SnappyGeneral purpose2-4xFast
LZ4General purpose2-5xExtremely fast

Performance Metrics

Write Performance

Write ModeThroughputLatency
Single point sync100k points/sec< 10μs
Batch sync (1000 points)1 million points/sec< 1ms
Async batch2 million points/sec< 1μs (enqueue)
Stream write1 million points/sec< 1μs

Query Performance

Query TypeData VolumeLatency
Single point query1 point< 1μs
Time range query10k points< 5ms
Time range query1 million points< 100ms
Aggregation query10 million points< 50ms
Downsampling query100 million points → 1000 points< 200ms
Multi-series query100 series × 10k points< 50ms

Storage Metrics

MetricValue
Raw data size1.2 GB/100 million points
After Gorilla compression100-150 MB/100 million points
Compression ratio8-12x
Index size50 MB/100 million points

Comparison with External Databases

Comparison with InfluxDB

DimensionInfluxDBEnerOS Time-Series Native
DeploymentIndependent processKernel component
Access methodHTTP/gRPCSystem call
Write latency1-5ms (network round trip)< 1μs (kernel mode)
Query latency5-50ms< 5ms
Throughput500k points/sec1 million points/sec
Multi-tenantWeakKernel native
Data consistencyEventual consistencyStrong consistency (with kernel state)
Operations costHigh (independent ops)Zero (kernel managed)
Failure impactIndependent failureSame lifecycle as kernel

Comparison with TimescaleDB

DimensionTimescaleDBEnerOS Time-Series Native
FoundationPostgreSQL extensionKernel component
SQL supportFull SQLQuery DSL
Write throughput300k points/sec1 million points/sec
Query latency10-100ms< 5ms
Transaction supportACIDACID
Deployment complexityHigh (requires PostgreSQL)Zero
Memory usageHigh (500MB+)Low (50MB)

Comparison with Prometheus

DimensionPrometheusEnerOS Time-Series Native
Design goalMonitoring metricsPower time-series data
Data modelMetric + labelsMeasurement + tags
Write methodPull modelPush (kernel direct write)
PrecisionMillisecondNanosecond
Retention policyFixed windowMulti-tier
Aggregation capabilityPromQLBuilt-in aggregation + downsampling
Applicable scenariosIT monitoringPower systems

Complete Usage Example

use eneros_timeseries::{TimeSeriesEngine, Point, Value, Aggregation, Duration, Downsample};

fn main() -> Result<(), Box<dyn std::error::Error>> {
    // 1. Open time-series engine
    let ts = TimeSeriesEngine::open("data/eneros.db")?;

    // 2. Configure retention policy
    let policy = eneros_timeseries::RetentionPolicy::new()
        .tier("raw", Duration::days(7), eneros_timeseries::Resolution::Raw)
        .tier("15min", Duration::days(365), eneros_timeseries::Resolution::Minutes(15));
    ts.set_retention_policy("bus_.*", policy)?;

    // 3. Write measurement data
    let now = eneros_timeseries::Timestamp::now();
    ts.batch_write(&[
        Point::new("bus_1_voltage", now, Value::Float(1.024)),
        Point::new("bus_1_frequency", now, Value::Float(50.001)),
        Point::new("bus_1_p", now, Value::Float(45.5)),
        Point::new("bus_1_q", now, Value::Float(12.3)),
    ])?;

    // 4. Query voltage for the last 1 hour
    let recent = ts.query()
        .point("bus_1_voltage")
        .range(now - Duration::hours(1), now)
        .execute()?;
    println!("Voltage points in last 1 hour: {}", recent.len());

    // 5. Aggregation query: 15-minute average for past 7 days
    let avg = ts.query()
        .point("bus_1_voltage")
        .range(now - Duration::days(7), now)
        .aggregation(Aggregation::Avg, Duration::minutes(15))
        .execute()?;
    println!("15-minute average points for past 7 days: {}", avg.len());

    // 6. Downsampling: reduce past 30 days data to 1000 points
    let downsampled = ts.query()
        .point("bus_1_voltage")
        .range(now - Duration::days(30), now)
        .downsample(Downsample::LTTB, 1000)
        .execute()?;
    println!("Points after downsampling: {}", downsampled.len());

    // 7. Statistics
    let stats = ts.stats("bus_1_voltage")?;
    println!("Series stats: {:?}", stats);

    Ok(())
}

Limitations and Trade-offs

Trade-offDescriptionMitigation Strategy
Single-node storageNot distributed by defaultSupports multi-region sync
No SQL supportUses query DSLProvides rich API
Memory usageKernel resident cacheConfigurable upper limit
Historical data migrationCross-engine migration requires toolsProvides CSV/Parquet import/export
Cross-node queriesRequires multi-region moduleeneros-multiregion support

Next Steps