Class: PWN::AI::Agent::Loop::Steering

Inherits:
Object
  • Object
show all
Defined in:
lib/pwn/ai/agent/loop.rb

Overview

Request-owned input broker; only the provider window is interruptible.

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

Instance Method Summary collapse

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

Raises:



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(messages, 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: 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