Class: JobWorkflow::Monitoring::ExecutionViewModel

Inherits:
Object
  • Object
show all
Defined in:
lib/job_workflow/monitoring/execution_view_model.rb,
sig/generated/job_workflow/monitoring/execution_view_model.rbs

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(job_id:, queue_name:, status:) ⇒ ExecutionViewModel

: (job_id: String, queue_name: String?, status: WorkflowStatus) -> void

Parameters:



15
16
17
18
19
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 15

def initialize(job_id:, queue_name:, status:)
  @job_id = job_id
  @queue_name = queue_name
  @status = status
end

Instance Attribute Details

#job_idString (readonly)

RBS:

  • @tasks: Array[Hash[Symbol, untyped]]

  • @failed_task_name: Symbol?

  • @dag_layout: Hash[Symbol, untyped]

Returns:

  • (String)


10
11
12
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 10

def job_id
  @job_id
end

#queue_nameString? (readonly)

Signature:

  • String?

Returns:

  • (String, nil)


11
12
13
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 11

def queue_name
  @queue_name
end

#statusWorkflowStatus (readonly)

Signature:

  • WorkflowStatus

Returns:



12
13
14
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 12

def status
  @status
end

Instance Method Details

#argumentsArguments

: () -> Arguments

Returns:



37
38
39
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 37

def arguments
  status.arguments
end

#callable_summary(value) ⇒ Object

: (untyped) -> untyped

Parameters:

  • (Object)

Returns:

  • (Object)


210
211
212
213
214
215
216
217
218
219
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 210

def callable_summary(value)
  case value
  when nil
    nil
  when Proc
    "proc"
  else
    value
  end
end

#completed_task?(task_name, task_outputs, task_job_statuses) ⇒ Boolean

: (Symbol, Array, Array) -> bool

Parameters:

Returns:

  • (Boolean)


101
102
103
104
105
106
107
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 101

def completed_task?(task_name, task_outputs, task_job_statuses)
  return true if !running? && task_outputs.any?
  return task_job_statuses.all?(&:succeeded?) if task_job_statuses.any?
  return true if status.completed_task_names.include?(task_name)

  task_outputs.any?
end

#current_task_nameSymbol?

: () -> Symbol?

Returns:

  • (Symbol, nil)


32
33
34
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 32

def current_task_name
  status.current_task_name
end

#current_task_running?(task_name) ⇒ Boolean

: (Symbol) -> bool

Parameters:

  • (Symbol)

Returns:

  • (Boolean)


116
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 116

def current_task_running?(task_name) = running? && current_task_name == task_name

#dag_layoutHash[Symbol, untyped]

: () -> Hash[Symbol, untyped]

Returns:

  • (Hash[Symbol, untyped])


65
66
67
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 65

def dag_layout
  @dag_layout ||= DagLayout.new(tasks:).to_h
end

#dependency_wait_configuration(task) ⇒ Hash[Symbol, untyped]

: (Task) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


200
201
202
203
204
205
206
207
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 200

def dependency_wait_configuration(task)
  {
    poll_timeout: task.dependency_wait.poll_timeout,
    poll_interval: task.dependency_wait.poll_interval,
    reschedule_delay: task.dependency_wait.reschedule_delay,
    polling_only: task.dependency_wait.polling_only?
  }
end

#each_progress(task_outputs, task_job_statuses) ⇒ Hash[Symbol, Integer]

: (Array, Array) -> Hash[Symbol, Integer]

Parameters:

Returns:

  • (Hash[Symbol, Integer])


227
228
229
230
231
232
233
234
235
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 227

def each_progress(task_outputs, task_job_statuses)
  {
    total: [task_outputs.size, task_job_statuses.size].max,
    succeeded: task_job_statuses.count { |task_job_status| task_job_status.succeeded? }, # rubocop:disable Style/SymbolProc
    failed: task_job_statuses.count { |task_job_status| task_job_status.failed? }, # rubocop:disable Style/SymbolProc
    pending: task_job_statuses.count { |task_status| task_status.status == :pending },
    running: task_job_statuses.count { |task_status| task_status.status == :running }
  }
