Class: Ace::Hitl::Lifecycle::Overseer
- Inherits:
-
Object
- Object
- Ace::Hitl::Lifecycle::Overseer
- Defined in:
- lib/ace/hitl/lifecycle/overseer.rb
Overview
The Root Overseer's bounded response channel — the reverse address of the HITL conversation (spec 8wm.t.y21 §6; ports of overseer-send/overseer-pending/overseer-ack). Responses are bounded, type-tagged, and free of internal identifiers; the transport that relays them is NOT part of this surface.
Constant Summary collapse
- MAX_RESPONSE =
1200- TYPE_TAGS =
%w[[decyzja] [pytanie] [info]].freeze
- TYPE_TAG =
/\A\[(decyzja|pytanie|info)\]/i- FULL_SHA =
/\b[0-9a-fA-F]{40,64}\b/- ATTEMPT_REFERENCE =
/\bA-[0-9a-fA-F]{24}\b/- WORK_REFERENCE =
/\bW[0-9]{3,}\b/- TASK_REFERENCE =
/\b[0-9a-z]{3,}\.t\.[0-9a-z.]+\b/i- SHA_WORD =
/\bSHA(?:-?256)?\b/i
Instance Attribute Summary collapse
-
#outbox_dir ⇒ Object
readonly
Returns the value of attribute outbox_dir.
Instance Method Summary collapse
-
#ack(message_id) ⇒ Object
Root-only acknowledgement: the transport removes one response exactly once it has been relayed.
-
#initialize(outbox_dir:, overseer_user: "mo", ownership: AtomicJson::DEFAULT_OWNERSHIP, identity: Identity) ⇒ Overseer
constructor
A new instance of Overseer.
-
#pending ⇒ Object
Root-only drain view of the outbox.
-
#send_response(reader:, reply_to: "") ⇒ Object
Queue one bounded, type-tagged response (read from the reader, e.g. stdin).
Constructor Details
#initialize(outbox_dir:, overseer_user: "mo", ownership: AtomicJson::DEFAULT_OWNERSHIP, identity: Identity) ⇒ Overseer
Returns a new instance of Overseer.
30 31 32 33 34 35 36 |
# File 'lib/ace/hitl/lifecycle/overseer.rb', line 30 def initialize(outbox_dir:, overseer_user: "mo", ownership: AtomicJson::DEFAULT_OWNERSHIP, identity: Identity) @outbox_dir = Pathname.new(outbox_dir) @overseer_user = overseer_user @ownership = ownership @identity = identity end |
Instance Attribute Details
#outbox_dir ⇒ Object (readonly)
Returns the value of attribute outbox_dir.
28 29 30 |
# File 'lib/ace/hitl/lifecycle/overseer.rb', line 28 def outbox_dir @outbox_dir end |
Instance Method Details
#ack(message_id) ⇒ Object
Root-only acknowledgement: the transport removes one response exactly once it has been relayed.
89 90 91 92 93 94 95 96 97 98 99 |
# File 'lib/ace/hitl/lifecycle/overseer.rb', line 89 def ack() require_root!("overseer-ack") id = .to_s raise StateError, "invalid message id" unless Kinds::REQUEST_ID.match?(id) path = @outbox_dir.join("#{id}.json") raise StateError, "unknown Overseer response" unless path.exist? path.unlink {"id" => id, "acknowledged" => true} end |
#pending ⇒ Object
Root-only drain view of the outbox.
82 83 84 85 |
# File 'lib/ace/hitl/lifecycle/overseer.rb', line 82 def pending require_root!("overseer-pending") @outbox_dir.glob("*.json").sort.filter_map { |path| AtomicJson.read(path) } end |
#send_response(reader:, reply_to: "") ⇒ Object
Queue one bounded, type-tagged response (read from the reader, e.g. stdin). Only the overseer user may respond.
40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 |
# File 'lib/ace/hitl/lifecycle/overseer.rb', line 40 def send_response(reader:, reply_to: "") unless @identity.username == @overseer_user raise PermissionError, "only the Root Overseer can send channel responses" end reply_to = reply_to.to_s if !reply_to.empty? && (!/\A\d+\z/.match?(reply_to) || reply_to.length > 32) raise StateError, "invalid source message id" end response = reader.call(MAX_RESPONSE + 1).to_s.strip if response.empty? || response.length > MAX_RESPONSE || response.include?("\u0000") raise StateError, "response must contain 1-#{MAX_RESPONSE} characters" end unless TYPE_TAG.match?(response) raise StateError, "response must open with a type tag: #{TYPE_TAGS.join(", ")}" end if FULL_SHA.match?(response) || ATTEMPT_REFERENCE.match?(response) || WORK_REFERENCE.match?(response) || TASK_REFERENCE.match?(response) || SHA_WORD.match?(response) raise StateError, "rewrite for Captain without SHA, Work/Attempt/task IDs, hashes, or internal references" end = "msg-#{SecureRandom.hex(8)}" value = { "id" => , "reply_to_message_id" => reply_to, "response" => response, "created_at" => Time.now.to_i, "requester" => @overseer_user } AtomicJson.call(@outbox_dir.join("#{}.json"), value, mode: 0o600, ownership_strategy: @ownership) response.clear { "id" => , "queued" => true, "reply_to_message_id" => reply_to } end |