Class: Broadlistening::Status

Inherits:
Object
  • Object
show all
Defined in:
lib/broadlistening/status.rb

Constant Summary collapse

LOCK_DURATION =

5分

300

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(output_dir) ⇒ Status

Returns a new instance of Status.



15
16
17
18
19
# File 'lib/broadlistening/status.rb', line 15

def initialize(output_dir)
  @output_dir = Pathname.new(output_dir)
  @status_file = @output_dir / "status.json"
  @data = load_or_initialize
end

Instance Attribute Details

#dataObject (readonly)

Returns the value of attribute data.



13
14
15
# File 'lib/broadlistening/status.rb', line 13

def data
  @data
end

#output_dirObject (readonly)

Returns the value of attribute output_dir.



13
14
15
# File 'lib/broadlistening/status.rb', line 13

def output_dir
  @output_dir
end

#status_fileObject (readonly)

Returns the value of attribute status_file.



13
14
15
# File 'lib/broadlistening/status.rb', line 13

def status_file
  @status_file
end

Instance Method Details

#complete_pipelineObject



70
71
72
73
74
75
76
# File 'lib/broadlistening/status.rb', line 70

def complete_pipeline
  merge_previous_jobs
  @data[:status] = "completed"
  @data[:end_time] = Time.now.iso8601
  @data.delete(:previous)
  save
end

#complete_step(step_name, params:, duration:, token_usage: 0) ⇒ Object



56
57
58
59
60
61
62
63
64
65
66
67
68
# File 'lib/broadlistening/status.rb', line 56

def complete_step(step_name, params:, duration:, token_usage: 0)
  @data[:completed_jobs] ||= []
  @data[:completed_jobs] << {
    step: step_name.to_s,
    completed: Time.now.iso8601,
    duration: duration,
    params: serialize_params(params),
    token_usage: token_usage
  }
  @data.delete(:current_job)
  @data.delete(:current_job_started)
  save
end

#error_pipeline(error) ⇒ Object



78
79
80
81
82
83
84
# File 'lib/broadlistening/status.rb', line 78

def error_pipeline(error)
  @data[:status] = "error"
  @data[:end_time] = Time.now.iso8601
  @data[:error] = "#{error.class}: #{error.message}"
  @data[:error_stack_trace] = error.backtrace&.join("\n")
  save
end

#load_or_initializeObject



21
22
23
24
25
26
27
28
29
30
31
# File 'lib/broadlistening/status.rb', line 21

def load_or_initialize
  if status_file.exist?
    JSON.parse(status_file.read, symbolize_names: true)
  else
    {
      status: "initialized",
      completed_jobs: [],
      previously_completed_jobs: []
    }
  end
end

#locked?Boolean

Returns:

  • (Boolean)


86
87
88
89
90
91
# File 'lib/broadlistening/status.rb', line 86

def locked?
  return false unless @data[:status] == "running"
  return false unless @data[:lock_until]

  Time.parse(@data[:lock_until]) > Time.now
end

#previous_completed_jobsObject



93
94
95
# File 'lib/broadlistening/status.rb', line 93

def previous_completed_jobs
  (@data[:completed_jobs] || []) + (@data[:previously_completed_jobs] || [])
end

#saveObject



33
34
35
36
# File 'lib/broadlistening/status.rb', line 33

def save
  FileUtils.mkdir_p(output_dir)
  status_file.write(JSON.pretty_generate(@data))
end

#start_pipeline(plan) ⇒ Object



38
39
40
41
42
43
44
45
46
47
# File 'lib/broadlistening/status.rb', line 38

def start_pipeline(plan)
  @data.merge!(
    status: "running",
    plan: plan.map { |p| serialize_plan_entry(p) },
    start_time: Time.now.iso8601,
    completed_jobs: [],
    lock_until: lock_time.iso8601
  )
  save
end

#start_step(step_name) ⇒ Object



49
50
51
52
53
54
# File 'lib/broadlistening/status.rb', line 49

def start_step(step_name)
  @data[:current_job] = step_name.to_s
  @data[:current_job_started] = Time.now.iso8601
  @data[:lock_until] = lock_time.iso8601
  save
end