1#include <kassert.h>
2#include <mem/alloc.h>
3#include <stat_series.h>
4#include <time/time.h>
5
6void stat_series_init(struct stat_series *s, struct stat_bucket *buckets,
7 uint32_t nbuckets, time_t bucket_us,
8 stat_series_callback bucket_reset, void *private) {
9 memset(buckets, 0, nbuckets * sizeof(struct stat_bucket));
10 s->buckets = buckets;
11 s->bucket_reset = bucket_reset;
12 s->nbuckets = nbuckets;
13 s->bucket_us = bucket_us;
14 s->current = 0;
15 s->last_update_us = time_get_us();
16 s->private = private;
17 spinlock_init(lock: &s->lock);
18
19 for (size_t i = 0; i < nbuckets; i++)
20 s->buckets[i].parent = s;
21}
22
23struct stat_series *stat_series_create(uint32_t nbuckets, time_t bucket_us,
24 stat_series_callback bucket_reset,
25 void *private) {
26 struct stat_bucket *buckets =
27 kmalloc(sizeof(struct stat_bucket) * nbuckets, ALLOC_FLAGS_ZERO);
28 if (!buckets)
29 return NULL;
30
31 struct stat_series *series =
32 kmalloc(sizeof(struct stat_series), ALLOC_FLAGS_ZERO);
33 if (!series) {
34 kfree(buckets);
35 return NULL;
36 }
37
38 stat_series_init(s: series, buckets, nbuckets, bucket_us, bucket_reset,
39 private);
40 return series;
41}
42
43void stat_series_reset(struct stat_series *s) {
44 enum irql irql = spin_lock(lock: &s->lock);
45
46 struct stat_bucket *bucket;
47
48 stat_series_for_each(s, bucket) {
49 atomic_store(&bucket->count, 0);
50 atomic_store(&bucket->sum, 0);
51 s->bucket_reset(bucket);
52 }
53
54 spin_unlock(lock: &s->lock, old: irql);
55}
56
57void stat_series_advance_internal(struct stat_series *s, time_t now_us,
58 bool already_locked) {
59 enum irql irql = IRQL_NONE;
60
61 if (!already_locked) {
62 /* Fast-path: check without taking lock */
63 time_t last = atomic_load(&s->last_update_us);
64 size_t delta = now_us - last;
65 uint32_t steps = delta / s->bucket_us;
66 if (steps == 0)
67 return; /* no rotation required, avoid locking */
68 } else {
69 SPINLOCK_ASSERT_HELD(&s->lock);
70 }
71
72 /* Take lock and re-check/recompute using the protected fields */
73 if (!already_locked) {
74 irql = spin_lock(lock: &s->lock);
75 } else {
76 SPINLOCK_ASSERT_HELD(&s->lock);
77 }
78
79 /* re-evaluate based on locked state (canonical update) */
80 size_t delta = now_us - s->last_update_us;
81 uint32_t steps = delta / s->bucket_us;
82 if (steps == 0)
83 goto out;
84
85 if (steps > s->nbuckets)
86 steps = s->nbuckets;
87
88 size_t series_current = s->current;
89
90 for (uint32_t i = 0; i < steps; i++) {
91 /* update canonical non-atomic s->current while holding lock */
92 size_t current = (series_current + 1) % s->nbuckets;
93 struct stat_bucket *bucket = &s->buckets[current];
94
95 atomic_store(&bucket->count, 0);
96 atomic_store(&bucket->sum, 0);
97 s->bucket_reset(bucket);
98 s->current = current; /* publish */
99 }
100
101 /* publish last_update_us atomically for lockless readers/writers */
102 atomic_store(&s->last_update_us, s->last_update_us + steps * s->bucket_us);
103
104out:
105 if (!already_locked)
106 spin_unlock(lock: &s->lock, old: irql);
107}
108
109void stat_series_advance(struct stat_series *s, time_t now_us) {
110 stat_series_advance_internal(s, now_us, /* already_locked = */ false);
111}
112
113void stat_series_record(struct stat_series *s, size_t value,
114 stat_series_callback callback) {
115 time_t now_us = time_get_us();
116
117 /* attempt to advance if needed */
118 stat_series_advance_internal(s, now_us, /* already_locked = */ false);
119
120 /* read canonical current published by advance (atomic load) */
121 uint32_t cur = atomic_load(&s->current); /* acquire ordering */
122 struct stat_bucket *b = &s->buckets[cur];
123
124 /* hot path: only atomic RMWs on bucket */
125 atomic_fetch_add(&b->count, 1);
126 atomic_fetch_add(&b->sum, value);
127 if (callback)
128 callback(b);
129}
130