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
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
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
end
end
|