first commit
This commit is contained in:
commit
39f5468436
5 files changed
+681
No files matched your search
@@ -0,0 +1,5 @@
|
||||
target/
|
||||
*.exe
|
||||
*.d
|
||||
*.pdb
|
||||
*.lock
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
[package]
|
||||
name = "actual"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
eframe = "0.29"
|
||||
egui = "0.29"
|
||||
tokio = { version = "1.38", features = ["rt-multi-thread", "macros"] }
|
||||
reqwest = { version = "0.12", features = ["json", "stream"] }
|
||||
serde = { version = "1.0", features = ["derive"] }
|
||||
serde_json = "1.0"
|
||||
futures-util = "0.3"
|
||||
anyhow = "1.0"
|
||||
egui_commonmark = "0.18"
|
||||
image = "0.25"
|
||||
Binary file not shown.
Binary file not shown.
|
After Width: | Height: | Size: 18 KiB |
+660
@@ -0,0 +1,660 @@
|
||||
use anyhow::{anyhow, Result};
|
||||
use eframe::egui;
|
||||
use egui::{FontData, FontDefinitions, FontFamily};
|
||||
use egui_commonmark::{CommonMarkCache, CommonMarkViewer};
|
||||
use futures_util::StreamExt;
|
||||
use image::GenericImageView;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::{
|
||||
sync::{Arc, Mutex},
|
||||
thread,
|
||||
};
|
||||
|
||||
fn load_icon() -> egui::IconData {
|
||||
let bytes = include_bytes!("../assets/icon.png");
|
||||
let img = image::load_from_memory(bytes).expect("icon.png must be a valid image");
|
||||
let img = img.to_rgba8();
|
||||
let (w, h) = img.dimensions();
|
||||
egui::IconData {
|
||||
rgba: img.into_raw(),
|
||||
width: w,
|
||||
height: h,
|
||||
}
|
||||
}
|
||||
|
||||
fn main() -> eframe::Result<()> {
|
||||
let icon = load_icon();
|
||||
let native_options = eframe::NativeOptions {
|
||||
viewport: egui::ViewportBuilder::default().with_icon(icon),
|
||||
..Default::default()
|
||||
};
|
||||
eframe::run_native(
|
||||
"Actual Computer",
|
||||
native_options,
|
||||
Box::new(|_cc| Ok(Box::new(ChatApp::default()))),
|
||||
)
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct ChatApp {
|
||||
// Connection config
|
||||
base_url: String, // e.g. http://127.0.0.1:1234/v1
|
||||
api_key: String, // if your server ignores it, can be empty
|
||||
model: String, // e.g. mdl_2eb3129c7de908b3
|
||||
|
||||
// UI state
|
||||
user_input: String,
|
||||
messages: Vec<(Role, String)>, // (role, content)
|
||||
status: String,
|
||||
|
||||
// Streaming output (written by worker thread)
|
||||
streaming_text: Arc<Mutex<String>>,
|
||||
is_streaming: bool,
|
||||
|
||||
// theme stuff
|
||||
styled: bool,
|
||||
|
||||
// markdown
|
||||
md_cache: CommonMarkCache,
|
||||
|
||||
// fuck you
|
||||
stream_done_rx: Option<std::sync::mpsc::Receiver<anyhow::Result<()>>>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
enum Role {
|
||||
User,
|
||||
Assistant,
|
||||
}
|
||||
|
||||
impl Role {
|
||||
fn as_str(&self) -> &'static str {
|
||||
match self {
|
||||
Role::User => "user",
|
||||
Role::Assistant => "assistant",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for Role {
|
||||
fn default() -> Self {
|
||||
Role::User
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for ChatApp {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
base_url: "http://127.0.0.1:8080/v1".to_string(),
|
||||
api_key: "".to_string(),
|
||||
model: "mdl_2eb3129c7de908b3".to_string(),
|
||||
user_input: "".to_string(),
|
||||
messages: vec![],
|
||||
status: "Idle".to_string(),
|
||||
streaming_text: Arc::new(Mutex::new(String::new())),
|
||||
is_streaming: false,
|
||||
styled: false,
|
||||
md_cache: CommonMarkCache::default(),
|
||||
stream_done_rx: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl eframe::App for ChatApp {
|
||||
fn update(&mut self, ctx: &egui::Context, _frame: &mut eframe::Frame) {
|
||||
if !self.styled {
|
||||
ChatApp::apply_terminal_glass_style(ctx);
|
||||
self.styled = true;
|
||||
}
|
||||
|
||||
if let Some(rx) = &self.stream_done_rx {
|
||||
if let Ok(_res) = rx.try_recv() {
|
||||
// stream is finished -> commit assistant message and re-enable sending
|
||||
self.finish_streaming(); // your existing method
|
||||
self.stream_done_rx = None;
|
||||
}
|
||||
}
|
||||
|
||||
// keep repainting while streaming so text updates
|
||||
if self.is_streaming {
|
||||
ctx.request_repaint();
|
||||
}
|
||||
|
||||
/*egui::CentralPanel::default().show(ctx, |ui| {
|
||||
ui.heading("Actual Computer");
|
||||
ui.add_space(8.0);
|
||||
|
||||
ui.horizontal(|ui| {
|
||||
ui.label("Base URL:");
|
||||
ui.text_edit_singleline(&mut self.base_url);
|
||||
});
|
||||
ui.horizontal(|ui| {
|
||||
ui.label("Model:");
|
||||
ui.text_edit_singleline(&mut self.model);
|
||||
});
|
||||
ui.horizontal(|ui| {
|
||||
ui.label("API key:");
|
||||
ui.add(egui::TextEdit::singleline(&mut self.api_key).password(true));
|
||||
});
|
||||
|
||||
ui.add_space(8.0);
|
||||
ui.separator();
|
||||
ui.add_space(8.0);
|
||||
|
||||
// Chat history
|
||||
egui::ScrollArea::vertical()
|
||||
.auto_shrink([false, false])
|
||||
.max_height(320.0)
|
||||
.show(ui, |ui| {
|
||||
for (role, content) in &self.messages {
|
||||
let prefix = match role {
|
||||
Role::User => "You",
|
||||
Role::Assistant => "Assistant",
|
||||
};
|
||||
ui.label(format!("{prefix}: {content}"));
|
||||
ui.add_space(4.0);
|
||||
}
|
||||
|
||||
// Streaming assistant text (in-progress)
|
||||
let cur = self.streaming_text.lock().unwrap().clone();
|
||||
if !cur.is_empty() {
|
||||
ui.label(format!("Assistant (streaming): {cur}"));
|
||||
}
|
||||
});
|
||||
|
||||
ui.add_space(8.0);
|
||||
ui.separator();
|
||||
|
||||
// Input area
|
||||
ui.label("Message:");
|
||||
let input = ui.add(
|
||||
egui::TextEdit::multiline(&mut self.user_input)
|
||||
.desired_rows(3)
|
||||
.hint_text("Type a message and press Send…"),
|
||||
);
|
||||
|
||||
ui.horizontal(|ui| {
|
||||
let send_clicked = ui
|
||||
.add_enabled(!self.is_streaming, egui::Button::new("Send"))
|
||||
.clicked();
|
||||
|
||||
let enter_pressed = input.lost_focus()
|
||||
&& ui.input(|i| i.key_pressed(egui::Key::Enter))
|
||||
&& !ui.input(|i| i.modifiers.shift);
|
||||
|
||||
if (send_clicked || enter_pressed) && !self.user_input.trim().is_empty() {
|
||||
let user_text = self.user_input.trim().to_string();
|
||||
self.user_input.clear();
|
||||
|
||||
// append user message to history
|
||||
self.messages.push((Role::User, user_text.clone()));
|
||||
|
||||
// clear streaming buffer
|
||||
*self.streaming_text.lock().unwrap() = String::new();
|
||||
|
||||
// start background streaming request
|
||||
self.status = "Streaming…".to_string();
|
||||
self.is_streaming = true;
|
||||
|
||||
let base_url = self.base_url.clone();
|
||||
let api_key = self.api_key.clone();
|
||||
let model = self.model.clone();
|
||||
let history = self
|
||||
.messages
|
||||
.iter()
|
||||
.map(|(r, c)| ChatMessage {
|
||||
role: r.as_str().to_string(),
|
||||
content: c.clone(),
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let streaming_text = self.streaming_text.clone();
|
||||
let streaming_text_clone = Arc::clone(&streaming_text);
|
||||
let done_flag = Arc::new(Mutex::new(false));
|
||||
let done_flag2 = done_flag.clone();
|
||||
|
||||
thread::spawn(move || {
|
||||
let rt = tokio::runtime::Runtime::new().expect("tokio runtime");
|
||||
let res = rt.block_on(stream_chat_openai_compat(
|
||||
&base_url,
|
||||
&api_key,
|
||||
&model,
|
||||
history,
|
||||
streaming_text_clone,
|
||||
));
|
||||
|
||||
// mark done; we can't mutate UI state directly from this thread,
|
||||
// but we'll encode done-ness in the buffer as a sentinel line on error.
|
||||
if let Err(e) = res {
|
||||
let mut buf = streaming_text.lock().unwrap();
|
||||
buf.push_str(&format!("\n\n[error: {e}]"));
|
||||
}
|
||||
|
||||
*done_flag2.lock().unwrap() = true;
|
||||
});
|
||||
|
||||
// We can't join the thread; instead we will check for completion by heuristic:
|
||||
// when the stream ends it stops appending; we’ll also toggle off when user clicks "Stop".
|
||||
// For simplicity, we implement a "Finish" check below using a small delay heuristic.
|
||||
self.status = "Streaming… (close window to quit)".to_string();
|
||||
|
||||
// store done flag in status? (kept minimal; see note below)
|
||||
// In a more complete app, you'd use channels to signal completion.
|
||||
std::mem::drop(done_flag);
|
||||
}
|
||||
|
||||
if ui
|
||||
.add_enabled(self.is_streaming, egui::Button::new("Stop (local)"))
|
||||
.clicked()
|
||||
{
|
||||
// This doesn't cancel the HTTP request (minimal app).
|
||||
// It only stops repainting and commits what we have.
|
||||
self.finish_streaming();
|
||||
}
|
||||
});
|
||||
|
||||
ui.add_space(4.0);
|
||||
ui.label(format!("Status: {}", self.status));
|
||||
|
||||
// If we are streaming but the server already ended (heuristic):
|
||||
// if buffer ends with "[DONE]" we finish automatically.
|
||||
if self.is_streaming {
|
||||
let cur = self.streaming_text.lock().unwrap().clone();
|
||||
if cur.contains("\n[DONE]") || cur.ends_with("[DONE]") {
|
||||
self.finish_streaming();
|
||||
}
|
||||
}
|
||||
});*/
|
||||
|
||||
egui::TopBottomPanel::top("top").show(ctx, |ui| {
|
||||
ui.horizontal(|ui| {
|
||||
ui.label(
|
||||
egui::RichText::new("Actual Computer")
|
||||
.monospace()
|
||||
.color(egui::Color32::from_rgb(40, 200, 160)),
|
||||
);
|
||||
ui.separator();
|
||||
ui.label(
|
||||
egui::RichText::new(&self.status)
|
||||
.monospace()
|
||||
.color(egui::Color32::from_gray(180)),
|
||||
);
|
||||
|
||||
ui.with_layout(egui::Layout::right_to_left(egui::Align::Center), |ui| {
|
||||
if ui.button("Clear").clicked() && !self.is_streaming {
|
||||
self.messages.clear();
|
||||
*self.streaming_text.lock().unwrap() = String::new();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
// Connection settings stay minimal but hidden
|
||||
egui::CollapsingHeader::new("Connection")
|
||||
.default_open(false)
|
||||
.show(ui, |ui| {
|
||||
ui.horizontal(|ui| {
|
||||
ui.label("Base URL");
|
||||
ui.text_edit_singleline(&mut self.base_url);
|
||||
});
|
||||
ui.horizontal(|ui| {
|
||||
ui.label("Model");
|
||||
ui.text_edit_singleline(&mut self.model);
|
||||
});
|
||||
ui.horizontal(|ui| {
|
||||
ui.label("API Key");
|
||||
ui.add(egui::TextEdit::singleline(&mut self.api_key).password(true));
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
egui::TopBottomPanel::bottom("composer").show(ctx, |ui| {
|
||||
ui.horizontal(|ui| {
|
||||
let send = ui
|
||||
.add_enabled(!self.is_streaming, egui::Button::new("SEND"))
|
||||
.clicked();
|
||||
|
||||
let edit = ui.add(
|
||||
egui::TextEdit::singleline(&mut self.user_input)
|
||||
.hint_text("Type here… (enter to send)"),
|
||||
);
|
||||
|
||||
let enter = edit.lost_focus() && ui.input(|i| i.key_pressed(egui::Key::Enter));
|
||||
if (send || enter) && !self.user_input.trim().is_empty() && !self.is_streaming {
|
||||
let user_text = self.user_input.trim().to_string();
|
||||
self.user_input.clear();
|
||||
|
||||
// append user message to history
|
||||
self.messages.push((Role::User, user_text.clone()));
|
||||
|
||||
// clear streaming buffer
|
||||
*self.streaming_text.lock().unwrap() = String::new();
|
||||
|
||||
// start background streaming request
|
||||
self.status = "Streaming…".to_string();
|
||||
self.is_streaming = true;
|
||||
|
||||
let base_url = self.base_url.clone();
|
||||
let api_key = self.api_key.clone();
|
||||
let model = self.model.clone();
|
||||
let history = self
|
||||
.messages
|
||||
.iter()
|
||||
.map(|(r, c)| ChatMessage {
|
||||
role: r.as_str().to_string(),
|
||||
content: c.clone(),
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
let streaming_text = self.streaming_text.clone();
|
||||
let streaming_text_clone = Arc::clone(&streaming_text);
|
||||
let done_flag = Arc::new(Mutex::new(false));
|
||||
let done_flag2 = done_flag.clone();
|
||||
|
||||
let (done_tx, done_rx) = std::sync::mpsc::channel();
|
||||
self.stream_done_rx = Some(done_rx);
|
||||
|
||||
thread::spawn(move || {
|
||||
let rt = tokio::runtime::Runtime::new().expect("tokio runtime");
|
||||
let res = rt.block_on(stream_chat_openai_compat(
|
||||
&base_url,
|
||||
&api_key,
|
||||
&model,
|
||||
history,
|
||||
streaming_text_clone,
|
||||
));
|
||||
|
||||
// still append error text if you want
|
||||
if let Err(e) = &res {
|
||||
let mut buf = streaming_text.lock().unwrap();
|
||||
buf.push_str(&format!("\n\n[error: {e}]"));
|
||||
}
|
||||
|
||||
// IMPORTANT: notify UI we're done (success or error)
|
||||
let _ = done_tx.send(res.map(|_| ()));
|
||||
});
|
||||
|
||||
// We can't join the thread; instead we will check for completion by heuristic:
|
||||
// when the stream ends it stops appending; we’ll also toggle off when user clicks "Stop".
|
||||
// For simplicity, we implement a "Finish" check below using a small delay heuristic.
|
||||
self.status = "Streaming… (close window to quit)".to_string();
|
||||
|
||||
// store done flag in status? (kept minimal; see note below)
|
||||
// In a more complete app, you'd use channels to signal completion.
|
||||
std::mem::drop(done_flag);
|
||||
}
|
||||
|
||||
if self.is_streaming {
|
||||
ui.add(egui::Spinner::new());
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
egui::CentralPanel::default().show(ctx, |ui| {
|
||||
egui::ScrollArea::vertical()
|
||||
.stick_to_bottom(true)
|
||||
.show(ui, |ui| {
|
||||
for (i, (role, content)) in self.messages.iter().enumerate() {
|
||||
message_panel(ui, &mut self.md_cache, *role, content, ("msg", i));
|
||||
}
|
||||
|
||||
// streaming message
|
||||
let cur = self.streaming_text.lock().unwrap().clone();
|
||||
if !cur.is_empty() {
|
||||
message_panel(ui, &mut self.md_cache, Role::Assistant, &cur, "streaming");
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
impl ChatApp {
|
||||
fn finish_streaming(&mut self) {
|
||||
let text = self.streaming_text.lock().unwrap().clone();
|
||||
let cleaned = text.replace("\n[Done]", "").trim().to_string();
|
||||
if !cleaned.is_empty() {
|
||||
self.messages.push((Role::Assistant, cleaned));
|
||||
}
|
||||
*self.streaming_text.lock().unwrap() = String::new();
|
||||
self.is_streaming = false;
|
||||
self.status = "Idle".to_string();
|
||||
}
|
||||
|
||||
fn apply_terminal_glass_style(ctx: &egui::Context) {
|
||||
// Dark visuals tuned for a "terminal glass" look
|
||||
let mut v = egui::Visuals::dark();
|
||||
v.override_text_color = Some(egui::Color32::from_gray(220));
|
||||
v.panel_fill = egui::Color32::from_rgb(10, 12, 16);
|
||||
v.window_fill = egui::Color32::from_rgb(10, 12, 16);
|
||||
v.faint_bg_color = egui::Color32::from_rgb(14, 16, 22);
|
||||
|
||||
// Accent (borders, selection, etc.)
|
||||
v.selection.bg_fill = egui::Color32::from_rgb(40, 200, 160);
|
||||
v.selection.stroke = egui::Stroke::new(1.0, egui::Color32::from_rgb(40, 200, 160));
|
||||
v.widgets.active.bg_fill = egui::Color32::from_rgb(14, 18, 26);
|
||||
v.widgets.hovered.bg_fill = egui::Color32::from_rgb(16, 22, 32);
|
||||
v.widgets.inactive.bg_fill = egui::Color32::from_rgb(12, 14, 20);
|
||||
|
||||
// Thin, crisp strokes
|
||||
v.widgets.noninteractive.fg_stroke = egui::Stroke::new(1.0, egui::Color32::from_gray(220));
|
||||
v.widgets.inactive.fg_stroke = egui::Stroke::new(1.0, egui::Color32::from_gray(220));
|
||||
v.widgets.hovered.fg_stroke = egui::Stroke::new(1.0, egui::Color32::from_rgb(40, 200, 160));
|
||||
v.widgets.active.fg_stroke = egui::Stroke::new(1.0, egui::Color32::from_rgb(40, 200, 160));
|
||||
|
||||
// Use square-ish corners for a "minimal tool" vibe
|
||||
v.window_rounding = egui::Rounding::same(6.0);
|
||||
v.menu_rounding = egui::Rounding::same(6.0);
|
||||
v.widgets.noninteractive.rounding = egui::Rounding::same(6.0);
|
||||
v.widgets.inactive.rounding = egui::Rounding::same(6.0);
|
||||
v.widgets.hovered.rounding = egui::Rounding::same(6.0);
|
||||
v.widgets.active.rounding = egui::Rounding::same(6.0);
|
||||
|
||||
ctx.set_visuals(v);
|
||||
|
||||
// Fonts: monospace-first. If you want a true bitmap look, add a bitmap TTF to assets.
|
||||
// Example choices: "PxPlus IBM VGA8", "Cozette", "Terminus", etc. (many are TTF/OTF variants).
|
||||
let mut fonts = FontDefinitions::default();
|
||||
|
||||
// OPTIONAL: include a custom font file.
|
||||
// Put a TTF at: assets/bitmap.ttf and uncomment below.
|
||||
//
|
||||
// fonts.font_data.insert(
|
||||
// "bitmap".to_string(),
|
||||
// FontData::from_static(include_bytes!("../assets/bitmap.ttf")),
|
||||
// );
|
||||
//
|
||||
// fonts
|
||||
// .families
|
||||
// .get_mut(&FontFamily::Monospace)
|
||||
// .unwrap()
|
||||
// .insert(0, "bitmap".to_string());
|
||||
//
|
||||
// fonts
|
||||
// .families
|
||||
// .get_mut(&FontFamily::Proportional)
|
||||
// .unwrap()
|
||||
// .insert(0, "bitmap".to_string());
|
||||
|
||||
fonts.font_data.insert(
|
||||
"bitmap".to_string(),
|
||||
FontData::from_static(include_bytes!("../assets/creep2.ttf")),
|
||||
);
|
||||
|
||||
// Even without a custom font, force monospace everywhere:
|
||||
// fonts.families.insert(FontFamily::Proportional, vec!["monospace".to_string()]);
|
||||
// fonts.families.insert(FontFamily::Monospace, vec!["monospace".to_string()]);
|
||||
|
||||
ctx.set_fonts(fonts);
|
||||
|
||||
// Slightly larger base size for readability
|
||||
ctx.set_style({
|
||||
let mut s = (*ctx.style()).clone();
|
||||
s.text_styles
|
||||
.insert(egui::TextStyle::Body, egui::FontId::monospace(15.0));
|
||||
s.text_styles
|
||||
.insert(egui::TextStyle::Button, egui::FontId::monospace(15.0));
|
||||
s.text_styles
|
||||
.insert(egui::TextStyle::Small, egui::FontId::monospace(13.0));
|
||||
s
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
fn message_panel(
|
||||
ui: &mut egui::Ui,
|
||||
cache: &mut egui_commonmark::CommonMarkCache,
|
||||
role: Role,
|
||||
text: &str,
|
||||
unique_id: impl std::hash::Hash,
|
||||
) {
|
||||
let accent = egui::Color32::from_rgb(40, 200, 160);
|
||||
let border = egui::Stroke::new(1.0, accent);
|
||||
let fill = match role {
|
||||
Role::User => egui::Color32::from_rgb(12, 18, 22),
|
||||
Role::Assistant => egui::Color32::from_rgb(12, 12, 18),
|
||||
};
|
||||
|
||||
egui::Frame::none()
|
||||
.fill(fill)
|
||||
.stroke(border)
|
||||
.inner_margin(egui::Margin::symmetric(10.0, 8.0))
|
||||
.rounding(egui::Rounding::same(6.0))
|
||||
.show(ui, |ui| {
|
||||
let header = match role {
|
||||
Role::User => "YOU",
|
||||
Role::Assistant => "AI",
|
||||
};
|
||||
ui.label(
|
||||
egui::RichText::new(header)
|
||||
.color(accent)
|
||||
.monospace()
|
||||
.size(12.0),
|
||||
);
|
||||
ui.add_space(6.0);
|
||||
|
||||
// IMPORTANT: give each markdown render a stable unique id
|
||||
ui.push_id(unique_id, |ui| {
|
||||
let viewer = egui_commonmark::CommonMarkViewer::new();
|
||||
viewer.show(ui, cache, text);
|
||||
});
|
||||
});
|
||||
|
||||
ui.add_space(8.0);
|
||||
}
|
||||
|
||||
#[derive(Serialize, Clone)]
|
||||
struct ChatMessage {
|
||||
role: String,
|
||||
content: String,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct ChatRequest {
|
||||
model: String,
|
||||
messages: Vec<ChatMessage>,
|
||||
stream: bool,
|
||||
|
||||
// Optional OpenAI-style params (keep minimal but matches your example shape)
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
temperature: Option<f32>,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
top_p: Option<f32>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
struct StreamChunk {
|
||||
choices: Vec<Choice>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
struct Choice {
|
||||
delta: Delta,
|
||||
#[allow(dead_code)]
|
||||
finish_reason: Option<String>,
|
||||
#[allow(dead_code)]
|
||||
index: usize,
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
struct Delta {
|
||||
#[serde(default)]
|
||||
content: Option<String>,
|
||||
#[serde(default)]
|
||||
role: Option<String>,
|
||||
}
|
||||
|
||||
// Streams an OpenAI-compatible chat.completions SSE response.
|
||||
// Expects lines like: `data: {...json...}` and `data: [DONE]`.
|
||||
async fn stream_chat_openai_compat(
|
||||
base_url: &str,
|
||||
api_key: &str,
|
||||
model: &str,
|
||||
messages: Vec<ChatMessage>,
|
||||
streaming_text: Arc<Mutex<String>>,
|
||||
) -> Result<()> {
|
||||
let url = format!("{}/chat/completions", base_url.trim_end_matches('/'));
|
||||
|
||||
let req_body = ChatRequest {
|
||||
model: model.to_string(),
|
||||
messages,
|
||||
stream: true,
|
||||
temperature: Some(0.7),
|
||||
top_p: Some(0.95),
|
||||
};
|
||||
|
||||
let client = reqwest::Client::new();
|
||||
let mut req = client.post(url).json(&req_body);
|
||||
|
||||
// OpenAI-compatible auth header (safe even if ignored by your server)
|
||||
if !api_key.trim().is_empty() {
|
||||
req = req.bearer_auth(api_key.trim());
|
||||
}
|
||||
|
||||
let resp = req.send().await?;
|
||||
if !resp.status().is_success() {
|
||||
let status = resp.status();
|
||||
let text = resp.text().await.unwrap_or_default();
|
||||
return Err(anyhow!("HTTP {status}: {text}"));
|
||||
}
|
||||
|
||||
let mut stream = resp.bytes_stream();
|
||||
|
||||
// We'll parse by lines; simplest: accumulate into a String and split on '\n'
|
||||
let mut buffer = String::new();
|
||||
|
||||
while let Some(item) = stream.next().await {
|
||||
let bytes = item?;
|
||||
let chunk_str = String::from_utf8_lossy(&bytes);
|
||||
buffer.push_str(&chunk_str);
|
||||
|
||||
while let Some(pos) = buffer.find('\n') {
|
||||
let mut line = buffer[..pos].to_string();
|
||||
buffer = buffer[pos + 1..].to_string();
|
||||
|
||||
line = line.trim().to_string();
|
||||
if line.is_empty() {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(rest) = line.strip_prefix("data:") {
|
||||
let data = rest.trim();
|
||||
|
||||
if data == "[DONE]" {
|
||||
let mut out = streaming_text.lock().unwrap();
|
||||
out.push_str("\n");
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Parse JSON chunk
|
||||
let parsed: StreamChunk = serde_json::from_str(data)?;
|
||||
if let Some(choice) = parsed.choices.get(0) {
|
||||
if let Some(token) = &choice.delta.content {
|
||||
let mut out = streaming_text.lock().unwrap();
|
||||
out.push_str(token);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
Reference in new issue
Block a user