Class: Step

Inherits:
Object
  • Object
show all
Defined in:
lib/rbbt/workflow/step.rb,
lib/rbbt/workflow/step/run.rb,
lib/rbbt/workflow/step/prepare.rb,
lib/rbbt/workflow/util/archive.rb,
lib/rbbt/workflow/step/accessor.rb,
lib/rbbt/workflow/util/provenance.rb,
lib/rbbt/workflow/step/dependencies.rb

Direct Known Subclasses

RemoteStep

Constant Summary collapse

RBBT_DEBUG_CLEAN =
ENV["RBBT_DEBUG_CLEAN"] == 'true'
MAIN_RSYNC_ARGS =
"-avztAXHP"
INFO_SERIALIZER =
begin
  if ENV["RBBT_INFO_SERIALIZER"]
    Kernel.const_get ENV["RBBT_INFO_SERIALIZER"]
  else
    Marshal
  end
end
STREAM_CACHE =
{}
STREAM_CACHE_MUTEX =
Mutex.new

Class Attribute Summary collapse

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(path, task = nil, inputs = nil, dependencies = nil, bindings = nil, clean_name = nil) ⇒ Step

Returns a new instance of Step.



46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
# File 'lib/rbbt/workflow/step.rb', line 46

def initialize(path, task = nil, inputs = nil, dependencies = nil, bindings = nil, clean_name = nil)
  path = Path.setup(Misc.sanitize_filename(path)) if String === path
  path = path.call if Proc === path

  @path = path
  @task = task
  @bindings = bindings
  @dependencies = case
                  when dependencies.nil? 
                    []
                  when Array === dependencies
                    dependencies
                  else
                    [dependencies]
                  end
  @mutex = Mutex.new
  @info_mutex = Mutex.new
  @inputs = inputs 
  NamedArray.setup @inputs, task.inputs.collect{|s| s.to_s} if task and task.respond_to? :inputs and task.inputs
  if Open.exists?(info_file) and (info[:path] != path)
    @relocated = true
  end
end

Class Attribute Details

.lock_dirObject

Returns the value of attribute lock_dir.



20
21
22
# File 'lib/rbbt/workflow/step.rb', line 20

def lock_dir
  @lock_dir
end

.log_relay_stepObject

Returns the value of attribute log_relay_step.



232
233
234
# File 'lib/rbbt/workflow/step.rb', line 232

def log_relay_step
  @log_relay_step
end

Instance Attribute Details

#bindingsObject

Returns the value of attribute bindings.



9
10
11
# File 'lib/rbbt/workflow/step.rb', line 9

def bindings
  @bindings
end

#clean_nameObject

Returns the value of attribute clean_name.



9
10
11
# File 'lib/rbbt/workflow/step.rb', line 9

def clean_name
  @clean_name
end

#dependenciesObject

Returns the value of attribute dependencies.



9
10
11
# File 'lib/rbbt/workflow/step.rb', line 9

def dependencies
  @dependencies
end

#duppedObject (readonly)

Returns the value of attribute dupped.



9
10
11
# File 'lib/rbbt/workflow/step/run.rb', line 9

def dupped
  @dupped
end

#exec(no_load = false) ⇒ Object

Returns the value of attribute exec.



12
13
14
# File 'lib/rbbt/workflow/step.rb', line 12

def exec
  @exec
end

#inputsObject

Returns the value of attribute inputs.



9
10
11
# File 'lib/rbbt/workflow/step.rb', line 9

def inputs
  @inputs
end

#mutexObject

Returns the value of attribute mutex.



14
15
16
# File 'lib/rbbt/workflow/step.rb', line 14

def mutex
  @mutex
end

#original_task_nameObject

Returns the value of attribute original_task_name.



15
16
17
# File 'lib/rbbt/workflow/step.rb', line 15

def original_task_name
  @original_task_name
end

#overridenObject

Returns the value of attribute overriden.



10
11
12
# File 'lib/rbbt/workflow/step.rb', line 10

def overriden
  @overriden
end

#pathObject

Returns the value of attribute path.



9
10
11
# File 'lib/rbbt/workflow/step.rb', line 9

def path
  @path
end

#pidObject

Returns the value of attribute pid.



11
12
13
# File 'lib/rbbt/workflow/step.rb', line 11

def pid
  @pid
end

#real_inputsObject

Returns the value of attribute real_inputs.



15
16
17
# File 'lib/rbbt/workflow/step.rb', line 15

def real_inputs
  @real_inputs
end

#relocatedObject

Returns the value of attribute relocated.



13
14
15
# File 'lib/rbbt/workflow/step.rb', line 13

def relocated
  @relocated
end

#resultObject

Returns the value of attribute result.



14
15
16
# File 'lib/rbbt/workflow/step.rb', line 14

def result
  @result
end

#saved_streamObject (readonly)

Returns the value of attribute saved_stream.



9
10
11
# File 'lib/rbbt/workflow/step/run.rb', line 9

def saved_stream
  @saved_stream
end

#seenObject

Returns the value of attribute seen.



14
15
16
# File 'lib/rbbt/workflow/step.rb', line 14

def seen
  @seen
end

#streamObject (readonly)

Returns the value of attribute stream.



9
10
11
# File 'lib/rbbt/workflow/step/run.rb', line 9

def stream
  @stream
end

#taskObject

Returns the value of attribute task.



9
10
11
# File 'lib/rbbt/workflow/step.rb', line 9

def task
  @task
end

#task_nameObject

Returns the value of attribute task_name.



10
11
12
# File 'lib/rbbt/workflow/step.rb', line 10

def task_name
  @task_name
end

#workflowObject

Returns the value of attribute workflow.



9
10
11
# File 'lib/rbbt/workflow/step.rb', line 9

def workflow
  @workflow
end

Class Method Details

.archive(files, target = nil, recursive = true) ⇒ Object



108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
# File 'lib/rbbt/workflow/util/archive.rb', line 108

def self.archive(files, target = nil, recursive = true)
  target = self.path + '.tar.gz' if target.nil?
  target = File.expand_path(target) if String === target

  job_files = job_files_for_archive files, recursive
  TmpFile.with_file do |tmpdir|
    job_files.each do |file|
      Step.link_job file, tmpdir
    end

    Misc.in_dir(tmpdir) do
      if File.directory?(target)
        CMD.cmd_log("rsync #{MAIN_RSYNC_ARGS} --copy-unsafe-links '#{ tmpdir }/' '#{ target }/'")
      else
        CMD.cmd_log("tar cvhzf '#{target}'  ./*")
      end
    end
    Log.debug "Archive finished at: #{target}"
  end
end

.clean(path) ⇒ Object



425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
# File 'lib/rbbt/workflow/step.rb', line 425

def self.clean(path)
  info_file = Step.info_file path
  pid_file = Step.pid_file path
  md5_file = Step.md5_file path
  files_dir = Step.files_dir path
  tmp_path = Step.tmp_path path

  if ! (Open.writable?(path) && Open.writable?(info_file))
    Log.warn "Could not clean #{path}: not writable"
    return 
  end

  if (Open.exists?(path) or Open.broken_link?(path)) or Open.exists?(pid_file) or Open.exists?(info_file) or Open.exists?(files_dir) or Open.broken_link?(files_dir)

    @result = nil
    @pid = nil

    Misc.insist do
      Open.rm info_file if Open.exists?(info_file)
      Open.rm md5_file if Open.exists?(md5_file)
      Open.rm path if (Open.exists?(path) || Open.broken_link?(path))
      Open.rm_rf files_dir if Open.exists?(files_dir) || Open.broken_link?(files_dir)
      Open.rm pid_file if Open.exists?(pid_file)
      Open.rm tmp_path if Open.exists?(tmp_path)
    end
  end
end

.dup_stream(stream) ⇒ Object



13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
# File 'lib/rbbt/workflow/step/dependencies.rb', line 13

def self.dup_stream(stream)
  case stream
  when IO, File, Step
    return stream if stream.respond_to?(:closed?) and stream.closed?
    return stream if stream.respond_to?(:done?) and stream.done?

    STREAM_CACHE_MUTEX.synchronize do
      stream_key = Misc.fingerprint(stream)
      current = STREAM_CACHE[stream_key]
      case current
      when nil, Step
        Log.medium "Not duplicating stream #{stream_key}"
        STREAM_CACHE[stream_key] = stream
      when File
        if Open.exists?(current.path)
          Log.medium "Reopening file #{stream_key}"
          Open.open(current.path)
        else
          new = Misc.dup_stream(current)
          Log.medium "Duplicating file #{stream_key} #{current.inspect} => #{Misc.fingerprint(new)}"
          new
        end
      else
        new = Misc.dup_stream(current)
        Log.medium "Duplicating stream #{stream_key} #{ Misc.fingerprint(stream) } => #{Misc.fingerprint(new)}"
        new
      end
    end
  when TSV::Dumper, TSV::Parser
    orig_stream = stream
    stream = stream.stream
    return stream if stream.closed?

    STREAM_CACHE_MUTEX.synchronize do
      if STREAM_CACHE[stream].nil?
        Log.high "Not duplicating #{Misc.fingerprint orig_stream} #{ stream.inspect }"
        STREAM_CACHE[stream] = stream
      else
        new = Misc.dup_stream(STREAM_CACHE[stream])
        Log.high "Duplicating #{Misc.fingerprint orig_stream} #{ stream.inspect } into #{new.inspect}"
        new
      end
    end
  else
    stream
  end
end

.files_dir(path) ⇒ Object



45
46
47
# File 'lib/rbbt/workflow/step/accessor.rb', line 45

def self.files_dir(path)
  path.nil? ? nil : path + '.files'
end

.info_file(path) ⇒ Object



49
50
51
# File 'lib/rbbt/workflow/step/accessor.rb', line 49

def self.info_file(path)
  path.nil? ? nil : path + '.info'
end

.job_files_for_archive(files, recursive = false, skip_overriden = false) ⇒ Object



54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
# File 'lib/rbbt/workflow/util/archive.rb', line 54

def self.job_files_for_archive(files, recursive = false, skip_overriden = false)
  job_files = Set.new

  jobs = files.collect do |file|  
    if Step === file
      file
    else
      file = file.sub(/\.info$/,'')
      Step.new(File.expand_path(file))
    end
  end.uniq

  jobs.each do |step|
    next unless File.exists?(step.path)
    next if skip_overriden && step.overriden

    job_files << step.path
    job_files << step.info_file if File.exists?(step.info_file)
    job_files << Step.md5_file(step.path) if File.exists?(Step.md5_file step.path)
    job_file_dir_content = Dir.glob(step.files_dir + '/**/*')
    job_files += job_file_dir_content
    job_files << step.files_dir if File.exists?(step.files_dir)
    rec_dependencies = Set.new

    next unless recursive

    deps = [step.path]
    seen = Set.new
    while deps.any?
      path = deps.shift

      dep = Workflow.load_step path
      seen << dep.path

      dep.load_dependencies_from_info

      dep.dependencies.each do |dep|
        next if seen.include? dep.path
        deps << dep.path
        rec_dependencies << dep.path
      end if dep.info[:dependencies]
    end

    rec_dependencies.each do |path|
      dep = Workflow.load_step path
      job_files << dep.path
      job_files << dep.files_dir if Dir.glob(dep.files_dir + '/*').any?
      job_files << dep.info_file if File.exists?(dep.info_file)
    end
  end

  job_files.to_a
end

.job_name_for_info_file(info_file, extension = nil) ⇒ Object



80
81
82
83
84
85
86
# File 'lib/rbbt/workflow/step/accessor.rb', line 80

