Class: Magick::PerformanceMetrics
- Inherits:
-
Object
- Object
- Magick::PerformanceMetrics
- Defined in:
- lib/magick/performance_metrics.rb
Defined Under Namespace
Classes: Metric
Constant Summary collapse
- METRICS_RING_CAP =
1_000
Instance Attribute Summary collapse
-
#redis_enabled ⇒ Object
readonly
Public accessor for redis_enabled.
Instance Method Summary collapse
- #average_duration(feature_name: nil, operation: nil) ⇒ Object
- #clear! ⇒ Object
- #enable_redis_tracking(enable: true) ⇒ Object
-
#ensure_async_processor! ⇒ Object
Restart the async processor after a fork.
- #flush_to_redis ⇒ Object
- #flush_to_redis_if_needed ⇒ Object
-
#force_flush_to_redis ⇒ Object
Force flush pending updates to Redis immediately Useful when you need up-to-date stats across processes.
-
#initialize(batch_size: 100, flush_interval: 60, redis_enabled: nil) ⇒ PerformanceMetrics
constructor
A new instance of PerformanceMetrics.
- #most_used_features(limit: 10) ⇒ Object
-
#process_async_record(feature_name_str, operation_str, duration, success) ⇒ Object
Internal: Process metrics from async queue (runs in background thread).
- #record(feature_name, operation, duration, success: true) ⇒ Object
-
#start_async_processor ⇒ Object
Start background thread to process async metrics.
-
#stop_async_processor ⇒ Object
Stop async processor (for cleanup).
- #usage_count(feature_name) ⇒ Object
Constructor Details
#initialize(batch_size: 100, flush_interval: 60, redis_enabled: nil) ⇒ PerformanceMetrics
Returns a new instance of PerformanceMetrics.
29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 |
# File 'lib/magick/performance_metrics.rb', line 29 def initialize(batch_size: 100, flush_interval: 60, redis_enabled: nil) @metrics = [] @mutex = Mutex.new @usage_count = Hash.new(0) # Memory-only counts (reset on each process boot) @pending_updates = Hash.new(0) # For Redis batching (reset on each process boot) @flushed_counts = Hash.new(0) # Track counts that have been flushed to Redis (to avoid double-counting) @batch_size = batch_size @flush_interval = flush_interval @last_flush = Time.now # If redis_enabled is explicitly set, use it; otherwise default to false # It will be enabled later via enable_redis_tracking if Redis adapter is available @redis_enabled = redis_enabled.nil? ? false : redis_enabled # Cache expensive checks for performance @_rails_events_enabled = defined?(Magick::Rails::Events) && Magick::Rails::Events.rails8? @_adapter_available = nil # Will be cached on first check @_redis_available = nil # Will be cached on first check # Async recording queue for non-blocking metrics @async_queue = Queue.new @async_thread = nil @async_enabled = true # Enable async by default for performance @owner_pid = Process.pid start_async_processor end |
Instance Attribute Details
#redis_enabled ⇒ Object (readonly)
Public accessor for redis_enabled
73 74 75 |
# File 'lib/magick/performance_metrics.rb', line 73 def redis_enabled @redis_enabled end |
Instance Method Details
#average_duration(feature_name: nil, operation: nil) ⇒ Object
249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 |
# File 'lib/magick/performance_metrics.rb', line 249 def average_duration(feature_name: nil, operation: nil) # Calculate from memory metrics (current process, not yet flushed). # Read under the mutex - the async processor appends to and evicts from # @metrics concurrently. memory_sum, memory_count = @mutex.synchronize do filtered = @metrics.select do |m| (feature_name.nil? || m.feature_name == feature_name.to_s) && (operation.nil? || m.operation == operation.to_s) && m.success end [filtered.sum(&:duration), filtered.length] end # Also get from Redis if available (persisted across processes) redis_sum = 0.0 redis_count = 0 begin redis = stats_redis_client if redis if feature_name && operation # Specific feature and operation sum_key = "magick:duration:sum:#{feature_name}:#{operation}" count_key = "magick:duration:count:#{feature_name}:#{operation}" redis_sum = redis.get(sum_key).to_f redis_count = redis.get(count_key).to_i elsif feature_name # All operations for this feature prefix = "magick:duration:sum:#{feature_name}:" scan_stats_keys(redis, "#{prefix}*").each do |sum_key| op = sum_key.to_s.sub(prefix, '') count_key = "magick:duration:count:#{feature_name}:#{op}" redis_sum += redis.get(sum_key).to_f redis_count += redis.get(count_key).to_i end else # All features and operations (not recommended, but supported) scan_stats_keys(redis, 'magick:duration:sum:*').each do |sum_key| count_key = sum_key.to_s.sub(':sum:', ':count:') redis_sum += redis.get(sum_key).to_f redis_count += redis.get(count_key).to_i end end end rescue StandardError # Silently fail end total_sum = memory_sum + redis_sum total_count = memory_count + redis_count return 0.0 if total_count == 0 total_sum / total_count.to_f end |
#clear! ⇒ Object
357 358 359 360 361 362 363 364 |
# File 'lib/magick/performance_metrics.rb', line 357 def clear! @mutex.synchronize do @metrics.clear @usage_count.clear @pending_updates.clear @flushed_counts.clear end end |
#enable_redis_tracking(enable: true) ⇒ Object
225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 |
# File 'lib/magick/performance_metrics.rb', line 225 def enable_redis_tracking(enable: true) old_value = @redis_enabled @redis_enabled = enable # Flush any pending updates when enabling if enable && !@mutex.synchronize { @pending_updates.empty? } begin flush_to_redis rescue StandardError => e # Don't fail if flush fails - the flag is still set if defined?(::Rails) && ::Rails.env.development? warn "Magick: Failed to flush stats when enabling Redis tracking: #{e.message}" end end end # Verify the value was set (for debugging) if !(@redis_enabled == enable) && defined?(::Rails) && ::Rails.env.development? warn "Magick: Failed to set redis_enabled to #{enable}, current value: #{@redis_enabled}" end true end |
#ensure_async_processor! ⇒ Object
Restart the async processor after a fork. Child processes inherit a dead thread reference + a queue that was populated in the parent; both must be recreated. The inherited thread (if alive in the parent's address space at fork time) is killed so it cannot keep polling a detached queue.
58 59 60 61 62 63 64 65 66 67 68 69 70 |
# File 'lib/magick/performance_metrics.rb', line 58 def ensure_async_processor! return if @owner_pid == Process.pid stale_thread = @async_thread stale_queue = @async_queue @async_queue = Queue.new @async_thread = nil @owner_pid = Process.pid start_async_processor stale_queue&.close if stale_queue.respond_to?(:close) stale_thread&.kill end |
#flush_to_redis ⇒ Object
186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 |
# File 'lib/magick/performance_metrics.rb', line 186 def flush_to_redis return if @mutex.synchronize { @pending_updates.empty? } # Never drain pending state when there is nowhere to write it. Without # this guard a deployment with no Redis credits @flushed_counts for # updates that were never persisted, and usage_count reads back zero. redis = stats_redis_client return if redis.nil? updates_to_flush = nil duration_stats_to_flush = {} @mutex.synchronize do return if @pending_updates.empty? updates_to_flush = @pending_updates.dup flushed_feature_names = updates_to_flush.keys.to_set @pending_updates.clear # Collect duration stats for flushed features # Group metrics by feature_name and operation, sum durations and count occurrences. # Each group carries its own metric objects so that only samples we # actually persist get dropped from memory afterwards. @metrics.each do |metric| next unless flushed_feature_names.include?(metric.feature_name) && metric.success key = "#{metric.feature_name}:#{metric.operation}" duration_stats_to_flush[key] ||= { sum: 0.0, count: 0, feature_name: metric.feature_name, operation: metric.operation, metrics: [] } duration_stats_to_flush[key][:sum] += metric.duration duration_stats_to_flush[key][:count] += 1 duration_stats_to_flush[key][:metrics] << metric end end return if updates_to_flush.nil? || updates_to_flush.empty? write_flush_to_redis(redis, updates_to_flush, duration_stats_to_flush) end |
#flush_to_redis_if_needed ⇒ Object
156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 |
# File 'lib/magick/performance_metrics.rb', line 156 def flush_to_redis_if_needed # Cache adapter availability check (expensive) if @_adapter_available.nil? adapter = Magick.adapter_registry || Magick.default_adapter_registry @_adapter_available = adapter @_redis_available = adapter.is_a?(Magick::Adapters::Registry) && adapter.redis_available? if adapter end return unless @_adapter_available return unless @_redis_available || @redis_enabled should_flush = false @mutex.synchronize do next if @pending_updates.empty? # Flush if we have enough pending updates (sum of all counts) or enough time has passed # Check total count of pending updates, not just number of keys total_pending_count = @pending_updates.values.sum should_flush = true if total_pending_count >= @batch_size || (Time.now - @last_flush) >= @flush_interval end flush_to_redis if should_flush end |
#force_flush_to_redis ⇒ Object
Force flush pending updates to Redis immediately Useful when you need up-to-date stats across processes
182 183 184 |
# File 'lib/magick/performance_metrics.rb', line 182 def force_flush_to_redis flush_to_redis end |
#most_used_features(limit: 10) ⇒ Object
333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 |
# File 'lib/magick/performance_metrics.rb', line 333 def most_used_features(limit: 10) # Combine memory and Redis counts. The memory side is read under the # mutex because the async processor writes @usage_count. combined_counts = @mutex.synchronize { @usage_count.dup } # Always check Redis if adapter is available (not just if @redis_enabled) # This ensures we get the full count even if redis_enabled flag wasn't set correctly begin redis = stats_redis_client if redis # Get all stats keys scan_stats_keys(redis, 'magick:stats:*').each do |key| feature_name = key.to_s.sub('magick:stats:', '') redis_count = redis.get(key).to_i combined_counts[feature_name] = (combined_counts[feature_name] || 0) + redis_count end end rescue StandardError # Silently fail end combined_counts.sort_by { |_name, count| -count }.first(limit).to_h end |
#process_async_record(feature_name_str, operation_str, duration, success) ⇒ Object
Internal: Process metrics from async queue (runs in background thread)
93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 |
# File 'lib/magick/performance_metrics.rb', line 93 def process_async_record(feature_name_str, operation_str, duration, success) # Minimize mutex lock time - only update counters pending_count = nil total_pending = nil @mutex.synchronize do # Ring buffer: append, then evict oldest. Capping by refusing new # samples would freeze @metrics on the first METRICS_RING_CAP entries # whenever there is no Redis to drain it, so averages would report # boot-time behaviour forever. @metrics << Metric.new(feature_name_str, operation_str, duration, success: success) @metrics.shift while @metrics.length > METRICS_RING_CAP @usage_count[feature_name_str] += 1 @pending_updates[feature_name_str] += 1 pending_count = @pending_updates[feature_name_str] total_pending = @pending_updates.values.sum end # Rails 8+ event for usage tracking (cached check) if @_rails_events_enabled Magick::Rails::Events.usage_tracked(feature_name_str, operation: operation_str, duration: duration, success: success) end # Batch flush check - only if we're close to batch size flush_to_redis_if_needed if pending_count >= @batch_size || total_pending >= @batch_size end |
#record(feature_name, operation, duration, success: true) ⇒ Object
75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 |
# File 'lib/magick/performance_metrics.rb', line 75 def record(feature_name, operation, duration, success: true) # Fast path: push to async queue (non-blocking, zero overhead in hot path) # Queue#<< is thread-safe and lock-free - extremely fast! return unless @async_enabled # Push to async queue - this is lock-free and extremely fast # Use non-blocking push (will raise if queue is full, but our queue is unbounded) begin @async_queue << [feature_name.to_s, operation.to_s, duration, success] rescue ThreadError, ClosedQueueError # Queue is closed or thread died, disable async @async_enabled = false end nil end |
#start_async_processor ⇒ Object
Start background thread to process async metrics
121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 |
# File 'lib/magick/performance_metrics.rb', line 121 def start_async_processor return if @async_thread&.alive? @async_thread = Thread.new do last_flush_check = Time.now loop do # Wait for metrics with timeout to allow periodic flush checks # Queue#pop with timeout returns nil on timeout, raises on closed queue begin data = @async_queue.pop(timeout: 1.0) rescue ThreadError => e # Queue closed or thread interrupted break if e..include?('queue closed') raise end if data feature_name_str, operation_str, duration, success = data process_async_record(feature_name_str, operation_str, duration, success) last_flush_check = Time.now elsif Time.now - last_flush_check >= 1.0 # Timeout - check if we need to flush based on time (every second) flush_to_redis_if_needed last_flush_check = Time.now end rescue StandardError => e # Log error but continue processing warn "Magick: Error in async metrics processor: #{e.message}" if defined?(::Rails) && ::Rails.env.development? sleep 0.1 # Brief pause on error end end @async_thread.abort_on_exception = false end |
#stop_async_processor ⇒ Object
Stop async processor (for cleanup)
367 368 369 370 371 372 |
# File 'lib/magick/performance_metrics.rb', line 367 def stop_async_processor @async_enabled = false @async_queue.close if @async_queue.respond_to?(:close) @async_thread&.kill @async_thread = nil end |
#usage_count(feature_name) ⇒ Object
305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 |
# File 'lib/magick/performance_metrics.rb', line 305 def usage_count(feature_name) feature_name_str = feature_name.to_s # Force flush any pending updates for this feature before reading to ensure accuracy # This ensures stats are synced across processes immediately. # With no Redis configured the flush is a no-op and the counts stay in memory. force_flush_to_redis if @mutex.synchronize { @pending_updates[feature_name_str] }.positive? # Memory count = total counts in current process minus what we've already flushed # This avoids double-counting with Redis memory_count = @mutex.synchronize do (@usage_count[feature_name_str] || 0) - (@flushed_counts[feature_name_str] || 0) end memory_count = 0 if memory_count.negative? # Safety check # Redis count = total counts from all processes (including this process's flushed counts) redis_count = 0 begin redis = stats_redis_client redis_count = redis.get("magick:stats:#{feature_name_str}").to_i if redis rescue StandardError # Silently fail end # Total = Redis (all processes, all time) + memory (current process, not yet flushed) redis_count + memory_count end |