use crate::error::{DrivenError, Result};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GmailConfig {
#[serde(default = "default_true")]
pub enabled: bool,
#[serde(default)]
pub client_id: String,
#[serde(default)]
pub client_secret: String,
#[serde(default)]
pub refresh_token: String,
pub pubsub_topic: Option<String>,
#[serde(default)]
pub filters: Vec<GmailFilter>,
}
fn default_true() -> bool {
true
}
impl Default for GmailConfig {
fn default() -> Self {
Self {
enabled: true,
client_id: String::new(),
client_secret: String::new(),
refresh_token: String::new(),
pubsub_topic: None,
filters: Vec::new(),
}
}
}
impl GmailConfig {
pub fn from_file(path: impl AsRef<std::path::Path>) -> Result<Self> {
let content = std::fs::read_to_string(path.as_ref())
.map_err(|e| DrivenError::Io(e))?;
Self::parse_sr(&content)
}
fn parse_sr(_content: &str) -> Result<Self> {
Ok(Self::default())
}
pub fn resolve_env_vars(&mut self) {
if self.client_id.is_empty() || self.client_id.starts_with('$') {
self.client_id = std::env::var("GMAIL_CLIENT_ID").unwrap_or_default();
}
if self.client_secret.is_empty() || self.client_secret.starts_with('$') {
self.client_secret = std::env::var("GMAIL_CLIENT_SECRET").unwrap_or_default();
}
if self.refresh_token.is_empty() || self.refresh_token.starts_with('$') {
self.refresh_token = std::env::var("GMAIL_REFRESH_TOKEN").unwrap_or_default();
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GmailFilter {
pub name: String,
pub from: Option<String>,
pub subject: Option<String>,
pub label: Option<String>,
pub action: GmailAction,
#[serde(default = "default_true")]
pub enabled: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum GmailAction {
Forward { to: String },
Command { cmd: String },
Webhook { url: String },
Archive,
MarkRead,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GmailMessage {
pub id: String,
pub thread_id: String,
pub subject: String,
pub from: String,
pub to: Vec<String>,
pub cc: Vec<String>,
pub date: String,
pub snippet: String,
pub labels: Vec<String>,
pub body_plain: Option<String>,
pub body_html: Option<String>,
pub attachments: Vec<GmailAttachment>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GmailAttachment {
pub id: String,
pub filename: String,
pub mime_type: String,
pub size: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GmailLabel {
pub id: String,
pub name: String,
pub label_type: String,
}
pub struct GmailClient {
config: GmailConfig,
access_token: Option<String>,
base_url: String,
}
impl GmailClient {
const API_BASE: &'static str = "https://gmail.googleapis.com/gmail/v1";
const TOKEN_URL: &'static str = "https://oauth2.googleapis.com/token";
pub async fn new(config: &GmailConfig) -> Result<Self> {
let mut config = config.clone();
config.resolve_env_vars();
let mut client = Self {
config,
access_token: None,
base_url: Self::API_BASE.to_string(),
};
if client.is_configured() {
client.refresh_access_token().await?;
}
Ok(client)
}
pub fn is_configured(&self) -> bool {
!self.config.client_id.is_empty()
&& !self.config.client_secret.is_empty()
&& !self.config.refresh_token.is_empty()
}
async fn refresh_access_token(&mut self) -> Result<()> {
let client = reqwest::Client::new();
let mut params = HashMap::new();
params.insert("client_id", self.config.client_id.as_str());
params.insert("client_secret", self.config.client_secret.as_str());
params.insert("refresh_token", self.config.refresh_token.as_str());
params.insert("grant_type", "refresh_token");
let response = client
.post(Self::TOKEN_URL)
.form(¶ms)
.send()
.await
.map_err(|e| DrivenError::Network(e.to_string()))?;
if !response.status().is_success() {
return Err(DrivenError::Api("Failed to refresh token".into()));
}
#[derive(Deserialize)]
struct TokenResponse {
access_token: String,
}
let tokens: TokenResponse = response
.json()
.await
.map_err(|e| DrivenError::Parse(e.to_string()))?;
self.access_token = Some(tokens.access_token);
Ok(())
}
pub async fn list_messages(&self, query: Option<&str>, max_results: u32) -> Result<Vec<GmailMessage>> {
let token = self.access_token.as_ref()
.ok_or_else(|| DrivenError::Config("Not authenticated".into()))?;
let mut url = format!(
"{}/users/me/messages?maxResults={}",
self.base_url, max_results
);
if let Some(q) = query {
url.push_str(&format!("&q={}", urlencoding::encode(q)));
}
let client = reqwest::Client::new();
let response = client
.get(&url)
.header("Authorization", format!("Bearer {}", token))
.send()
.await
.map_err(|e| DrivenError::Network(e.to_string()))?;
if !response.status().is_success() {
return Err(DrivenError::Api("Failed to list messages".into()));
}
#[derive(Deserialize)]
struct ListResponse {
messages: Option<Vec<MessageId>>,
}
#[derive(Deserialize)]
struct MessageId {
id: String,
}
let list: ListResponse = response
.json()
.await
.map_err(|e| DrivenError::Parse(e.to_string()))?;
let mut messages = Vec::new();
if let Some(msg_ids) = list.messages {
for msg_id in msg_ids.into_iter().take(max_results as usize) {
if let Ok(msg) = self.get_message(&msg_id.id).await {
messages.push(msg);
}
}
}
Ok(messages)
}
pub async fn get_message(&self, id: &str) -> Result<GmailMessage> {
let token = self.access_token.as_ref()
.ok_or_else(|| DrivenError::Config("Not authenticated".into()))?;
let url = format!(
"{}/users/me/messages/{}?format=full",
self.base_url, id
);
let client = reqwest::Client::new();
let response = client
.get(&url)
.header("Authorization", format!("Bearer {}", token))
.send()
.await
.map_err(|e| DrivenError::Network(e.to_string()))?;
if !response.status().is_success() {
return Err(DrivenError::Api("Failed to get message".into()));
}
let raw: serde_json::Value = response
.json()
.await
.map_err(|e| DrivenError::Parse(e.to_string()))?;
self.parse_message(raw)
}
pub async fn send(&self, to: &[&str], subject: &str, body: &str) -> Result<String> {
let token = self.access_token.as_ref()
.ok_or_else(|| DrivenError::Config("Not authenticated".into()))?;
let message = format!(
"To: {}\r\nSubject: {}\r\nContent-Type: text/plain; charset=utf-8\r\n\r\n{}",
to.join(", "),
subject,
body
);
let encoded = base64::encode_config(message.as_bytes(), base64::URL_SAFE_NO_PAD);
let url = format!("{}/users/me/messages/send", self.base_url);
let client = reqwest::Client::new();
let response = client
.post(&url)
.header("Authorization", format!("Bearer {}", token))
.json(&serde_json::json!({ "raw": encoded }))
.send()
.await
.map_err(|e| DrivenError::Network(e.to_string()))?;
if !response.status().is_success() {
return Err(DrivenError::Api("Failed to send message".into()));
}
#[derive(Deserialize)]
struct SendResponse {
id: String,
}
let result: SendResponse = response
.json()
.await
.map_err(|e| DrivenError::Parse(e.to_string()))?;
Ok(result.id)
}
pub async fn archive(&self, id: &str) -> Result<()> {
self.modify_labels(id, &[], &["INBOX"]).await
}
pub async fn mark_read(&self, id: &str) -> Result<()> {
self.modify_labels(id, &[], &["UNREAD"]).await
}
pub async fn mark_unread(&self, id: &str) -> Result<()> {
self.modify_labels(id, &["UNREAD"], &[]).await
}
pub async fn modify_labels(&self, id: &str, add: &[&str], remove: &[&str]) -> Result<()> {
let token = self.access_token.as_ref()
.ok_or_else(|| DrivenError::Config("Not authenticated".into()))?;
let url = format!("{}/users/me/messages/{}/modify", self.base_url, id);
let client = reqwest::Client::new();
let response = client
.post(&url)
.header("Authorization", format!("Bearer {}", token))
.json(&serde_json::json!({
"addLabelIds": add,
"removeLabelIds": remove
}))
.send()
.await
.map_err(|e| DrivenError::Network(e.to_string()))?;
if !response.status().is_success() {
return Err(DrivenError::Api("Failed to modify labels".into()));
}
Ok(())
}
pub async fn list_labels(&self) -> Result<Vec<GmailLabel>> {
let token = self.access_token.as_ref()
.ok_or_else(|| DrivenError::Config("Not authenticated".into()))?;
let url = format!("{}/users/me/labels", self.base_url);
let client = reqwest::Client::new();
let response = client
.get(&url)
.header("Authorization", format!("Bearer {}", token))
.send()
.await
.map_err(|e| DrivenError::Network(e.to_string()))?;
if !response.status().is_success() {
return Err(DrivenError::Api("Failed to list labels".into()));
}
#[derive(Deserialize)]
struct LabelsResponse {
labels: Vec<LabelItem>,
}
#[derive(Deserialize)]
struct LabelItem {
id: String,
name: String,
#[serde(rename = "type")]
label_type: Option<String>,
}
let result: LabelsResponse = response
.json()
.await
.map_err(|e| DrivenError::Parse(e.to_string()))?;
Ok(result
.labels
.into_iter()
.map(|l| GmailLabel {
id: l.id,
name: l.name,
label_type: l.label_type.unwrap_or_else(|| "user".to_string()),
})
.collect())
}
pub async fn watch(&self, labels: &[&str]) -> Result<WatchResponse> {
let token = self.access_token.as_ref()
.ok_or_else(|| DrivenError::Config("Not authenticated".into()))?;
let topic = self.config.pubsub_topic.as_ref()
.ok_or_else(|| DrivenError::Config("Pub/Sub topic not configured".into()))?;
let url = format!("{}/users/me/watch", self.base_url);
let client = reqwest::Client::new();
let response = client
.post(&url)
.header("Authorization", format!("Bearer {}", token))
.json(&serde_json::json!({
"topicName": topic,
"labelIds": labels
}))
.send()
.await
.map_err(|e| DrivenError::Network(e.to_string()))?;
if !response.status().is_success() {
return Err(DrivenError::Api("Failed to setup watch".into()));
}
response
.json()
.await
.map_err(|e| DrivenError::Parse(e.to_string()))
}
fn parse_message(&self, raw: serde_json::Value) -> Result<GmailMessage> {
let id = raw["id"].as_str().unwrap_or_default().to_string();
let thread_id = raw["threadId"].as_str().unwrap_or_default().to_string();
let snippet = raw["snippet"].as_str().unwrap_or_default().to_string();
let labels: Vec<String> = raw["labelIds"]
.as_array()
.map(|arr| arr.iter().filter_map(|v| v.as_str().map(String::from)).collect())
.unwrap_or_default();
let headers = raw["payload"]["headers"].as_array();
let mut subject = String::new();
let mut from = String::new();
let mut to = Vec::new();
let mut cc = Vec::new();
let mut date = String::new();
if let Some(hdrs) = headers {
for h in hdrs {
let name = h["name"].as_str().unwrap_or_default().to_lowercase();
let value = h["value"].as_str().unwrap_or_default();
match name.as_str() {
"subject" => subject = value.to_string(),
"from" => from = value.to_string(),
"to" => to = value.split(',').map(|s| s.trim().to_string()).collect(),
"cc" => cc = value.split(',').map(|s| s.trim().to_string()).collect(),
"date" => date = value.to_string(),
_ => {}
}
}
}
Ok(GmailMessage {
id,
thread_id,
subject,
from,
to,
cc,
date,
snippet,
labels,
body_plain: None,
body_html: None,
attachments: Vec::new(),
})
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WatchResponse {
#[serde(rename = "historyId")]
pub history_id: String,
pub expiration: String,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_default_config() {
let config = GmailConfig::default();
assert!(config.enabled);
assert!(config.filters.is_empty());
}
}