stuff
All checks were successful
/ upload (release) Successful in 40s

This commit is contained in:
pavel 2026-02-17 19:07:06 +01:00
commit 9aca4cadf0
13 changed files with 483 additions and 21 deletions

View file

@ -3,6 +3,7 @@ pub mod tools;
use chrono::Utc;
use sea_orm::DatabaseConnection;
use std::sync::Arc;
use std::time::{Duration, Instant};
use self::api::{ChatRequest, ChatResponse, Message, Tool};
@ -13,6 +14,8 @@ pub struct Agent {
url: String,
zen_api_key: Option<String>,
tavily_api_key: Option<String>,
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
pub user_sub: Option<String>,
pub messages: Vec<Message>,
tools: Option<Vec<Tool>>,
logs: String,
@ -26,6 +29,8 @@ impl Agent {
db: DatabaseConnection,
zen_api_key: Option<String>,
tavily_api_key: Option<String>,
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
user_sub: Option<String>,
initial_message: String,
) -> AppResult<Self> {
let intro = format!(
@ -50,13 +55,22 @@ impl Agent {
},
];
Self::with_messages(db, zen_api_key, tavily_api_key, messages)
Self::with_messages(
db,
zen_api_key,
tavily_api_key,
calendar_client,
user_sub,
messages,
)
}
pub fn with_messages(
db: DatabaseConnection,
zen_api_key: Option<String>,
tavily_api_key: Option<String>,
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
user_sub: Option<String>,
messages: Vec<Message>,
) -> AppResult<Self> {
let tools = Some(tools::get_tools());
@ -72,6 +86,8 @@ impl Agent {
url: "https://opencode.ai/zen/v1/chat/completions".to_string(),
zen_api_key,
tavily_api_key,
calendar_client,
user_sub,
messages,
tools,
logs: String::new(),
@ -166,12 +182,15 @@ impl Agent {
for tool_call in tool_calls {
self.log(&format!("Calling tool: {}", tool_call.function.name));
let (tool_message, is_final, tool_answer) =
tools::handle_tool_call(tool_call, &self.tavily_api_key, &self.db)
.await
.map_err(|e| {
AppError::Internal(format!("Tool execution failed: {}", e))
})?;
let (tool_message, is_final, tool_answer) = tools::handle_tool_call(
tool_call,
&self.tavily_api_key,
&self.db,
&self.calendar_client,
self.user_sub.as_deref(),
)
.await
.map_err(|e| AppError::Internal(format!("Tool execution failed: {}", e)))?;
if let Some(ans) = tool_answer {
self.answer = Some(ans.clone());

View file

@ -3,6 +3,7 @@ use sea_orm::{
ColumnTrait, Condition, DatabaseConnection, EntityTrait, QueryFilter, QueryOrder, QuerySelect,
};
use serde::Deserialize;
use std::sync::Arc;
use uuid::Uuid;
#[derive(Deserialize)]
@ -96,6 +97,47 @@ pub fn get_tools() -> Vec<Tool> {
}),
},
},
Tool {
tool_type: "function".to_string(),
function: FunctionDefinition {
name: "calendar_list_events".to_string(),
description: "List calendar events".to_string(),
parameters: serde_json::json!({
"type": "object",
"properties": {
"upcoming": {
"type": "boolean",
"description": "If true, only upcoming events will be listed"
}
}
}),
},
},
Tool {
tool_type: "function".to_string(),
function: FunctionDefinition {
name: "calendar_create_event".to_string(),
description: "Create a new calendar event".to_string(),
parameters: serde_json::json!({
"type": "object",
"properties": {
"name": {
"type": "string",
"description": "Name of the event"
},
"from": {
"type": "string",
"description": "Start time in ISO 8601 format (e.g., 2023-10-27T10:00:00Z)"
},
"to": {
"type": "string",
"description": "End time in ISO 8601 format (e.g., 2023-10-27T11:00:00Z)"
}
},
"required": ["name", "from", "to"]
}),
},
},
]
}
@ -103,6 +145,8 @@ pub async fn handle_tool_call(
tool_call: &ToolCall,
tavily_api_key: &Option<String>,
db: &DatabaseConnection,
calendar: &Arc<crate::domain::calendar::CalendarClient>,
user_sub: Option<&str>,
) -> Result<(Message, bool, Option<String>), Box<dyn std::error::Error>> {
let mut answer = None;
let name = &tool_call.function.name;
@ -198,6 +242,34 @@ pub async fn handle_tool_call(
));
}
(out, false)
} else if name == "calendar_list_events" {
let args: serde_json::Value = serde_json::from_str(&tool_call.function.arguments)?;
let upcoming = args["upcoming"].as_bool();
match calendar
.list_events(user_sub.map(|s| s.to_string()), upcoming)
.await
{
Ok(events) => {
tracing::info!("{:#?}", events);
(serde_json::to_string(&events)?, false)
}
Err(e) => (format!("Error listing events: {}", e), false),
}
} else if name == "calendar_create_event" {
let args: serde_json::Value = serde_json::from_str(&tool_call.function.arguments)?;
let name_val = args["name"].as_str().unwrap_or_default();
let from_val = args["from"].as_str().unwrap_or_default();
let to_val = args["to"].as_str().unwrap_or_default();
match calendar
.create_event(user_sub.map(|s| s.to_string()), name_val, from_val, to_val)
.await
{
Ok(event) => (
format!("Event created: {}", serde_json::to_string(&event)?),
false,
),
Err(e) => (format!("Error creating event: {}", e), false),
}
} else {
(format!("Error: Unknown tool {}", name), false)
};

