Class: Ace::Hitl::Lifecycle::Store
- Inherits:
-
Object
- Object
- Ace::Hitl::Lifecycle::Store
- Defined in:
- lib/ace/hitl/lifecycle/store.rb
Overview
The generic, file-backed HITL request lifecycle (spec 8wm.t.y21 §3): create/pending/states/deliver/consume/cancel. Time never ends a request (W651): only an answer, its consumption, or an explicit audited cancellation removes one. All terminal transitions share one per-request flock, the same transaction boundary the lab-side stop protocol uses.
Constant Summary collapse
- MAX_ANSWER =
4096- MAX_ANSWER_BYTES =
IO#read(limit) is byte-oriented, so the character limit is enforced on a decoded string with this separate byte bound (review 8wq2zttz on PR#336); 4 bytes per character is the UTF-8 worst case.
MAX_ANSWER * 4
- MAX_PLAN_QUESTION =
240- MAX_OPTIONS =
8- MAX_OPTION =
80- PUBLIC_MODE =
0o440- ANSWER_MODE =
0o400- REQUEST_MODE =
0o600- CONSUME_POLL_SECONDS =
1- ROOT_MODE =
Pinned directory modes (spec §2): everything 0700 except the public projection (0755). Enforced umask-proof at first create (review F6 on W696).
0o700- DIR_MODES =
{ "requests" => 0o700, "secrets" => 0o700, "answers" => 0o700, "public" => 0o755, "effects" => 0o700 }.freeze
- DEFAULT_ADMIN_USER =
"lab-admin"- DEFAULT_GROUP =
"lab-control"
Instance Attribute Summary collapse
-
#binding ⇒ Object
readonly
Returns the value of attribute binding.
-
#escalation_sink ⇒ Object
The escalation spool seam (spec 8wm.t.y21 §1): the generic core records deduped escalation state; the wake/spool glue stays lab-side and wires in through this callable.
-
#root ⇒ Object
readonly
Returns the value of attribute root.
Instance Method Summary collapse
- #answer_path(value) ⇒ Object
- #answers_dir ⇒ Object
-
#cancel(id, reason: "") ⇒ Object
The only way to abandon a request: explicit, audited cancellation (W651).
-
#consume(id, timeout: 0) ⇒ Object
Consumes the answer for one requester-owned request.
-
#create(id:, work:, attempt:, plan:, question:, ace_hitl_id:, project: "ace", harness: DEFAULT_ADMIN_USER, kind: "text", options: [], effect: nil) ⇒ Object
---- requester side ------------------------------------------------.
-
#deliver(id, answer_reader) ⇒ Object
Answer one pending request as the host broker.
- #effects_dir ⇒ Object
-
#ensure_layout! ⇒ Object
Provision the store layout with the pinned directory modes (spec §2; review F6 on W696).
- #group_id ⇒ Object
-
#initialize(root:, binding:, admin_user: DEFAULT_ADMIN_USER, group: DEFAULT_GROUP, ownership: AtomicJson::DEFAULT_OWNERSHIP, identity: Identity, poll_seconds: CONSUME_POLL_SECONDS) ⇒ Store
constructor
A new instance of Store.
-
#initialize_projection!(value) ⇒ Object
First projection of a fresh incarnation, run under the same per-request flock every transition uses (review 8wq2zttv on PR#336).
- #load_request(request_id) ⇒ Object
- #ownership_for(value) ⇒ Object
-
#pending ⇒ Object
Answerable requests; never purges or cancels anything (W651).
- #public_dir ⇒ Object
- #public_path(request_id) ⇒ Object
- #read_answer(value) ⇒ Object
- #remove_request(value, keep_public: false) ⇒ Object
- #request_path(request_id) ⇒ Object
-
#requests_dir ⇒ Object
---- shared internals ----------------------------------------------.
- #require_root!(operation) ⇒ Object
- #safe_id(value) ⇒ Object
- #secrets_dir ⇒ Object
-
#states ⇒ Object
All public lifecycle records.
-
#update_public(value, state, audit: nil) ⇒ Object
The public lifecycle projection: merge-on-write — every write merges over the existing projection (state, timestamps and audit fields win; effect/escalation fields recorded by an earlier write survive every later lifecycle write, review F-R2 on W696), 0440, owner = requester, group = the control group.
-
#with_request_lock(request_id) ⇒ Object
The per-request flock: the transaction boundary shared with deliver/consume/cancel and the lab-side stop protocol.
Constructor Details
#initialize(root:, binding:, admin_user: DEFAULT_ADMIN_USER, group: DEFAULT_GROUP, ownership: AtomicJson::DEFAULT_OWNERSHIP, identity: Identity, poll_seconds: CONSUME_POLL_SECONDS) ⇒ Store
Returns a new instance of Store.
58 59 60 61 62 63 64 65 66 67 68 69 70 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 58 def initialize(root:, binding:, admin_user: DEFAULT_ADMIN_USER, group: DEFAULT_GROUP, ownership: AtomicJson::DEFAULT_OWNERSHIP, identity: Identity, poll_seconds: CONSUME_POLL_SECONDS) raise ArgumentError, "a binding policy is required (fail closed without one)" unless binding @root = Pathname.new(root) @binding = binding @admin_user = admin_user @group = group @ownership = ownership @identity = identity @poll_seconds = poll_seconds end |
Instance Attribute Details
#binding ⇒ Object (readonly)
Returns the value of attribute binding.
51 52 53 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 51 def binding @binding end |
#escalation_sink ⇒ Object
The escalation spool seam (spec 8wm.t.y21 §1): the generic core records deduped escalation state; the wake/spool glue stays lab-side and wires in through this callable. No default spool.
56 57 58 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 56 def escalation_sink @escalation_sink end |
#root ⇒ Object (readonly)
Returns the value of attribute root.
51 52 53 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 51 def root @root end |
Instance Method Details
#answer_path(value) ⇒ Object
319 320 321 322 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 319 def answer_path(value) directory = (value["sensitive"] == true) ? secrets_dir : answers_dir directory.join("#{safe_id(value["id"])}.answer") end |
#answers_dir ⇒ Object
307 308 309 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 307 def answers_dir @root.join("answers") end |
#cancel(id, reason: "") ⇒ Object
The only way to abandon a request: explicit, audited cancellation (W651). Time never cancels anything.
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 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 198 def cancel(id, reason: "") request_id = safe_id(id) cancelled_by = @identity.username audit = {"cancelled_by" => cancelled_by, "reason" => reason.to_s.strip.empty? ? "unspecified" : reason.to_s.strip} locked_value = nil with_request_lock(request_id) do value = load_request(request_id) # Ownership is re-checked against the LOCKED record: an # unlocked check would race a concurrent recreate of the id # (review 8wq2zttu on PR#336). unless @identity.root? || @identity.username == value["requester"].to_s raise PermissionError, "only the requesting role can cancel this request" end update_public(value, "cancelled", audit: audit) remove_request(value, keep_public: true) locked_value = value end { "id" => request_id, "work" => locked_value["work"], "attempt" => locked_value["attempt"], "cancelled" => true, "cancelled_by" => cancelled_by, "reason" => audit["reason"] } end |
#consume(id, timeout: 0) ⇒ Object
Consumes the answer for one requester-owned request. Without a positive timeout the wait is indefinite; a timeout bounds ONLY this local wait and never cancels the request (W651).
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 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 156 def consume(id, timeout: 0) request_id = safe_id(id) value = load_request(request_id) requester_gate!(value) deadline = timeout.positive? ? Time.now.to_i + timeout : nil loop do answer = nil with_request_lock(request_id) do value = load_request(request_id) requester_gate!(value) verify_active!(value) answer = read_answer(value) unless answer next end # The terminal transition commits under the same per-request # flock as deliver/cancel (spec §3): a concurrent cancel # must not interleave between reading the answer and # removing the request (review F5 on W696). update_public(value, "consumed") remove_request(value, keep_public: true) end if answer return { "id" => request_id, "work" => value["work"], "attempt" => value["attempt"], "answer" => answer, "sensitive" => value["sensitive"] == true } end break if deadline && Time.now.to_i > deadline sleep(@poll_seconds) end raise StateError, "timed out waiting for HITL answer; the request remains pending" end |
#create(id:, work:, attempt:, plan:, question:, ace_hitl_id:, project: "ace", harness: DEFAULT_ADMIN_USER, kind: "text", options: [], effect: nil) ⇒ Object
---- requester side ------------------------------------------------
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 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/ace/hitl/lifecycle/store.rb', line 74 def create(id:, work:, attempt:, plan:, question:, ace_hitl_id:, project: "ace", harness: DEFAULT_ADMIN_USER, kind: "text", options: [], effect: nil) requester = @identity.username request_id = safe_id(id) unless Kinds::WORK_ID.match?(work) raise StateError, "invalid work id" end unless attempt && Kinds::ATTEMPT_ID.match?(attempt) raise StateError, "HITL request requires the exact active Attempt id" end unless Kinds::SAFE_LABEL.match?(project) && Kinds::SAFE_LABEL.match?(harness) raise StateError, "invalid project or harness label" end plan = plan.to_s.strip question = question.to_s.strip if plan.empty? || plan.length > MAX_PLAN_QUESTION || question.empty? || question.length > MAX_PLAN_QUESTION raise StateError, "plan and question must contain 1-#{MAX_PLAN_QUESTION} characters" end = Array().map(&:strip) if .length > MAX_OPTIONS || .any? { |item| item.empty? || item.length > MAX_OPTION } raise StateError, "options must contain at most eight 1-#{MAX_OPTION} character values" end if kind == "secret" raise StateError, "secret HITL must use the approved otp kind" end unless Kinds.valid?(kind) raise StateError, "unknown HITL kind: #{kind}" end if Kinds.secret?(kind) unless .empty? raise StateError, "OTP requests must not offer choices" end validate_request_binding(work, attempt, project, requester) elsif requester != @admin_user validate_request_binding(work, attempt, project, requester) end value = { "id" => request_id, "work" => work, "attempt" => attempt, "project" => project, "harness" => harness, "kind" => kind, "sensitive" => Kinds.secret?(kind), "plan" => plan, "question" => question, "options" => , "ace_hitl_id" => ace_hitl_id, "requester" => requester, "created_at" => Time.now.to_i } Effects.validate_declaration!(effect) if effect value["effect"] = Effects.normalized_declaration(effect) if effect value["incarnation"] = SecureRandom.hex(8) request_path = requests_dir.join("#{request_id}.json") raise StateError, "HITL request already exists" if request_path.exist? ensure_layout! begin # link(2) is the create-once commit point: rename(2) would # silently let a concurrent duplicate-id create overwrite the # winner (review F8 on W696). AtomicJson.call(request_path, value, mode: REQUEST_MODE, ownership_strategy: @ownership, exclusive: true) rescue Errno::EEXIST raise StateError, "HITL request already exists" end initialize_projection!(value) { "id" => request_id, "work" => value["work"], "attempt" => attempt, "requested" => true } end |
#deliver(id, answer_reader) ⇒ Object
Answer one pending request as the host broker. The answer is ALWAYS relayed unchanged; the declared effect callback (if any) then executes inside the same locked critical section.
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 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 230 def deliver(id, answer_reader) require_root!("deliver") request_id = safe_id(id) value = load_request(request_id) if value["sensitive"] == true && !Kinds.secret?(value["kind"].to_s) raise StateError, "secret HITL answers are forbidden" end raise StateError, "HITL request already has an answer" if answer_path(value).exist? answer = read_bounded_answer(answer_reader) Kinds.check_answer!(value["kind"].to_s, answer) requester_uid, requester_gid = @identity.user_ids(value["requester"]) # Snapshot of the incarnation the unlocked checks and the # dropped credentials were derived from: a concurrent # cancel+recreate may replace the record before this broker # takes the lock (review 8wq2zttu on PR#336). incarnation = [value["created_at"], value["kind"], value["sensitive"], value["requester"]] with_request_lock(request_id) do value = load_request(request_id) if [value["created_at"], value["kind"], value["sensitive"], value["requester"]] != incarnation raise StateError, "HITL request was cancelled and recreated while the answer was read" end raise StateError, "HITL request already has an answer" if answer_path(value).exist? begin binding.require_active(work: value["work"], attempt: value["attempt"]) rescue BindingError update_public(value, "cancelled") remove_request(value, keep_public: true) raise end write_answer(value, answer, requester_uid, requester_gid) update_public(value, "answer-delivered") Effects.run( self, value, answer, requester_uid: requester_uid, requester_gid: requester_gid ) end { "id" => request_id, "work" => value["work"], "attempt" => value["attempt"], "delivered" => true, "sensitive" => value["sensitive"] == true } ensure answer&.clear end |
#effects_dir ⇒ Object
315 316 317 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 315 def effects_dir @root.join("effects") end |
#ensure_layout! ⇒ Object
Provision the store layout with the pinned directory modes (spec §2; review F6 on W696). mkdir(2) bakes in the process umask, so every directory is chmod'd explicitly after mkdir — the same umask-proof pattern as the atomic writer.
332 333 334 335 336 337 338 339 340 341 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 332 def ensure_layout! FileUtils.mkdir_p(@root) File.chmod(ROOT_MODE, @root) DIR_MODES.each do |name, mode| dir = @root.join(name) FileUtils.mkdir_p(dir) File.chmod(mode, dir) end nil end |
#group_id ⇒ Object
452 453 454 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 452 def group_id @group_id ||= @identity.group_id(@group) end |
#initialize_projection!(value) ⇒ Object
First projection of a fresh incarnation, run under the same per-request flock every transition uses (review 8wq2zttv on PR#336). The projection-and-effects cleanup of a consumed or cancelled predecessor happens here, keyed by the incarnation token: a predecessor's artifacts never leak into the new incarnation (review F-B on W696), and a broker that won the lock first and already delivered can never be regressed to "created" — a current-incarnation projection is left untouched.
412 413 414 415 416 417 418 419 420 421 422 423 424 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 412 def initialize_projection!(value) request_id = safe_id(value["id"].to_s) with_request_lock(request_id) do current = AtomicJson.read(public_path(request_id)) predecessor = current.nil? || current["incarnation"] != value["incarnation"] if predecessor public_path(request_id).unlink if public_path(request_id).exist? effects_dir.join("#{request_id}.json").unlink if effects_dir.join("#{request_id}.json").exist? update_public(value, "created") end end nil end |
#load_request(request_id) ⇒ Object
426 427 428 429 430 431 432 433 434 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 426 def load_request(request_id) path = request_path(request_id) raise StateError, "unknown or invalid HITL request" unless path.exist? value = AtomicJson.read(path) raise StateError, "unknown or invalid HITL request" unless value.is_a?(Hash) value end |
#ownership_for(value) ⇒ Object
445 446 447 448 449 450 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 445 def ownership_for(value) uid, _ = @identity.user_ids(value["requester"].to_s) AtomicJson::Ownership.new(uid: uid, gid: group_id) rescue Lifecycle::Error nil end |
#pending ⇒ Object
Answerable requests; never purges or cancels anything (W651).
279 280 281 282 283 284 285 286 287 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 279 def pending require_root!("pending") requests_dir.glob("*.json").sort.filter_map do |path| value = AtomicJson.read(path) next unless value.is_a?(Hash) value if !answer_path(value).exist? end end |
#public_dir ⇒ Object
311 312 313 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 311 def public_dir @root.join("public") end |
#public_path(request_id) ⇒ Object
343 344 345 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 343 def public_path(request_id) public_dir.join("#{safe_id(request_id)}.json") end |
#read_answer(value) ⇒ Object
436 437 438 439 440 441 442 443 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 436 def read_answer(value) path = answer_path(value) return nil unless path.exist? answer = path.read Kinds.check_answer!(value["kind"].to_s, answer) answer end |
#remove_request(value, keep_public: false) ⇒ Object
378 379 380 381 382 383 384 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 378 def remove_request(value, keep_public: false) request_id = safe_id(value["id"].to_s) request_path(request_id).unlink if request_path(request_id).exist? public_path(request_id).unlink if !keep_public && public_path(request_id).exist? answer_path(value).unlink if answer_path(value).exist? nil end |
#request_path(request_id) ⇒ Object
324 325 326 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 324 def request_path(request_id) requests_dir.join("#{safe_id(request_id)}.json") end |
#requests_dir ⇒ Object
---- shared internals ----------------------------------------------
299 300 301 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 299 def requests_dir @root.join("requests") end |
#require_root!(operation) ⇒ Object
456 457 458 459 460 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 456 def require_root!(operation) unless @identity.root? raise PermissionError, "#{operation} is a host-broker operation" end end |
#safe_id(value) ⇒ Object
462 463 464 465 466 467 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 462 def safe_id(value) id = value.to_s raise StateError, "invalid HITL request id" unless Kinds::REQUEST_ID.match?(id) id end |
#secrets_dir ⇒ Object
303 304 305 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 303 def secrets_dir @root.join("secrets") end |
#states ⇒ Object
All public lifecycle records.
290 291 292 293 294 295 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 290 def states require_root!("states") public_dir.glob("*.json").sort.filter_map do |path| AtomicJson.read(path) end end |
#update_public(value, state, audit: nil) ⇒ Object
The public lifecycle projection: merge-on-write — every write merges over the existing projection (state, timestamps and audit fields win; effect/escalation fields recorded by an earlier write survive every later lifecycle write, review F-R2 on W696), 0440, owner = requester, group = the control group. Never carries answer content.
353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 353 def update_public(value, state, audit: nil) request_id = safe_id(value["id"].to_s) public = { "id" => request_id, "work" => value["work"].to_s, "attempt" => value["attempt"].to_s, "project" => value["project"].to_s, "harness" => value["harness"].to_s, "kind" => value["kind"].to_s, "state" => state, "created_at" => Integer(value["created_at"]), "updated_at" => Time.now.to_i } public["incarnation"] = value["incarnation"] if value["incarnation"] public["effect_state"] = value["effect_state"] if value["effect_state"] public.update(audit) if audit existing = AtomicJson.read(public_path(request_id)) public = existing.merge(public) if existing.is_a?(Hash) ownership = ownership_for(value) AtomicJson.call( public_path(request_id), public, mode: PUBLIC_MODE, ownership: ownership, ownership_strategy: @ownership ) end |
#with_request_lock(request_id) ⇒ Object
The per-request flock: the transaction boundary shared with deliver/consume/cancel and the lab-side stop protocol. A record that vanished under a concurrent cancel is a StateError, never a raw Errno escape (review F5 on W696).
390 391 392 393 394 395 396 397 398 399 400 401 402 |
# File 'lib/ace/hitl/lifecycle/store.rb', line 390 def with_request_lock(request_id) path = request_path(request_id) begin File.open(path, "r") do |file| file.flock(File::LOCK_EX) yield path ensure file.flock(File::LOCK_UN) end rescue Errno::ENOENT raise StateError, "unknown or invalid HITL request" end end |