Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
304 changes: 241 additions & 63 deletions .rubocop_todo.yml

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion Gemfile
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ group :postgresql, optional: ENV.key?('CI') && ENV['DB'] != 'postgresql' do
end

group :lint do
gem 'theforeman-rubocop', '~> 0.0.4'
gem 'theforeman-rubocop', '~> 0.1.0'
end

group :memory_watcher do
Expand Down
2 changes: 1 addition & 1 deletion lib/dynflow/action/rescue.rb
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ module Action::Rescue
end

SuggestedStrategy = Algebrick.type do
fields! action: Action,
fields! action: Action,
strategy: Strategy
end

Expand Down
6 changes: 3 additions & 3 deletions lib/dynflow/action/v2/with_sub_plans.rb
Original file line number Diff line number Diff line change
Expand Up @@ -140,9 +140,9 @@ def recalculate_counts
failed = sub_plans_count('state' => %w(paused stopped), 'result' => %w(error warning)) - cancelled_scheduled_plans
success = sub_plans_count('state' => 'stopped', 'result' => 'success')
output.update(:pending_count => total - failed - success - cancelled_scheduled_plans,
:failed_count => failed - output.fetch(:resumed_count, 0),
:success_count => success,
:cancelled_count => cancelled)
:failed_count => failed - output.fetch(:resumed_count, 0),
:success_count => success,
:cancelled_count => cancelled)
end

def counts_set?
Expand Down
4 changes: 2 additions & 2 deletions lib/dynflow/action/with_polling_sub_plans.rb
Original file line number Diff line number Diff line change
Expand Up @@ -65,9 +65,9 @@ def recalculate_counts
total = sub_plans_count
failed = sub_plans_count('state' => %w(paused stopped), 'result' => 'error')
success = sub_plans_count('state' => 'stopped', 'result' => 'success')
output.update(:total_count => total - output.fetch(:resumed_count, 0),
output.update(:total_count => total - output.fetch(:resumed_count, 0),
:pending_count => total - failed - success,
:failed_count => failed - output.fetch(:resumed_count, 0),
:failed_count => failed - output.fetch(:resumed_count, 0),
:success_count => success)
end
end
Expand Down
2 changes: 1 addition & 1 deletion lib/dynflow/action/with_sub_plans.rb
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ def backtrace

SubPlanFinished = Algebrick.type do
fields! :execution_plan_id => String,
:success => type { variants TrueClass, FalseClass }
:success => type { variants TrueClass, FalseClass }
end

def run(event = nil)
Expand Down
6 changes: 3 additions & 3 deletions lib/dynflow/active_job/queue_adapter.rb
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@ def enqueue_at(job, timestamp)
job.provider_job_id = job.job_id
::Rails.application.dynflow.world
.delay_with_options(id: job.provider_job_id,
action_class: JobWrapper,
delay_options: { :start_at => Time.at(timestamp) },
args: [job.serialize])
action_class: JobWrapper,
delay_options: { :start_at => Time.at(timestamp) },
args: [job.serialize])
end
end

Expand Down
6 changes: 3 additions & 3 deletions lib/dynflow/clock.rb
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,9 @@ class Clock < Actor
include Algebrick::Types

Timer = Algebrick.type do
fields! who: Object, # to ping back
when: Time, # to deliver
what: Maybe[Object], # to send
fields! who: Object, # to ping back
when: Time, # to deliver
what: Maybe[Object], # to send
where: Symbol # it should be delivered, which method
end

Expand Down
10 changes: 5 additions & 5 deletions lib/dynflow/delayed_plan.rb
Original file line number Diff line number Diff line change
Expand Up @@ -59,11 +59,11 @@ def execute(future = Concurrent::Promises.resolvable_future)

def to_hash
recursive_to_hash :execution_plan_uuid => @execution_plan_uuid,
:start_at => @start_at,
:start_before => @start_before,
:serialized_args => @args_serializer.serialized_args,
:args_serializer => @args_serializer.class.name,
:frozen => @frozen
:start_at => @start_at,
:start_before => @start_before,
:serialized_args => @args_serializer.serialized_args,
:args_serializer => @args_serializer.class.name,
:frozen => @frozen
end

