Class: Tengine::Job::Runtime::Edge
- Inherits:
-
Object
- Object
- Tengine::Job::Runtime::Edge
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
39
|
# File 'lib/tengine/job/runtime/edge.rb', line 39
def alive?; !!phase_entry[:alive]; end
|
#alive_or_closing? ⇒ Boolean
40
|
# File 'lib/tengine/job/runtime/edge.rb', line 40
def alive_or_closing?; alive? || closing?; end
|
#alive_or_closing_or_closed? ⇒ Boolean
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
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)
je.close(nil)
closed_edges << e
end
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.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
when :suspended, :keeping then
raise Tengine::Job::Runtime::Edge::StatusError, "#{self.class.name}#complete not available on #{phase_key.inspect} at #{self.inspect}"
end
end
|
#destination ⇒ Object
53
54
55
|
# File 'lib/tengine/job/runtime/edge.rb', line 53
def destination
owner.children.detect{|c| c.id == destination_id}
end
|
#inspect ⇒ Object
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_message ⇒ Object
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
|
#origin ⇒ Object
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("#{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}"
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
|