Class: PWN::AI::Agent::Loop::Steering
- Inherits:
-
Object
- Object
- PWN::AI::Agent::Loop::Steering
- Defined in:
- lib/pwn/ai/agent/loop.rb
Overview
Request-owned input broker; only the provider window is interruptible.
Direct Known Subclasses
Plugins::REPL::AIConsole::Control, Plugins::REPL::AISwarm::JobControl
Defined Under Namespace
Classes: ModelCancelled, Restart
Constant Summary collapse
- INPUT_LOCK =
Mutex.new
- INPUT_OWNERS =
rubocop:disable Style/MutableConstant -- guarded by INPUT_LOCK
{}
Instance Attribute Summary collapse
-
#reader ⇒ Object
readonly
Returns the value of attribute reader.
Instance Method Summary collapse
- #checkpoint(messages:, phase: :boundary) ⇒ Object
-
#initialize(input:, output:) ⇒ Steering
constructor
A new instance of Steering.
- #model(messages:, &block) ⇒ Object
- #submit(text) ⇒ Object
- #with_reader ⇒ Object
Constructor Details
#initialize(input:, output:) ⇒ Steering
Returns a new instance of Steering.
82 83 84 85 86 87 88 89 |
# File 'lib/pwn/ai/agent/loop.rb', line 82 def initialize(input:, output:) @input = input @output = output @owner = Thread.current @mutex = Mutex.new @pending = [] @phase = :boundary end |
Instance Attribute Details
#reader ⇒ Object (readonly)
Returns the value of attribute reader.
80 81 82 |
# File 'lib/pwn/ai/agent/loop.rb', line 80 def reader @reader end |
Instance Method Details
#checkpoint(messages:, phase: :boundary) ⇒ Object
113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 |
# File 'lib/pwn/ai/agent/loop.rb', line 113 def checkpoint(messages:, phase: :boundary) instructions = @mutex.synchronize do if @pending.empty? @phase = phase @model_signal_sent = false if phase == :model nil else @phase = :boundary @pending.shift(@pending.length) end end return unless instructions yield if block_given? raise Restart.new(, instructions) end |
#model(messages:, &block) ⇒ Object
130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 |
# File 'lib/pwn/ai/agent/loop.rb', line 130 def model(messages:, &block) # Drain an in-flight signal inside this rescue, never in a tool. Thread.handle_interrupt(ModelCancelled => :never) do checkpoint(messages: , phase: :model) Thread.handle_interrupt(ModelCancelled => :immediate, &block) ensure @mutex.synchronize { @phase = :boundary } begin Thread.handle_interrupt(ModelCancelled => :immediate) { nil } rescue ModelCancelled nil end end rescue ModelCancelled nil end |
#submit(text) ⇒ Object
91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 |
# File 'lib/pwn/ai/agent/loop.rb', line 91 def submit(text) @mutex.synchronize do if text.strip.empty? @output.puts('[pwn-ai] Usage: /steer <instruction>') return :empty end if @phase == :finished @output.puts('[pwn-ai] /steer: no active request') return :idle end @pending << text notice = @phase == :tool ? 'wait until the current tool finishes; already-started work is not undone' : 'discarding the pending model result at the next safe boundary' @output.puts("[pwn-ai] steering queued: #{notice}") @output.flush if @phase == :model && !@model_signal_sent @model_signal_sent = true @owner.raise(ModelCancelled) end end :queued end |
#with_reader ⇒ Object
147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 |
# File 'lib/pwn/ai/agent/loop.rb', line 147 def with_reader require 'io/console' input_id = @input.fileno claimed = INPUT_LOCK.synchronize do raise ArgumentError, 'steering input is already owned by an active request' if INPUT_OWNERS.key?(input_id) INPUT_OWNERS[input_id] = self end original_flags = @input.fcntl(Fcntl::F_GETFL) source = @input.dup forwarded, @forward = IO.pipe @stop_read, @stop_write = IO.pipe # Preserve descriptor identity for STDIN and child-process prompts. # Non-command lines go to that descriptor, never into Pry's buffer. @input.reopen(forwarded) forwarded.close run = proc do @reader = Thread.new { read_lines(source) } yield ensure @stop_write.write('.') unless @stop_write.closed? @reader&.join end source.tty? ? source.cooked(&run) : run.call ensure if source && !source.closed? @input.reopen(source) @input.fcntl(Fcntl::F_SETFL, original_flags) end [source, forwarded, @forward, @stop_read, @stop_write].compact.each { |io| io.close unless io.closed? } INPUT_LOCK.synchronize { INPUT_OWNERS.delete(input_id) } if claimed end |