Class: Cosmo::API::Busy
- Inherits:
-
Object
- Object
- Cosmo::API::Busy
- Defined in:
- lib/cosmo/api/busy.rb,
sig/cosmo/api/busy.rbs
Constant Summary collapse
- TTL =
70- HEARTBEAT =
30- BUCKET =
"cosmo_jobs_busy"
Class Method Summary collapse
Instance Method Summary collapse
- #add(message) ⇒ void
- #delete(message) ⇒ void
- #heartbeat_loop ⇒ void
-
#initialize ⇒ Busy
constructor
A new instance of Busy.
- #list(limit: 25) ⇒ Array[Hash[Symbol, untyped]]
- #size ⇒ ::Integer
- #with(message) { ... } ⇒ void
- #worker_id ⇒ ::String
Constructor Details
Class Method Details
.instance ⇒ Busy
12 13 14 |
# File 'lib/cosmo/api/busy.rb', line 12 def self.instance @instance ||= new end |
Instance Method Details
#add(message) ⇒ void
This method returns an undefined value.
28 29 30 31 32 33 34 |
# File 'lib/cosmo/api/busy.rb', line 28 def add() @thread ||= Thread.new { heartbeat_loop } seq = ..sequence.stream value = Utils::Json.dump({ data: .data, stream: ..stream, worker: worker_id, started_at: Time.now.to_i }) @messages[seq] = value @kv.set(seq, value) end |
#delete(message) ⇒ void
This method returns an undefined value.
36 37 38 39 40 |
# File 'lib/cosmo/api/busy.rb', line 36 def delete() seq = ..sequence.stream @messages.delete(seq) @kv.purge(seq) end |
#heartbeat_loop ⇒ void
This method returns an undefined value.
52 53 54 55 56 57 58 59 |
# File 'lib/cosmo/api/busy.rb', line 52 def heartbeat_loop loop do sleep(HEARTBEAT) @messages.dup.each { |seq, value| @kv.set(seq, value) rescue StandardError } rescue StandardError => e Logger.debug "Busy heartbeat error: #{e.class} #{e.}" end end |
#list(limit: 25) ⇒ Array[Hash[Symbol, untyped]]
42 43 44 |
# File 'lib/cosmo/api/busy.rb', line 42 def list(limit: 25) @kv.keys(limit:).filter_map { Utils::Json.parse(@kv.get(_1)&.value) }.map { _1.merge(data: Utils::Json.parse(_1[:data])) } end |
#size ⇒ ::Integer
46 47 48 |
# File 'lib/cosmo/api/busy.rb', line 46 def size @kv.size end |
#with(message) { ... } ⇒ void
This method returns an undefined value.
21 22 23 24 25 26 |
# File 'lib/cosmo/api/busy.rb', line 21 def with() add() yield ensure delete() end |
#worker_id ⇒ ::String
61 62 63 |
# File 'lib/cosmo/api/busy.rb', line 61 def worker_id @worker_id ||= "#{Socket.gethostname}-#{Process.pid}" end |