| 1 | #include <kassert.h> |
| 2 | #include <mem/alloc.h> |
| 3 | #include <stat_series.h> |
| 4 | #include <time/time.h> |
| 5 | |
| 6 | void 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 | |
| 23 | struct 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 | |
| 43 | void 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 | |
| 57 | void 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 | |
| 104 | out: |
| 105 | if (!already_locked) |
| 106 | spin_unlock(lock: &s->lock, old: irql); |
| 107 | } |
| 108 | |
| 109 | void stat_series_advance(struct stat_series *s, time_t now_us) { |
| 110 | stat_series_advance_internal(s, now_us, /* already_locked = */ false); |
| 111 | } |
| 112 | |
| 113 | void 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 | |