Class: CloseYourIt::Subscribers::JobPerformance

Inherits:
Object
  • Object
show all
Defined in:
lib/closeyourit/subscribers/job_performance.rb

Overview

Misura durata di esecuzione e attesa in coda (queue latency) dei background job — ActiveJob e Sidekiq — e, oltre le soglie configurate, emette metriche performance_issue (subtype slow_job e job_queue_latency). Logica PURA e SENZA STATO condiviso: tutto arriva per parametri, quindi job concorrenti sullo stesso thread/processo non si contaminano. Il wiring ad ActiveSupport::Notifications e al middleware Sidekiq vive altrove (Railtie / JobMetricsMiddleware). Rispetta il master switch monitor_jobs, le soglie e il jobs_sample_rate.

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(configuration = nil) ⇒ JobPerformance

Returns a new instance of JobPerformance.



15
16
17
# File 'lib/closeyourit/subscribers/job_performance.rb', line 15

def initialize(configuration = nil)
  @configuration = configuration
end

Class Method Details

.latency_ms(enqueued_at, now:) ⇒ Object

Normalizza enqueued_at (Time, epoch Numerico in secondi come Sidekiq, o String ISO8601) in attesa (ms) rispetto a now. nil o non parsabile → nil: nessuna metrica di attesa (adapter che non popola l'istante di enqueue). Clamp a 0 se negativa (clock skew, enqueue "nel futuro"): lo schema di ingest esige duration_ms >= 0.



69
70
71
72
73
74
75
# File 'lib/closeyourit/subscribers/job_performance.rb', line 69

def self.latency_ms(enqueued_at, now:)
  started = to_time(enqueued_at)
  return nil if started.nil?

  ms = (now - started) * 1000.0
  ms.negative? ? 0.0 : ms
end

.parse_time(value) ⇒ Object



85
86
87
88
89
# File 'lib/closeyourit/subscribers/job_performance.rb', line 85

def self.parse_time(value)
  Time.parse(value)
rescue ArgumentError
  nil
end

.to_time(value) ⇒ Object



77
78
79
80
81
82
83
# File 'lib/closeyourit/subscribers/job_performance.rb', line 77

def self.to_time(value)
  case value
  when Time    then value
  when Numeric then Time.at(value)
  when String  then parse_time(value)
  end
end

Instance Method Details

#active_job_performed(job, duration_ms) ⇒ Object

Hook perform.active_job: a fine esecuzione la durata è event.duration (ms).



54
55
56
57
58
59
60
61
62
63
# File 'lib/closeyourit/subscribers/job_performance.rb', line 54

def active_job_performed(job, duration_ms)
  record(
    job_class: job.class.name,
    queue: (job.queue_name if job.respond_to?(:queue_name)),
    adapter: "active_job",
    duration_ms: duration_ms,
    attempt: (job.executions if job.respond_to?(:executions)),
    trace_id: (job.job_id if job.respond_to?(:job_id))
  )
end

#active_job_started(job, now: Time.now.utc) ⇒ Object

Hook perform_start.active_job: l'attesa in coda è nota appena il job parte (now - enqueued_at).



42
43
44
45
46
47
48
49
50
51
# File 'lib/closeyourit/subscribers/job_performance.rb', line 42

def active_job_started(job, now: Time.now.utc)
  record(
    job_class: job.class.name,
    queue: (job.queue_name if job.respond_to?(:queue_name)),
    adapter: "active_job",
    queue_latency_ms: self.class.latency_ms(enqueued_at(job), now: now),
    attempt: (job.executions if job.respond_to?(:executions)),
    trace_id: (job.job_id if job.respond_to?(:job_id))
  )
end

#record(job_class:, queue: nil, adapter: nil, duration_ms: nil, queue_latency_ms: nil, attempt: nil, trace_id: nil) ⇒ Object

Punto unico di decisione: dai valori misurati (durata e/o attesa) costruisce 0..2 metriche, applica le soglie (stretto >, così X non genera rumore e X+1 sì) e il sampling, poi le spedisce fire-and-forget. duration_ms e queue_latency_ms sono opzionali: ActiveJob li fornisce da due hook distinti (perform_start → attesa, perform → durata), Sidekiq entrambi in una sola chiamata. Il sampling è applicato SOLO ai candidati già oltre soglia (i job normali non consumano né generano nulla).



25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
# File 'lib/closeyourit/subscribers/job_performance.rb', line 25

def record(job_class:, queue: nil, adapter: nil, duration_ms: nil, queue_latency_ms: nil,
           attempt: nil, trace_id: nil)
  config = configuration
  return unless config.monitor_jobs

  common = { job_class: job_class, queue: queue, adapter: adapter, attempt: attempt, trace_id: trace_id }
  events = []
  events << build(config, "slow_job", duration_ms, common) if slow?(config, duration_ms)
  events << build(config, "job_queue_latency", queue_latency_ms, common) if waited?(config, queue_latency_ms)
  return if events.empty?
  return unless sampled?(config)

  events.each { |event| CloseYourIt.capture_event(event) }
  nil
end