Module: ClaudeAgentServer::Services::SseStream

Defined in:
lib/claude_agent_server/services/sse_stream.rb

Defined Under Namespace

Classes: StreamBody, StreamWriter

Class Method Summary collapse

Class Method Details

.format_sse(event, data, id: nil) ⇒ Object



40
41
42
43
44
45
46
# File 'lib/claude_agent_server/services/sse_stream.rb', line 40

def format_sse(event, data, id: nil)
  parts = []
  parts << "id: #{id}" unless id.nil?
  parts << "event: #{event}"
  parts << "data: #{JSON.generate(data)}"
  "#{parts.join("\n")}\n\n"
end

.stream_query(prompt:, options:) ⇒ Object



10
11
12
13
14
15
16
17
18
19
20
21
22
23
# File 'lib/claude_agent_server/services/sse_stream.rb', line 10

def stream_query(prompt:, options:)
  index = 0
  StreamBody.new do |stream|
    QueryExecutor.stream(prompt: prompt, options: options) do |message|
      serialized = MessageSerializer.serialize(message)
      event_type = serialized[:type] || 'message'
      stream.write(format_sse(event_type, serialized, id: index))
      index += 1
    end
    stream.write(format_sse('done', { status: 'complete' }))
  rescue IOError, Errno::EPIPE
    # Client disconnected — exit cleanly
  end
end

.stream_session(session_entry, last_event_id: nil) ⇒ Object



25
26
27
28
29
30
31
32
33
34
35
36
37
38
# File 'lib/claude_agent_server/services/sse_stream.rb', line 25

def stream_session(session_entry, last_event_id: nil)
  offset = last_event_id ? last_event_id.to_i + 1 : 0

  StreamBody.new do |stream|
    session_entry.subscribe(offset: offset) do |event|
      serialized = MessageSerializer.serialize(event.message)
      event_type = serialized[:type] || 'message'
      stream.write(format_sse(event_type, serialized, id: event.index))
    end
    stream.write(format_sse('done', { status: 'complete' }))
  rescue IOError, Errno::EPIPE
    # Client disconnected — session stays alive
  end
end