View file

@ -190,4 +190,27 @@ impl Authenticator {
Ok(res)
}
pub async fn client_credentials(
&self,
scope: &str,
) -> Result<serde_json::Value, Box<dyn std::error::Error>> {
let params = [
("grant_type", "client_credentials"),
("client_id", &self.client_id),
("client_secret", &self.client_secret),
("scope", scope),
];
let res = self
.client
.post(&self.token_url)
.form(&params)
.send()
.await?
.json()
.await?;
Ok(res)
}
}

292
src/domain/calendar/mod.rs Normal file
View file

@ -0,0 +1,292 @@
use crate::domain::auth::Authenticator;
use crate::error::{AppError, AppResult};
use chrono::{DateTime, Duration, Utc};
use reqwest::Client;
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use tokio::sync::RwLock;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CalendarEvent {
pub id: Option<i64>,
pub name: String,
pub from: DateTime<Utc>,
pub to: DateTime<Utc>,
pub user_sub: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CreateEventRequest {
pub name: String,
pub from: String,
pub to: String,
pub user_sub: Option<String>,
}
struct TokenState {
access_token: String,
expires_at: DateTime<Utc>,
}
pub struct CalendarClient {
base_url: String,
client: Client,
authenticator: Arc<Authenticator>,
token_state: RwLock<Option<TokenState>>,
}
impl CalendarClient {
pub fn new(base_url: String, authenticator: Arc<Authenticator>) -> Self {
Self {
base_url: base_url.trim_end_matches('/').to_string(),
client: Client::new(),
authenticator,
token_state: RwLock::new(None),
}
}
async fn get_token(&self) -> AppResult<String> {
{
let state = self.token_state.read().await;
if let Some(token) = &*state {
if token.expires_at > Utc::now() + Duration::seconds(30) {
tracing::debug!("Using cached Calendar API token");
return Ok(token.access_token.clone());
}
}
}
let mut state = self.token_state.write().await;
// Double check after acquiring write lock
if let Some(token) = &*state {
if token.expires_at > Utc::now() + Duration::seconds(30) {
return Ok(token.access_token.clone());
}
}
tracing::info!("Refreshing Calendar API token via Client Credentials flow");
let token_data = self
.authenticator
.client_credentials("profile")
.await
.map_err(|e| AppError::Internal(format!("Failed to get client credentials: {}", e)))?;
let access_token = token_data["access_token"]
.as_str()
.ok_or_else(|| AppError::Internal("Missing access_token in response".into()))?
.to_string();
let expires_in = token_data["expires_in"].as_i64().unwrap_or(3600);
let expires_at = Utc::now() + Duration::seconds(expires_in);
*state = Some(TokenState {
access_token: access_token.clone(),
expires_at,
});
Ok(access_token)
}
pub async fn list_events(
&self,
user_sub: Option<String>,
upcoming: Option<bool>,
) -> AppResult<Vec<CalendarEvent>> {
let token = self.get_token().await?;
let mut url = format!("{}/service/v1/events", self.base_url);
let mut params = Vec::new();
if let Some(uid) = &user_sub {
params.push(format!("user_sub={}", uid));
}
if let Some(u) = upcoming {
params.push(format!("upcoming={}", u));
}
if !params.is_empty() {
url.push_str("?");
url.push_str(&params.join("&"));
}
tracing::info!(method = "GET", %url, "Sending Calendar API request");
let res = self
.client
.get(&url)
.bearer_auth(token)
.send()
.await
.map_err(AppError::Network)?;
let status = res.status();
tracing::info!(%status, %url, "Received Calendar API response");
if !status.is_success() {
let error_body = res.text().await.unwrap_or_default();
tracing::error!(%status, %url, body = %error_body, "Calendar API request failed");
return Err(AppError::Internal(format!(
"Failed to list events: {} - {}",
status, error_body
)));
}
let body = res
.json()
.await
.map_err(|e| AppError::Internal(e.to_string()));
tracing::info!(%status, %url, body = ?body, "Received Calendar API response");
body
}
pub async fn create_event(
&self,
user_sub: Option<String>,
name: &str,
from: &str,
to: &str,
) -> AppResult<CalendarEvent> {
let token = self.get_token().await?;
let url = format!("{}/service/v1/events", self.base_url);
let request = CreateEventRequest {
name: name.to_string(),
from: from.to_string(),
to: to.to_string(),
user_sub,
};
tracing::info!(method = "POST", %url, "Sending Calendar API request");
let res = self
.client
.post(&url)
.bearer_auth(token)
.json(&request)
.send()
.await
.map_err(AppError::Network)?;
let status = res.status();
tracing::info!(%status, %url, "Received Calendar API response");
if !status.is_success() {
let error_body = res.text().await.unwrap_or_default();
tracing::error!(%status, %url, body = %error_body, "Calendar API request failed");
return Err(AppError::Internal(format!(
"Failed to create event: {} - {}",
status, error_body
)));
}
res.json()
.await
.map_err(|e| AppError::Internal(e.to_string()))
}
#[allow(dead_code)]
pub async fn get_event(&self, id: i32) -> AppResult<CalendarEvent> {
let token = self.get_token().await?;
let url = format!("{}/service/v1/events/{}", self.base_url, id);
tracing::info!(method = "GET", %url, "Sending Calendar API request");
let res = self
.client
.get(&url)
.bearer_auth(token)
.send()
.await
.map_err(AppError::Network)?;
let status = res.status();
tracing::info!(%status, %url, "Received Calendar API response");
if !status.is_success() {
let error_body = res.text().await.unwrap_or_default();
tracing::error!(%status, %url, body = %error_body, "Calendar API request failed");
return Err(AppError::Internal(format!(
"Failed to get event: {} - {}",
status, error_body
)));
}
res.json()
.await
.map_err(|e| AppError::Internal(e.to_string()))
}
#[allow(dead_code)]
pub async fn update_event(
&self,
id: i32,
user_sub: Option<String>,
name: &str,
from: &str,
to: &str,
) -> AppResult<CalendarEvent> {
let token = self.get_token().await?;
let url = format!("{}/service/v1/events/{}", self.base_url, id);
let request = CreateEventRequest {
name: name.to_string(),
from: from.to_string(),
to: to.to_string(),
user_sub,
};
tracing::info!(method = "PUT", %url, "Sending Calendar API request");
let res = self
.client
.put(&url)
.bearer_auth(token)
.json(&request)
.send()
.await
.map_err(AppError::Network)?;
let status = res.status();
tracing::info!(%status, %url, "Received Calendar API response");
if !status.is_success() {
let error_body = res.text().await.unwrap_or_default();
tracing::error!(%status, %url, body = %error_body, "Calendar API request failed");
return Err(AppError::Internal(format!(
"Failed to update event: {} - {}",
status, error_body
)));
}
res.json()
.await
.map_err(|e| AppError::Internal(e.to_string()))
}
#[allow(dead_code)]
pub async fn delete_event(&self, id: i32) -> AppResult<()> {
let token = self.get_token().await?;
let url = format!("{}/service/v1/events/{}", self.base_url, id);
tracing::info!(method = "DELETE", %url, "Sending Calendar API request");
let res = self
.client
.delete(&url)
.bearer_auth(token)
.send()
.await
.map_err(AppError::Network)?;
let status = res.status();
tracing::info!(%status, %url, "Received Calendar API response");
if !status.is_success() {
let error_body = res.text().await.unwrap_or_default();
tracing::error!(%status, %url, body = %error_body, "Calendar API request failed");
return Err(AppError::Internal(format!(
"Failed to delete event: {} - {}",
status, error_body
)));
}
Ok(())
}
}

View file

@ -1,4 +1,5 @@
pub mod agent;
pub mod auth;
pub mod calendar;
pub mod notifications;
pub mod tasks;

View file

@ -56,6 +56,7 @@ pub async fn execute_agent_run(
db: &DatabaseConnection,
_scheduler: &Arc<Scheduler>,
config: &Arc<Config>,
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
task_id: Uuid,
goal: String,
) -> AppResult<TaskResponse> {
@ -92,6 +93,8 @@ pub async fn execute_agent_run(
db.clone(),
config.zen_api_key.clone(),
config.tavily_api_key.clone(),
calendar_client.clone(),
None,
goal.clone(),
)?;