use std::ffi::{CStr, CString, c_char};
use std::future::Future;
use std::path::Path;
use std::sync::Arc;
use anyhow::{Result, anyhow, bail};
use serde_json::json;
use tdlib_rs::{enums, functions, types};
use tokio::sync::mpsc::UnboundedSender;
use tokio::sync::mpsc::error::SendError;
use crate::config::{ApiKeys, Config};
pub enum TgEvent {
Update(Box<enums::Update>),
Error(String),
ChatsLoaded {
all: bool,
failed: bool,
},
History {
chat_id: i64,
page: Page,
messages: Option<Vec<types::Message>>,
},
Found {
chat_id: i64,
query: String,
found: Option<Found>,
},
Replied {
chat_id: i64,
message_id: i64,
replied: Option<Box<types::Message>>,
},
Deletable {
chat_id: i64,
message_id: i64,
deletable: Option<Deletable>,
},
Downloaded {
file_id: i32,
path: Option<String>,
},
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Page {
Latest,
Older(i64),
Newer(i64),
Around(i64),
}
pub struct Found {
pub ids: Vec<i64>,
pub total: i32,
pub next_from: i64,
}
#[derive(Clone, Copy)]
pub struct Deletable {
pub for_everyone: bool,
pub for_me: bool,
}
const DOWNLOAD_PRIORITY: i32 = 16;
const QUIET_DOWNLOAD_PRIORITY: i32 = 8;
const NOTIFICATION_GROUPS: i64 = 5;
unsafe extern "C" {
fn td_execute(request: *const c_char) -> *const c_char;
}
fn log_to_file(path: &Path) -> Result<()> {
let requests = [
json!({
"@type": "setLogStream",
"log_stream": {
"@type": "logStreamFile",
"path": path.to_string_lossy(),
"max_file_size": 10 * 1024 * 1024,
"redirect_stderr": false,
},
}),
json!({ "@type": "setLogVerbosityLevel", "new_verbosity_level": 2 }),
];
for request in requests {
let request = CString::new(request.to_string())?;
let response = unsafe {
let response = td_execute(request.as_ptr());
if response.is_null() {
String::new()
} else {
CStr::from_ptr(response).to_string_lossy().into_owned()
}
};
if !response.contains(r#""@type":"ok""#) {
bail!("TDLib log setup failed: {response}");
}
}
Ok(())
}
pub type Tagged = (i32, TgEvent);
#[derive(Clone)]
struct Events {
client_id: i32,
tx: UnboundedSender<Tagged>,
}
impl Events {
fn send(&self, event: TgEvent) -> Result<(), SendError<Tagged>> {
self.tx.send((self.client_id, event))
}
}
#[derive(Clone)]
pub struct Tg {
client_id: i32,
tx: Events,
config: Arc<Config>,
}
impl Tg {
pub async fn start(config: Config, tx: UnboundedSender<Tagged>) -> Result<Self> {
log_to_file(&config.data_dir.join("tdlib.log"))?;
let client_id = tdlib_rs::create_client();
std::thread::spawn({
let tx = tx.clone();
move || {
loop {
if let Some((update, client_id)) = tdlib_rs::receive()
&& tx
.send((client_id, TgEvent::Update(Box::new(update))))
.is_err()
{
break;
}
}
}
});
functions::get_option("version".into(), client_id)
.await
.map_err(|e| anyhow!("TDLib didn't start: {}", e.message))?;
Ok(Self {
client_id,
tx: Events { client_id, tx },
config: Arc::new(config),
})
}
pub fn client_id(&self) -> i32 {
self.client_id
}
pub fn set_tdlib_parameters(&self, keys: ApiKeys) {
let config = Arc::clone(&self.config);
let client_id = self.client_id;
self.spawn(async move {
let dir = |name: &str| config.data_dir.join(name).to_string_lossy().into_owned();
functions::set_tdlib_parameters(
false,
dir("db"),
dir("files"),
String::new(), true, true, true, false, keys.id,
keys.hash,
"en".into(),
"Terminal".into(), std::env::consts::OS.into(),
env!("CARGO_PKG_VERSION").into(),
client_id,
)
.await
});
}
pub fn send_phone_number(&self, phone: String) {
self.spawn(functions::set_authentication_phone_number(
phone,
None,
self.client_id,
));
}
pub fn send_code(&self, code: String) {
self.spawn(functions::check_authentication_code(code, self.client_id));
}
pub fn send_password(&self, password: String) {
self.spawn(functions::check_authentication_password(
password,
self.client_id,
));
}
pub fn send_email(&self, email: String) {
self.spawn(functions::set_authentication_email_address(
email,
self.client_id,
));
}
pub fn request_qr_code(&self) {
self.spawn(functions::request_qr_code_authentication(
Vec::new(),
self.client_id,
));
}
pub fn send_email_code(&self, code: String) {
self.spawn(functions::check_authentication_email_code(
enums::EmailAddressAuthentication::Code(types::EmailAddressAuthenticationCode { code }),
self.client_id,
));
}
pub fn load_chats(&self, limit: i32) {
let tx = self.tx.clone();
let client_id = self.client_id;
tokio::spawn(async move {
let (all, failed) =
match functions::load_chats(Some(enums::ChatList::Main), limit, client_id).await {
Ok(()) => (false, false),
Err(e) if e.code == 404 => (true, false),
Err(e) => {
let _ = tx.send(TgEvent::Error(e.message));
(false, true)
}
};
let _ = tx.send(TgEvent::ChatsLoaded { all, failed });
});
}
pub fn open_chat(&self, chat_id: i64) {
self.spawn(functions::open_chat(chat_id, self.client_id));
}
pub fn close_chat(&self, chat_id: i64) {
self.spawn(functions::close_chat(chat_id, self.client_id));
}
pub fn load_history(&self, chat_id: i64, page: Page, limit: i32) {
let (from, offset) = match page {
Page::Latest => (0, 0),
Page::Older(id) => (id, 0),
Page::Newer(id) => (id, 1 - limit),
Page::Around(id) => (id, -limit / 2),
};
let tx = self.tx.clone();
let client_id = self.client_id;
tokio::spawn(async move {
let result =
functions::get_chat_history(chat_id, from, offset, limit, false, client_id).await;
let messages = match result {
Ok(enums::Messages::Messages(page)) => {
Some(page.messages.into_iter().flatten().collect())
}
Err(e) => {
let _ = tx.send(TgEvent::Error(e.message));
None
}
};
let _ = tx.send(TgEvent::History {
chat_id,
page,
messages,
});
});
}
pub fn search_messages(&self, chat_id: i64, query: String, from: i64, limit: i32) {
let tx = self.tx.clone();
let client_id = self.client_id;
tokio::spawn(async move {
let result = functions::search_chat_messages(
chat_id,
None,
query.clone(),
None,
from,
0,
limit,
None,
client_id,
)
.await;
let found = match result {
Ok(enums::FoundChatMessages::FoundChatMessages(f)) => Some(Found {
ids: f.messages.iter().map(|m| m.id).collect(),
total: f.total_count,
next_from: f.next_from_message_id,
}),
Err(e) => {
let _ = tx.send(TgEvent::Error(e.message));
None
}
};
let _ = tx.send(TgEvent::Found {
chat_id,
query,
found,
});
});
}
pub fn get_replied_message(&self, chat_id: i64, message_id: i64) {
let tx = self.tx.clone();
let client_id = self.client_id;
tokio::spawn(async move {
let replied = functions::get_replied_message(chat_id, message_id, client_id)
.await
.ok()
.map(|enums::Message::Message(m)| Box::new(m));
let _ = tx.send(TgEvent::Replied {
chat_id,
message_id,
replied,
});
});
}
pub fn check_deletable(&self, chat_id: i64, message_id: i64) {
let tx = self.tx.clone();
let client_id = self.client_id;
tokio::spawn(async move {
let result = functions::get_message_properties(chat_id, message_id, client_id).await;
let deletable = match result {
Ok(enums::MessageProperties::MessageProperties(p)) => Some(Deletable {
for_everyone: p.can_be_deleted_for_all_users,
for_me: p.can_be_deleted_only_for_self,
}),
Err(e) => {
let _ = tx.send(TgEvent::Error(e.message));
None
}
};
let _ = tx.send(TgEvent::Deletable {
chat_id,
message_id,
deletable,
});
});
}
pub fn delete_message(&self, chat_id: i64, message_id: i64, revoke: bool) {
self.spawn(functions::delete_messages(
chat_id,
vec![message_id],
revoke,
self.client_id,
));
}
pub fn send_text(&self, chat_id: i64, text: String, reply_to: Option<i64>) {
let content = enums::InputMessageContent::InputMessageText(types::InputMessageText {
text: types::FormattedText {
text,
entities: Vec::new(),
},
link_preview_options: None,
clear_draft: true,
});
let reply_to = reply_to.map(|message_id| {
enums::InputMessageReplyTo::Message(types::InputMessageReplyToMessage {
message_id,
quote: None,
checklist_task_id: 0,
})
});
self.spawn(functions::send_message(
chat_id,
None,
reply_to,
None,
content,
self.client_id,
));
}
pub fn download(&self, file_id: i32) {
self.fetch_file(file_id, DOWNLOAD_PRIORITY, true);
}
pub fn download_quiet(&self, file_id: i32) {
self.fetch_file(file_id, QUIET_DOWNLOAD_PRIORITY, false);
}
fn fetch_file(&self, file_id: i32, priority: i32, report_errors: bool) {
let tx = self.tx.clone();
let client_id = self.client_id;
tokio::spawn(async move {
let result = functions::download_file(file_id, priority, 0, 0, true, client_id).await;
let path = match result {
Ok(enums::File::File(f)) if f.local.is_downloading_completed => Some(f.local.path),
Ok(_) => None,
Err(e) => {
if report_errors {
let _ = tx.send(TgEvent::Error(e.message));
}
None
}
};
let _ = tx.send(TgEvent::Downloaded { file_id, path });
});
}
pub fn set_online(&self, online: bool) {
let value = enums::OptionValue::Boolean(types::OptionValueBoolean { value: online });
self.set_option("online", value);
}
pub fn enable_notifications(&self) {
let value = enums::OptionValue::Integer(types::OptionValueInteger {
value: NOTIFICATION_GROUPS,
});
self.set_option("notification_group_count_max", value);
}
fn set_option(&self, name: &str, value: enums::OptionValue) {
self.spawn(functions::set_option(
name.into(),
Some(value),
self.client_id,
));
}
pub fn view_messages(&self, chat_id: i64, message_ids: Vec<i64>) {
self.spawn(functions::view_messages(
chat_id,
message_ids,
None,
true,
self.client_id,
));
}
pub fn close(&self) {
let client_id = self.client_id;
self.spawn(async move {
let offline = enums::OptionValue::Boolean(types::OptionValueBoolean { value: false });
let _ = functions::set_option("online".into(), Some(offline), client_id).await;
functions::close(client_id).await
});
}
pub fn log_out(&self) {
self.spawn(functions::log_out(self.client_id));
}
pub fn reopen(&mut self) {
self.client_id = tdlib_rs::create_client();
self.tx.client_id = self.client_id;
self.spawn(functions::get_option("version".into(), self.client_id));
}
fn spawn<T: Send + 'static>(
&self,
request: impl Future<Output = Result<T, types::Error>> + Send + 'static,
) {
let tx = self.tx.clone();
tokio::spawn(async move {
if let Err(e) = request.await {
let _ = tx.send(TgEvent::Error(e.message));
}
});
}
}