Module: Ace::Hitl::Lifecycle::Effects
- Defined in:
- lib/ace/hitl/lifecycle/effects.rb
Overview
The requester-declared answer-effect layer (spec 8wm.t.y21 §5; recovered from the deployed 8wl.t.ga9 contract). A request may declare a callback that executes exactly once when the answer arrives — always AFTER the answer is relayed — AS THE REQUESTER, through exec-style argv (a shell never sees the answer). The answer content and the substituted argv never appear in any log; full attempts live only in the root-only effects log.
Constant Summary collapse
- DEFAULT_TIMEOUT_S =
120- TERMINATE_GRACE_SECONDS =
0.2- OUTCOME_OK =
"callback-ok"- OUTCOME_ESCALATED =
"callback-escalated"- DeclarationError =
A Lifecycle::Error so every lifecycle caller's typed rescue (notably the provider seam's orphan-event ask wrapper) catches declaration violations instead of a raw backtrace escaping to the CLI (review F-R1 on W696).
Class.new(Lifecycle::Error)
Class Method Summary collapse
-
.drop_child_groups(gid) ⇒ Object
Supplementary-group isolation for the effect child (spec §5; review F7 on W696): spawn has no groups hook, so root swaps the process supplementary list for exactly the requester's group for the duration of the block and restores it afterwards; the child inherits the dropped list.
-
.escalate_once(store, value, attempt) ⇒ Object
The deduped escalation: recorded once per request in the effects log; the wake/spool glue stays behind the sink seam.
-
.execute(declaration, answer, requester_uid, requester_gid, spawner = Process, group_dropper: nil) ⇒ Object
Exec-style spawn, never a shell.
-
.fullmatch?(source, answer) ⇒ Boolean
Python re.fullmatch equivalent: the whole answer must match.
-
.normalized_declaration(effect) ⇒ Object
The persisted record shape, byte-compatible with the deployed contract: cwd:, match:, timeout_s:.
- .outcome_of(attempt) ⇒ Object
-
.record_effect_attempt(store, value, attempt) ⇒ Object
The root-only effects log: redacted attempts (no answer, no argv) plus the deduped escalation marker.
-
.run(store, value, answer, requester_uid:, requester_gid:, identity: Identity, spawner: Process, group_dropper: nil) ⇒ Object
Executed inside deliver's locked critical section, after the answer is relayed.
- .sleep_after_terminate ⇒ Object
-
.validate_declaration!(effect) ⇒ Object
Validated by Atoms::HitlEffectValidator at the CLI boundary; the store re-applies the declaration checks so direct API use cannot bypass the bounds.
Class Method Details
.drop_child_groups(gid) ⇒ Object
Supplementary-group isolation for the effect child (spec §5; review F7 on W696): spawn has no groups hook, so root swaps the process supplementary list for exactly the requester's group for the duration of the block and restores it afterwards; the child inherits the dropped list. A non-root process has nothing to drop and would hit EPERM, so it yields unchanged.
236 237 238 239 240 241 242 243 244 245 246 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 236 def drop_child_groups(gid) return yield unless Process.euid.zero? previous = Process.groups Process::Sys.setgroups([gid]) begin yield ensure Process::Sys.setgroups(previous) end end |
.escalate_once(store, value, attempt) ⇒ Object
The deduped escalation: recorded once per request in the effects log; the wake/spool glue stays behind the sink seam.
141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 141 def escalate_once(store, value, attempt) return if value["escalation_spooled"] sink = store.escalation_sink value["escalation_spooled"] = true return unless sink sink.call( request_id: value["id"], work: value["work"], attempt: value["attempt"], project: value["project"], outcome: attempt["outcome"] ) end |
.execute(declaration, answer, requester_uid, requester_gid, spawner = Process, group_dropper: nil) ⇒ Object
Exec-style spawn, never a shell. answer substitutes once per element as a plain string replace. Output is discarded (redaction). Process.spawn applies only setgid/setuid, so the supplementary-group drop happens around the spawn window (see drop_child_groups) and the forked child inherits the dropped list. Returns [exit_status, timed_out].
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 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 170 def execute(declaration, answer, requester_uid, requester_gid, spawner = Process, group_dropper: nil) group_dropper ||= method(:drop_child_groups) # Block form: the answer is inserted literally. The # replacement-string form would interpret backslash sequences # in the answer as backreferences (review 8wq2ztty on PR#336). argv = declaration["argv"].map { |element| element.gsub("{answer}") { answer } } timeout = declaration["timeout_s"] || DEFAULT_TIMEOUT_S deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout status = nil timed_out = false group_dropper.call(requester_gid) do # [cmd, cmd] forces exec/argv semantics: a single-element # argv would otherwise be handed to a shell as a command # string (review 8wq2ztu2 on PR#336), and a substituted # answer must never become shell syntax. # # pgroup: true makes the child a process group leader so the # timeout can signal its whole group (review 8wq2zttw on # PR#336). pid = spawner.spawn( [argv.first, argv.first], *argv.drop(1), chdir: declaration["cwd"], gid: requester_gid, uid: requester_uid, pgroup: true, out: File::NULL, err: File::NULL ) loop do _, status = spawner.waitpid2(pid, Process::WNOHANG) break if status if Process.clock_gettime(Process::CLOCK_MONOTONIC) > deadline timed_out = true # -pid addresses the process group: a callback that # forked leaves no descendants running after the # timeout. ESRCH means the group is already gone; macOS # additionally raises EPERM when signaling a group whose # only remaining members are zombies — both are fine, # the direct child is reaped by the wait below. begin spawner.kill("TERM", -pid) rescue Errno::ESRCH, Errno::EPERM end sleep_after_terminate begin spawner.kill("KILL", -pid) rescue Errno::ESRCH, Errno::EPERM end _, status = spawner.waitpid2(pid) break end sleep 0.05 end end [status&.exitstatus, timed_out] end |
.fullmatch?(source, answer) ⇒ Boolean
Python re.fullmatch equivalent: the whole answer must match.
158 159 160 161 162 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 158 def fullmatch?(source, answer) Regexp.new("\\A(?:#{source})\\z").match?(answer) rescue RegexpError false end |
.normalized_declaration(effect) ⇒ Object
The persisted record shape, byte-compatible with the deployed contract: cwd:, match:, timeout_s:.
57 58 59 60 61 62 63 64 65 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 57 def normalized_declaration(effect) timeout = effect[:timeout_s] || effect[:effect_timeout] || DEFAULT_TIMEOUT_S { "argv" => Array(effect[:effect_args] || effect[:argv]).map(&:to_s), "cwd" => (effect[:cwd] || effect[:effect_cwd]).to_s, "match" => effect[:match]&.to_s, "timeout_s" => Integer(timeout) } end |
.outcome_of(attempt) ⇒ Object
120 121 122 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 120 def outcome_of(attempt) (attempt["match_ok"] && !attempt["timed_out"] && attempt["exit_status"] == 0) ? "ok" : "escalated" end |
.record_effect_attempt(store, value, attempt) ⇒ Object
The root-only effects log: redacted attempts (no answer, no argv) plus the deduped escalation marker.
126 127 128 129 130 131 132 133 134 135 136 137 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 126 def record_effect_attempt(store, value, attempt) store.effects_dir.mkpath path = store.effects_dir.join("#{value["id"]}.json") record = AtomicJson.read(path) record = {"id" => value["id"], "attempts" => [], "escalated" => nil} unless record.is_a?(Hash) record["attempts"] << attempt if attempt["outcome"] != "ok" && record["escalated"].nil? record["escalated"] = {"at" => Time.now.to_i, "outcome" => attempt["outcome"]} end AtomicJson.call(path, record, mode: 0o600) nil end |
.run(store, value, answer, requester_uid:, requester_gid:, identity: Identity, spawner: Process, group_dropper: nil) ⇒ Object
Executed inside deliver's locked critical section, after the answer is relayed. Exactly one attempt; failure escalates once.
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 107 108 109 110 111 112 113 114 115 116 117 118 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 69 def run(store, value, answer, requester_uid:, requester_gid:, identity: Identity, spawner: Process, group_dropper: nil) declaration = value["effect"] return nil unless declaration.is_a?(Hash) && Array(declaration["argv"]).any? if !identity.root? && requester_uid != identity.euid raise Lifecycle::PermissionError, "effect callback requires the requester's identity and root authority to drop to it" end match_ok = declaration["match"].nil? || fullmatch?(declaration["match"], answer) attempt = { "at" => Time.now.to_i, "match_ok" => match_ok, "timed_out" => false, "exit_status" => nil, "duration_s" => nil } if match_ok start = Process.clock_gettime(Process::CLOCK_MONOTONIC) begin status, timed_out = execute( declaration, answer, requester_uid, requester_gid, spawner, group_dropper: group_dropper ) rescue SystemCallError => e # A spawn/wait failure (missing binary, non-executable # argv, EPERM) must never escape deliver AFTER the answer # was relayed: the operator signal would be silently lost. # It is an escalation outcome exactly like a timeout or a # nonzero exit (spec §5; review F2 on W696). attempt["error"] = e.class.name status = nil timed_out = false end attempt["timed_out"] = timed_out attempt["exit_status"] = status attempt["duration_s"] = (Process.clock_gettime(Process::CLOCK_MONOTONIC) - start).round(3) end attempt["outcome"] = outcome_of(attempt) record_effect_attempt(store, value, attempt) effect_state = (attempt["outcome"] == "ok") ? OUTCOME_OK : OUTCOME_ESCALATED value["effect_state"] = effect_state store.update_public(value, "answer-delivered") if effect_state == OUTCOME_ESCALATED escalate_once(store, value, attempt) end effect_state end |
.sleep_after_terminate ⇒ Object
248 249 250 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 248 def sleep_after_terminate sleep(TERMINATE_GRACE_SECONDS) end |
.validate_declaration!(effect) ⇒ Object
Validated by Atoms::HitlEffectValidator at the CLI boundary; the store re-applies the declaration checks so direct API use cannot bypass the bounds. Declarations are never rewritten.
34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 |
# File 'lib/ace/hitl/lifecycle/effects.rb', line 34 def validate_declaration!(effect) match = effect[:match]&.to_s cwd = (effect[:cwd] || effect[:effect_cwd])&.to_s # cwd is a required declaration field (spec §5: {argv:, cwd:, # match:, timeout_s:}). A cwd-less declaration would otherwise # pass validation and fail at answer time as a misleading # Errno::ENOENT spawn escalation (review F-A on W696). if cwd.nil? || cwd.strip.empty? raise DeclarationError, "--effect-cwd is required for an effect callback (absolute path to an existing directory)" end Atoms::HitlEffectValidator.validate!( match: match, effect_args: Array(effect[:effect_args] || effect[:argv]).map(&:to_s), effect_cwd: cwd, effect_timeout: (effect[:timeout_s] || effect[:effect_timeout]).to_s ) rescue Atoms::HitlEffectValidator::ValidationError => e raise DeclarationError, e. end |