Skip to content

Commit 5ff8e21

Browse files
author
MPCoreDeveloper
committed
feat(phase8.2): Add Bucket Storage System - TimeSeriesBucket, BucketPartitioner, BucketManager, TimeSeriesTable - all 24 tests passing
1 parent 7e8574f commit 5ff8e21

5 files changed

Lines changed: 1704 additions & 0 deletions

File tree

Lines changed: 389 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,389 @@
1+
// <copyright file="BucketManager.cs" company="MPCoreDeveloper">
2+
// Copyright (c) 2025-2026 MPCoreDeveloper and GitHub Copilot. All rights reserved.
3+
// Licensed under the MIT License. See LICENSE file in the project root for full license information.
4+
// </copyright>
5+
6+
namespace SharpCoreDB.TimeSeries;
7+
8+
using System;
9+
using System.Collections.Concurrent;
10+
using System.Collections.Generic;
11+
using System.Linq;
12+
using System.Threading;
13+
14+
/// <summary>
15+
/// Manages time-series bucket lifecycle.
16+
/// C# 14: Lock class, primary constructors, modern patterns.
17+
///
18+
/// ✅ SCDB Phase 8.2: Bucket Storage System
19+
///
20+
/// Purpose:
21+
/// - Automatic bucket creation and rotation
22+
/// - Tier management (Hot → Warm → Cold)
23+
/// - Bucket lookup and retrieval
24+
/// - Integration with Phase 8.1 compression
25+
/// </summary>
26+
public sealed class BucketManager : IDisposable
27+
{
28+
private readonly ConcurrentDictionary<string, TimeSeriesBucket> _buckets = new();
29+
private readonly ConcurrentDictionary<string, TimeSeriesBatch> _hotData = new();
30+
private readonly ConcurrentDictionary<string, TimeSeriesCompression.CompressedData> _warmTimestamps = new();
31+
private readonly ConcurrentDictionary<string, TimeSeriesCompression.CompressedData> _warmValues = new();
32+
private readonly Lock _bucketLock = new();
33+
private readonly BucketGranularity _defaultGranularity;
34+
private readonly TimeSpan _hotToColdThreshold;
35+
private bool _disposed;
36+
37+
/// <summary>
38+
/// Initializes a new instance of the <see cref="BucketManager"/> class.
39+
/// </summary>
40+
/// <param name="defaultGranularity">Default bucket granularity.</param>
41+
/// <param name="hotToWarmThreshold">Time after which hot data becomes warm (default: 1 hour).</param>
42+
public BucketManager(
43+
BucketGranularity defaultGranularity = BucketGranularity.Hour,
44+
TimeSpan? hotToWarmThreshold = null)
45+
{
46+
_defaultGranularity = defaultGranularity;
47+
_hotToColdThreshold = hotToWarmThreshold ?? TimeSpan.FromHours(1);
48+
}
49+
50+
/// <summary>Gets the number of buckets.</summary>
51+
public int BucketCount => _buckets.Count;
52+
53+
/// <summary>Gets the default granularity.</summary>
54+
public BucketGranularity DefaultGranularity => _defaultGranularity;
55+
56+
/// <summary>
57+
/// Gets or creates a bucket for the given timestamp.
58+
/// </summary>
59+
public TimeSeriesBucket GetOrCreateBucket(string tableName, DateTime timestamp)
60+
{
61+
ArgumentException.ThrowIfNullOrWhiteSpace(tableName);
62+
63+
var bucketId = BucketPartitioner.GetBucketId(tableName, timestamp, _defaultGranularity);
64+
65+
return _buckets.GetOrAdd(bucketId, _ =>
66+
{
67+
var bucketStart = BucketPartitioner.GetBucketStart(timestamp, _defaultGranularity);
68+
var bucketEnd = BucketPartitioner.GetBucketEnd(bucketStart, _defaultGranularity);
69+
70+
return new TimeSeriesBucket
71+
{
72+
BucketId = bucketId,
73+
TableName = tableName,
74+
StartTime = bucketStart,
75+
EndTime = bucketEnd,
76+
Granularity = _defaultGranularity,
77+
Tier = BucketTier.Hot
78+
};
79+
});
80+
}
81+
82+
/// <summary>
83+
/// Inserts data points into the appropriate bucket.
84+
/// </summary>
85+
public void Insert(string tableName, IEnumerable<TimeSeriesDataPoint> points)
86+
{
87+
ArgumentException.ThrowIfNullOrWhiteSpace(tableName);
88+
ArgumentNullException.ThrowIfNull(points);
89+
90+
// Group points by bucket
91+
var grouped = points.GroupBy(p =>
92+
BucketPartitioner.GetBucketId(tableName, p.Timestamp, _defaultGranularity));
93+
94+
foreach (var group in grouped)
95+
{
96+
var bucketId = group.Key;
97+
var bucket = GetOrCreateBucket(tableName, group.First().Timestamp);
98+
99+
// Add to hot data
100+
var batch = TimeSeriesBatch.FromDataPoints(group);
101+
InsertBatch(bucketId, bucket, batch);
102+
}
103+
}
104+
105+
/// <summary>
106+
/// Inserts a batch of data into a bucket.
107+
/// </summary>
108+
public void InsertBatch(string tableName, DateTime timestamp, TimeSeriesBatch batch)
109+
{
110+
ArgumentException.ThrowIfNullOrWhiteSpace(tableName);
111+
ArgumentNullException.ThrowIfNull(batch);
112+
113+
var bucket = GetOrCreateBucket(tableName, timestamp);
114+
InsertBatch(bucket.BucketId, bucket, batch);
115+
}
116+
117+
private void InsertBatch(string bucketId, TimeSeriesBucket bucket, TimeSeriesBatch batch)
118+
{
119+
lock (_bucketLock)
120+
{
121+
// Merge with existing hot data
122+
if (_hotData.TryGetValue(bucketId, out var existing))
123+
{
124+
// Merge batches
125+
var mergedTimestamps = existing.Timestamps.Concat(batch.Timestamps).OrderBy(t => t).ToArray();
126+
var mergedValues = new double[mergedTimestamps.Length];
127+
128+
// Build lookup for merging
129+
var existingDict = existing.Timestamps.Zip(existing.Values).ToDictionary(x => x.First, x => x.Second);
130+
var newDict = batch.Timestamps.Zip(batch.Values).ToDictionary(x => x.First, x => x.Second);
131+
132+
for (int i = 0; i < mergedTimestamps.Length; i++)
133+
{
134+
var ts = mergedTimestamps[i];
135+
mergedValues[i] = newDict.TryGetValue(ts, out var v) ? v : existingDict[ts];
136+
}
137+
138+
batch = new TimeSeriesBatch
139+
{
140+
Timestamps = mergedTimestamps,
141+
Values = mergedValues
142+
};
143+
}
144+
145+
_hotData[bucketId] = batch;
146+
147+
// Update bucket metadata
148+
bucket.RowCount = batch.Count;
149+
bucket.UncompressedSize = batch.Count * (sizeof(long) + sizeof(double));
150+
151+
if (batch.Count > 0)
152+
{
153+
bucket.MinTimestamp = new DateTime(batch.Timestamps.Min(), DateTimeKind.Utc);
154+
bucket.MaxTimestamp = new DateTime(batch.Timestamps.Max(), DateTimeKind.Utc);
155+
}
156+
157+
bucket.ModifiedAt = DateTime.UtcNow;
158+
}
159+
}
160+
161+
/// <summary>
162+
/// Queries data points within a time range.
163+
/// </summary>
164+
public IEnumerable<TimeSeriesDataPoint> Query(
165+
string tableName,
166+
DateTime startTime,
167+
DateTime endTime)
168+
{
169+
ArgumentException.ThrowIfNullOrWhiteSpace(tableName);
170+
171+
var bucketIds = BucketPartitioner.GetBucketIdsForRange(
172+
tableName, startTime, endTime, _defaultGranularity);
173+
174+
foreach (var bucketId in bucketIds)
175+
{
176+
if (!_buckets.TryGetValue(bucketId, out var bucket))
177+
{
178+
continue;
179+
}
180+
181+
TimeSeriesBatch? batch = null;
182+
183+
if (bucket.Tier == BucketTier.Hot)
184+
{
185+
// Read from hot data
186+
_hotData.TryGetValue(bucketId, out batch);
187+
}
188+
else if (bucket.Tier == BucketTier.Warm)
189+
{
190+
// Decompress from warm data
191+
batch = DecompressBucket(bucketId);
192+
}
193+
194+
if (batch == null)
195+
{
196+
continue;
197+
}
198+
199+
// Filter by time range
200+
var startTicks = startTime.Ticks;
201+
var endTicks = endTime.Ticks;
202+
203+
for (int i = 0; i < batch.Count; i++)
204+
{
205+
var ts = batch.Timestamps[i];
206+
if (ts >= startTicks && ts < endTicks)
207+
{
208+
yield return new TimeSeriesDataPoint
209+
{
210+
Timestamp = new DateTime(ts, DateTimeKind.Utc),
211+
Value = batch.Values[i]
212+
};
213+
}
214+
}
215+
}
216+
}
217+
218+
/// <summary>
219+
/// Compresses a hot bucket to warm tier using Phase 8.1 codecs.
220+
/// </summary>
221+
public bool CompressBucket(string bucketId)
222+
{
223+
if (!_buckets.TryGetValue(bucketId, out var bucket))
224+
{
225+
return false;
226+
}
227+
228+
if (bucket.Tier != BucketTier.Hot)
229+
{
230+
return false; // Already compressed
231+
}
232+
233+
if (!_hotData.TryGetValue(bucketId, out var batch))
234+
{
235+
return false;
236+
}
237+
238+
lock (_bucketLock)
239+
{
240+
// Compress using Phase 8.1 codecs
241+
var compressedTimestamps = TimeSeriesCompression.CompressTimestamps(batch.Timestamps);
242+
var compressedValues = TimeSeriesCompression.CompressValues(batch.Values);
243+
244+
_warmTimestamps[bucketId] = compressedTimestamps;
245+
_warmValues[bucketId] = compressedValues;
246+
247+
// Update bucket metadata
248+
bucket.Tier = BucketTier.Warm;
249+
bucket.CompressedSize = compressedTimestamps.Data.Length + compressedValues.Data.Length;
250+
bucket.IsSealed = true;
251+
bucket.ModifiedAt = DateTime.UtcNow;
252+
253+
// Remove hot data
254+
_hotData.TryRemove(bucketId, out _);
255+
}
256+
257+
return true;
258+
}
259+
260+
/// <summary>
261+
/// Decompresses a warm bucket back to a batch.
262+
/// </summary>
263+
private TimeSeriesBatch? DecompressBucket(string bucketId)
264+
{
265+
if (!_warmTimestamps.TryGetValue(bucketId, out var compressedTs) ||
266+
!_warmValues.TryGetValue(bucketId, out var compressedVals))
267+
{
268+
return null;
269+
}
270+
271+
var timestamps = TimeSeriesCompression.DecompressTimestamps(compressedTs);
272+
var values = TimeSeriesCompression.DecompressValues(compressedVals);
273+
274+
return new TimeSeriesBatch
275+
{
276+
Timestamps = timestamps,
277+
Values = values
278+
};
279+
}
280+
281+
/// <summary>
282+
/// Compresses all eligible hot buckets.
283+
/// </summary>
284+
public int CompressEligibleBuckets()
285+
{
286+
var now = DateTime.UtcNow;
287+
int compressed = 0;
288+
289+
foreach (var bucket in _buckets.Values.Where(b => b.Tier == BucketTier.Hot))
290+
{
291+
// Compress if bucket end time is past the threshold
292+
if (now - bucket.EndTime > _hotToColdThreshold)
293+
{
294+
if (CompressBucket(bucket.BucketId))
295+
{
296+
compressed++;
297+
}
298+
}
299+
}
300+
301+
return compressed;
302+
}
303+
304+
/// <summary>
305+
/// Gets all buckets for a table.
306+
/// </summary>
307+
public IEnumerable<TimeSeriesBucket> GetBuckets(string tableName)
308+
{
309+
ArgumentException.ThrowIfNullOrWhiteSpace(tableName);
310+
311+
return _buckets.Values
312+
.Where(b => b.TableName == tableName)
313+
.OrderBy(b => b.StartTime);
314+
}
315+
316+
/// <summary>
317+
/// Gets a bucket by ID.
318+
/// </summary>
319+
public TimeSeriesBucket? GetBucket(string bucketId)
320+
{
321+
return _buckets.TryGetValue(bucketId, out var bucket) ? bucket : null;
322+
}
323+
324+
/// <summary>
325+
/// Gets statistics for all buckets.
326+
/// </summary>
327+
public BucketManagerStats GetStats()
328+
{
329+
var buckets = _buckets.Values.ToList();
330+
331+
return new BucketManagerStats
332+
{
333+
TotalBuckets = buckets.Count,
334+
HotBuckets = buckets.Count(b => b.Tier == BucketTier.Hot),
335+
WarmBuckets = buckets.Count(b => b.Tier == BucketTier.Warm),
336+
ColdBuckets = buckets.Count(b => b.Tier == BucketTier.Cold),
337+
TotalRows = buckets.Sum(b => b.RowCount),
338+
TotalCompressedSize = buckets.Sum(b => b.CompressedSize),
339+
TotalUncompressedSize = buckets.Sum(b => b.UncompressedSize)
340+
};
341+
}
342+
343+
/// <summary>
344+
/// Disposes the bucket manager.
345+
/// </summary>
346+
public void Dispose()
347+
{
348+
if (_disposed) return;
349+
350+
_buckets.Clear();
351+
_hotData.Clear();
352+
_warmTimestamps.Clear();
353+
_warmValues.Clear();
354+
355+
_disposed = true;
356+
}
357+
}
358+
359+
/// <summary>
360+
/// Bucket manager statistics.
361+
/// </summary>
362+
public sealed record BucketManagerStats
363+
{
364+
/// <summary>Total number of buckets.</summary>
365+
public int TotalBuckets { get; init; }
366+
367+
/// <summary>Number of hot buckets.</summary>
368+
public int HotBuckets { get; init; }
369+
370+
/// <summary>Number of warm buckets.</summary>
371+
public int WarmBuckets { get; init; }
372+
373+
/// <summary>Number of cold buckets.</summary>
374+
public int ColdBuckets { get; init; }
375+
376+
/// <summary>Total row count across all buckets.</summary>
377+
public long TotalRows { get; init; }
378+
379+
/// <summary>Total compressed size in bytes.</summary>
380+
public long TotalCompressedSize { get; init; }
381+
382+
/// <summary>Total uncompressed size in bytes.</summary>
383+
public long TotalUncompressedSize { get; init; }
384+
385+
/// <summary>Overall compression ratio.</summary>
386+
public double CompressionRatio => TotalCompressedSize > 0
387+
? (double)TotalUncompressedSize / TotalCompressedSize
388+
: 1.0;
389+
}

0 commit comments

Comments
 (0)