Class: Ace::Hitl::Hermes::Organisms::HermesBox

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

Instance Method Summary collapse

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.message_path(@channel.folder, id)})"
  end
  Molecules::HermesContract.verify_folder!(@channel.folder)
  path = Molecules::HermesContract.message_path(@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.message_path(@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).

Returns:

  • (Boolean)


188
189
190
191
192
193
# File 'lib/ace/hitl/hermes/organisms/hermes_box.rb', line 188

def identical?(id, bytes)
  path = Molecules::HermesContract.message_path(@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)
  messages = []
  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)
      message = Molecules::HermesMessage.from_hash(hash, filename_id: id)
    rescue Error => e
      quarantined << quarantine(path, id, e.message)
      next
    ensure
      file.close
    end

    notify(:question_received, address(message.id)) if message.question?
    messages << message
  end

  PollResult.new(messages: 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
    message = build_message(
      kind: kind, id: id || @id_generator.call,
      sender: sender, timestamp: 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 = message.to_json
    Molecules::HermesFormats.decode!(envelope)
    Molecules::HermesMessage.from_hash(
      JSON.parse(envelope), filename_id: message.id
    )
    begin
      Molecules::HermesAtomicWriter.write(
        Molecules::HermesContract.message_path(@channel.folder, message.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(message.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(message.id)) if message.answer?
    return message
  end
end