Class: LittleGhost::Providers::Bedrock::StreamNormalizer

Inherits:
Object
  • Object
show all
Defined in:
lib/little_ghost/providers/bedrock.rb

Overview

:nodoc:

Instance Method Summary collapse

Constructor Details

#initialize(model:) ⇒ StreamNormalizer

Returns a new instance of StreamNormalizer.



356
357
358
359
360
361
362
363
364
365
# File 'lib/little_ghost/providers/bedrock.rb', line 356

def initialize(model:)
  @model = model
  @message_id = nil
  @text = +""
  @reasoning_blocks = {}
  @tool_calls = {}
  @usage = Usage.new
  @stop_reason = nil
  @finished = false
end

Instance Method Details

#consume(event) ⇒ Object



367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
# File 'lib/little_ghost/providers/bedrock.rb', line 367

def consume(event)
  if event["event_type"]
    type = event["event_type"].to_s
    payload = event.except("event_type")
  else
    type, payload = event.first
  end
  case type
  when "message_start"
    @message_id = payload["role"]
    [StreamEvent.build(:message_start, id: nil, model: @model)]
  when "content_block_start"
    content_start(payload)
  when "content_block_delta"
    content_delta(payload)
  when "content_block_stop"
    content_stop(payload)
  when "message_stop"
    @stop_reason = normalize_stop(payload["stop_reason"])
    @terminal = true
    []
  when "metadata"
    (payload)
  when *TRANSIENT_STREAM_ERRORS, "validation_exception"
    message = payload.is_a?(Hash) ? payload["message"].to_s : ""
    message = "Bedrock returned #{type}" if message.empty?
    raise StreamError.new(message, event_type: type)
  else
    []
  end
end

#finishObject



399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
# File 'lib/little_ghost/providers/bedrock.rb', line 399

def finish
  return [] if @finished
  raise ProtocolError, "Bedrock stream ended before message_stop" unless @terminal

  @finished = true
  blocks = []
  @reasoning_blocks.sort.each do |_index, reasoning|
    if !reasoning[:redacted_content].empty?
      blocks << Content::Reasoning.new(redacted_content: reasoning[:redacted_content])
    elsif !reasoning[:text].empty?
      signature = reasoning[:signature]
      blocks << Content::Reasoning.new(
        text: reasoning[:text],
        signature: signature.empty? ? nil : signature
      )
    end
  end
  blocks << Content::Text.new(text: @text) unless @text.empty?
  @tool_calls.sort.each do |_index, tool|
    input = tool[:arguments].empty? ? {} : JSON.parse(tool[:arguments])
    blocks << Content::ToolUse.new(id: tool[:id], name: tool[:name], input:)
  end
  response = ModelResponse.new(
    message: Message.new(role: :assistant, content: blocks),
    stop_reason: @stop_reason || (@tool_calls.empty? ? :end_turn : :tool_use),
    usage: @usage,
    metadata: {model: @model}
  )
  [StreamEvent.build(:message_stop, response:)]
rescue JSON::ParserError, ArgumentError => error
  raise MalformedToolCallError, "Bedrock returned an invalid tool call: #{error.message}"
end