Class: Cosmo::API::Busy

Inherits:
Object
  • Object
show all
Defined in:
lib/cosmo/api/busy.rb,
sig/cosmo/api/busy.rbs

Constant Summary collapse

TTL =

Returns:

  • (::Integer)
70
HEARTBEAT =

Returns:

  • (::Integer)
30
BUCKET =

Returns:

  • (::String)
"cosmo_jobs_busy"

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initializeBusy

Returns a new instance of Busy.



16
17
18
19
# File 'lib/cosmo/api/busy.rb', line 16

def initialize
  @messages = {}
  @kv = KV.new(BUCKET, { ttl: TTL })
end

Class Method Details

.instanceBusy

Returns:



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.

Parameters:

  • message (Object)


28
29
30
31
32
33
34
# File 'lib/cosmo/api/busy.rb', line 28

def add(message)
  @thread ||= Thread.new { heartbeat_loop }
  seq = message..sequence.stream
  value = Utils::Json.dump({ data: message.data, stream: message..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.

Parameters:

  • message (Object)


36
37
38
39
40
# File 'lib/cosmo/api/busy.rb', line 36

def delete(message)
  seq = message..sequence.stream
  @messages.delete(seq)
  @kv.purge(seq)
end

#heartbeat_loopvoid

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.message}"
  end
end

#list(limit: 25) ⇒ Array[Hash[Symbol, untyped]]

Parameters:

  • limit: (::Integer) (defaults to: 25)

Returns:

  • (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

Returns:

  • (::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.

Parameters:

  • message (Object)

Yields:

Yield Returns:

  • (void)


21
22
23
24
25
26
# File 'lib/cosmo/api/busy.rb', line 21

def with(message)
  add(message)
  yield
ensure
  delete(message)
end

#worker_id::String

Returns:

  • (::String)


61
62
63
# File 'lib/cosmo/api/busy.rb', line 61

def worker_id
  @worker_id ||= "#{Socket.gethostname}-#{Process.pid}"
end