Class: DecisionAgent::Monitoring::MetricsCollector
- Inherits:
-
Object
- Object
- DecisionAgent::Monitoring::MetricsCollector
- Includes:
- MonitorMixin
- Defined in:
- lib/decision_agent/monitoring/metrics_collector.rb
Overview
Thread-safe metrics collector for decision analytics
Instance Attribute Summary collapse
-
#metrics ⇒ Object
readonly
Returns the value of attribute metrics.
-
#storage_adapter ⇒ Object
readonly
Returns the value of attribute storage_adapter.
-
#window_size ⇒ Object
readonly
Returns the value of attribute window_size.
Instance Method Summary collapse
-
#add_observer(&block) ⇒ Object
Register observer for real-time updates.
-
#cleanup_old_metrics_from_storage(older_than:) ⇒ Object
Cleanup old metrics from persistent storage.
-
#clear! ⇒ Object
Clear all metrics.
-
#initialize(window_size: 3600, storage: :auto, cleanup_threshold: 100) ⇒ MetricsCollector
constructor
A new instance of MetricsCollector.
-
#metrics_count ⇒ Object
Get current metrics count.
-
#record_decision(decision, context, duration_ms: nil) ⇒ Object
Record a decision for analytics.
-
#record_error(error, context: {}) ⇒ Object
Record error.
-
#record_evaluation(evaluation) ⇒ Object
Record individual evaluation metrics.
-
#record_performance(operation:, duration_ms:, success: true, metadata: {}) ⇒ Object
Record performance metrics.
-
#statistics(time_range: nil) ⇒ Object
Get aggregated statistics.
-
#time_series(metric_type:, bucket_size: 60, time_range: 3600) ⇒ Object
Get time-series data for graphing.
Constructor Details
#initialize(window_size: 3600, storage: :auto, cleanup_threshold: 100) ⇒ MetricsCollector
Returns a new instance of MetricsCollector.
21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 21 def initialize(window_size: 3600, storage: :auto, cleanup_threshold: 100) super() @window_size = window_size # Default: 1 hour window @cleanup_threshold = cleanup_threshold # Cleanup every N records @cleanup_counter = 0 @storage_adapter = initialize_storage_adapter(storage, window_size) # Legacy in-memory metrics for backward compatibility with observers @metrics = { decisions: [], evaluations: [], performance: [], errors: [] } @observers = [] freeze_config end |
Instance Attribute Details
#metrics ⇒ Object (readonly)
Returns the value of attribute metrics.
19 20 21 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 19 def metrics @metrics end |
#storage_adapter ⇒ Object (readonly)
Returns the value of attribute storage_adapter.
19 20 21 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 19 def storage_adapter @storage_adapter end |
#window_size ⇒ Object (readonly)
Returns the value of attribute window_size.
19 20 21 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 19 def window_size @window_size end |
Instance Method Details
#add_observer(&block) ⇒ Object
Register observer for real-time updates
230 231 232 233 234 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 230 def add_observer(&block) synchronize do @observers << block end end |
#cleanup_old_metrics_from_storage(older_than:) ⇒ Object
Cleanup old metrics from persistent storage
264 265 266 267 268 269 270 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 264 def cleanup_old_metrics_from_storage(older_than:) synchronize do return 0 unless @storage_adapter.respond_to?(:cleanup) @storage_adapter.cleanup(older_than: older_than) end end |
#clear! ⇒ Object
Clear all metrics
237 238 239 240 241 242 243 244 245 246 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 237 def clear! synchronize do @metrics.each_value(&:clear) # Also clear storage adapter if using MemoryAdapter if @storage_adapter.is_a?(Storage::MemoryAdapter) # Clear all by using a very large time period (100 years in seconds) @storage_adapter.cleanup(older_than: 100 * 365 * 24 * 60 * 60) end end end |
#metrics_count ⇒ Object
Get current metrics count
249 250 251 252 253 254 255 256 257 258 259 260 261 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 249 def metrics_count synchronize do # Use in-memory metrics for MemoryAdapter (to maintain backward compatibility) # Only delegate to ActiveRecordAdapter for persistent storage use_storage = @storage_adapter.respond_to?(:metrics_count) && !@storage_adapter.is_a?(Storage::MemoryAdapter) return @storage_adapter.metrics_count if use_storage # Use in-memory @metrics.transform_values(&:size) end end |
#record_decision(decision, context, duration_ms: nil) ⇒ Object
Record a decision for analytics
40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 40 def record_decision(decision, context, duration_ms: nil) synchronize do metric = { timestamp: Time.now.utc, decision: decision.decision, confidence: decision.confidence, evaluations_count: decision.evaluations.size, context_size: context.to_h.size, duration_ms: duration_ms, evaluator_names: decision.evaluations.map(&:evaluator_name).uniq } # Store in-memory for observers (backward compatibility) @metrics[:decisions] << metric maybe_cleanup_old_metrics! # Persist to storage adapter @storage_adapter.record_decision( decision.decision, context.to_h, confidence: decision.confidence, evaluations_count: decision.evaluations.size, duration_ms: duration_ms, status: determine_decision_status(decision) ) notify_observers(:decision, metric) metric end end |
#record_error(error, context: {}) ⇒ Object
Record error
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 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 127 def record_error(error, context: {}) synchronize do metric = { timestamp: Time.now.utc, error_class: error.class.name, error_message: error., context: context } # Store in-memory for observers (backward compatibility) @metrics[:errors] << metric maybe_cleanup_old_metrics! # Persist to storage adapter @storage_adapter.record_error( error.class.name, message: error., stack_trace: error.backtrace, severity: determine_error_severity(error), context: context ) notify_observers(:error, metric) metric end end |
#record_evaluation(evaluation) ⇒ Object
Record individual evaluation metrics
72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 72 def record_evaluation(evaluation) synchronize do metric = { timestamp: Time.now.utc, decision: evaluation.decision, weight: evaluation.weight, evaluator_name: evaluation.evaluator_name } # Store in-memory for observers (backward compatibility) @metrics[:evaluations] << metric maybe_cleanup_old_metrics! # Persist to storage adapter @storage_adapter.record_evaluation( evaluation.evaluator_name, score: evaluation.weight, success: evaluation.weight.positive?, details: { decision: evaluation.decision } ) notify_observers(:evaluation, metric) metric end end |
#record_performance(operation:, duration_ms:, success: true, metadata: {}) ⇒ Object
Record performance metrics
99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 99 def record_performance(operation:, duration_ms:, success: true, metadata: {}) synchronize do metric = { timestamp: Time.now.utc, operation: operation, duration_ms: duration_ms, success: success, metadata: } # Store in-memory for observers (backward compatibility) @metrics[:performance] << metric maybe_cleanup_old_metrics! # Persist to storage adapter @storage_adapter.record_performance( operation, duration_ms: duration_ms, status: success ? "success" : "failure", metadata: ) notify_observers(:performance, metric) metric end end |
#statistics(time_range: nil) ⇒ Object
Get aggregated statistics
155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 155 def statistics(time_range: nil) synchronize do # Use in-memory metrics for MemoryAdapter (to maintain backward compatibility) # Only delegate to ActiveRecordAdapter for persistent storage use_storage = time_range && @storage_adapter.respond_to?(:statistics) && !@storage_adapter.is_a?(Storage::MemoryAdapter) if use_storage stats = @storage_adapter.statistics(time_range: time_range) return stats.merge(timestamp: Time.now.utc, storage: @storage_adapter.class.name) if stats end # Use in-memory metrics range_start = time_range ? Time.now.utc - time_range : nil decisions = filter_by_time(@metrics[:decisions], range_start) evaluations = filter_by_time(@metrics[:evaluations], range_start) performance = filter_by_time(@metrics[:performance], range_start) errors = filter_by_time(@metrics[:errors], range_start) { summary: { total_decisions: decisions.size, total_evaluations: evaluations.size, total_errors: errors.size, time_range: range_start ? "Last #{time_range}s" : "All time" }, decisions: compute_decision_stats(decisions), evaluations: compute_evaluation_stats(evaluations), performance: compute_performance_stats(performance), errors: compute_error_stats(errors), timestamp: Time.now.utc, storage: "memory (fallback)" } end end |
#time_series(metric_type:, bucket_size: 60, time_range: 3600) ⇒ Object
Get time-series data for graphing
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 224 225 226 227 |
# File 'lib/decision_agent/monitoring/metrics_collector.rb', line 194 def time_series(metric_type:, bucket_size: 60, time_range: 3600) synchronize do # Use in-memory metrics for MemoryAdapter (to maintain backward compatibility) # Only delegate to ActiveRecordAdapter for persistent storage use_storage = @storage_adapter.respond_to?(:time_series) && !@storage_adapter.is_a?(Storage::MemoryAdapter) if use_storage series = @storage_adapter.time_series(metric_type, bucket_size: bucket_size, time_range: time_range) return series if series && series[:timestamps] end # Use in-memory metrics data = @metrics[metric_type] || [] range_start = Time.now.utc - time_range buckets = {} data.each do |metric| next if metric[:timestamp] < range_start bucket_key = (metric[:timestamp].to_i / bucket_size) * bucket_size buckets[bucket_key] ||= [] buckets[bucket_key] << metric end buckets.sort.map do |, metrics| { timestamp: Time.at().utc, count: metrics.size, metrics: metrics } end end end |