def self.job_name_for_info_file(info_file, extension = nil)
  if extension and not extension.empty?
    info_file.sub(/\.#{extension}\.info$/,'')
  else
    info_file.sub(/\.info$/,'')
  end
end


5
6
7
8
9
10
11
12
13
14
15
16
17
18
# File 'lib/rbbt/workflow/util/archive.rb', line 5

def self.link_job(path, target_dir, task = nil, workflow = nil)
  Path.setup(target_dir)

  name = File.basename(path)
  task = File.basename(File.dirname(path)) if task.nil?
  workflow = File.basename(File.dirname(File.dirname(path))) if workflow.nil?

  return if target_dir[workflow][task][name].exists? || File.symlink?(target_dir[workflow][task][name].find)
  Log.debug "Linking #{ path }"
  FileUtils.mkdir_p target_dir[workflow][task] unless target_dir[workflow][task].exists?
  FileUtils.ln_s path, target_dir[workflow][task][name].find if File.exists?(path)
  FileUtils.ln_s path + '.files', target_dir[workflow][task][name].find + '.files' if File.exists?(path + '.files')
  FileUtils.ln_s path + '.info', target_dir[workflow][task][name].find + '.info' if File.exists?(path + '.info')
end

.load_serialized_info(io) ⇒ Object



16
17
18
# File 'lib/rbbt/workflow/step/accessor.rb', line 16

def self.load_serialized_info(io)
  IndiferentHash.setup(INFO_SERIALIZER.load(io))
end

.log(status, message, path, &block) ⇒ Object



395
396
397
398
399
400
401
402
403
404
405
# File 'lib/rbbt/workflow/step/accessor.rb', line 395

def self.log(status, message, path, &block)
  if block
    if Hash === message
      log_progress(status, message, path, &block)
    else
      log_block(status, message, path, &block)
    end
  else
    log_string(status, message, path)
  end
end

.log_block(status, message, path, &block) ⇒ Object



323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
# File 'lib/rbbt/workflow/step/accessor.rb', line 323

def self.log_block(status, message, path, &block)
  start = Time.now
  status = status.to_s
  status_color = self.status_color status

  Log.info do 
    now = Time.now
    str = Log.color :reset
    str << "#{ Log.color status_color, status}"
    str << ": #{ message }" if message and message != :result
    str << " -- #{Log.color :blue, path.to_s}" if path
    str << " #{Log.color :yellow, Process.pid}"
    str
  end
  res = yield
  eend = Time.now
  Log.info do 
    now = Time.now
    str = "#{ Log.color :cyan, status.to_s } +#{Log.color :green, "%.2f" % (eend - start)}"
    str << ": #{ res }" if message == :result
    str << " -- #{Log.color :blue, path.to_s}" if path
    str << " #{Log.color :yellow, Process.pid}"
    str
  end
  res
end

.log_progress(status, options = {}, path = nil, &block) ⇒ Object



365
366
367
368
369
370
371
372
373
374
375
376
377
378
# File 'lib/rbbt/workflow/step/accessor.rb', line 365

def self.log_progress(status, options = {}, path = nil, &block)
  options = Misc.add_defaults options, :severity => Log::INFO, :file => path
  max = Misc.process_options options, :max
  Log::ProgressBar.with_bar(max, options) do |bar|
    begin
      res = yield bar
      raise KeepBar.new res if IO === res
      res
    rescue
      Log.exception $!
      raise $!
    end
  end
end

.log_string(status, message, path) ⇒ Object



350
351
352
353
354
355
356
357
358
359
360
361
362
363
# File 'lib/rbbt/workflow/step/accessor.rb', line 350

def self.log_string(status, message, path)
  Log.info do 

    status = status.to_s
    status_color = self.status_color status

    str = Log.color :reset
    str << "#{ Log.color status_color, status}"
    str << ": #{ message }" if message
    str << " -- #{Log.color :blue, path.to_s}" if path
    str << " #{Log.color :yellow, Process.pid}"
    str
  end
end

.md5_file(path) ⇒ Object



61
62
63
# File 'lib/rbbt/workflow/step/accessor.rb', line 61

def self.md5_file(path)
  path.nil? ? nil : path + '.md5'
end

.migrate(path, search_path, options = {}) ⇒ Object



129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
# File 'lib/rbbt/workflow/util/archive.rb', line 129

def self.migrate(path, search_path, options = {})
  resource=Rbbt

  orig_path = path
  other_rsync_args = options[:rsync]

  recursive = options[:recursive]
  recursive = false if recursive.nil?

  paths = if options[:source]
            Misc.ssh_run(options[:source], <<-EOF).split("\n")
require 'rbbt-util'
require 'rbbt/workflow'

path = "#{path}"
recursive = #{ recursive.to_s }

if File.exists?(path)
path = #{resource.to_s}.identify(path)
else
path = Path.setup(path)
end

files = path.glob_all

files = Step.job_files_for_archive(files, recursive)

puts files * "\n"
            EOF

          else
            if File.exists?(path)
              path = resource.identify(path)
              raise "Resource #{resource} could not identify #{orig_path}" if path.nil?
            else
              path = Path.setup(path)
            end
            files = path.glob_all
            files = Step.job_files_for_archive(files, recursive)
            files
          end


  target = if options[:target] 
             target = Misc.ssh_run(options[:target], <<-EOF).split("\n").first
require 'rbbt-util'
path = "var/jobs"
resource = #{resource.to_s}
search_path = "#{search_path}"
puts resource[path].find(search_path)
             EOF
           else
             resource['var/jobs'].find(search_path)
           end

  subpath_files = {}
  paths.sort.each do |path|
    parts = path.split("/")
    subpath = parts[0..-4] * "/" + "/"

    if subpath_files.keys.any? && subpath.start_with?(subpath_files.keys.last)
      subpath = subpath_files.keys.last
    end

    source = path[subpath.length..-1]

    subpath_files[subpath] ||= []
    subpath_files[subpath] << source
  end

  synced_files = []
  subpath_files.each do |subpath, files|
    if options[:target]
      CMD.cmd("ssh #{options[:target]} mkdir -p '#{File.dirname(target)}'")
    else
      Open.mkdir File.dirname(target)
    end

    if options[:source]
      source = [options[:source], subpath] * ":"
    else
      source = subpath
    end
    target = [options[:target], target] * ":" if options[:target]

    next if File.exists?(source) && File.exists?(target) && File.expand_path(source) == File.expand_path(target)

    files_and_dirs = Set.new( files )
    files.each do |file|
      synced_files << File.join(subpath, file)

      parts = file.split("/")[0..-2].reject{|p| p.empty?}
      while parts.any?
        files_and_dirs <<  parts * "/"
        parts.pop
      end
    end

    TmpFile.with_file(files_and_dirs.sort_by{|l| l.length}.to_a * "\n") do |tmp_include_file|
      test_str = options[:test] ? '-nv' : ''

      cmd = "rsync #{MAIN_RSYNC_ARGS} --progress #{test_str} --files-from='#{tmp_include_file}' #{source}/ #{target}/ #{other_rsync_args}"

      #cmd << " && rm -Rf #{source}" if options[:delete]
      if options[:print]
        ppp Open.read(tmp_include_file)
        puts cmd 
      else
        CMD.cmd_log(cmd, :log => Log::INFO)
      end
    end
  end

  if options[:delete] && synced_files.any?
    puts Log.color :magenta, "About to erase these files:"
    synced_files.each do |p|
      puts Log.color :red, p
    end 

    if options[:non_interactive]
      response = 'yes'
    else
      puts Log.color :magenta, "Type 'yes' if you are sure:"
      response = STDIN.gets.chomp
    end

    if response == 'yes'
      synced_files.each do |p|
        Open.rm p
      end
    end
  end
end

.pid_file(path) ⇒ Object



65
66
67
# File 'lib/rbbt/workflow/step/accessor.rb', line 65

def self.pid_file(path)
  path.nil? ? nil : path + '.pid'
end

.prepare_dependencies(jobs, tasks, cpus) ⇒ Object



2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
# File 'lib/rbbt/workflow/step/prepare.rb', line 2

def self.prepare_dependencies(jobs, tasks, cpus)
  deps = []

  jobs = [jobs] unless Array === jobs
  tasks = [tasks] unless Array === tasks
  tasks = tasks.collect{|t| t.to_s}

  jobs.each do |job|
    job.rec_dependencies.each do |dep|
      next if dep.done?
      dep.clean if dep.error? && dep.recoverable_error?
      deps << dep if tasks.include?(dep.task_name.to_s) or tasks.include?([dep.workflow.to_s, dep.task_name] * "#")
    end
  end

  cpus = jobs.length if cpus.to_s == "max"
  cpus = cpus.to_i if String === cpus
  TSV.traverse deps.collect{|dep| dep.path}, :type => :array, :cpus => cpus, :bar => "Prepare dependencies #{Misc.fingerprint tasks} for #{Misc.fingerprint jobs}" do |path|
    dep = deps.select{|dep| dep.path == path}.first
    dep.produce
    nil
  end
end

.prepare_for_execution(job) ⇒ Object

Raises:



73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
# File 'lib/rbbt/workflow/step/dependencies.rb', line 73

def self.prepare_for_execution(job)
  return if job.done? && ! job.dirty?

  status = job.status.to_s

  if defined?(WorkflowRemoteClient) && WorkflowRemoteClient::RemoteStep === job 
    return unless (status == 'done' or status == 'error' or status == 'aborted')
  else
    return if status == 'streaming' and job.running?
  end

  canfail = nil
  job.status_lock.synchronize do
    status = job.status.to_s

    if (status == 'error' && (job.recoverable_error? || job.dirty?)) ||
        (job.noinfo? && Open.exists?(job.pid_file)) ||
        job.aborted? ||
        (job.done? && ! job.updated?)  || (job.error? && ! job.updated?) ||
        (job.done? && job.dirty?)  || (job.error? && job.dirty?) ||
        (!(job.noinfo? || job.done? || job.error? || job.aborted? || job.running?))

      if ! (job.resumable? && (job.updated? && ! job.dirty?))
        Log.high "About to clean -- status: #{status}, present #{File.exists?(job.path)}, " +
          %w(done? error? recoverable_error? noinfo? updated? dirty? aborted? running? resumable?).
          collect{|v| [v, job.send(v)]*": "} * ", " if RBBT_DEBUG_CLEAN

        job.clean
      end
      job.set_info :status, :cleaned
    end

    job.dup_inputs unless status == 'done' or job.started?
    job.init_info(status == 'noinfo') unless status == 'waiting' || status == 'done' || job.started? || ! Workflow.job_path?(job.path)

    canfail = ComputeDependency === job && job.canfail?
  end

  raise DependencyError, job if job.error? and not canfail
end

.prov_report(step, offset = 0, task = nil, seen = [], expand_repeats = false) ⇒ Object



72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
# File 'lib/rbbt/workflow/util/provenance.rb', line 72

def self.prov_report(step, offset = 0, task = nil, seen = [], expand_repeats = false)
  info = step.info  || {}
  info[:task_name] = task
  path  = step.path
  status = info[:status] || :missing
  status = "remote" if Open.remote?(path) || Open.ssh?(path)
  name = info[:name] || File.basename(path)
  status = :unsync if status == :done and not Open.exist?(path)
  status = :notfound if status == :noinfo and not Open.exist?(path)

  str = " " * offset
  str << prov_report_msg(status, name, path, info)
  step.dependencies.reverse.each do |dep|
    path = dep.path
    new = ! seen.include?(path)
    if new
      seen << path
      str << prov_report(dep, offset + 1, task, seen, expand_repeats)
    else
      if expand_repeats
        str << Log.color(:green, Log.uncolor(prov_report(dep, offset+1, task)))
      else
        info = dep.info  || {}
        status = info[:status] || :missing
        status = "remote" if Open.remote?(path) || Open.ssh?(path)
        name = info[:name] || File.basename(path)
        status = :unsync if status == :done and not Open.exist?(path)
        status = :notfound if status == :noinfo and not Open.exist?(path)

        str << Log.color(status == :notfound ? :blue : :green, " " * (offset + 1) + Log.uncolor(prov_report_msg(status, name, path, info)))
      end
    end
  end if step.dependencies
  str
end

.prov_report_msg(status, name, path, info = nil) ⇒ Object



24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
# File 'lib/rbbt/workflow/util/provenance.rb', line 24

def self.prov_report_msg(status, name, path, info = nil)
  parts = path.sub(/\{.*/,'').split "/"

  parts.pop
  
  task = Log.color(:yellow, parts.pop)
  workflow = Log.color(:magenta, parts.pop)
  #if status.to_s == 'noinfo' && parts.last != 'jobs'
  if ! Workflow.job_path?(path)
    task, status, workflow = Log.color(:yellow, info[:task_name]), Log.color(:green, "file"), Log.color(:magenta, "-")
  end

  path_mtime = begin
                 Open.mtime(path)
               rescue Exception
                 nil
               end
  str = if ! (Open.remote?(path) || Open.ssh?(path)) && (Open.exists?(path) && $main_mtime && path_mtime && ($main_mtime - path_mtime) < -2)
          prov_status_msg(status.to_s) << " " << [workflow, task, path].compact * " " << " (#{Log.color(:red, "Mtime out of sync") })"
        else
          prov_status_msg(status.to_s) << " " << [workflow, task, path].compact * " " 
        end

  if $inputs and $inputs.any? 
    job_inputs = Workflow.load_step(path).recursive_inputs.to_hash
    IndiferentHash.setup(job_inputs)

    $inputs.each do |input|
      value = job_inputs[input]
      next if  value.nil?
      value_str = Misc.fingerprint(value)
      str << "\t#{Log.color :magenta, input}=#{value_str}"
    end
  end

  if $info_fields and $info_fields.any?
    $info_fields.each do |field|
      IndiferentHash.setup(info)
      value = info[field]
      next if value.nil?
      value_str = Misc.fingerprint(value)
      str << "\t#{Log.color :magenta, field}=#{value_str}"
    end
  end

  str << "\n"
end

.prov_status_msg(status) ⇒ Object



2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
# File 'lib/rbbt/workflow/util/provenance.rb', line 2

def self.prov_status_msg(status)
  color = case status.to_sym
          when :error, :aborted, :missing, :dead, :unsync
            :red
          when :streaming, :started
            :cyan
          when :done, :noinfo
            :green
          when :dependencies, :waiting, :setup
            :yellow
          when :notfound, :cleaned
            :blue
          else
            if status.to_s.index ">"
              :cyan
            else
              :cyan
            end
          end
  Log.color(color, status.to_s)
end

.purge(path, recursive = false, skip_overriden = true) ⇒ Object



263
264
265
266
267
268
269
270
271
272
273
274
275
# File 'lib/rbbt/workflow/util/archive.rb', line 263

def self.purge(path, recursive = false, skip_overriden = true)
  path = [path] if String === path
  job_files = job_files_for_archive path, recursive, skip_overriden

  job_files.each do |file|
    begin
      Log.debug "Purging #{file}"
      Open.rm_rf file if Open.exists?(file)
    rescue
      Log.warn "Could not erase '#{file}': #{$!.message}"
    end
  end
end

.purge_stream_cacheObject



6
7
8
9
10
11
# File 'lib/rbbt/workflow/step/dependencies.rb', line 6

def self.purge_stream_cache
  Log.debug "Purging dup. stream cache"
  STREAM_CACHE_MUTEX.synchronize do
    STREAM_CACHE.clear
  end
end

.save_inputs(inputs, input_types, dir) ⇒ Object



88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
# File 'lib/rbbt/workflow/step/accessor.rb', line 88

def self.save_inputs(inputs, input_types, dir)
  inputs.each do |name,value|
    type = input_types[name]
    type = type.to_s if type
    path = File.join(dir, name.to_s)

    Log.debug "Saving job input #{name} (#{type}) into #{path}"
    case
    when Step === value
      Open.ln_s(value.path, path)
    when type.to_s == "file"
      if String === value && File.exists?(value)
        Open.ln_s(value, path)
      else
        Open.write(path + '.yaml', value.to_yaml)
      end
    when Array === value
      Open.write(path, value.collect{|v| Step === v ? v.path : v.to_s} * "\n")
    when IO === value
      if value.filename && String === value.filename && File.exists?(value.filename)
        Open.ln_s(value.filename, path)
      else
        Open.write(path, value)
      end
    else
      Open.write(path, value.to_s)
    end
  end.any?
end

.save_job_inputs(job, dir, options = nil) ⇒ Object



118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
# File 'lib/rbbt/workflow/step/accessor.rb', line 118

def self.save_job_inputs(job, dir, options = nil)
  options = IndiferentHash.setup options.dup if options

  task_name = Symbol === job.overriden ? job.overriden : job.task_name
  workflow = job.workflow
  workflow = Kernel.const_get workflow if String === workflow
  if workflow
    task_info = workflow.task_info(task_name)
    input_types = task_info[:input_types]
    task_inputs = task_info[:inputs]
    input_defaults = task_info[:input_defaults]
  else
    task_info = input_types = task_inputs = input_defaults = {}
  end

  inputs = {}
  real_inputs = job.real_inputs || job.info[:real_inputs]
  job.recursive_inputs.zip(job.recursive_inputs.fields).each do |value,name|
    next unless task_inputs.include? name.to_sym
    next unless real_inputs.include? name.to_sym
    next if options && ! options.include?(name)
    next if value.nil?
    next if input_defaults[name] == value
    inputs[name] = value
  end

  if options && options.include?('override_dependencies')
    inputs.merge!(:override_dependencies => open[:override_dependencies])
    input_types = IndiferentHash.setup(input_types.merge(:override_dependencies => :array))
  end
  save_inputs(inputs, input_types, dir)

  inputs.keys
end

.serialize_info(info) ⇒ Object



11
12
13
14
# File 'lib/rbbt/workflow/step/accessor.rb', line 11

def self.serialize_info(info)
  info = info.clean_version if IndiferentHash === info
  INFO_SERIALIZER.dump(info)
end

.status_color(status) ⇒ Object



309
310
311
312
313
314
315
316
317
318
319
320
321
# File 'lib/rbbt/workflow/step/accessor.rb', line 309

def self.status_color(status)
  status = status.split(">").last
  case status
  when "starting"
    :yellow
  when "error", "aborted"
    :red
  when "done"
    :green
  else
    :cyan
  end
end

.step_info(path) ⇒ Object



69
70
71
72
73
74
75
76
77
78
# File 'lib/rbbt/workflow/step/accessor.rb', line 69

def self.step_info(path)
  begin
    Open.open(info_file(path), :mode => 'rb') do |f|
      self.load_serialized_info(f)
    end
  rescue Exception
    Log.exception $!
    {}
  end
end

.tmp_path(path) ⇒ Object



53
54
55
56
57
58
59
# File 'lib/rbbt/workflow/step/accessor.rb', line 53

def self.tmp_path(path)
  path = path.find if Path === path
  path = File.expand_path(path)
  dir = File.dirname(path)
  filename = File.basename(path)
  File.join(dir, '.' << filename)
end

.wait_for_jobs(jobs) ⇒ Object



21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
# File 'lib/rbbt/workflow/step/accessor.rb', line 21

def self.wait_for_jobs(jobs)
  jobs = [jobs] if Step === jobs
  begin
    threads = []

    threads = jobs.collect do |j| 
      Thread.new do
        begin
          j.join unless j.done?
        rescue Exception
          Log.error "Exception waiting for job: #{Log.color :blue, j.path}"
          raise $!
        end
      end
    end

    threads.each{|t| t.join }
  rescue Exception
    threads.each{|t| t.exit }
    jobs.each do |j| j.abort end
    raise $!
  end
end

Instance Method Details

#_abortObject



628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
# File 'lib/rbbt/workflow/step/run.rb', line 628

def _abort
  return if @aborted
  @aborted = true
  Log.medium{"#{Log.color :red, "Aborting"} #{Log.color :blue, path}"}
  doretry = true
  begin
    return if done?
    abort_pid if running?
    kill_children
    abort_stream
    stop_dependencies
  rescue Aborted, Interrupt
    Log.medium{"#{Log.color :red, "Aborting ABORTED RETRY"} #{Log.color :blue, path}"}
    if doretry
      doretry = false
      retry
    end
    raise $!
  rescue Exception
    if doretry
      doretry = false
      retry
    end
  ensure
    _clean_finished
  end
end

#_clean_finishedObject



617
618
619
620
621
622
623
624
625
626
# File 'lib/rbbt/workflow/step/run.rb', line 617

def _clean_finished
  if Open.exists?(path) && status != :done
    Log.warn "Aborted job had finished. Removing result -- #{ path }"
    begin
      Open.rm path
    rescue Exception
      Log.warn "Exception removing result of aborted job: #{$!.message}"
    end
  end
end

#_execObject



93
94
95
96
97
98
99
100
101
102
103
# File 'lib/rbbt/workflow/step/run.rb', line 93

def _exec
  resolve_input_steps
  rewind_inputs
  @exec = true if @exec.nil?
  begin
    old = Signal.trap("INT"){ Thread.current.raise Aborted }
    @task.exec_in((bindings ? bindings : self), *@inputs)
  ensure
    Signal.trap("INT", old)
  end
end

#abortObject



656
657
658
659
660
661
# File 'lib/rbbt/workflow/step/run.rb', line 656

def abort
  return if done? and (status == :done or status == :noinfo)
  _abort
  log(:aborted, "Job aborted") unless aborted? or error?
  self
end

#abort_pidObject



576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
# File 'lib/rbbt/workflow/step/run.rb', line 576

def abort_pid
  @pid ||= info[:pid] || Open.read(pid_file)

  case @pid
  when nil
    Log.medium "Could not abort #{path}: no pid"
    false
  when Process.pid
    Log.medium "Could not abort #{path}: same process"
    false
  else
    Log.medium "Aborting pid #{path}: #{ @pid } #{Process.pid}"
    begin
      Process.kill("TERM", @pid.to_i)
      s = Process.waitpid2 @pid.to_i
      Log.medium "Aborted pid #{path} #{s}"
    rescue Exception
      Log.debug("Aborted job #{@pid} was not killed: #{$!.message}")
    end
    true
  end
end

#abort_streamObject



599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
# File 'lib/rbbt/workflow/step/run.rb', line 599

def abort_stream
  stream = @result if IO === @result
  @saved_stream = nil
  if stream and stream.respond_to? :abort and not stream.aborted?
    doretry = true
    begin
      Log.medium "Aborting job stream #{stream.inspect} -- #{Log.color :blue, path}"
      stream.abort 
    rescue Aborted, Interrupt
      Log.medium "Aborting job stream #{stream.inspect} ABORTED RETRY -- #{Log.color :blue, path}"
      if doretry
        doretry = false
        retry
      end
    end
  end
end

#aborted?Boolean

Returns:

  • (Boolean)


546
547
548
549
# File 'lib/rbbt/workflow/step/accessor.rb', line 546

def aborted?
  status = self.status
  status == :aborted || ((status != :cleaned && status != :noinfo && status != :setup && status != :noinfo) && nopid?)
end

#accessObject



666
667
668
# File 'lib/rbbt/workflow/step/accessor.rb', line 666

def access
  CMD.cmd("touch -c -h -a #{self.path} #{self.info_file}")
end

#archive(target = nil) ⇒ Object



20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
# File 'lib/rbbt/workflow/util/archive.rb', line 20

def archive(target = nil)
  target = self.path + '.tar.gz' if target.nil?
  target = File.expand_path(target)
  TmpFile.with_file do |tmpdir|
    Step.link_job self.path, tmpdir
    rec_dependencies = Set.new
    deps = [self.path]
    seen = Set.new
    while deps.any?
      path = deps.shift
      dep = Workflow.load_step path
      seen << dep.path
      dep.dependencies.each do |dep|
        next if seen.include? dep.path
        deps << dep.path
        rec_dependencies << dep.path
      end if dep.dependencies
    end

    rec_dependencies.each do |path|
      Step.link_job path, tmpdir
    end

    Misc.in_dir(tmpdir) do
      if File.directory?(target)
        CMD.cmd_log("rsync #{MAIN_RSYNC_ARGS} --copy-unsafe-links '#{ tmpdir }/' '#{ target }/'")
      else
        CMD.cmd_log("tar cvhzf '#{target}'  ./*")
      end
    end
    Log.debug "Archive finished at: #{target}"
  end
end

#archive_depsObject



127
128
129
130
# File 'lib/rbbt/workflow/step.rb', line 127

def archive_deps
  self.set_info :archived_info, archived_info
  self.set_info :archived_dependencies, info[:dependencies]
end

#archived_infoObject



132
133
134
135
136
137
138
139
140
141
142
# File 'lib/rbbt/workflow/step.rb', line 132

def archived_info
  return info[:archived_info] if info[:archived_info]

  archived_info = {}
  dependencies.each do |dep|
    archived_info[dep.path] = dep.info
    archived_info.merge!(dep.archived_info)
  end if dependencies

  archived_info
end

#archived_inputsObject



144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
# File 'lib/rbbt/workflow/step.rb', line 144

def archived_inputs
  return {} unless info[:archived_dependencies]
  archived_info = self.archived_info

  all_inputs = IndiferentHash.setup({})
  deps = info[:archived_dependencies].collect{|p| p.last}
  seen = []
  while path = deps.pop
    dep_info = archived_info[path]
    if dep_info
      dep_info[:inputs].each do |k,v|
        all_inputs[k] = v unless all_inputs.include?(k)
      end if dep_info[:inputs]
      deps.concat(dep_info[:dependencies].collect{|p| p.last } - seen) if dep_info[:dependencies]
      deps.concat(dep_info[:archived_dependencies].collect{|p| p.last } - seen) if dep_info[:archived_dependencies]
    end
    seen << path
  end

  all_inputs
end

#canfail_pathsObject



317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
# File 'lib/rbbt/workflow/step/dependencies.rb', line 317

def canfail_paths
  return Set.new if done? && ! Open.exists?(info_file)

  @canfail_paths ||= begin 
                       if info[:canfail] 
                         paths = info[:canfail].uniq
                         paths = Workflow.relocate_array self.path, paths if relocated
                         Set.new(paths)
                       else
                         canfail_paths = Set.new
                         all_deps = dependencies || []
                         all_deps.each do |dep|
                           next if canfail_paths.include? dep.path
                           canfail_paths += dep.canfail_paths
                           next unless ComputeDependency === dep && dep.canfail?
                           canfail_paths << dep.path
                           canfail_paths += dep.rec_dependencies.collect{|d| d.path }
                         end
                         canfail_paths
                         begin
                           set_info :canfail, canfail_paths.to_a
                         rescue Errno::EROFS
                         end
                         canfail_paths
                       end
                     end
end

#checksObject



138
139
140
# File 'lib/rbbt/workflow/step/run.rb', line 138

def checks
  (dependency_checks + input_checks).uniq
end

#child(&block) ⇒ Object



342
343
344
345
346
347
348
349
350
351
352
# File 'lib/rbbt/workflow/step.rb', line 342

def child(&block)
  child_pid = Process.fork &block
  children_pids = info[:children_pids]
  if children_pids.nil?
    children_pids = [child_pid]
  else
    children_pids << child_pid
  end
  set_info :children_pids, children_pids
  child_pid
end

#cleanObject



460
461
462
463
464
465
466
467
468
469
470
# File 'lib/rbbt/workflow/step.rb', line 460

def clean
  status = []
  status << "dirty" if done? && dirty?
  status << "not running" if ! done? && ! running? 
  status.unshift " " if status.any?
  Log.high "Cleaning step: #{path}#{status * " "}"
  Log.stack caller if RBBT_DEBUG_CLEAN
  abort if ! done? && running?
  Step.clean(path)
  self
end

#cmd(*args) ⇒ Object



354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
# File 'lib/rbbt/workflow/step.rb', line 354

def cmd(*args)
  all_args = *args

  all_args << {} unless Hash === all_args.last

  level = all_args.last[:log] || 0
  level = 0 if TrueClass === level
  level = 10 if FalseClass === level
  level = level.to_i

  all_args.last[:log] = true
  all_args.last[:pipe] = true

  io = CMD.cmd(*all_args)
  child_pid = io.pids.first

  children_pids = info[:children_pids]
  if children_pids.nil?
    children_pids = [child_pid]
  else
    children_pids << child_pid
  end
  set_info :children_pids, children_pids

  while c = io.getc
    STDERR << c if Log.severity <= level
    if c == "\n"
      if pid
        Log.logn "STDOUT [#{pid}]: ", level
      else
        Log.logn "STDOUT: ", level
      end
    end
  end 

  io.join

  nil
end

#config(key, *tokens) ⇒ Object



651
652
653
654
655
656
657
658
659
660
661
662
663
664
# File 'lib/rbbt/workflow/step/accessor.rb', line 651

def config(key, *tokens)
  options = tokens.pop if Hash === tokens.last
  options ||= {}

  new_tokens = []
  if workflow
    workflow_name = workflow.to_s
    new_tokens << ("workflow:" << workflow_name)
    new_tokens << ("task:" << workflow_name << "#" << task_name.to_s)
  end
  new_tokens << ("task:" << task_name.to_s)

  Rbbt::Config.get(key, tokens + new_tokens, options)
end

#copy_files_dirObject



115
116
117
118
119
120
121
122
123
124
125
# File 'lib/rbbt/workflow/step.rb', line 115

def copy_files_dir
  if File.symlink?(self.files_dir)
    begin
      realpath = Open.realpath(self.files_dir)
      Open.rm self.files_dir
      Open.cp realpath, self.files_dir
    rescue
      Log.warn "Copy files_dir for #{self.workflow_short_path}: " + $!.message
    end
  end
end

#dependency_checksObject



122
123
124
125
126
127
128
129
130
131
# File 'lib/rbbt/workflow/step/run.rb', line 122

def dependency_checks
  return [] if ENV["RBBT_UPDATE"] != "true"

  rec_dependencies(true).
    reject{|dependency| (defined?(WorkflowRemoteClient) && WorkflowRemoteClient::RemoteStep === dependency) || Open.remote?(dependency.path) }.
    reject{|dependency| dependency.error? }.
    #select{|dependency| Open.exists?(dependency.path) || ((Open.exists?(dependency.info_file) && (dependency.status == :cleaned) || dependency.status == :waiting)) }.
    #select{|dependency| dependency.updatable? }.
    collect{|dependency| Workflow.relocate_dependency(self, dependency)}
end

#dirty?Boolean

Returns:

  • (Boolean)


481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
# File 'lib/rbbt/workflow/step/accessor.rb', line 481

def dirty?
  return true if Open.exists?(pid_file) && ! ( Open.exists?(info_file) || done? )
  return false unless done? || status == :done
  return false unless ENV["RBBT_UPDATE"] == "true"

  status = self.status

  if done? and not (status == :done or status == :ending or status == :producing) and not status == :noinfo
    return true 
  end

  if status == :done and not done?
    return true 
  end

  if dirty_files.any?
    Log.low "Some dirty files found for #{self.path}: #{Misc.fingerprint dirty_files}"
    true
  else
    ! self.updated?
  end
end

#dirty_filesObject



468
469
470
471
472
473
474
475
476
477
478
479
# File 'lib/rbbt/workflow/step/accessor.rb', line 468

def dirty_files
  rec_dependencies = self.rec_dependencies(true)
  return [] if rec_dependencies.empty?
  canfail_paths = self.canfail_paths

  dirty_files = rec_dependencies.reject{|dep|
    (defined?(WorkflowRemoteClient) && WorkflowRemoteClient::RemoteStep === dep) || 
      ! Open.exists?(dep.info_file) ||
      (dep.path && (Open.exists?(dep.path) || Open.remote?(dep.path))) || 
      ((dep.error? || dep.aborted?) && (! dep.recoverable_error? || canfail_paths.include?(dep.path)))
  }
end

#done?Boolean

Returns:

  • (Boolean)


504
505
506
# File 'lib/rbbt/workflow/step/accessor.rb', line 504

def done?
  path and Open.exists? path
end

#dup_inputsObject



61
62
63
64
65
66
67
68
69
70
71
# File 'lib/rbbt/workflow/step/dependencies.rb', line 61

def dup_inputs
  return if @inputs.nil?
  return if @dupped or ENV["RBBT_NO_STREAM"] == 'true'
  return if ComputeDependency === self and self.compute == :produce
  Log.low "Dupping inputs for #{path}"
  dupped_inputs = @inputs.collect do |input|
    Step.dup_stream input
  end
  @inputs.replace dupped_inputs
  @dupped = true
end

#error?Boolean

Returns:

  • (Boolean)


538
539
540
# File 'lib/rbbt/workflow/step/accessor.rb', line 538

def error?
  status == :error
end

#exception(ex, msg = nil) ⇒ Object



415
416
417
418
419
420
421
422
423
424
425
426
427
# File 'lib/rbbt/workflow/step/accessor.rb', line 415

def exception(ex, msg = nil)
  ex_class = ex.class.to_s
  backtrace = ex.backtrace if ex.respond_to?(:backtrace)
  message = ex.message if ex.respond_to?(:message)
  set_info :backtrace, backtrace
  set_info :exception, {:class => ex_class, :message => message, :backtrace => backtrace}
  if msg.nil?
    log :error, "#{ex_class} -- #{message}"
  else
    log :error, "#{msg} -- #{message}"
  end
  self._abort
end

#execute_and_dup(step, dep_step, log = true) ⇒ Object



204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
# File 'lib/rbbt/workflow/step/dependencies.rb', line 204

def execute_and_dup(step, dep_step, log = true)
  dup = step.result.nil?
  execute_dependency(step, log)
  if dup and step.streaming? and not step.result.nil?
    if dep_step[step.path] and dep_step[step.path].length > 1
      stream = step.result
      other_steps = dep_step[step.path]
      return unless other_steps.length > 1
      log_dependency_exec(step, "duplicating #{other_steps.length}") 
      copies = Misc.tee_stream_thread_multiple(stream, other_steps.length)
      copies.extend StreamArray
      step.instance_variable_set("@result", copies)
    end
  end
end

#execute_dependency(dependency, log = true) ⇒ Object



133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
# File 'lib/rbbt/workflow/step/dependencies.rb', line 133

def execute_dependency(dependency, log = true)
  task_name = self.task_name
  canfail_paths = self.canfail_paths
  already_failed = []
  begin

    dependency.resolve_input_steps

    if dependency.done?
      dependency.inputs.each do |v|
        Misc.consume_stream(v) if IO === v
        Misc.consume_stream(TSV.get_stream v) if Step === v and not v.done?  and  v.streaming?
      end
      log_dependency_exec(dependency, :done) if log
      return
    end

    dependency.status_lock.synchronize do
      if dependency.aborted? || (dependency.error? && dependency.recoverable_error? && ! canfail_paths.include?(dependency.path) && ! already_failed.include?(dependency.path)) || (!Open.remote?(dependency.path) && dependency.missing?)
        if dependency.resumable?
          dependency.status = :resume
        else
          Log.warn "Cleaning dep. on exec #{Log.color :blue, dependency.path} (missing: #{dependency.missing?}; error #{dependency.error?})"
          dependency.clean
          already_failed << dependency.path
          raise TryAgain
        end
      end
    end

    if dependency.status == :resume || ! (dependency.started? || dependency.error?)
      log_dependency_exec(dependency, :starting)
      dependency.run(true)
      raise TryAgain
    end

    dependency.grace

    if dependency.error?
      log_dependency_exec(dependency, :error)
      raise DependencyError, [dependency.path, dependency.messages.last] * ": " if dependency.error?
    end

    if dependency.streaming?
      log_dependency_exec(dependency, :streaming) if log
      return
    end

    begin
      log_dependency_exec(dependency, :joining)
      dependency.join
      raise TryAgain unless dependency.done?
    rescue Aborted
      raise TryAgain
    end

  rescue TryAgain
    Log.low "Retrying dep. #{Log.color :yellow, dependency.task_name.to_s} -- [#{dependency.status}] #{(dependency.messages || ["No message"]).last}"
    retry
  rescue Aborted, Interrupt
    Log.error "Aborted dep. #{Log.color :red, dependency.task_name.to_s}"
    raise $!
  rescue Interrupt
    Log.error "Interrupted while in dep. #{Log.color :red, dependency.task_name.to_s}"
    raise $!
  rescue Exception
    Log.error "Exception in dep. #{ Log.color :red, dependency.task_name.to_s } -- #{$!.message}"
    raise $! unless canfail_paths.include? dependency.path
  end
end

#file(name) ⇒ Object



568
569
570
# File 'lib/rbbt/workflow/step/accessor.rb', line 568

def file(name)
  Path.setup(File.join(files_dir, name.to_s), workflow, self)
end

#filesObject



561
562
563
564
565
566
# File 'lib/rbbt/workflow/step/accessor.rb', line 561

def files
  files = Dir.glob(File.join(files_dir, '**', '*')).reject{|path| File.directory? path}.collect do |path| 
    Misc.path_relative_to(files_dir, path) 
  end
  files
end

#files_dirObject

{{{ INFO



553
554
555
# File 'lib/rbbt/workflow/step/accessor.rb', line 553

def files_dir
  @files_dir ||= Step.files_dir path
end

#fork(no_load = false, semaphore = nil) ⇒ Object



497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
# File 'lib/rbbt/workflow/step/run.rb', line 497

def fork(no_load = false, semaphore = nil)
  raise "Can not fork: Step is waiting for proces #{@pid} to finish" if not @pid.nil? and not Process.pid == @pid and Misc.pid_exists?(@pid) and not done? and info[:forked]
  sout, sin = Misc.pipe if no_load == :stream
  @pid = Process.fork do
    Signal.trap(:TERM) do
      raise Aborted, "Recieved TERM Signal on forked process #{Process.pid}"
    end
    sout.close if sout
    Misc.pre_fork
    Open.mkdir File.dirname(path) unless Open.exist?(File.dirname(path))
    Open.write(pid_file, Process.pid.to_s) unless Open.exists?(path) or Open.exists?(pid_file)

    if semaphore
      init_info
      log :queue, "Queued over semaphore: #{semaphore}"
      ret = RbbtSemaphore.wait_semaphore(semaphore)
      raise SemaphoreInterrupted if ret == -1
    end

    begin
      begin
        @forked = true
        res = run no_load
        set_info :forked, true
        if sin
          io = TSV.get_stream res
          if io.respond_to? :setup
            io.setup(sin) 
            sin.pair = io
            io.pair = sin
          end
          begin
            Misc.consume_stream(io, false, sin)
          rescue 
            Log.warn "Could not consume stream (#{io.closed? ? 'closed' : 'open'}) into pipe for forked job: #{self.path}"
            Misc.consume_stream(io) unless io.closed?
          end
        end
      rescue Aborted, Interrupt
        Log.debug{"Forked process aborted: #{path}"}
        log :aborted, "Job aborted (#{Process.pid})"
        raise $!
      rescue Exception
        Log.debug("Exception '#{$!.message}' caught on forked process: #{path}")
        raise $!
      ensure
        join_stream
      end

      begin
        children_pids = info[:children_pids]
        if children_pids
          children_pids.each do |pid|
            if Misc.pid_exists? pid
              begin
                Process.waitpid pid
              rescue Errno::ECHILD
                Log.low "Waiting on #{ pid } failed: #{$!.message}"
              end
            end
          end
          set_info :children_done, Time.now
        end
      rescue Exception
        Log.debug("Exception waiting for children: #{$!.message}")
        RbbtSemaphore.post_semaphore(semaphore) if semaphore
        Kernel.exit! -1
      end
      #set_info :pid, nil
    ensure
      RbbtSemaphore.post_semaphore(semaphore) if semaphore
    end
  end
  sin.close if sin
  @result = sout if sout 
  Process.detach(@pid)
  self
end

#get_exceptionObject



429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
# File 'lib/rbbt/workflow/step/accessor.rb', line 429

def get_exception
  if info[:exception].nil?
    return Aborted if aborted?
    return Exception.new(messages.last) if error?
    Exception.new "" 
  else
    ex_class, ex_message, ex_backtrace = info[:exception].values_at :class, :message, :backtrace
    begin
      klass = Kernel.const_get(ex_class)
      ex = klass.new ex_message
      ex.set_backtrace ex_backtrace unless ex_backtrace.nil? or ex_backtrace.empty?
      ex
    rescue
      Log.exception $!
      Exception.new ex_message
    end
  end
end

#get_streamObject



11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
# File 'lib/rbbt/workflow/step/run.rb', line 11

def get_stream
  @mutex.synchronize do
    Log.low "Getting stream from #{path} #{!@saved_stream} [#{object_id}-#{Misc.fingerprint(@result)}]"
    begin
      if IO === @result 
        return nil if @saved_stream
        @saved_stream = @result 
      elsif StreamArray === @result and @result.any?
        @saved_stream = @result.pop 
      else
        nil
      end
    end
  end
end

#graceObject



684
685
686
687
688
689
# File 'lib/rbbt/workflow/step/run.rb', line 684

def grace
  until done? || result || error? || aborted? || streaming? || waiting? 
    sleep 1 
  end
  self
end

#info(check_lock = true) ⇒ Object



203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
# File 'lib/rbbt/workflow/step/accessor.rb', line 203

def info(check_lock = true)
  return {:status => :noinfo} if info_file.nil? or not Open.exists? info_file
  begin
    Misc.insist do
      begin
        return @info_cache if @info_cache and @info_cache_time and Open.ctime(info_file) < @info_cache_time 
      rescue Exception
        raise $!
      end

      begin
        @info_cache = Misc.insist(3, 1.6, info_file) do
          Misc.insist(2, 1, info_file) do
            Misc.insist(3, 0.2, info_file) do
              raise TryAgain, "Info locked" if check_lock and info_lock.locked?
              info_lock.lock if check_lock and false
              begin
                Open.open(info_file, :mode => 'rb') do |file|
                  Step.load_serialized_info(file)
                end
              ensure
                info_lock.unlock if check_lock and false
              end
            end
          end
        end
        @info_cache_time = Time.now
        @info_cache
      end
    end
  rescue Exception
    Log.debug{"Error loading info file: " + info_file}
    Log.exception $!
    Open.rm info_file
    Misc.sensiblewrite(info_file, Step.serialize_info({:status => :error, :messages => ["Info file lost"]}))
    raise $!
  end
end

#info_fileObject

{{{ INFO



177
178
179
# File 'lib/rbbt/workflow/step/accessor.rb', line 177

def info_file
  @info_file ||= Step.info_file(path)
end

#info_lockObject



185
186
187
188
189
190
191
192
# File 'lib/rbbt/workflow/step/accessor.rb', line 185

def info_lock
  @info_lock = begin
                 path = Persist.persistence_path(info_file + '.lock', {:dir => Step.lock_dir})
                 #Lockfile.new path, :refresh => false, :dont_use_lock_id => true
                 Lockfile.new path
               end if @info_lock.nil?
               @info_lock
end

#init_info(force = false) ⇒ Object



242
243
244
245
246
247
248
249
250
251
# File 'lib/rbbt/workflow/step/accessor.rb', line 242

def init_info(force = false)
  return nil if @exec || info_file.nil? || (Open.exists?(info_file) && ! force)
  Open.lock(info_file, :lock => info_lock) do
    i = {:status => :waiting, :pid => Process.pid, :path => path, :real_inputs => real_inputs}
    i[:dependencies] = dependencies.collect{|dep| [dep.task_name, dep.name, dep.path]} if dependencies
    Misc.sensiblewrite(info_file, Step.serialize_info(i), :force => true, :lock => false)
    @info_cache = IndiferentHash.setup(i)
    @info_cache_time = Time.now
  end
end

#input_checksObject



133
134
135
136
# File 'lib/rbbt/workflow/step/run.rb', line 133

def input_checks
  (inputs.select{|i| Step === i } + inputs.select{|i| Path === i && Step === i.resource}.collect{|i| i.resource})
    #select{|dependency| dependency.updatable? }
end

#input_dependenciesObject



129
130
131
# File 'lib/rbbt/workflow/step/dependencies.rb', line 129

def input_dependencies
  (inputs.flatten.select{|i| Step === i} + inputs.flatten.select{|dep| Path === dep && Step === dep.resource}.collect{|dep| dep.resource})
end

#joinObject



691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
# File 'lib/rbbt/workflow/step/run.rb', line 691

def join

  grace if Open.exists?(info_file) 

  if streaming?
    join_stream 
  end

  return self if not Open.exists? info_file

  return self if info[:joined]

  pid = @pid 

  Misc.insist [0.1, 0.2, 0.5, 1] do
    pid ||= info[:pid]
  end

  begin

    if pid.nil? or Process.pid == pid
      dependencies.each{|dep| dep.join }
    else
      begin
        pid = pid.to_i if String === pid
        Log.debug{"Waiting for pid: #{pid}"}
        Process.waitpid pid
      rescue Errno::ECHILD
        Log.debug{"Process #{ pid } already finished: #{ path }"}
      end if Misc.pid_exists? pid
      pid = nil
      dependencies.each{|dep| dep.join }
    end

    until (Open.exists?(path) && (status == :done || status == :noinfo)) or error? or aborted? or waiting?
      sleep 1
      join_stream if streaming?
    end

    self
  ensure
    begin
      set_info :joined, true 
    rescue
    end if Open.exists?(info_file) && writable?
    @result = nil
  end
end

#join_streamObject



663
664
665
666
667
668
669
670
671
672
673
674
675
# File 'lib/rbbt/workflow/step/run.rb', line 663

def join_stream
  stream = get_stream if @result
  @result = nil
  if stream
    begin
      Misc.consume_stream stream 
      stream.join if stream.respond_to? :join
    rescue Exception
      stream.abort $!
      self._abort
    end
  end
end

#kill_childrenObject



197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
# File 'lib/rbbt/workflow/step/run.rb', line 197

def kill_children
  begin
    children_pids = info[:children_pids]
    if children_pids and children_pids.any?
      Log.medium("Killing children: #{ children_pids * ", " }")
      children_pids.each do |pid|
        Log.medium("Killing child #{ pid }")
        begin
          Process.kill "TERM", pid.to_i
        rescue Exception
          Log.medium("Exception killing child #{ pid }: #{$!.message}")
        end
      end
    end
  rescue
    Log.medium("Exception finding children")
  end
end

#knowledge_base(organism = nil) ⇒ Object



729
730
731
732
733
734
# File 'lib/rbbt/workflow/step/accessor.rb', line 729

def knowledge_base(organism = nil)
  @_kb ||= begin
             kb_dir = self.file('knowledge_base')
             KnowledgeBase.new kb_dir, organism
           end
end

#loadObject



395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
# File 'lib/rbbt/workflow/step.rb', line 395

def load
  res = begin
          @result = nil if IO === @result && @result.closed?
          if @result && @path != @result
            res = @result
          else
            join if not done?
            res = @path.exists? ? Persist.load_file(@path, result_type) : exec
          end

          if result_description
            entity_info = info.dup
            entity_info.merge! info[:inputs] if info[:inputs]
            res = prepare_result res, result_description, entity_info 
          end

          res
        rescue IOError
          if @result
            @result = nil
            retry
          end
          raise $!
        ensure
          @result = nil if IO === @result
        end

  res
end

#load_dependencies_from_infoObject



89
90
91
92
93
94
95
96
97
98
99
100
101
102
# File 'lib/rbbt/workflow/step.rb', line 89

def load_dependencies_from_info
  relocated = nil
  @dependencies = (self.info[:dependencies] || []).collect do |task,name,dep_path|
    if Open.exists?(dep_path) || Open.exists?(dep_path + '.info')
      Workflow._load_step dep_path
    else
      next if FalseClass === relocated
      new_path = Workflow.relocate(path, dep_path)
      relocated = true if Open.exists?(new_path) || Open.exists?(new_path + '.info')
      Workflow._load_step new_path
    end
  end.compact
  @relocated = relocated
end

#load_file(name, type = nil, options = {}) ⇒ Object



588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
# File 'lib/rbbt/workflow/step/accessor.rb', line 588

def load_file(name, type = nil, options = {})
  if type.nil? and name =~ /.*\.(\w+)$/
    extension = name.match(/.*\.(\w+)$/)[1]
    case extension
    when "tc"
      type = :tc
    when "tsv"
      type = :tsv
    when "list", "ary", "array"
      type = :array
    when "yaml"
      type = :yaml
    when "marshal"
      type = :marshal
    else
      type = :other
    end
  else
    type ||= :other
  end

  case type.to_sym
  when :tc
    Persist.open_tokyocabinet(file(name), false)
  when :tsv
    TSV.open Open.open(file(name)), options
  when :array
    #Open.read(file(name)).split /\n|,\s*/
    Open.read(file(name)).split "\n"
  when :yaml
    YAML.load(Open.open(file(name)))
  when :marshal
    Marshal.load(Open.open(file(name)))
  else
    Open.read(file(name))
  end
end

#load_inputs_from_infoObject



75
76
77
78
79
80
81
82
83
84
85
86
87
# File 'lib/rbbt/workflow/step.rb', line 75

def load_inputs_from_info
  if info[:inputs]
    info_inputs = info[:inputs]
    if task && task.respond_to?(:inputs) && task.inputs
      IndiferentHash.setup info_inputs
      @inputs = NamedArray.setup info_inputs.values_at(*task.inputs.collect{|name| name.to_s}), task.inputs
    else
      @inputs = NamedArray.setup info_inputs.values, info_inputs.keys
    end
  else
    nil
  end
end

#log(status, message = nil, &block) ⇒ Object



407
408
409
410
411
412
413
# File 'lib/rbbt/workflow/step/accessor.rb', line 407

def log(status, message = nil, &block)
  self.status = status
  if message
    self.message Log.uncolor(message)
  end
  Step.log(status, message, path, &block)
end

#log_dependency_exec(dependency, action) ⇒ Object



114
115
116
117
118
119
120
121
122
123
124
125
126
127
# File 'lib/rbbt/workflow/step/dependencies.rb', line 114

def log_dependency_exec(dependency, action)
  task_name = self.task_name

  str = Log.color(:reset, "")
  str << Log.color(:yellow, task_name.to_s || "") 
  str << " "
  str << Log.color(:magenta, action.to_s)
  str << " "
  str << Log.color(:yellow, dependency.task_name.to_s || "")
  str << " -- "
  str << "#{Log.color :blue, dependency.path}"

  Log.info str
end

#log_progress(status, options = {}, &block) ⇒ Object



380
381
382
# File 'lib/rbbt/workflow/step/accessor.rb', line 380

def log_progress(status, options = {}, &block)
  Step.log_progress(status, options, file(:progress), &block)
end

#merge_info(hash) ⇒ Object



268
269
270
271
272
273
274
275
276
277
278
279
280
281
# File 'lib/rbbt/workflow/step/accessor.rb', line 268

def merge_info(hash)
  return nil if @exec or info_file.nil?
  return nil if ! writable?
  value = Annotated.purge value if defined? Annotated
  Open.lock(info_file, :lock => info_lock) do
    i = info(false)
    i.merge! hash
    dump = Step.serialize_info(i)
    @info_cache = IndiferentHash.setup(i)
    Misc.sensiblewrite(info_file, dump, :force => true, :lock => false) if Open.exists?(info_file)
    @info_cache_time = Time.now
    value
  end
end

#message(message) ⇒ Object



304
305
306
307
# File 'lib/rbbt/workflow/step/accessor.rb', line 304

def message(message)
  message = Log.uncolor(message)
  set_info(:messages, (messages || []) << message)
end

#messagesObject



296
297
298
299
300
301
302
# File 'lib/rbbt/workflow/step/accessor.rb', line 296

def messages
  if messages = info[:messages]
    messages
  else
    set_info(:messages, []) if self.respond_to?(:set_info)
  end
end

#missing?Boolean

Returns:

  • (Boolean)


534
535
536
# File 'lib/rbbt/workflow/step/accessor.rb', line 534

def missing?
  status == :done && ! Open.exists?(path)
end

#monitor_stream(stream, options = {}, &block) ⇒ Object



677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
# File 'lib/rbbt/workflow/step/accessor.rb', line 677

def monitor_stream(stream, options = {}, &block)
  case options[:bar] 
  when TrueClass
    bar = progress_bar 
  when Hash
    bar = progress_bar options[:bar]
  when Numeric
    bar = progress_bar :max => options[:bar]
  else
    bar = options[:bar]
  end

  out = if bar.nil?
          Misc.line_monitor_stream stream, &block
        elsif (block.nil? || block.arity == 0)
          Misc.line_monitor_stream stream do
            bar.tick
          end
        elsif block.arity == 1
          Misc.line_monitor_stream stream do |line|
            bar.tick
            block.call line
          end
        elsif block.arity == 2
          Misc.line_monitor_stream stream do |line|
            block.call line, bar
          end
        end

  ConcurrentStream.setup(out, :abort_callback => Proc.new{
    Log::ProgressBar.remove_bar(bar, true) if bar
  }, :callback => Proc.new{
    Log::ProgressBar.remove_bar(bar) if bar
  })

  bgzip = (options[:compress] || options[:gzip]).to_s == 'bgzip'
  bgzip = true if options[:bgzip]

  gzip = true if options[:compress] || options[:gzip]
  if bgzip
    Open.bgzip(out)
  elsif gzip
    Open.gzip(out)
  else
    out
  end
end

#nameObject



153
154
155
# File 'lib/rbbt/workflow/step/accessor.rb', line 153

def name
  @name ||= path.sub(/.*\/#{Regexp.quote task_name.to_s}\/(.*)/, '\1')
end

#noinfo?Boolean

Returns:

  • (Boolean)


512
513
514
# File 'lib/rbbt/workflow/step/accessor.rb', line 512

def noinfo?
  status == :noinfo
end

#nopid?Boolean

Returns:

  • (Boolean)


542
543
544
# File 'lib/rbbt/workflow/step/accessor.rb', line 542

def nopid?
  ! Open.exists?(pid_file) && ! (status.nil? || status == :aborted || status == :done || status == :error || status == :cleaned)
end

#out_of_dateObject



151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
# File 'lib/rbbt/workflow/step/run.rb', line 151

def out_of_date

  checks = self.checks
  return [] if checks.empty?
  outdated_time  = []
  outdated_dep  = []
  canfail_paths = self.canfail_paths
  this_mtime = Open.mtime(self.path) if Open.exists?(self.path)

  #outdated_time = checks.select{|dep| dep.updatable? && dep.done? && Persist.newer?(path, dep.path) }
  outdated_time = checks.select{|dep| dep.done? && Persist.newer?(path, dep.path) }
  outdated_dep = checks.reject{|dep| dep.done? || (dep.error? && ! dep.recoverable_error? && canfail_paths.include?(dep.path)) }

  #checks.each do |dep| 
  #  next unless dep.updatable?
  #  dep_done = dep.done?

  #  begin
  #    if this_mtime && dep_done && Open.exists?(dep.path) && (Open.mtime(dep.path) > this_mtime + 1)
  #      outdated_time << dep
  #    end
  #  rescue
  #  end

  #  # Is this pointless? this would mean some dep got updated after a later
  #  # dep but but before this one.
  #  #if (! dep.done? && ! canfail_paths.include?(dep.path)) || ! dep.updated?

  #  if (! dep_done && ! canfail_paths.include?(dep.path))
  #    outdated_dep << dep
  #  end
  #end

  Log.high "Some newer files found: #{Misc.fingerprint outdated_time}" if outdated_time.any?
  Log.high "Some outdated files found: #{Misc.fingerprint outdated_dep}" if outdated_dep.any?

  outdated_time + outdated_dep
end

#persist_checksObject



142
143
144
145
146
147
148
149
# File 'lib/rbbt/workflow/step/run.rb', line 142

def persist_checks
  canfail_paths = self.canfail_paths
  checks.collect do |dep| 
    path = dep.path
    next if ! dep.done? && canfail_paths.include?(path)
    path 
  end.compact
end

#pid_fileObject



181
182
183
# File 'lib/rbbt/workflow/step/accessor.rb', line 181

def pid_file
  @pid_file ||= Step.pid_file(path)
end

#prepare_result(value, description = nil, entity_info = nil) ⇒ Object



268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
# File 'lib/rbbt/workflow/step.rb', line 268

def prepare_result(value, description = nil, entity_info = nil)
  res = case 
  when IO === value
    begin
      res = case result_type
            when :array
              array = []
              while line = value.gets
                array << line.chomp
              end
              array
            when :tsv
              begin
                TSV.open(value)
              rescue IOError
                TSV.setup({})
              end
            else
              value.read
            end
      value.join if value.respond_to? :join
      res
    rescue Exception
      value.abort if value.respond_to? :abort
      self.abort
      raise $!
    end
  when (not defined? Entity or description.nil? or not Entity.formats.include? description)
    value
  when (Annotated === value and info.empty?)
    value
  when Annotated === value
    annotations = value.annotations
    entity_info ||= begin 
                      entity_info = info.dup
                      entity_info.merge! info[:inputs] if info[:inputs]
                      entity_info
                    end
    entity_info.each do |k,v|
      value.send("#{h}=", v) if annotations.include? k
    end
                      
    value
  else
    entity_info ||= begin 
                      entity_info = info.dup
                      entity_info.merge! info[:inputs] if info[:inputs]
                      entity_info
                    end
    Entity.formats[description].setup(value, entity_info.merge(:format => description))
  end

  if Annotated === res
    dep_hash = nil
    res.annotations.each do |a|
      a = a.to_s
      varname = "@" + a
      next unless res.instance_variable_get(varname).nil? 

      dep_hash ||= begin
                     h = {}
                     rec_dependencies.each{|dep| h[dep.task_name.to_s] ||= dep }
                     h
                   end
      dep = dep_hash[a]
      next if dep.nil?
      res.send(a.to_s+"=", dep.load)
    end 
  end

  res
end

#produce(force = false, dofork = false) ⇒ Object



459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
# File 'lib/rbbt/workflow/step/run.rb', line 459

def produce(force=false, dofork=false)
  return self if done? and not dirty?

  self.status_lock.synchronize do
    if error? || aborted? || stalled?
      if stalled?
        Log.warn "Aborting stalled job #{self.path}"
        abort
      end
      if force or aborted? or recoverable_error?
        clean
      else
        e = get_exception
        if e
          Log.error "Raising exception in produced job #{self.path}: #{e.message}" 
          raise e
        else
          raise "Error in job: #{self.path}"
        end
      end
    end
  end

  update if done?

  if dofork
    fork(true) unless started?

    join unless done?
  else
    run(true) unless started?

    join unless done?
  end

  self
end

#progress_bar(msg = "Progress", options = nil) ⇒ Object



384
385
386
387
388
389
390
391
392
393
# File 'lib/rbbt/workflow/step/accessor.rb', line 384

def progress_bar(msg = "Progress", options = nil)
  if Hash === msg and options.nil?
    options = msg
    msg = nil
  end
  options = {} if options.nil?

  max = options[:max]
  Log::ProgressBar.new_bar(max, {:desc => msg, :file => file(:progress)}.merge(options))
end

#provenanceObject



626
627
628
629
630
631
632
633
634
635
636
637
# File 'lib/rbbt/workflow/step/accessor.rb', line 626

def provenance
  provenance = {}
  dependencies.each do |dep|
    next unless dep.path.exists?
    if Open.exists? dep.info_file
      provenance[dep.path] = dep.provenance if Open.exists? dep.path
    else
      provenance[dep.path] = nil
    end
  end
  {:inputs => info[:inputs], :provenance => provenance}
end

#provenance_pathsObject



639
640
641
642
643
644
645
# File 'lib/rbbt/workflow/step/accessor.rb', line 639

def provenance_paths
  provenance = {}
  dependencies.each do |dep|
    provenance[dep.path] = dep.provenance_paths if Open.exists? dep.path
  end
  provenance
end

#rec_accessObject



670
671
672
673
674
675
# File 'lib/rbbt/workflow/step/accessor.rb', line 670

def rec_access
  access
  rec_dependencies.each do |dep|
    dep.access
  end
end

#rec_dependencies(connected = false, seen = []) ⇒ Object

connected = true means that dependency searching ends when a result is done but dependencies are absent, meanining that the file could have been dropped in



475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
# File 'lib/rbbt/workflow/step.rb', line 475

def rec_dependencies(connected = false, seen = [])

  # A step result with no info_file means that it was manually
  # placed. In that case, do not consider its dependencies
  return [] if ! (defined? WorkflowRemoteClient && WorkflowRemoteClient::RemoteStep === self) && ! Open.exists?(self.info_file) && Open.exists?(self.path.to_s) 

  return [] if dependencies.nil? or dependencies.empty?

  new_dependencies = []
  archived_deps = self.info[:archived_info] ? self.info[:archived_info].keys : []

  dependencies.each{|step| 
    #next if self.done? && Open.exists?(info_file) && info[:dependencies] && info[:dependencies].select{|task,name,path| path == step.path }.empty?
    next if archived_deps.include? step.path
    next if seen.include? step.path
    next if self.done? && connected && ! updatable?

    r = step.rec_dependencies(connected, new_dependencies.collect{|d| d.path})
    new_dependencies.concat r
    new_dependencies << step
  }
  new_dependencies.uniq
end

#recoverable_error?Boolean

Returns:

  • (Boolean)


448
449
450
451
452
453
454
455
456
457
458
# File 'lib/rbbt/workflow/step/accessor.rb', line 448

def recoverable_error?
  return true if aborted?
  return false unless error?
  begin
    return true unless info[:exception]
    klass = Kernel.const_get(info[:exception][:class])
    ! (klass <= RbbtException)
  rescue Exception
    true
  end
end

#recursive_cleanObject



503
504
505
506
507
508
509
# File 'lib/rbbt/workflow/step.rb', line 503

def recursive_clean
  dependencies.each do |step| 
    step.recursive_clean 
  end
  clean if Open.exists?(self.info_file)
  self
end

#recursive_inputsObject



171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
# File 'lib/rbbt/workflow/step.rb', line 171

def recursive_inputs
  if NamedArray === inputs
    i = {}
    inputs.zip(inputs.fields).each do |v,f|
      i[f] = v
    end
  else
    i = {}
  end
  rec_dependencies.each do |dep|
    next unless NamedArray === dep.inputs

    dep.inputs.zip(dep.inputs.fields).each do |v,f|
      if i.include?(f) && i[f] != v
        Log.debug "Conflict in #{ f }: #{[Misc.fingerprint(i[f]), Misc.fingerprint(v)] * " <-> "}"
      else 
        i[f] = v
      end
    end

    dep.archived_inputs.each do |k,v|
      i[k] = v unless i.include? k
    end
  end

  self.archived_inputs.each do |k,v|
    i[k] = v unless i.include? k
  end

  #dependencies.each do |dep|
  #  di = dep.recursive_inputs
  #  next unless NamedArray === di
  #  di.fields.zip(di).each do |k,v|
  #    i[k] = v unless i.include? k
  #  end
  #end
  
  v = i.values
  NamedArray.setup v, i.keys 
  v
end

#relay_log(step) ⇒ Object



235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
# File 'lib/rbbt/workflow/step.rb', line 235

def relay_log(step)
  return self unless Task === self.task and not self.task.name.nil?
  if not self.respond_to? :original_log
    class << self
      attr_accessor :relay_step
      alias original_log log 
      def log(status, message = nil)
        self.status = status
        message Log.uncolor message
        relay_step.log([task.name.to_s, status.to_s] * ">", message.nil? ? nil : message ) unless (relay_step.done? or relay_step.error? or relay_step.aborted?)
      end
    end
  end
  @relay_step = step
  self
end

#relocated?Boolean

Returns:

  • (Boolean)


725
726
727
# File 'lib/rbbt/workflow/step/accessor.rb', line 725

def relocated?
  done? && info[:path] && info[:path] != path
end

#resolve_input_stepsObject



27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
# File 'lib/rbbt/workflow/step/run.rb', line 27

def resolve_input_steps
  step = false
  pos = 0

  input_options = Workflow === workflow ? workflow.task_info(task_name)[:input_options] : {}
  new_inputs = inputs.collect do |i| 
    begin
      if Step === i
        if i.error?
          e = i.get_exception
          if e
            raise e
          else
            raise DependencyError, "Error in dep. #{Log.blue e.path}"
          end
        end
        step = true
        i.produce unless i.done? || i.error? || i.started?
        if i.done?
          if (task.input_options[task.inputs[pos]] || {})[:stream]
            TSV.get_stream i
          else
            if (task.input_options[task.inputs[pos]] || {})[:nofile]
              i.path
            else
              i.load
            end
          end
        elsif i.streaming? and (task.input_options[task.inputs[pos]] || {})[:stream]
          TSV.get_stream i
        else
          i.join
          if (task.input_options[task.inputs[pos]] || {})[:stream]
            TSV.get_stream i
          else
            if (task.input_options[task.inputs[pos]] || {})[:nofile]
              i.path
            else
              i.load
            end
          end
        end
      else
        i
      end
    ensure
      pos += 1
    end
  end
  @inputs.replace new_inputs if step
end

#result_descriptionObject



260
261
262
263
264
265
266
# File 'lib/rbbt/workflow/step.rb', line 260

def result_description
  @result_description ||= if @task.nil?
                     info[:result_description]
                   else
                     @task.result_description
                   end
end

#result_typeObject



252
253
254
255
256
257
258
# File 'lib/rbbt/workflow/step.rb', line 252

def result_type
  @result_type ||= if @task.nil?
                     info[:result_type] || :binary
                   else
                     @task.result_type || info[:result_type] || :string
                   end
end

#resumable?Boolean

Returns:

  • (Boolean)


647
648
649
# File 'lib/rbbt/workflow/step/accessor.rb', line 647

def resumable?
  task && task.resumable
end

#rewind_inputsObject



79
80
81
82
83
84
85
86
87
88
89
90
91
# File 'lib/rbbt/workflow/step/run.rb', line 79

def rewind_inputs
  return if @inputs.nil?
  Log.debug "Rewinding inputs for #{path}"
  @inputs.each do |input|
    next unless input.respond_to? :rewind
    begin
      input.rewind
      input.first_line = nil if TSV::Parser === input
      Log.debug "Rewinded #{Misc.fingerprint input}"
    rescue
    end
  end
end

#run(no_load = false) ⇒ Object



216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
# File 'lib/rbbt/workflow/step/run.rb', line 216

def run(no_load = false)
  result = nil

  begin
    time_elapsed = total_time_elapsed = nil
    res = @mutex.synchronize do
      no_load = :stream if no_load

      Open.write(pid_file, Process.pid.to_s) unless Open.exists?(path) or Open.exists?(pid_file)
      result_type = @task.result_type if @task
      result_type = info[:result_type] if result_type.nil?
      result = Persist.persist "Job", result_type, :file => path, :check => persist_checks, :no_load => no_load do 
        if Step === Step.log_relay_step and not self == Step.log_relay_step
          relay_log(Step.log_relay_step) unless self.respond_to? :relay_step and self.relay_step
        end

        Open.write(pid_file, Process.pid.to_s) unless Open.exists? pid_file

        @exec = false
        init_info(true)

        log :setup, "#{Log.color :green, "Setup"} step #{Log.color :yellow, task.name.to_s || ""}"

        merge_info({
          :issued => (issue_time = Time.now),
          :name => name,
          :pid => Process.pid.to_s,
          :pid_hostname => Socket.gethostname,
          :clean_name => clean_name,
          :workflow => (@workflow || @task.workflow).to_s,
          :task_name => @task.name,
          :result_type => @task.result_type,
          :result_description => @task.result_description,
          :dependencies => dependencies.collect{|dep| [dep.task_name, dep.name, dep.path]},
          :versions => Rbbt.versions
        })

        new_inputs = []
        @inputs.each_with_index do |input,i|
          name = @task.inputs[i]
          type = @task.input_types[name]

          if type == :directory
            directory_inputs = file('directory_inputs')
            input_source = directory_inputs['.source'][name].find
            input_dir = directory_inputs[name].find

            case input
            when Path
              if input.directory?
                new_inputs << input
              else
                input.open do |io|
                  begin
                    Misc.untar(io, input_source)
                  rescue
                    raise ParameterException, "Error unpackaging tar directory input '#{name}':\n\n#{$!.message}"
                  end
                end
                tar_1 = input_source.glob("*")
                raise ParameterException, "When using tar.gz files for directories, the directory must be the single first level entry" if tar_1.length != 1
                FileUtils.ln_s Misc.path_relative_to(directory_inputs, tar_1.first), input_dir
                new_inputs << input_dir
              end
            when File, IO, Tempfile
              begin
                Misc.untar(Open.gunzip(input), input_source)
              rescue
                raise ParameterException, "Error unpackaging tar directory input '#{name}':\n\n#{$!.message}"
              end
              tar_1 = input_source.glob("*")
              raise ParameterException, "When using tar.gz files for directories, the directory must be the single first level entry" if tar_1.length != 1
              FileUtils.ln_s Misc.path_relative_to(directory_inputs, tar_1.first), input_dir
              new_inputs << input_dir
            else
              raise ParameterException, "Format of directory input '#{name}' not understood: #{Misc.fingerprint input}"
            end
          else
            new_inputs << input
          end
        end if @inputs

        @inputs = new_inputs if @inputs

        if @inputs && ! task.inputs.nil?
          info_inputs = @inputs.collect do |i| 
            if Path === i 
              i.to_s
            else 
              i 
            end
          end
          set_info :inputs, Misc.remove_long_items(Misc.zip2hash(task.inputs, info_inputs)) 
        end

        begin
          run_dependencies
        rescue Exception
          Open.rm pid_file if Open.exists?(pid_file)
          stop_dependencies
          raise $!
        end

        set_info :started, (start_time = Time.now)
        log :started, "Starting step #{Log.color :yellow, task.name.to_s || ""}"

        config_keys_pre = Rbbt::Config::GOT_KEYS.dup
        begin

          result = _exec
        rescue Aborted, Interrupt
          log(:aborted, "Aborted")
          raise $!
        rescue Exception
          backtrace = $!.backtrace

          # HACK: This fixes an strange behaviour in 1.9.3 where some
          # backtrace strings are coded in ASCII-8BIT
          backtrace = backtrace.collect{|l| l.dup.force_encoding("UTF-8")} if String.instance_methods.include? :force_encoding
          set_info :backtrace, backtrace 
          log(:error, "#{$!.class}: #{$!.message}")
          stop_dependencies
          raise $!
        end

        if not no_load or ENV["RBBT_NO_STREAM"] == "true" 
          result = prepare_result result, @task.description, info if IO === result 
          result = prepare_result result.stream, @task.description, info if TSV::Dumper === result 
        end

        stream = case result
                 when IO
                   result
                 when TSV::Dumper
                   result.stream
                 end

        if stream
          log :streaming, "Streaming step #{Log.color :yellow, task.name.to_s || ""}"

          callback = Proc.new do
            if AbortedStream === stream
              if stream.exception
                raise stream.exception 
              else
                raise Aborted
              end
            end
            begin
              status = self.status
              if status != :done and status != :error and status != :aborted
                Misc.insist do
                  merge_info({
                    :done => (done_time = Time.now),
                    :total_time_elapsed => (total_time_elapsed = done_time - issue_time),
                    :time_elapsed => (time_elapsed = done_time - start_time),
                    :versions => Rbbt.versions
                  })
                  log :done, "Completed step #{Log.color :yellow, task.name.to_s || ""} in #{time_elapsed.to_i}+#{(total_time_elapsed - time_elapsed).to_i} sec."
                end
              end
            rescue
              Log.exception $!
            ensure
              Step.purge_stream_cache
              Open.rm pid_file if Open.exist?(pid_file)
            end
          end

          abort_callback = Proc.new do |exception|
            begin
              if exception
                self.exception exception
              else
                log :aborted, "#{Log.color :red, "Aborted"} step #{Log.color :yellow, task.name.to_s || ""}" if status == :streaming
              end
              _clean_finished
            rescue
              stop_dependencies
              Open.rm pid_file if Open.exist?(pid_file)
            end
          end

          ConcurrentStream.setup stream, :callback => callback, :abort_callback => abort_callback

          if AbortedStream === stream 
            exception = stream.exception || Aborted.new("Aborted stream: #{Misc.fingerprint stream}")
            self.exception exception
            _clean_finished
            raise exception
          end
        else
          merge_info({
            :done => (done_time = Time.now),
            :total_time_elapsed => (total_time_elapsed = done_time - issue_time),
            :time_elapsed => (time_elapsed = done_time - start_time),
            :versions => Rbbt.versions
          })
          log :ending
          Step.purge_stream_cache
          Open.rm pid_file if Open.exist?(pid_file)
        end

        set_info :dependencies, dependencies.collect{|dep| [dep.task_name, dep.name, dep.path]}

        config_keys = Rbbt::Config::GOT_KEYS[config_keys_pre.length..-1]
        set_info :config_keys, config_keys

        if result.nil? && File.exists?(self.tmp_path) && ! File.exists?(self.path)
          Open.mv self.tmp_path, self.path
        end
        result
      end # END PERSIST
      log :done, "Completed step #{Log.color :yellow, task.name.to_s || ""} in #{time_elapsed.to_i}+#{(total_time_elapsed - time_elapsed).to_i} sec." unless stream or time_elapsed.nil?

      if no_load
        @result ||= result
        self
      else
        Step.purge_stream_cache
        @result = prepare_result result, @task.result_description
      end
    end # END SYNC
    res
  rescue DependencyError
    exception $!
  rescue LockInterrupted
    raise $!
  rescue Aborted, Interrupt
    abort
    stop_dependencies
    raise $!
  rescue Exception
    exception $!
    stop_dependencies
    raise $!
  ensure 
    no_load = false unless IO === result
    Open.rm pid_file if Open.exist?(pid_file) unless no_load
    #set_info :pid, nil unless no_load
  end
end

#run_compute_dependencies(type, list, dep_step = {}) ⇒ Object



220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
# File 'lib/rbbt/workflow/step/dependencies.rb', line 220

def run_compute_dependencies(type, list, dep_step = {})
  if Array === type
    type, *rest = type
  end

  canfail = (rest && rest.include?(:canfail)) || type == :canfail

  case type
  when :canfail
    list.each do |dep|
      begin
        dep.produce
      rescue RbbtException
        Log.warn "Allowing failing of #{dep.path}: #{dep.messages.last if dep.messages}"
      end
      nil
    end
  when :produce, :no_dup
    list.each do |step|
      Misc.insist do
        begin
          step.produce
        rescue RbbtException
          raise $! unless canfail || step.canfail?
        rescue Exception
          step.exception $!
          if step.recoverable_error?
            raise $!
          else
            raise StopInsist.new($!)
          end
        end
      end
    end
    nil
  when :bootstrap
    cpus = rest.nil? ? nil : rest.first 

    cpus = config('dep_cpus', 'bootstrap', :default => [5, list.length / 2].min) if cpus.nil? || cpus.to_i == 0

    respawn = rest && rest.include?(:respawn)
    respawn = false if rest && rest.include?(:norespawn)
    respawn = rest && rest.include?(:always_respawn)
    respawn = :always if respawn.nil?

    Misc.bootstrap(list, cpus, :bar => "Bootstrapping dependencies for #{self.short_path} [#{cpus}]", :respawn => respawn) do |dep|
      begin
        Signal.trap(:INT) do
          dep.abort
          raise Aborted
        end

        Misc.insist do
          begin
            dep.produce 
            Log.warn "Error in bootstrap dependency #{dep.path}: #{dep.messages.last}" if dep.error? or dep.aborted?

          rescue Aborted
            ex = $!
            begin
              dep.abort
              Log.warn "Aborted bootstrap dependency #{dep.path}: #{dep.messages.last}" if dep.error? or dep.aborted?
            rescue
            end
            raise StopInsist.new(ex)

          rescue RbbtException
            if canfail || dep.canfail?
              Log.warn "Allowing failing of #{dep.path}: #{dep.messages.last}"
            else
              Log.warn "NOT Allowing failing of #{dep.path}: #{dep.messages.last}"
              dep.exception $!
              if dep.recoverable_error?
                begin
                  dep.abort
                rescue
                end
                raise $!
              else
                raise StopInsist.new($!)
              end
            end
          end
        end
      rescue
        dep.abort
        raise $!
      end
      nil
    end
  else
    list.each do |step|
      execute_and_dup(step, dep_step, false)
    end
  end
end

#run_dependenciesObject



345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
# File 'lib/rbbt/workflow/step/dependencies.rb', line 345

def run_dependencies
  dep_step = {}

  rec_dependencies = self.rec_dependencies(true)

  return if rec_dependencies.empty?

  all_deps = rec_dependencies + [self]

  compute_deps = rec_dependencies.collect do |dep|
    next unless ComputeDependency === dep
    dep.rec_dependencies + dep.inputs.flatten.select{|i| Step === i}
  end.compact.flatten.uniq

  canfail_paths = self.canfail_paths

  seen_paths = Set.new
  all_deps.uniq.each do |step|
    next if seen_paths.include? step.path
    seen_paths << step.path
    begin
      Step.prepare_for_execution(step) unless step == self
    rescue DependencyError
      raise $! unless canfail_paths.include? step.path
    end
    next unless step.dependencies and step.dependencies.any?
    (step.dependencies + step.input_dependencies).each do |step_dep|
      next if step_dep.done? or step_dep.running? or (ComputeDependency === step_dep and (step_dep.compute == :nodup or step_dep.compute == :ignore))
      dep_step[step_dep.path] ||= []
      dep_step[step_dep.path] << step
    end
  end

  produced = []
  (dependencies + input_dependencies).each do |dep|
    next unless ComputeDependency === dep
    if dep.compute == :produce
      dep.produce 
      produced << dep.path
    end
  end

  self.dup_inputs

  required_dep_paths = []
  dep_step.each do |path,list|
    #required_dep_paths << path if list.length > 1
    required_dep_paths << path if (list & dependencies).any?
  end

  required_dep_paths.concat dependencies.collect{|dep| dep.path}

  required_dep_paths.concat input_dependencies.collect{|dep| dep.path}

  required_dep_paths.concat(dependencies.collect do |dep| 
    [dep.path] + dep.input_dependencies
  end.flatten)

  log :dependencies, "Dependencies for step #{Log.color :yellow, task.name.to_s || ""}"

  pre_deps = []
  compute_pre_deps = {}
  last_deps = []
  compute_last_deps = {}
  seen_paths = Set.new
  rec_dependencies.uniq.each do |step| 
    next if seen_paths.include? step.path
    seen_paths << step.path
    next unless required_dep_paths.include? step.path
    if step.inputs.flatten.select{|i| Step === i}.any?
      if ComputeDependency === step
        next if produced.include? step.path 
        compute_last_deps[step.compute] ||= []
        compute_last_deps[step.compute] << step
      else
        last_deps << step
      end
    else
      if ComputeDependency === step
        next if produced.include? step.path 
        compute_pre_deps[step.compute] ||= []
        compute_pre_deps[step.compute] << step
      else
        pre_deps << step #if dependencies.include?(step)
      end
    end
  end

  Log.medium "Computing pre dependencies: #{Misc.fingerprint(compute_pre_deps)} - #{Log.color :blue, self.path}" if compute_pre_deps.any?
  compute_pre_deps.each do |type,list|
    run_compute_dependencies(type, list, dep_step)
  end

  Log.medium "Processing pre dependencies: #{Misc.fingerprint(pre_deps)} - #{Log.color :blue, self.path}" if pre_deps.any?
  pre_deps.each do |step|
    next if compute_deps.include? step
    begin
      execute_and_dup(step, dep_step, false)
    rescue Exception
      raise $! unless canfail_paths.include?(step.path)
    end
  end

  Log.medium "Processing last dependencies: #{Misc.fingerprint(last_deps)} - #{Log.color :blue, self.path}" if last_deps.any?
  last_deps.each do |step|
    next if compute_deps.include? step
    begin Exception
      execute_and_dup(step, dep_step) 
    rescue 
      raise $! unless canfail_paths.include? step.path
    end
  end

  Log.medium "Computing last dependencies: #{Misc.fingerprint(compute_last_deps)} - #{Log.color :blue, self.path}" if compute_last_deps.any?
  compute_last_deps.each do |type,list|
    run_compute_dependencies(type, list, dep_step)
  end

  dangling_deps = all_deps.reject{|dep| dep.done? || canfail_paths.include?(dep.path) }.
    select{|dep| dep.waiting? }

  Log.medium "Aborting (actually not) waiting dangling dependencies #{Misc.fingerprint dangling_deps}" if dangling_deps.any?
  #dangling_deps.each{|dep| dep.abort }

end

#running?Boolean

Returns:

  • (Boolean)


516
517
518
519
520
521
522
523
524
525
526
527
528
# File 'lib/rbbt/workflow/step/accessor.rb', line 516

def running? 
  return false if ! (started? || status == :ending)
  return nil unless Open.exist?(self.pid_file)
  pid = Open.read(self.pid_file).to_i

  return false if done? or error? or aborted? 

  if Misc.pid_exists?(pid) 
    pid
  else
    done? or error? or aborted? 
  end
end

#save_file(name, content) ⇒ Object



572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
# File 'lib/rbbt/workflow/step/accessor.rb', line 572

def save_file(name, content)
  content = case
            when String === content
              content
            when Array === content
              content * "\n"
            when TSV === content
              content.to_s
            when Hash === content
              content.collect{|*p| p * "\t"} * "\n"
            else
              content.to_s
            end
  Open.write(file(name), content)
end

#set_info(key, value) ⇒ Object



253
254
255
256
257
258
259
260
261
262
263
264
265
266
# File 'lib/rbbt/workflow/step/accessor.rb', line 253

def set_info(key, value)
  return nil if @exec or info_file.nil?
  return nil if ! writable?
  value = Annotated.purge value if defined? Annotated
  Open.lock(info_file, :lock => info_lock) do
    i = info(false).dup
    i[key] = value 
    dump = Step.serialize_info(i)
    @info_cache = IndiferentHash.setup(i)
    Misc.sensiblewrite(info_file, dump, :force => true, :lock => false) if Open.exists?(info_file)
    @info_cache_time = Time.now
    value
  end
end

#short_pathObject



158
159
160
# File 'lib/rbbt/workflow/step/accessor.rb', line 158

def short_path
  [task_name, name] * "/"
end

#soft_graceObject



677
678
679
680
681
682
# File 'lib/rbbt/workflow/step/run.rb', line 677

def soft_grace
  until done? or (Open.exist?(info_file) && info[:status] != :noinfo)
    sleep 1 
  end
  self
end

#stalled?Boolean

Returns:

  • (Boolean)


530
531
532
# File 'lib/rbbt/workflow/step/accessor.rb', line 530

def stalled?
  started? && ! (done? || running? || done? || error? || aborted?)
end

#started?Boolean

Returns:

  • (Boolean)


460
461
462
# File 'lib/rbbt/workflow/step/accessor.rb', line 460

def started?
  Open.exists?(path) or (Open.exists?(pid_file) && Open.exists?(info_file))
end

#statusObject



283
284
285
286
287
288
289
290
# File 'lib/rbbt/workflow/step/accessor.rb', line 283

def status
  begin
    info[:status]
  rescue Exception
    Log.error "Exception reading status: #{$!.message}" 
    :error
  end
end

#status=(status) ⇒ Object



292
293
294
# File 'lib/rbbt/workflow/step/accessor.rb', line 292

def status=(status)
  set_info(:status, status)
end

#status_lockObject



194
195
196
197
198
199
200
201
# File 'lib/rbbt/workflow/step/accessor.rb', line 194

def status_lock
  return @mutex
  #@status_lock = begin
  #               path = Persist.persistence_path(info_file + '.status.lock', {:dir => Step.lock_dir})
  #               Lockfile.new path, :refresh => false, :dont_use_lock_id => true
  #             end if @status_lock.nil?
  #@status_lock
end

#step(name) ⇒ Object



511
512
513
514
515
516
517
518
519
520
521
522
523
524
# File 'lib/rbbt/workflow/step.rb', line 511

def step(name)
  @steps ||= {}
  @steps[name] ||= begin
                     deps = rec_dependencies.select{|step| 
                       step.task_name.to_sym == name.to_sym
                     }
                     raise "Dependency step not found: #{ name }" if deps.empty?
                     if (deps & self.dependencies).any?
                       (deps & self.dependencies).first
                     else
                       deps.first
                     end
                   end
end

#stop_dependenciesObject



471
472
473
474
475
476
477
478
479
480
481
482
483
484
# File 'lib/rbbt/workflow/step/dependencies.rb', line 471

def stop_dependencies
  return if dependencies.nil?
  dependencies.each do |dep|
    if dep.nil?
      Log.warn "Dependency is nil #{Misc.fingerprint step} -- #{Misc.fingerprint dependencies}"
      next
    end

    next if dep.done? or dep.aborted?

    dep.abort if dep.running?
  end
  kill_children
end

#streaming?Boolean

Returns:

  • (Boolean)


508
509
510
# File 'lib/rbbt/workflow/step/accessor.rb', line 508

def streaming?
  (IO === @result) or (not @saved_stream.nil?) or status == :streaming
end

#task_signatureObject



171
172
173
# File 'lib/rbbt/workflow/step/accessor.rb', line 171

def task_signature
  [workflow.to_s, task_name] * "#"
end

#tmp_pathObject



557
558
559
# File 'lib/rbbt/workflow/step/accessor.rb', line 557

def tmp_path
  @tmp_path ||= Step.tmp_path path
end

#updatable?Boolean

Returns:

  • (Boolean)


114
115
116
117
118
119
120
# File 'lib/rbbt/workflow/step/run.rb', line 114

def updatable?
  return true if ENV["RBBT_UPDATE_ALL_JOBS"] == 'true'
  return false unless ENV["RBBT_UPDATE"] == "true"
  return false unless Open.exists?(info_file)
  return true if status != :noinfo && ! (relocated? && done?)
  false
end

#updateObject



453
454
455
456
457
458
# File 'lib/rbbt/workflow/step.rb', line 453

def update
  if dirty?
    dependencies.collect{|d| d.update } if dependencies
    clean
  end
end

#updated?Boolean

Returns:

  • (Boolean)


190
191
192
193
194
195
# File 'lib/rbbt/workflow/step/run.rb', line 190

def updated?
  return true if ENV["RBBT_UPDATE"] != "true"
  return true unless (done? || error? || ! writable?)

  @updated ||= out_of_date.empty?
end

#waiting?Boolean

Returns:

  • (Boolean)


464
465
466
# File 'lib/rbbt/workflow/step/accessor.rb', line 464

def waiting?
  Open.exists?(info_file) and not started?
end

#workflow_short_pathObject



162
163
164
165
# File 'lib/rbbt/workflow/step/accessor.rb', line 162

def workflow_short_path
  return short_path unless workflow
  workflow.to_s + "#" + short_path
end

#writable?Boolean

Returns:

  • (Boolean)


499
500
501
# File 'lib/rbbt/workflow/step.rb', line 499

def writable?
  Open.writable?(self.path) && Open.writable?(self.info_file)
end