Class: Ace::Hitl::Hermes::Organisms::HermesBox
- Inherits:
-
Object
- Object
- Ace::Hitl::Hermes::Organisms::HermesBox
- Defined in:
- lib/ace/hitl/hermes/organisms/hermes_box.rb
Overview
The folder interface per channel (spec 8wm.t.vs1 §9): publish (atomic write + collision retry), poll (validate + quarantine), ack (deletion IS the ACK), age (retry clock), identical? (redelivery identity). The Box is role-agnostic: the hermes side polls for questions, the lab side (labd) polls for answers.
Defined Under Namespace
Classes: PollResult, QuarantinedFile
Instance Attribute Summary collapse
-
#channel ⇒ Object
readonly
Returns the value of attribute channel.
Instance Method Summary collapse
-
#ack(id) ⇒ Object
Deletion is the ACK; idempotent (an absent file reports :already_acked, never an error).
- #address(id) ⇒ Object
-
#age(id) ⇒ Object
Seconds since the message file was written (mtime); nil when the file is gone (already acked).
-
#identical?(id, bytes) ⇒ Boolean
True when the file for
idcurrently holds exactlybytes- the identity check every redelivery must pass (retries never duplicate answers). -
#initialize(channel:, registry: nil, id_generator: -> { SecureRandom.hex(6) }, notifier: nil, euid_provider: -> { Process.euid }) ⇒ HermesBox
constructor
A new instance of HermesBox.
-
#poll ⇒ Object
Validate + collect everything currently in the folder.
-
#publish(kind:, body:, sender:, timestamp:, id: nil) ⇒ Object
Atomic write of a validated message.
Constructor Details
#initialize(channel:, registry: nil, id_generator: -> { SecureRandom.hex(6) }, notifier: nil, euid_provider: -> { Process.euid }) ⇒ HermesBox
Returns a new instance of HermesBox.
29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 29 def initialize(channel:, registry: nil, id_generator: -> { SecureRandom.hex(6) }, notifier: nil, euid_provider: -> { Process.euid }) @channel = if channel.is_a?(Molecules::HermesChannels::Channel) channel elsif registry registry.resolve(channel) else raise ContractError, "hermes box needs a Channel or a registry to resolve the channel name" end @id_generator = id_generator @notifier = notifier || ->(_line) {} @euid_provider = euid_provider end |
Instance Attribute Details
#channel ⇒ Object (readonly)
Returns the value of attribute channel.
27 28 29 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 27 def channel @channel end |
Instance Method Details
#ack(id) ⇒ Object
Deletion is the ACK; idempotent (an absent file reports :already_acked, never an error). The delete is a folder write, so the no-root gate applies exactly like on publish.
160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 160 def ack(id) if @euid_provider.call.zero? raise RootUserError, "hermes writes without root: refusing ack deletion as euid 0 " \ "(#{Molecules::HermesContract.(@channel.folder, id)})" end Molecules::HermesContract.verify_folder!(@channel.folder) path = Molecules::HermesContract.(@channel.folder, id) if File.exist?(path) File.delete(path) notify(:acked, address(id)) :acked else :already_acked end end |
#address(id) ⇒ Object
45 46 47 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 45 def address(id) @channel.address(id) end |
#age(id) ⇒ Object
Seconds since the message file was written (mtime); nil when the file is gone (already acked). The undeleted-file retry clock (spec §8).
180 181 182 183 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 180 def age(id) path = Molecules::HermesContract.(@channel.folder, id) File.exist?(path) ? Time.now - File.mtime(path) : nil end |
#identical?(id, bytes) ⇒ Boolean
True when the file for id currently holds exactly bytes -
the identity check every redelivery must pass (retries never
duplicate answers).
188 189 190 191 192 193 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 188 def identical?(id, bytes) path = Molecules::HermesContract.(@channel.folder, id) File.exist?(path) && read_capped(path) == bytes rescue Errno::ENOENT false end |
#poll ⇒ Object
Validate + collect everything currently in the folder.
Valid files become Messages; files failing the fail-closed
gate are quarantined (never delivered). Foreign names
(dotfiles, tmp leftovers, non-<token>.json) are ignored.
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 152 153 154 155 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 105 def poll Molecules::HermesContract.verify_folder!(@channel.folder) = [] quarantined = [] Dir.children(@channel.folder).sort.each do |name| id = Molecules::HermesContract.parse_file_name(name) next if id.nil? # dotfiles, tmp leftovers, foreign names path = File.join(@channel.folder, name) # Opened with O_NOFOLLOW and verified regular on the open # descriptor: a top-level symlink must never be followed - # it could import content from outside the channel or # reverse a quarantine's terminal state (review 8wq2ztu3 # on PR#336). begin file = File.open(path, File::RDONLY | File::NOFOLLOW) rescue Errno::ENOENT next # concurrently acked between listing and reading rescue Errno::ELOOP quarantined << quarantine(path, id, "message file is a symlink") next end begin begin regular = file.stat.file? rescue Errno::ENOENT next # concurrently acked after the open end unless regular quarantined << quarantine(path, id, "message file is not a regular file") next end bytes = file.read(Molecules::HermesContract::MAX_BYTES + 1) || "" hash = Molecules::HermesFormats.decode!(bytes) = Molecules::HermesMessage.from_hash(hash, filename_id: id) rescue Error => e quarantined << quarantine(path, id, e.) next ensure file.close end notify(:question_received, address(.id)) if .question? << end PollResult.new(messages: , quarantined: quarantined) end |
#publish(kind:, body:, sender:, timestamp:, id: nil) ⇒ Object
Atomic write of a validated message. A GENERATED id is
regenerated (bounded, notified) on a filename collision; an
EXPLICITLY supplied id never gets regenerated - it asserts
correlation (e.g. the answer of question q-1 is q-1.json),
so a collision there fails loudly instead of silently breaking
the question -> answer pairing. Returns the published Message.
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 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 |
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 55 def publish(kind:, body:, sender:, timestamp:, id: nil) explicit_id = !id.nil? attempts = 0 loop do = ( kind: kind, id: id || @id_generator.call, sender: sender, timestamp: , body: body ) # Producer fail-closed gate: the exact BYTES about to be # written must survive the consumers' decode gate (UTF-8 + # 64 KiB cap + schema id), and the decoded envelope must # satisfy the exact message contract, BEFORE any disk # write - so an envelope its own poll would terminally # quarantine can never be published (review 8wq2zttx on # PR#336). envelope = .to_json Molecules::HermesFormats.decode!(envelope) Molecules::HermesMessage.from_hash( JSON.parse(envelope), filename_id: .id ) begin Molecules::HermesAtomicWriter.write( Molecules::HermesContract.(@channel.folder, .id), envelope, euid_provider: @euid_provider ) rescue CollisionError raise if explicit_id attempts += 1 unless Molecules::HermesRetryPolicy.collision_retry?(attempts) raise end notify(:retry_scheduled, address(.id), attempt: attempts, max_attempts: Molecules::HermesRetryPolicy::DEFAULT_MAX_ATTEMPTS, policy: "collision") id = nil # force a fresh id on the next iteration next end notify(:answer_written, address(.id)) if .answer? return end end |