# Retrieves arguments from the serializer
Expand Down
10 changes: 5 additions & 5 deletions lib/dynflow/director.rb
Original file line number Diff line number Diff line change
Expand Up @@ -12,12 +12,12 @@ class Director
include Algebrick::TypeCheck

Event = Algebrick.type do
fields! request_id: String,
fields! request_id: String,
execution_plan_id: String,
step_id: Integer,
event: Object,
result: Concurrent::Promises::ResolvableFuture,
optional: Algebrick::Types::Boolean
step_id: Integer,
event: Object,
result: Concurrent::Promises::ResolvableFuture,
optional: Algebrick::Types::Boolean
end

UnprocessableEvent = Class.new(Dynflow::Error)
Expand Down
8 changes: 4 additions & 4 deletions lib/dynflow/dispatcher.rb
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,10 @@ module Dispatcher
Request = Algebrick.type do
Event = type do
fields! execution_plan_id: String,
step_id: Integer,
event: Object,
time: type { variants Time, NilClass },
optional: Algebrick::Types::Boolean
step_id: Integer,
event: Object,
time: type { variants Time, NilClass },
optional: Algebrick::Types::Boolean
end

Execution = type do
Expand Down
4 changes: 2 additions & 2 deletions lib/dynflow/execution_history.rb
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@ class ExecutionHistory
include Enumerable

Event = Algebrick.type do
fields! time: Integer,
name: String,
fields! time: Integer,
name: String,
world_id: type { variants String, NilClass }
end

Expand Down
8 changes: 4 additions & 4 deletions lib/dynflow/execution_plan/output_reference.rb
Original file line number Diff line number Diff line change
Expand Up @@ -53,11 +53,11 @@ def [](subkey)
end

def to_hash
recursive_to_hash class: self.class.to_s,
recursive_to_hash class: self.class.to_s,
execution_plan_id: execution_plan_id,
step_id: step_id,
action_id: action_id,
subkeys: subkeys
step_id: step_id,
action_id: action_id,
subkeys: subkeys
end

def to_s
Expand Down
26 changes: 13 additions & 13 deletions lib/dynflow/execution_plan/steps/abstract.rb
Original file line number Diff line number Diff line change
Expand Up @@ -87,19 +87,19 @@ def to_s

def to_hash
recursive_to_hash execution_plan_uuid: execution_plan_id,
id: id,
state: state,
class: self.class.to_s,
action_class: action_class.to_s,
action_id: action_id,
error: error,
started_at: started_at,
ended_at: ended_at,
execution_time: execution_time,
real_time: real_time,
progress_done: progress_done,
progress_weight: progress_weight,
queue: queue
id: id,
state: state,
class: self.class.to_s,
action_class: action_class.to_s,
action_id: action_id,
error: error,
started_at: started_at,
ended_at: ended_at,
execution_time: execution_time,
real_time: real_time,
progress_done: progress_done,
progress_weight: progress_weight,
queue: queue
end

def progress_done
Expand Down
6 changes: 3 additions & 3 deletions lib/dynflow/execution_plan/steps/error.rb
Original file line number Diff line number Diff line change
Expand Up @@ -42,10 +42,10 @@ def self.new_from_hash(hash)
end

def to_hash
recursive_to_hash class: self.class.name,
recursive_to_hash class: self.class.name,
exception_class: exception_class.to_s,
message: message,
backtrace: backtrace
message: message,
backtrace: backtrace
end

def to_s
Expand Down
4 changes: 2 additions & 2 deletions lib/dynflow/executors/parallel.rb
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ def initialize(world,
queues_options: { :default => { :pool_size => 5 } })
@world = world
@logger = world.logger
@core = executor_class.spawn name: 'parallel-executor-core',
args: [world, heartbeat_interval, queues_options],
@core = executor_class.spawn name: 'parallel-executor-core',
args: [world, heartbeat_interval, queues_options],
initialized: @core_initialized = Concurrent::Promises.resolvable_future
end

Expand Down
4 changes: 2 additions & 2 deletions lib/dynflow/middleware/register.rb
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@ class Middleware::Register

