Class: Ace::Hitl::Lifecycle::Store

Inherits:
Object
  • Object
show all
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

Instance Method Summary collapse

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.

Raises:

  • (ArgumentError)


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).

Raises:



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 ------------------------------------------------

Raises:



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
  options = Array(options).map(&:strip)
  if options.length > MAX_OPTIONS || 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 options.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" => 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

Raises:



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

Raises:



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