Class: Tengine::Job::Runtime::Edge

Inherits:
Object
  • Object
show all
Includes:
Mongoid::Document, Mongoid::Timestamps, Core::SelectableAttr, Signal::Transmittable, Structure::Visitor::Accepter
Defined in:
lib/tengine/job/runtime/edge.rb

Overview

Vertexとともにジョブネットを構成するグラフの「辺」を表すモデル Tengine::Job::Runtime::Jobnetにembeddedされます。

Defined Under Namespace

Classes: Closer, StatusError

Instance Method Summary collapse

Instance Method Details

#alive?Boolean

Returns:



39
# File 'lib/tengine/job/runtime/edge.rb', line 39

def alive?; !!phase_entry[:alive]; end

#alive_or_closing?Boolean

Returns:



40
# File 'lib/tengine/job/runtime/edge.rb', line 40

def alive_or_closing?; alive? || closing?; end

#alive_or_closing_or_closed?Boolean

Returns:



41
# File 'lib/tengine/job/runtime/edge.rb', line 41

def alive_or_closing_or_closed?; alive? || closing? || closed?; end

#close(signal) ⇒ Object



125
126
127
128
129
130
131
132
133
134
# File 'lib/tengine/job/runtime/edge.rb', line 125

def close(signal)
  case phase_key
  when :active, :suspended, :keeping, :transmitting then
    self.phase_key = :closing
  when :closing, :closed then
    # ignored
  else
    Tengine.logger.warn "#{object_id} #{inspect} wasn't closed"
  end
end

#close_followings(signal, options = {}) ⇒ Object



136
137
138
139
140
# File 'lib/tengine/job/runtime/edge.rb', line 136

def close_followings(signal, options = {})
  v = Tengine::Job::Runtime::Edge::Closer.new(signal, options)
  accept_visitor(v)
  v.closed_edges
end

#close_followings_and_trasmit(signal) ⇒ Object

ownerのupdate_with_lockを使っています。



143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
# File 'lib/tengine/job/runtime/edge.rb', line 143

def close_followings_and_trasmit(signal)
  jobnet = signal.cache(self.owner)
  closing_edges = nil
  closed_edges = []
  jobnet.update_with_lock do
    closing_edges = self.close_followings(signal)
    closing_edges.each do |e|
      next unless e.owner.id == jobnet.id
      je = signal.cache(e) # jobnet単位で保存するので、jobnetオブエジェクトに紐付けられたものを見つける
      je.close(nil)
      closed_edges << e
    end
    # jobnetオブジェクトのedgesに含まれないエッジについては、そのowner毎にまとめて保存する
    self.transmit(signal)

    Tengine.logger.debug "<" * 100
    Tengine.logger.debug "#{__FILE__}##{__LINE__}"
    jobnet.edges.each do |edge|
      Tengine.logger.debug "#{edge.object_id} #{edge.inspect} BEFORE end of block for update_with_lock"
    end
  end
  # signal.cache_list
  signal.changed_vertecs.each(&:save!)
end

#complete(signal) ⇒ Object



98
99
100
101
102
103
104
105
106
107
108
# File 'lib/tengine/job/runtime/edge.rb', line 98

def complete(signal)
  case phase_key
  when :transmitting then
    self.phase_key = :transmitted
  when :active, :closed then
    # IG
  when :suspended, :keeping then
    # N/A
    raise Tengine::Job::Runtime::Edge::StatusError, "#{self.class.name}#complete not available on #{phase_key.inspect} at #{self.inspect}"
  end
end

#destinationObject



53
54
55
# File 'lib/tengine/job/runtime/edge.rb', line 53

def destination
  owner.children.detect{|c| c.id == destination_id}
end

#inspectObject



61
62
63
# File 'lib/tengine/job/runtime/edge.rb', line 61

def inspect
  "#<#{self.class.name} #{phase_key.inspect} #{name_for_message}>"
end

#name_for_messageObject



57
58
59
# File 'lib/tengine/job/runtime/edge.rb', line 57

def name_for_message
  "edge(#{id.to_s}) from #{origin ? origin.name_path : 'no origin'} to #{destination ? destination.name_path : 'no destination'}"
end

#originObject



49
50
51
# File 'lib/tengine/job/runtime/edge.rb', line 49

def origin
  owner.children.detect{|c| c.id == origin_id}
end

#phase_key=(phase_key) ⇒ Object



168
169
170
171
172
# File 'lib/tengine/job/runtime/edge.rb', line 168

def phase_key=(phase_key)
  # Tengine.logger.debug("edge phase changed. <#{self.id.to_s}> #{self.phase_name} -> #{Tengine::Job::Runtime::Edge.phase_name_by_key(phase_key)}")
  Tengine.logger.debug("#{object_id} edge phase changed. <#{inspect}> #{self.phase_name} -> #{Tengine::Job::Runtime::Edge.phase_name_by_key(phase_key)}")
  self.write_attribute(:phase_cd, Tengine::Job::Runtime::Edge.phase_id_by_key(phase_key))
end

#reset(signal) ⇒ Object



110
111
112
113
114
115
116
117
118
119
120
121
122
123
# File 'lib/tengine/job/runtime/edge.rb', line 110

def reset(signal)
  # 全てのステータスから遷移する
  if d = destination
    if signal.execution.in_scope?(d)
      self.phase_key = :active

      signal.call_later do
        d.reset(signal)
      end
    end
  else
    raise "destination not found: #{destination_id.inspect} from #{origin.inspect}"
  end
end

#transmit(signal) ⇒ Object



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
# File 'lib/tengine/job/runtime/edge.rb', line 66

def transmit(signal)
  case phase_key
  when :active then
    d = destination
    if signal.execution.in_scope?(d)
      self.phase_key = :transmitting
      signal.leave(self)
    else
      Tengine.logger.info("#{d.name_path} will not be executed, becauase it is out of execution scope.")
    end
  when :suspended then
    self.phase_key = :keeping
  when :closing then

Tengine.logger.debug "c" * 100
Tengine.logger.debug "#{object_id} #{inspect}"
# Tengine.logger.debug caller[0, 20].join("\n  ")

# binding.pry

    self.phase_key = :closed
    signal.paths << self
    signal.with_paths_backup do
      if destination.is_a?(Tengine::Job::Runtime::NamedVertex)
        signal.cache(destination.next_edges.first).transmit(signal)
      else
        signal.leave(self)
      end
    end
  end
end