def initialize
@rules = Hash.new do |h, k|
h[k] = { before: [],
after: [],
h[k] = { before: [],
after: [],
replace: [] }
end
end
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,14 +4,14 @@

Sequel.migration do
helper = MsgpackMigrationHelper.new({
:dynflow_actions => [:data, :input, :output],
:dynflow_coordinator_records => [:data],
:dynflow_delayed_plans => [:serialized_args, :data],
:dynflow_envelopes => [:data],
:dynflow_execution_plans => [:run_flow, :finalize_flow, :execution_history, :step_ids],
:dynflow_steps => [:error, :children],
:dynflow_output_chunks => [:chunk]
})
:dynflow_actions => [:data, :input, :output],
:dynflow_coordinator_records => [:data],
:dynflow_delayed_plans => [:serialized_args, :data],
:dynflow_envelopes => [:data],
:dynflow_execution_plans => [:run_flow, :finalize_flow, :execution_history, :step_ids],
:dynflow_steps => [:error, :children],
:dynflow_output_chunks => [:chunk]
})

up do
helper.up(self)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,9 @@

Sequel.migration do
helper = MsgpackMigrationHelper.new({
:dynflow_execution_plans => [:data],
:dynflow_steps => [:data]
})
:dynflow_execution_plans => [:data],
:dynflow_steps => [:data]
})

up do
helper.up(self)
Expand Down
28 changes: 14 additions & 14 deletions test/persistence_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,8 @@ module PersistenceTest

def prepare_plans
execution_plans_data.map do |h|
h.merge result: nil, started_at: Time.now.utc - 20, ended_at: Time.now.utc - 10,
real_time: 0.0, execution_time: 0.0
h.merge result: nil, started_at: Time.now.utc - 20, ended_at: Time.now.utc - 10,
real_time: 0.0, execution_time: 0.0
end
end

Expand Down Expand Up @@ -173,8 +173,8 @@ def self.it_acts_as_persistence_adapter

loaded_plans = adapter.find_execution_plan_statuses(filters: { state: ['paused'] })
_(loaded_plans).must_equal({ "plan1" => { :state => "paused", :result => nil },
"plan3" => { :state => "paused", :result => nil },
"plan4" => { :state => "paused", :result => nil } })
"plan3" => { :state => "paused", :result => nil },
"plan4" => { :state => "paused", :result => nil } })

loaded_plans = adapter.find_execution_plan_statuses(filters: { state: ['stopped'] })
_(loaded_plans).must_equal({ "plan2" => { :state => "stopped", :result => nil } })
Expand All @@ -184,14 +184,14 @@ def self.it_acts_as_persistence_adapter

loaded_plans = adapter.find_execution_plan_statuses(filters: { state: ['stopped', 'paused'] })
_(loaded_plans).must_equal({ "plan1" => { :state => "paused", :result => nil },
"plan2" => { :state => "stopped", :result => nil },
"plan3" => { :state => "paused", :result => nil }, "plan4" => { :state => "paused", :result => nil } })
"plan2" => { :state => "stopped", :result => nil },
"plan3" => { :state => "paused", :result => nil }, "plan4" => { :state => "paused", :result => nil } })

loaded_plans = adapter.find_execution_plan_statuses(filters: { 'state' => ['stopped', 'paused'] })
_(loaded_plans).must_equal({ "plan1" => { :state => "paused", :result => nil },
"plan2" => { :state => "stopped", :result => nil },
"plan3" => { :state => "paused", :result => nil },
"plan4" => { :state => "paused", :result => nil } })
"plan2" => { :state => "stopped", :result => nil },
"plan3" => { :state => "paused", :result => nil },
"plan4" => { :state => "paused", :result => nil } })

