Class: JobWorkflow::Monitoring::ExecutionViewModel
- Inherits:
-
Object
- Object
- JobWorkflow::Monitoring::ExecutionViewModel
- Defined in:
- lib/job_workflow/monitoring/execution_view_model.rb,
sig/generated/job_workflow/monitoring/execution_view_model.rbs
Instance Attribute Summary collapse
- #job_id ⇒ String readonly
- #queue_name ⇒ String? readonly
- #status ⇒ WorkflowStatus readonly
Instance Method Summary collapse
-
#arguments ⇒ Arguments
: () -> Arguments.
-
#callable_summary(value) ⇒ Object
: (untyped) -> untyped.
- #completed_task?(task_name, task_outputs, task_job_statuses) ⇒ Boolean
-
#current_task_name ⇒ Symbol?
: () -> Symbol?.
-
#current_task_running?(task_name) ⇒ Boolean
: (Symbol) -> bool.
-
#dag_layout ⇒ Hash[Symbol, untyped]
: () -> Hash[Symbol, untyped].
-
#dependency_wait_configuration(task) ⇒ Hash[Symbol, untyped]
: (Task) -> Hash[Symbol, untyped].
- #each_progress(task_outputs, task_job_statuses) ⇒ Hash[Symbol, Integer]
-
#enqueue_configuration(task) ⇒ Hash[Symbol, untyped]
: (Task) -> Hash[Symbol, untyped].
-
#failed_task_name ⇒ Symbol?
: () -> Symbol?.
-
#filtered_arguments ⇒ Hash[untyped, untyped]
: () -> Hash[untyped, untyped].
-
#initialize(job_id:, queue_name:, status:) ⇒ ExecutionViewModel
constructor
: (job_id: String, queue_name: String?, status: WorkflowStatus) -> void.
-
#job_class_name ⇒ String
: () -> String.
-
#mission_control_job_path ⇒ String?
: () -> String?.
-
#mission_control_job_path_for(job_id, status = nil) ⇒ String?
: (String?, Symbol?) -> String?.
-
#output_configuration(task) ⇒ Array[Hash[Symbol, untyped]]
: (Task) -> Array[Hash[Symbol, untyped]].
-
#primitive_summary(value) ⇒ Object
: (untyped) -> untyped.
-
#retry_configuration(task) ⇒ Hash[Symbol, untyped]
: (Task) -> Hash[Symbol, untyped].
-
#running? ⇒ Boolean
: () -> bool.
-
#sub_task_jobs_view(task_job_statuses) ⇒ Array[Hash[Symbol, untyped]]
: (Array) -> Array[Hash[Symbol, untyped]].
-
#task_configuration(task) ⇒ Hash[Symbol, untyped]
: (Task) -> Hash[Symbol, untyped].
-
#task_configuration_view(task) ⇒ Hash[Symbol, untyped]
: (Task) -> Hash[Symbol, untyped].
-
#task_outputs_view(task_outputs) ⇒ Array[Hash[Symbol, untyped]]
: (Array) -> Array[Hash[Symbol, untyped]].
-
#task_running?(task_name, task_job_statuses) ⇒ Boolean
: (Symbol, Array) -> bool.
- #task_runtime_view(task_name, task_outputs, task_job_statuses) ⇒ Hash[Symbol, untyped]
-
#task_state(task_name) ⇒ [ Array[TaskOutput], Array[TaskJobStatus] ]
: (Symbol) -> [Array[TaskOutput], Array[TaskJobStatus]].
- #task_status(task_name, task_outputs, task_job_statuses) ⇒ Symbol
-
#task_view_model(task) ⇒ Hash[Symbol, untyped]
: (Task) -> Hash[Symbol, untyped].
-
#tasks ⇒ Array[Hash[Symbol, untyped]]
: () -> Array[Hash[Symbol, untyped]].
-
#throttle_configuration(task) ⇒ Hash[Symbol, untyped]
: (Task) -> Hash[Symbol, untyped].
-
#to_h ⇒ Hash[Symbol, untyped]
: () -> Hash[Symbol, untyped].
-
#workflow_status ⇒ Symbol
: () -> Symbol.
Constructor Details
#initialize(job_id:, queue_name:, status:) ⇒ ExecutionViewModel
: (job_id: String, queue_name: String?, status: WorkflowStatus) -> void
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_id ⇒ String (readonly)
10 11 12 |
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 10 def job_id @job_id end |
#queue_name ⇒ String? (readonly)
11 12 13 |
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 11 def queue_name @queue_name end |
#status ⇒ WorkflowStatus (readonly)
12 13 14 |
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 12 def status @status end |
Instance Method Details
#arguments ⇒ Arguments
: () -> Arguments
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
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
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_name ⇒ Symbol?
: () -> Symbol?
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
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_layout ⇒ Hash[Symbol, untyped]
: () -> 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]
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]
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]
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_name ⇒ Symbol?
: () -> Symbol?
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_arguments ⇒ Hash[untyped, untyped]
: () -> 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_name ⇒ String
: () -> 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_path ⇒ String?
: () -> String?
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?
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]]
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
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]
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
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]]
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]
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]
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]]
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
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]
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]]
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
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]
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 |
#tasks ⇒ Array[Hash[Symbol, untyped]]
: () -> 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]
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_h ⇒ Hash[Symbol, untyped]
: () -> 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_status ⇒ Symbol
: () -> Symbol
27 28 29 |
# File 'lib/job_workflow/monitoring/execution_view_model.rb', line 27 def workflow_status status.status end |