Class: Kettle::Family::ReleaseStateEventTape

Inherits:
Object
  • Object
show all
Defined in:
lib/kettle/family/release_state_event_tape.rb

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(root:, stream: nil, clock: nil, wall_clock: nil) ⇒ ReleaseStateEventTape

Returns a new instance of ReleaseStateEventTape.



12
13
14
15
16
17
18
19
20
21
# File 'lib/kettle/family/release_state_event_tape.rb', line 12

def initialize(root:, stream: nil, clock: nil, wall_clock: nil)
  @directory = File.join(root, "tmp", "kettle-family", "release-state-#{Time.now.strftime("%Y%m%d-%H%M%S")}-#{Process.pid}")
  FileUtils.mkdir_p(@directory)
  @stream = stream
  @clock = clock || lambda { Process.clock_gettime(Process::CLOCK_MONOTONIC) }
  @wall_clock = wall_clock || lambda { Time.now.utc.iso8601 }
  @started_at = @clock.call
  @sequence_by_member = Hash.new(0)
  @mutex = Mutex.new
end

Instance Attribute Details

#directoryObject (readonly)

Returns the value of attribute directory.



10
11
12
# File 'lib/kettle/family/release_state_event_tape.rb', line 10

def directory
  @directory
end

Instance Method Details

#call(event) ⇒ Object



23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
# File 'lib/kettle/family/release_state_event_tape.rb', line 23

def call(event)
  payload = event.merge(
    "event_version" => 1,
    "type" => "release_state",
    "timestamp" => @wall_clock.call,
    "elapsed_seconds" => @clock.call - @started_at
  )
  member = payload.fetch("member", "family").to_s
  @mutex.synchronize do
    @sequence_by_member[member] += 1
    payload["sequence"] = @sequence_by_member.fetch(member)
    line = JSON.generate(payload)
    File.open(path_for(member), "a") { |file| file.puts(line) }
    @stream&.puts(line)
    @stream.flush if @stream&.respond_to?(:flush)
  end
  payload
rescue
  nil
end