end

#enqueue_configuration(task) ⇒ Hash[Symbol, untyped]

: (Task) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


168
169
170
171
172
173
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 168

def enqueue_configuration(task)
  {
    enabled: primitive_summary(task.enqueue.condition),
    queue: task.enqueue.queue
  }
end

#failed_task_nameSymbol?

: () -> Symbol?

Returns:

  • (Symbol, nil)


52
53
54
55
56
57
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 52

def failed_task_name
  @failed_task_name ||= begin
    failed_task = tasks.find { |task| task[:status] == :failed }
    failed_task&.fetch(:name)
  end
end

#filtered_argumentsHash[untyped, untyped]

: () -> Hash[untyped, untyped]

Returns:

  • (Hash[untyped, untyped])


42
43
44
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 42

def filtered_arguments
  ParameterFilter.filter(arguments.to_h)
end

#job_class_nameString

: () -> String

Returns:

  • (String)


22
23
24
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 22

def job_class_name
  status.job_class_name
end

#mission_control_job_pathString?

: () -> String?

Returns:

  • (String, nil)


60
61
62
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 60

def mission_control_job_path
  JobWorkflow::Monitoring.mission_control_job_path(job_id, status: workflow_status)
end

#mission_control_job_path_for(job_id, status = nil) ⇒ String?

: (String?, Symbol?) -> String?

Parameters:

  • (String, nil)
  • (Symbol, nil)

Returns:

  • (String, nil)


257
258
259
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 257

def mission_control_job_path_for(job_id, status = nil)
  JobWorkflow::Monitoring.mission_control_job_path(job_id, status:)
end

#output_configuration(task) ⇒ Array[Hash[Symbol, untyped]]

: (Task) -> Array[Hash[Symbol, untyped]]

Parameters:

Returns:

  • (Array[Hash[Symbol, untyped]])


176
177
178
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 176

def output_configuration(task)
  task.output.map { |output| { name: output.name, type: output.type } }
end

#primitive_summary(value) ⇒ Object

: (untyped) -> untyped

Parameters:

  • (Object)

Returns:

  • (Object)


222
223
224
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 222

def primitive_summary(value)
  value.is_a?(Proc) ? "proc" : value
end

#retry_configuration(task) ⇒ Hash[Symbol, untyped]

: (Task) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


181
182
183
184
185
186
187
188
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 181

def retry_configuration(task)
  {
    count: task.task_retry.count,
    strategy: task.task_retry.strategy,
    base_delay: task.task_retry.base_delay,
    jitter: task.task_retry.jitter
  }
end

#running?Boolean

: () -> bool

Returns:

  • (Boolean)


70
71
72
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 70

def running?
  workflow_status == :running
end

#sub_task_jobs_view(task_job_statuses) ⇒ Array[Hash[Symbol, untyped]]

: (Array) -> Array[Hash[Symbol, untyped]]

Parameters:

Returns:

  • (Array[Hash[Symbol, untyped]])


245
246
247
248
249
250
251
252
253
254
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 245

def sub_task_jobs_view(task_job_statuses)
  Array(task_job_statuses).map do |task_job_status|
    {
      job_id: task_job_status.job_id,
      each_index: task_job_status.each_index,
      status: task_job_status.status,
      mission_control_job_path: mission_control_job_path_for(task_job_status.job_id, task_job_status.status)
    }
  end
end

#task_configuration(task) ⇒ Hash[Symbol, untyped]

: (Task) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


152
153
154
155
156
157
158
159
160
161
162
163
164
165
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 152

def task_configuration(task)
  {
    job_name: task.job_name,
    each: callable_summary(task.each),
    condition: callable_summary(task.condition),
    enqueue: enqueue_configuration(task),
    outputs: output_configuration(task),
    retry: retry_configuration(task),
    throttle: throttle_configuration(task),
    timeout: task.timeout,
    dependency_wait: dependency_wait_configuration(task),
    dry_run: callable_summary(task.dry_run_config.evaluator)
  }
end

#task_configuration_view(task) ⇒ Hash[Symbol, untyped]

: (Task) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


129
130
131
132
133
134
135
136
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 129