loaded_plans = adapter.find_execution_plan_statuses(filters: { label: ['test1'], :delayed => true })
_(loaded_plans).must_equal({})
Expand Down Expand Up @@ -347,11 +347,11 @@ def self.it_acts_as_persistence_adapter
start_time = Time.now.utc
prepare_and_save_plans
adapter.save_delayed_plan('plan1', :execution_plan_uuid => 'plan1', :frozen => false, :start_at => format_time(start_time + 60),
:start_before => format_time(start_time - 60))
:start_before => format_time(start_time - 60))
adapter.save_delayed_plan('plan2', :execution_plan_uuid => 'plan2', :frozen => false, :start_at => format_time(start_time - 60))
adapter.save_delayed_plan('plan3', :execution_plan_uuid => 'plan3', :frozen => false, :start_at => format_time(start_time + 60))
adapter.save_delayed_plan('plan4', :execution_plan_uuid => 'plan4', :frozen => false, :start_at => format_time(start_time - 60),
:start_before => format_time(start_time - 60))
:start_before => format_time(start_time - 60))
plans = adapter.find_ready_delayed_plans(start_time)
_(plans.length).must_equal 3
_(plans.map { |plan| plan[:execution_plan_uuid] }).must_equal %w(plan2 plan4 plan1)
Expand All @@ -362,9 +362,9 @@ def self.it_acts_as_persistence_adapter
prepare_and_save_plans

adapter.save_delayed_plan('plan1', :execution_plan_uuid => 'plan1', :frozen => false, :start_at => format_time(start_time + 60),
:start_before => format_time(start_time - 60))
:start_before => format_time(start_time - 60))
adapter.save_delayed_plan('plan2', :execution_plan_uuid => 'plan2', :frozen => true, :start_at => format_time(start_time + 60),
:start_before => format_time(start_time - 60))
:start_before => format_time(start_time - 60))

plans = adapter.find_ready_delayed_plans(start_time)
_(plans.length).must_equal 1
Expand Down Expand Up @@ -519,7 +519,7 @@ def self.it_acts_as_persistence_adapter

envelopes.each { |e| adapter.push_envelope(e) }
adapter.insert_coordinator_record({ "class" => "Dynflow::Coordinator::ExecutorWorld",
"id" => executor_world_id, "meta" => {}, "active" => true })
"id" => executor_world_id, "meta" => {}, "active" => true })

assert_equal 1, adapter.prune_undeliverable_envelopes
assert_equal 0, adapter.prune_undeliverable_envelopes
Expand Down
6 changes: 4 additions & 2 deletions test/support/test_execution_log.rb
Original file line number Diff line number Diff line change
Expand Up @@ -24,11 +24,13 @@ def size
end

def self.setup
@run, @finalize = self.new, self.new
@run = new
@finalize = new
end

def self.teardown
@run, @finalize = nil, nil
@run = nil
@finalize = nil
end

def self.run
Expand Down
6 changes: 3 additions & 3 deletions test/testing_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -63,19 +63,19 @@ module Dynflow

3.times { progress_action_time action }
_(action.output).must_equal('task' => { 'progress' => 30, 'done' => false },
'poll_attempts' => { 'total' => 2, 'failed' => 0 })
'poll_attempts' => { 'total' => 2, 'failed' => 0 })
_(action.run_progress).must_equal 0.3

run_action action, Dynflow::Action::Polling::Poll
run_action action, Dynflow::Action::Polling::Poll
_(action.output).must_equal('task' => { 'progress' => 50, 'done' => false },
'poll_attempts' => { 'total' => 4, 'failed' => 0 })
'poll_attempts' => { 'total' => 4, 'failed' => 0 })
_(action.run_progress).must_equal 0.5

5.times { progress_action_time action }

_(action.output).must_equal('task' => { 'progress' => 100, 'done' => true },
'poll_attempts' => { 'total' => 9, 'failed' => 0 })
'poll_attempts' => { 'total' => 9, 'failed' => 0 })
_(action.run_progress).must_equal 1
end

Expand Down
4 changes: 2 additions & 2 deletions test/world_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,8 @@ module WorldTest
registered_world = world.coordinator.find_worlds(false, id: world.id).first
registered_world.meta.delete('last_seen')
_(registered_world.meta).must_equal('hostname' => Socket.gethostname, 'pid' => Process.pid,
'queues' => { 'default' => { 'pool_size' => 5 },
'slow' => { 'pool_size' => 1 } })
'queues' => { 'default' => { 'pool_size' => 5 },
'slow' => { 'pool_size' => 1 } })
end

it 'is configurable' do
Expand Down
Loading