Class: PhaseoAgentSdk::StreamResult

Inherits:
Object
  • Object
show all
Defined in:
lib/phaseo_agent_sdk.rb

Instance Method Summary collapse

Constructor Details

#initialize(&work) ⇒ StreamResult

Returns a new instance of StreamResult.



35
36
37
38
39
40
# File 'lib/phaseo_agent_sdk.rb', line 35

def initialize(&work)
  @events=[];@mutex=Mutex.new;@condition=ConditionVariable.new;@done=false
  @worker=Thread.new do
	begin;@result=work.call(method(:push));rescue Exception=>error;@error=error;ensure;@mutex.synchronize{@done=true;@condition.broadcast};end
  end
end

Instance Method Details

#cancelObject



43
# File 'lib/phaseo_agent_sdk.rb', line 43

def cancel;@worker.kill if @worker&.alive?;@mutex.synchronize{@done=true;@condition.broadcast};end

#full_streamObject



44
45
46
# File 'lib/phaseo_agent_sdk.rb', line 44

def full_stream
  Enumerator.new do|yielded|;index=0;loop do;event=nil;finished=false;@mutex.synchronize do;@condition.wait(@mutex) while index>=@events.length&&!@done;event=@events[index] if index<@events.length;index+=1 if event;finished=@done&&event.nil?;end;break if finished;yielded<<event if event;end;raise @error if @error;end
end

#item_streamObject



49
# File 'lib/phaseo_agent_sdk.rb', line 49

def item_stream = full_stream.lazy.select{|event|event[:type]=="response.item"}.map{|event|event[:item]}

#push(event) ⇒ Object



41
# File 'lib/phaseo_agent_sdk.rb', line 41

def push(event)=@mutex.synchronize{@events<<event;@condition.broadcast}

#reasoning_streamObject



48
# File 'lib/phaseo_agent_sdk.rb', line 48

def reasoning_stream = full_stream.lazy.select{|event|event[:type]=="response.reasoning.delta"}.map{|event|event[:delta].to_s}

#resultObject

Raises:

  • (@error)


42
# File 'lib/phaseo_agent_sdk.rb', line 42

def result;@worker.join;raise @error if @error;@result;end

#text_streamObject



47
# File 'lib/phaseo_agent_sdk.rb', line 47

def text_stream = full_stream.lazy.select{|event|event[:type]=="response.output_text.delta"}.map{|event|event[:delta].to_s}