def task_configuration_view(task)
  {
    name: task.task_name,
    depends_on: task.depends_on,
    each: task.each?,
    configuration: task_configuration(task)
  }
end

#task_outputs_view(task_outputs) ⇒ Array[Hash[Symbol, untyped]]

: (Array) -> Array[Hash[Symbol, untyped]]

Parameters:

Returns:

  • (Array[Hash[Symbol, untyped]])


238
239
240
241
242
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 238

def task_outputs_view(task_outputs)
  task_outputs.map do |output|
    { each_index: output.each_index, data: ParameterFilter.filter(output.data) }
  end
end

#task_running?(task_name, task_job_statuses) ⇒ Boolean

: (Symbol, Array) -> bool

Parameters:

Returns:

  • (Boolean)


110
111
112
113
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 110

def task_running?(task_name, task_job_statuses)
  current_task_running?(task_name) ||
    (!task_job_statuses.empty? && task_job_statuses.any? { |task_status| !task_status.finished? })
end

#task_runtime_view(task_name, task_outputs, task_job_statuses) ⇒ Hash[Symbol, untyped]

: (Symbol, Array, Array) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


139
140
141
142
143
144
145
146
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 139

def task_runtime_view(task_name, task_outputs, task_job_statuses)
  {
    status: task_status(task_name, task_outputs, task_job_statuses),
    each_progress: each_progress(task_outputs, task_job_statuses),
    outputs: task_outputs_view(task_outputs),
    sub_task_jobs: sub_task_jobs_view(task_job_statuses)
  }
end

#task_state(task_name) ⇒ [ Array[TaskOutput], Array[TaskJobStatus] ]

: (Symbol) -> [Array[TaskOutput], Array[TaskJobStatus]]

Parameters:

  • (Symbol)

Returns:



149
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 149

def task_state(task_name) = [status.output.fetch_all(task_name:), status.job_status.fetch_all(task_name:)]

#task_status(task_name, task_outputs, task_job_statuses) ⇒ Symbol

: (Symbol, Array, Array) -> Symbol

Parameters:

Returns:

  • (Symbol)


92
93
94
95
96
97
98
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 92

def task_status(task_name, task_outputs, task_job_statuses)
  return :failed if task_job_statuses.any?(&:failed?)
  return :succeeded if completed_task?(task_name, task_outputs, task_job_statuses)
  return :running if task_running?(task_name, task_job_statuses)

  :pending
end

#task_view_model(task) ⇒ Hash[Symbol, untyped]

: (Task) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


119
120
121
122
123
124
125
126
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 119

def task_view_model(task)
  task_name = task.task_name
  task_outputs, task_job_statuses = task_state(task_name)

  task_configuration_view(task).merge(
    task_runtime_view(task_name, task_outputs, task_job_statuses)
  )
end

#tasksArray[Hash[Symbol, untyped]]

: () -> Array[Hash[Symbol, untyped]]

Returns:

  • (Array[Hash[Symbol, untyped]])


47
48
49
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 47

def tasks
  @tasks ||= status.context.workflow.tasks.map { |task| task_view_model(task) }
end

#throttle_configuration(task) ⇒ Hash[Symbol, untyped]

: (Task) -> Hash[Symbol, untyped]

Parameters:

Returns:

  • (Hash[Symbol, untyped])


191
192
193
194
195
196
197
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 191

def throttle_configuration(task)
  {
    key: task.throttle.key,
    limit: task.throttle.limit,
    ttl: task.throttle.ttl
  }
end

#to_hHash[Symbol, untyped]

: () -> Hash[Symbol, untyped]

Returns:

  • (Hash[Symbol, untyped])


75
76
77
78
79
80
81
82
83
84
85
86
87
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 75

def to_h
  {
    job_id:,
    queue_name:,
    job_class_name:,
    status: workflow_status,
    current_task_name:,
    failed_task_name:,
    arguments: filtered_arguments,
    tasks:,
    mission_control_job_path:
  }
end

#workflow_statusSymbol

: () -> Symbol

Returns:

  • (Symbol)


27
28
29
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 27

def workflow_status
  status.status
end