Class: DecisionAgent::Monitoring::MetricsCollector

Inherits:
Object
  • Object
show all
Includes:
MonitorMixin
Defined in:
lib/decision_agent/monitoring/metrics_collector.rb

Overview

Thread-safe metrics collector for decision analytics

Instance Attribute Summary collapse

Instance Method Summary collapse

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

#metricsObject (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_adapterObject (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_sizeObject (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_countObject

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.message,
      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.message,
      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 |timestamp, metrics|
      {
        timestamp: Time.at(timestamp).utc,
        count: metrics.size,
        metrics: metrics
      }
    end
  end
end