1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
|
mod markdown;
mod stream;
pub use self::markdown::{MarkdownRender, RenderOptions};
use self::stream::{markdown_stream, raw_stream};
use crate::client::Client;
use crate::config::GlobalConfig;
use crate::utils::AbortSignal;
use anyhow::{Context, Result};
use crossbeam::channel::{unbounded, Sender};
use crossbeam::sync::WaitGroup;
use is_terminal::IsTerminal;
use nu_ansi_term::{Color, Style};
use std::io::stdout;
use std::thread::spawn;
pub fn render_stream(
input: &str,
client: &dyn Client,
config: &GlobalConfig,
abort: AbortSignal,
) -> Result<String> {
let wg = WaitGroup::new();
let wg_cloned = wg.clone();
let render_options = config.read().get_render_options()?;
let mut stream_handler = {
let (tx, rx) = unbounded();
let abort_clone = abort.clone();
let highlight = config.read().highlight;
spawn(move || {
let run = move || {
if stdout().is_terminal() {
let mut render = MarkdownRender::init(render_options)?;
markdown_stream(&rx, &mut render, &abort)
} else {
raw_stream(&rx, &abort)
}
};
if let Err(err) = run() {
render_error(err, highlight);
}
drop(wg_cloned);
});
ReplyHandler::new(tx, abort_clone)
};
let ret = client.send_message_streaming(input, &mut stream_handler);
wg.wait();
let output = stream_handler.get_buffer().to_string();
match ret {
Ok(_) => {
println!();
Ok(output)
}
Err(err) => {
if !output.is_empty() {
println!();
}
Err(err)
}
}
}
pub fn render_error(err: anyhow::Error, highlight: bool) {
let err = format!("{err:?}");
if highlight {
let style = Style::new().fg(Color::Red);
eprintln!("{}", style.paint(err));
} else {
eprintln!("{err}");
}
}
pub struct ReplyHandler {
sender: Sender<ReplyEvent>,
buffer: String,
abort: AbortSignal,
}
impl ReplyHandler {
pub fn new(sender: Sender<ReplyEvent>, abort: AbortSignal) -> Self {
Self {
sender,
abort,
buffer: String::new(),
}
}
pub fn text(&mut self, text: &str) -> Result<()> {
debug!("ReplyText: {}", text);
if text.is_empty() {
return Ok(());
}
self.buffer.push_str(text);
let ret = self
.sender
.send(ReplyEvent::Text(text.to_string()))
.with_context(|| "Failed to send ReplyEvent:Text");
self.safe_ret(ret)?;
Ok(())
}
pub fn done(&mut self) -> Result<()> {
debug!("ReplyDone");
let ret = self
.sender
.send(ReplyEvent::Done)
.with_context(|| "Failed to send ReplyEvent::Done");
self.safe_ret(ret)?;
Ok(())
}
pub fn get_buffer(&self) -> &str {
&self.buffer
}
pub fn get_abort(&self) -> AbortSignal {
self.abort.clone()
}
fn safe_ret(&self, ret: Result<()>) -> Result<()> {
if ret.is_err() && self.abort.aborted() {
return Ok(());
}
ret
}
}
pub enum ReplyEvent {
Text(String),
Done,
}
|