feat: store chat_id to allow multiple chats, use enums in tauri
communication
This commit is contained in:
@@ -1,5 +0,0 @@
|
||||
CREATE TABLE IF NOT EXISTS messages (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
text TEXT NOT NULL,
|
||||
is_user BOOL NOT NULL
|
||||
);
|
||||
12
crates/daemon/migrations/20260222_1_init.sql
Normal file
12
crates/daemon/migrations/20260222_1_init.sql
Normal file
@@ -0,0 +1,12 @@
|
||||
CREATE TABLE IF NOT EXISTS messages (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
chat_id INTEGER NOT NULL,
|
||||
timestamp DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
text TEXT NOT NULL,
|
||||
is_user BOOL NOT NULL,
|
||||
|
||||
UNIQUE(id, chat_id)
|
||||
);
|
||||
|
||||
CREATE INDEX idx_message_timestamp ON messages(timestamp);
|
||||
CREATE INDEX idx_message_chat_id ON messages(chat_id);
|
||||
@@ -1,22 +1,29 @@
|
||||
use anyhow::Result;
|
||||
use directories::ProjectDirs;
|
||||
use sqlx::sqlite::SqliteConnectOptions;
|
||||
use sqlx::Row;
|
||||
use sqlx::SqlitePool;
|
||||
use sqlx::{Row, SqlitePool};
|
||||
use tokio::fs;
|
||||
use tonic::async_trait;
|
||||
|
||||
#[derive(Debug, sqlx::FromRow)]
|
||||
pub struct ChatMessageData {
|
||||
pub id: i64,
|
||||
pub chat_id: i64,
|
||||
pub text: String,
|
||||
pub is_user: bool,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
pub trait ChatRepository {
|
||||
async fn save_message(&self, text: &str, is_user: &bool) -> Result<ChatMessageData>;
|
||||
async fn get_latest_messages(&self) -> Result<Vec<ChatMessageData>>;
|
||||
async fn save_message(
|
||||
&self,
|
||||
text: &str,
|
||||
is_user: &bool,
|
||||
chat_id: &i32,
|
||||
) -> Result<ChatMessageData>;
|
||||
async fn get_latest_messages(&self, chat_id: &i32, count: &i32)
|
||||
-> Result<Vec<ChatMessageData>>;
|
||||
async fn get_chat_ids(&self) -> Result<Box<[i32]>>;
|
||||
}
|
||||
|
||||
pub struct SqliteChatRepository {
|
||||
@@ -51,31 +58,46 @@ impl SqliteChatRepository {
|
||||
|
||||
#[async_trait]
|
||||
impl ChatRepository for SqliteChatRepository {
|
||||
async fn save_message(&self, text: &str, is_user: &bool) -> Result<ChatMessageData> {
|
||||
async fn save_message(
|
||||
&self,
|
||||
text: &str,
|
||||
is_user: &bool,
|
||||
chat_id: &i32,
|
||||
) -> Result<ChatMessageData> {
|
||||
let result = sqlx::query_as::<_, ChatMessageData>(
|
||||
r#"
|
||||
INSERT INTO messages (text, is_user)
|
||||
VALUES (?, ?)
|
||||
RETURNING id, text, is_user
|
||||
INSERT INTO messages (text, is_user, chat_id)
|
||||
VALUES (?, ?, ?)
|
||||
RETURNING id, chat_id, text, is_user
|
||||
"#,
|
||||
)
|
||||
.bind(text)
|
||||
.bind(is_user)
|
||||
.bind(chat_id)
|
||||
.fetch_one(&self.pool)
|
||||
.await
|
||||
.inspect_err(|e| println!("sql error: {}", e))?;
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
async fn get_latest_messages(&self) -> Result<Vec<ChatMessageData>> {
|
||||
async fn get_latest_messages(
|
||||
&self,
|
||||
chat_id: &i32,
|
||||
count: &i32,
|
||||
) -> Result<Vec<ChatMessageData>> {
|
||||
// From all chat ids get the latest id.
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT * FROM (
|
||||
SELECT id, text, is_user
|
||||
FROM messages
|
||||
ORDER BY id DESC
|
||||
LIMIT 10
|
||||
) AS subquery ORDER BY id ASC"#,
|
||||
format!(
|
||||
r#"
|
||||
SELECT * FROM (
|
||||
SELECT id, chat_id, text, is_user
|
||||
FROM messages
|
||||
WHERE chat_id = {chat_id}
|
||||
ORDER BY id DESC
|
||||
LIMIT {count}
|
||||
) AS subquery ORDER BY id ASC;"#
|
||||
)
|
||||
.as_str(),
|
||||
)
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
@@ -85,11 +107,27 @@ impl ChatRepository for SqliteChatRepository {
|
||||
.into_iter()
|
||||
.map(|row| ChatMessageData {
|
||||
id: row.get(0),
|
||||
text: row.get(1),
|
||||
is_user: row.get(2),
|
||||
chat_id: row.get(1),
|
||||
text: row.get(2),
|
||||
is_user: row.get(3),
|
||||
})
|
||||
.collect();
|
||||
|
||||
Ok(messages)
|
||||
}
|
||||
|
||||
async fn get_chat_ids(&self) -> Result<Box<[i32]>> {
|
||||
let rows = sqlx::query("SELECT DISTINCT(chat_id) FROM messages ORDER BY chat_id DESC")
|
||||
.fetch_all(&self.pool)
|
||||
.await
|
||||
.inspect_err(|e| println!("sql error: {}", e))?;
|
||||
let ids: Vec<i32> = rows
|
||||
.into_iter()
|
||||
.map(|row| {
|
||||
let i: i32 = row.get(0);
|
||||
i
|
||||
})
|
||||
.collect();
|
||||
Ok(ids.into_boxed_slice())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,7 +28,10 @@ impl DaemonServer {
|
||||
impl AiDaemon for DaemonServer {
|
||||
async fn chat(&self, request: Request<CRequest>) -> Result<Response<CResponse>, Status> {
|
||||
let r = request.into_inner();
|
||||
let mut messages = gather_history(self.repo.clone())
|
||||
let chat_id = get_chat_id(self.repo.clone(), r.chat_id)
|
||||
.await
|
||||
.map_err(|e| Status::new(Code::Internal, e.to_string()))?;
|
||||
let mut messages = gather_history(self.repo.clone(), &chat_id)
|
||||
.await
|
||||
.map_err(|e| Status::new(Code::Internal, e.to_string()))?;
|
||||
messages.push(ChatMessage::user(r.text()));
|
||||
@@ -42,7 +45,7 @@ impl AiDaemon for DaemonServer {
|
||||
let user_message = message_to_dto(
|
||||
&self
|
||||
.repo
|
||||
.save_message(r.text(), &true)
|
||||
.save_message(r.text(), &true, &0)
|
||||
.await
|
||||
.map_err(|e| Status::new(Code::Internal, e.to_string()))?,
|
||||
);
|
||||
@@ -52,12 +55,12 @@ impl AiDaemon for DaemonServer {
|
||||
};
|
||||
|
||||
println!("User: {}", r.text());
|
||||
println!("AI: {}", response_text.clone());
|
||||
println!("AI: {}", response_text);
|
||||
|
||||
let ai_message = message_to_dto(
|
||||
&self
|
||||
.repo
|
||||
.save_message(response_text, &false)
|
||||
.save_message(response_text, &false, &0)
|
||||
.await
|
||||
.map_err(|e| Status::new(Code::Internal, e.to_string()))?,
|
||||
);
|
||||
@@ -70,11 +73,14 @@ impl AiDaemon for DaemonServer {
|
||||
|
||||
async fn chat_history(
|
||||
&self,
|
||||
_: Request<ChatHistoryRequest>,
|
||||
request: Request<ChatHistoryRequest>,
|
||||
) -> Result<Response<ChatHistoryResponse>, Status> {
|
||||
let chat_id = get_chat_id(self.repo.clone(), request.into_inner().chat_id)
|
||||
.await
|
||||
.map_err(|e| Status::new(Code::Internal, e.to_string()))?;
|
||||
let messages = self
|
||||
.repo
|
||||
.get_latest_messages()
|
||||
.get_latest_messages(&chat_id, &20)
|
||||
.await
|
||||
.map_err(|e| Status::new(Code::Internal, e.to_string()))?;
|
||||
|
||||
@@ -106,8 +112,11 @@ pub fn message_to_dto(msg: &ChatMessageData) -> CMessage {
|
||||
}
|
||||
}
|
||||
|
||||
async fn gather_history(repo: Arc<dyn ChatRepository + Send + Sync>) -> Result<Vec<ChatMessage>> {
|
||||
let messages = repo.get_latest_messages().await?;
|
||||
async fn gather_history(
|
||||
repo: Arc<dyn ChatRepository + Send + Sync>,
|
||||
chat_id: &i32,
|
||||
) -> Result<Vec<ChatMessage>> {
|
||||
let messages = repo.get_latest_messages(chat_id, &10).await?;
|
||||
Ok(messages
|
||||
.iter()
|
||||
.map(|m| match m.is_user {
|
||||
@@ -116,3 +125,13 @@ async fn gather_history(repo: Arc<dyn ChatRepository + Send + Sync>) -> Result<V
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
|
||||
async fn get_chat_id(
|
||||
repo: Arc<dyn ChatRepository + Send + Sync>,
|
||||
chat_id: Option<i64>,
|
||||
) -> Result<i32> {
|
||||
Ok(match chat_id {
|
||||
Some(i) => i as i32,
|
||||
None => repo.get_chat_ids().await?.get(0).copied().unwrap_or(0),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -3,7 +3,6 @@ mod daemongrpc;
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use genai::chat::{ChatMessage, ChatRequest};
|
||||
use genai::Client;
|
||||
use shared::ai::ai_daemon_server::AiDaemonServer;
|
||||
use tonic::transport::Server;
|
||||
@@ -11,20 +10,6 @@ use tonic::transport::Server;
|
||||
use chatpersistence::SqliteChatRepository;
|
||||
use daemongrpc::DaemonServer;
|
||||
|
||||
async fn prompt_ollama(
|
||||
client: &Client,
|
||||
model: &str,
|
||||
prompt: &str,
|
||||
) -> Result<String, Box<dyn std::error::Error>> {
|
||||
let chat_req = ChatRequest::new(vec![ChatMessage::user(prompt)]);
|
||||
let chat_res = client.exec_chat(model, chat_req, None).await?;
|
||||
let output = chat_res
|
||||
.first_text()
|
||||
.unwrap_or("No response content!")
|
||||
.to_string();
|
||||
Ok(output)
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let chat_repo = SqliteChatRepository::new().await?;
|
||||
|
||||
@@ -13,6 +13,26 @@ pub mod chatmessage {
|
||||
pub chat_id: Option<i64>,
|
||||
pub history: Vec<Message>,
|
||||
}
|
||||
|
||||
pub enum TauriCommand {
|
||||
Chat,
|
||||
ChatHistory,
|
||||
DaemonState,
|
||||
ToggleDarkMode,
|
||||
TogglePopup,
|
||||
}
|
||||
|
||||
impl TauriCommand {
|
||||
pub fn as_str(&self) -> &'static str {
|
||||
match self {
|
||||
TauriCommand::TogglePopup => "toggle_popup",
|
||||
TauriCommand::Chat => "chat",
|
||||
TauriCommand::ChatHistory => "chat_history",
|
||||
TauriCommand::DaemonState => "daemon_state",
|
||||
TauriCommand::ToggleDarkMode => "toggle_dark_mode",
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub mod daemon {
|
||||
|
||||
Reference in New Issue
Block a user