diff --git a/frontend/public/sw.js b/frontend/public/sw.js index c5efb91..e976c34 100644 --- a/frontend/public/sw.js +++ b/frontend/public/sw.js @@ -1,31 +1,77 @@ -const CACHE_NAME = 'agency-cache-v1'; +const CACHE_NAME = 'agency-cache-v4'; const ASSETS = [ '/', '/index.html', - '/src/main.js', - '/src/style.css', '/manifest.json', '/icon-192.png', '/icon-512.png' ]; +// Force immediate update to the latest SW self.addEventListener('install', (event) => { event.waitUntil( caches.open(CACHE_NAME).then((cache) => { return cache.addAll(ASSETS); - }) + }).then(() => self.skipWaiting()) + ); +}); + +// Clean up old caches and take control of all clients immediately +self.addEventListener('activate', (event) => { + event.waitUntil( + caches.keys().then((cacheNames) => { + return Promise.all( + cacheNames.map((cacheName) => { + if (cacheName !== CACHE_NAME) { + console.log('Deleting old cache:', cacheName); + return caches.delete(cacheName); + } + }) + ); + }).then(() => self.clients.claim()) ); }); self.addEventListener('fetch', (event) => { + // Only intercept http/https requests + if (!event.request.url.startsWith('http')) return; + event.respondWith( - caches.match(event.request).then((response) => { - return response || fetch(event.request); + caches.match(event.request).then((cachedResponse) => { + if (cachedResponse) { + return cachedResponse; + } + + return fetch(event.request).catch((error) => { + // If network fetch fails and it's a navigation request, return index.html + if (event.request.mode === 'navigate') { + return caches.match('/index.html'); + } + + // For assets, return a failure response instead of throwing. + // Re-throwing (or returning a rejected promise) causes the browser to show + // the "unexpected error" interception UI. + console.warn('Fetch failed for:', event.request.url, error); + + return new Response('Network error occurred', { + status: 503, + statusText: 'Service Unavailable', + headers: new Headers({ 'Content-Type': 'text/plain' }) + }); + }); }) ); }); + self.addEventListener('push', (event) => { - const data = event.data ? event.data.json() : { title: 'Notification', body: 'New update from Agency' }; + let data = { title: 'Notification', body: 'New update from Agency' }; + try { + if (event.data) { + data = event.data.json(); + } + } catch (e) { + console.error('Error parsing push data:', e); + } const options = { body: data.body, @@ -34,7 +80,9 @@ self.addEventListener('push', (event) => { vibrate: [100, 50, 100], data: { dateOfArrival: Date.now(), - primaryKey: '1' + primaryKey: '1', + taskId: data.task_id, + runId: data.run_id } }; @@ -45,7 +93,27 @@ self.addEventListener('push', (event) => { self.addEventListener('notificationclick', (event) => { event.notification.close(); + + const taskId = event.notification.data.taskId; + const runId = event.notification.data.runId; + + let url = '/'; + if (taskId && runId) { + url = `/?taskId=${taskId}&runId=${runId}`; + } + event.waitUntil( - clients.openWindow('/') + clients.matchAll({ type: 'window', includeUncontrolled: true }).then((windowClients) => { + // Check if there is already a window open and focus it, or open a new one + for (let client of windowClients) { + if ('focus' in client) { + // Navigate the existing client to the new URL if it's the same app + return client.navigate(url).then(c => c.focus()); + } + } + if (clients.openWindow) { + return clients.openWindow(url); + } + }) ); }); diff --git a/frontend/src/main.js b/frontend/src/main.js index 8fe9419..ffae909 100644 --- a/frontend/src/main.js +++ b/frontend/src/main.js @@ -159,6 +159,16 @@ async function fetchTasks() { const response = await fetchWithAuth(`${API_URL}/tasks`); const newTasks = await response.json(); + // Check for deep link in URL + const params = new URLSearchParams(window.location.search); + const urlTaskId = params.get('taskId'); + const urlRunId = params.get('runId'); + + if (urlTaskId && !state.selectedTaskId) { + state.selectedTaskId = urlTaskId; + state.selectedRunId = urlRunId; + } + // Check if we should follow the latest run let newSelectedRunId = state.selectedRunId; if (state.selectedTaskId) { @@ -179,6 +189,12 @@ async function fetchTasks() { tasks: newTasks, selectedRunId: newSelectedRunId }); + + // If we just loaded from a deep link, clear the params and select it + if (urlTaskId) { + window.history.replaceState({}, document.title, "/"); + selectTask(urlTaskId, urlRunId); + } } catch (error) { console.error('Error fetching tasks:', error); } @@ -806,6 +822,16 @@ async function initializeApp() { if (hasSession) { appEl.classList.remove('hidden'); loginOverlay.classList.add('hidden'); + + // Handle deep links from notifications + const params = new URLSearchParams(window.location.search); + const taskId = params.get('taskId'); + const runId = params.get('runId'); + if (taskId) { + state.selectedTaskId = taskId; + state.selectedRunId = runId; + } + await fetchTasks(); connectWebSocket(); } else { @@ -847,6 +873,13 @@ async function setupPush() { const vapidResponse = await fetch(`${API_URL}/notifications/vapid-key`); const { publicKey } = await vapidResponse.json(); + // Always clear existing subscription to ensure we use latest VAPID key + const existingSub = await state.swRegistration.pushManager.getSubscription(); + if (existingSub) { + await existingSub.unsubscribe(); + console.log('Unsubscribed existing push subscription'); + } + const subscription = await state.swRegistration.pushManager.subscribe({ userVisibleOnly: true, applicationServerKey: urlBase64ToUint8Array(publicKey) @@ -870,7 +903,11 @@ async function setupPush() { } function b64(buffer) { - return btoa(String.fromCharCode.apply(null, new Uint8Array(buffer))); + const binary = String.fromCharCode.apply(null, new Uint8Array(buffer)); + return btoa(binary) + .replace(/\+/g, '-') + .replace(/\//g, '_') + .replace(/=/g, ''); } function urlBase64ToUint8Array(base64String) { diff --git a/frontend/vite.config.js b/frontend/vite.config.js index a71c6f4..a06ef36 100644 --- a/frontend/vite.config.js +++ b/frontend/vite.config.js @@ -4,7 +4,7 @@ export default defineConfig({ server: { proxy: { '/api': { - target: 'http://localhost:3000', + target: 'http://localhost:3001', changeOrigin: true, ws: true, } diff --git a/openapi/calendar/openapi.json b/openapi/calendar/openapi.json new file mode 100644 index 0000000..8c8908c --- /dev/null +++ b/openapi/calendar/openapi.json @@ -0,0 +1,593 @@ +{ + "openapi": "3.1.0", + "info": { + "title": "calendar", + "description": "", + "license": { + "name": "" + }, + "version": "0.1.0" + }, + "paths": { + "/auth/me": { + "get": { + "tags": [ + "crate::handlers::auth" + ], + "operationId": "me", + "responses": { + "200": { + "description": "Current user profile", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/CurrentUser" + } + } + } + }, + "401": { + "description": "Unauthorized" + } + }, + "security": [ + { + "oidc": [] + } + ] + } + }, + "/events": { + "get": { + "tags": [ + "crate::handlers::event" + ], + "operationId": "list_events", + "parameters": [ + { + "name": "upcoming", + "in": "query", + "required": false, + "schema": { + "type": [ + "boolean", + "null" + ] + } + } + ], + "responses": { + "200": { + "description": "List of events", + "content": { + "application/json": { + "schema": { + "type": "array", + "items": { + "$ref": "#/components/schemas/Model" + } + } + } + } + }, + "401": { + "description": "Unauthorized" + } + }, + "security": [ + { + "oidc": [] + } + ] + }, + "post": { + "tags": [ + "crate::handlers::event" + ], + "operationId": "create_event", + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/CreateEventRequest" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Event created successfully", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Model" + } + } + } + }, + "400": { + "description": "Invalid request payload" + }, + "401": { + "description": "Unauthorized" + } + }, + "security": [ + { + "oidc": [] + } + ] + } + }, + "/events/{id}": { + "get": { + "tags": [ + "crate::handlers::event" + ], + "operationId": "get_event", + "parameters": [ + { + "name": "id", + "in": "path", + "description": "Event database id", + "required": true, + "schema": { + "type": "integer", + "format": "int32" + } + } + ], + "responses": { + "200": { + "description": "Event details", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Model" + } + } + } + }, + "401": { + "description": "Unauthorized" + }, + "404": { + "description": "Event not found" + } + }, + "security": [ + { + "oidc": [] + } + ] + }, + "put": { + "tags": [ + "crate::handlers::event" + ], + "operationId": "update_event", + "parameters": [ + { + "name": "id", + "in": "path", + "description": "Event database id", + "required": true, + "schema": { + "type": "integer", + "format": "int32" + } + } + ], + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/CreateEventRequest" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Event updated successfully", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Model" + } + } + } + }, + "400": { + "description": "Invalid request payload" + }, + "401": { + "description": "Unauthorized" + }, + "404": { + "description": "Event not found" + } + }, + "security": [ + { + "oidc": [] + } + ] + }, + "delete": { + "tags": [ + "crate::handlers::event" + ], + "operationId": "delete_event", + "parameters": [ + { + "name": "id", + "in": "path", + "description": "Event database id", + "required": true, + "schema": { + "type": "integer", + "format": "int32" + } + } + ], + "responses": { + "204": { + "description": "Event deleted successfully" + }, + "401": { + "description": "Unauthorized" + }, + "404": { + "description": "Event not found" + } + }, + "security": [ + { + "oidc": [] + } + ] + } + }, + "/service/v1/events": { + "get": { + "tags": [ + "crate::handlers::service" + ], + "operationId": "service_list_events", + "parameters": [ + { + "name": "user_id", + "in": "query", + "required": false, + "schema": { + "type": [ + "integer", + "null" + ], + "format": "int32" + } + }, + { + "name": "upcoming", + "in": "query", + "required": false, + "schema": { + "type": [ + "boolean", + "null" + ] + } + } + ], + "responses": { + "200": { + "description": "List of events", + "content": { + "application/json": { + "schema": { + "type": "array", + "items": { + "$ref": "#/components/schemas/Model" + } + } + } + } + }, + "401": { + "description": "Unauthorized" + } + }, + "security": [ + { + "oidc": [] + } + ] + }, + "post": { + "tags": [ + "crate::handlers::service" + ], + "operationId": "service_create_event", + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ServiceCreateEventRequest" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Event created successfully", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Model" + } + } + } + }, + "400": { + "description": "Invalid request payload" + }, + "401": { + "description": "Unauthorized" + } + }, + "security": [ + { + "oidc": [] + } + ] + } + }, + "/service/v1/events/{id}": { + "get": { + "tags": [ + "crate::handlers::service" + ], + "operationId": "service_get_event", + "parameters": [ + { + "name": "id", + "in": "path", + "description": "Event database id", + "required": true, + "schema": { + "type": "integer", + "format": "int32" + } + } + ], + "responses": { + "200": { + "description": "Event details", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Model" + } + } + } + }, + "401": { + "description": "Unauthorized" + }, + "404": { + "description": "Event not found" + } + }, + "security": [ + { + "oidc": [] + } + ] + }, + "put": { + "tags": [ + "crate::handlers::service" + ], + "operationId": "service_update_event", + "parameters": [ + { + "name": "id", + "in": "path", + "description": "Event database id", + "required": true, + "schema": { + "type": "integer", + "format": "int32" + } + } + ], + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ServiceCreateEventRequest" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Event updated successfully", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Model" + } + } + } + }, + "400": { + "description": "Invalid request payload" + }, + "401": { + "description": "Unauthorized" + }, + "404": { + "description": "Event not found" + } + }, + "security": [ + { + "oidc": [] + } + ] + }, + "delete": { + "tags": [ + "crate::handlers::service" + ], + "operationId": "service_delete_event", + "parameters": [ + { + "name": "id", + "in": "path", + "description": "Event database id", + "required": true, + "schema": { + "type": "integer", + "format": "int32" + } + } + ], + "responses": { + "204": { + "description": "Event deleted successfully" + }, + "401": { + "description": "Unauthorized" + }, + "404": { + "description": "Event not found" + } + }, + "security": [ + { + "oidc": [] + } + ] + } + } + }, + "components": { + "schemas": { + "CreateEventRequest": { + "type": "object", + "required": [ + "name", + "from", + "to" + ], + "properties": { + "from": { + "type": "string" + }, + "name": { + "type": "string" + }, + "to": { + "type": "string" + } + } + }, + "CurrentUser": { + "type": "object", + "required": [ + "id", + "sub", + "email", + "name" + ], + "properties": { + "email": { + "type": "string" + }, + "id": { + "type": "integer", + "format": "int32" + }, + "name": { + "type": "string" + }, + "sub": { + "type": "string" + } + } + }, + "Model": { + "type": "object", + "required": [ + "id", + "name", + "from", + "to" + ], + "properties": { + "from": { + "type": "string", + "format": "date-time" + }, + "id": { + "type": "integer", + "format": "int64" + }, + "name": { + "type": "string" + }, + "to": { + "type": "string", + "format": "date-time" + }, + "user_id": { + "type": [ + "integer", + "null" + ], + "format": "int32" + } + } + }, + "ServiceCreateEventRequest": { + "type": "object", + "required": [ + "name", + "from", + "to" + ], + "properties": { + "from": { + "type": "string" + }, + "name": { + "type": "string" + }, + "to": { + "type": "string" + }, + "user_id": { + "type": [ + "integer", + "null" + ], + "format": "int32" + } + } + } + } + }, + "tags": [ + { + "name": "calendar", + "description": "Calendar Management API" + } + ] +} \ No newline at end of file diff --git a/src/config.rs b/src/config.rs index b69a4c0..b81a727 100644 --- a/src/config.rs +++ b/src/config.rs @@ -15,6 +15,7 @@ pub struct Config { pub agent_max_turns: u32, pub agent_max_duration_secs: u64, pub vapid_private_key: String, + pub calendar_api_url: String, } impl Config { @@ -58,6 +59,9 @@ impl Config { let vapid_private_key = env::var("VAPID_PRIVATE_KEY") .map_err(|_| AppError::Config("VAPID_PRIVATE_KEY must be set".into()))?; + let calendar_api_url = + env::var("CALENDAR_API_URL").unwrap_or_else(|_| "http://localhost:8000".to_string()); + Ok(Config { database_url, port, @@ -71,6 +75,7 @@ impl Config { agent_max_turns, agent_max_duration_secs, vapid_private_key, + calendar_api_url, }) } } diff --git a/src/domain/agent/api.rs b/src/domain/agent/api.rs index f886c34..cbd93a8 100644 --- a/src/domain/agent/api.rs +++ b/src/domain/agent/api.rs @@ -75,6 +75,8 @@ pub async fn perform_search( tracing::info!(query = %query, "Performing Tavily web search"); let client = reqwest::Client::builder() .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .pool_idle_timeout(std::time::Duration::from_secs(60)) .build()?; let response = client .post("https://api.tavily.com/search") diff --git a/src/domain/agent/mod.rs b/src/domain/agent/mod.rs index 1b6dde8..8f738cc 100644 --- a/src/domain/agent/mod.rs +++ b/src/domain/agent/mod.rs @@ -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, tavily_api_key: Option, + calendar_client: Arc, + pub user_sub: Option, pub messages: Vec, tools: Option>, logs: String, @@ -26,6 +29,8 @@ impl Agent { db: DatabaseConnection, zen_api_key: Option, tavily_api_key: Option, + calendar_client: Arc, + user_sub: Option, initial_message: String, ) -> AppResult { let intro = format!( @@ -50,19 +55,31 @@ 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, tavily_api_key: Option, + calendar_client: Arc, + user_sub: Option, messages: Vec, ) -> AppResult { let tools = Some(tools::get_tools()); let client = reqwest::Client::builder() .timeout(std::time::Duration::from_secs(120)) + .connect_timeout(std::time::Duration::from_secs(10)) + .tcp_keepalive(std::time::Duration::from_secs(30)) + .pool_idle_timeout(std::time::Duration::from_secs(60)) .build() .map_err(|e| AppError::Internal(format!("Failed to build HTTP client: {}", e)))?; @@ -72,6 +89,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(), @@ -103,6 +122,8 @@ impl Agent { return Err(AppError::Internal("Agent run exceeded max turns".into())); } + tracing::info!("Turn {}", turns); + turns += 1; let current_role = self .messages @@ -114,17 +135,39 @@ impl Agent { turns, current_role )); - let assistant_message = self.execute_turn().await?; + let assistant_message = + match tokio::time::timeout(Duration::from_secs(180), self.execute_turn()).await { + Ok(res) => res?, + Err(_) => { + tracing::error!("Agent execution turn timed out after 180s"); + return Err(AppError::Internal("Agent execution turn timed out".into())); + } + }; if let Some(tool_calls) = &assistant_message.tool_calls { + tracing::info!("Assistant tool calls: {:#?}", tool_calls); if tool_calls.iter().any(|tc| tc.function.name == "answer") { finished = true; } } if self.answer.is_some() { + tracing::info!("Answer: {}", self.answer.as_ref().unwrap()); finished = true; } + + if !finished { + self.messages.push(Message { + role: "system".to_string(), + content: Some( + "continue, use the finish tool to submit your final answer".to_string(), + ), + tool_calls: None, + tool_call_id: None, + }); + } + + tracing::info!("Finished turn"); } self.log("\n--- Execution Finished ---"); @@ -132,10 +175,12 @@ impl Agent { } pub async fn execute_turn(&mut self) -> AppResult { - let max_sub_turns = 10; + let max_sub_turns = 20; let mut sub_turns = 0; loop { + tracing::info!("Sub turn {}", sub_turns); + sub_turns += 1; if sub_turns > max_sub_turns { return Err(AppError::Internal( @@ -151,27 +196,34 @@ impl Agent { .message .clone(); + tracing::info!("Assistant message: {:#?}", assistant_message); + self.messages.push(assistant_message.clone()); if let Some(content) = &assistant_message.content { + tracing::info!("Assistant content: {}", content); if !content.is_empty() { self.log(&format!("\nAssistant: {}", content)); } } if let Some(tool_calls) = &assistant_message.tool_calls { + tracing::info!("Assistant tool calls: {:#?}", tool_calls); let mut is_final_cycle = false; let mut final_answer = None; 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()); @@ -190,6 +242,7 @@ impl Agent { } if is_final_cycle { + tracing::info!("Final answer: {}", final_answer.as_ref().unwrap()); return Ok(Message { role: "assistant".to_string(), content: final_answer.or(assistant_message.content), @@ -212,7 +265,11 @@ impl Agent { tools: self.tools.clone(), }; - let mut request_builder = self.client.post(&self.url).json(&request); + let mut request_builder = self + .client + .post(&self.url) + .json(&request) + .timeout(Duration::from_secs(60)); if let Some(key) = &self.zen_api_key { request_builder = request_builder.header("Authorization", format!("Bearer {}", key)); @@ -222,10 +279,12 @@ impl Agent { let response = request_builder.send().await.map_err(|e| { let duration = start.elapsed(); let is_timeout = e.is_timeout(); + let is_connect = e.is_connect(); tracing::error!( - "Network error after {:?} during LLM call (Timeout: {}): {:?}", + "Network error after {:?} during LLM call (Timeout: {}, Connect: {}): {:?}", duration, is_timeout, + is_connect, e ); AppError::Network(e) diff --git a/src/domain/agent/tools.rs b/src/domain/agent/tools.rs index bbdff4c..2b257c0 100644 --- a/src/domain/agent/tools.rs +++ b/src/domain/agent/tools.rs @@ -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_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, db: &DatabaseConnection, + calendar: &Arc, + user_sub: Option<&str>, ) -> Result<(Message, bool, Option), Box> { 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) }; diff --git a/src/domain/auth.rs b/src/domain/auth.rs index c70e0c7..a180185 100644 --- a/src/domain/auth.rs +++ b/src/domain/auth.rs @@ -46,7 +46,10 @@ pub struct JwksVerifier { impl JwksVerifier { pub async fn new(issuer: String, audience: String) -> Result> { - let client = Client::new(); + let client = Client::builder() + .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .build()?; let discovery_url = format!( "{}/.well-known/openid-configuration", issuer.trim_end_matches('/') @@ -123,7 +126,10 @@ impl Authenticator { client_id: String, client_secret: String, ) -> Result> { - let client = Client::new(); + let client = Client::builder() + .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .build()?; let discovery_url = format!( "{}/.well-known/openid-configuration", issuer.trim_end_matches('/') @@ -190,4 +196,27 @@ impl Authenticator { Ok(res) } + + pub async fn client_credentials( + &self, + scope: &str, + ) -> Result> { + 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(¶ms) + .send() + .await? + .json() + .await?; + + Ok(res) + } } diff --git a/src/domain/calendar/mod.rs b/src/domain/calendar/mod.rs new file mode 100644 index 0000000..c6f49ae --- /dev/null +++ b/src/domain/calendar/mod.rs @@ -0,0 +1,298 @@ +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, + pub name: String, + pub from: DateTime, + pub to: DateTime, + pub user_sub: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CreateEventRequest { + pub name: String, + pub from: String, + pub to: String, + pub user_sub: Option, +} +struct TokenState { + access_token: String, + expires_at: DateTime, +} + +pub struct CalendarClient { + base_url: String, + client: Client, + authenticator: Arc, + token_state: RwLock>, +} + +impl CalendarClient { + pub fn new(base_url: String, authenticator: Arc) -> Self { + let client = Client::builder() + .timeout(std::time::Duration::from_secs(30)) + .connect_timeout(std::time::Duration::from_secs(10)) + .build() + .unwrap_or_else(|_| Client::new()); + + Self { + base_url: base_url.trim_end_matches('/').to_string(), + client, + authenticator, + token_state: RwLock::new(None), + } + } + + async fn get_token(&self) -> AppResult { + { + 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, + upcoming: Option, + ) -> AppResult> { + 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(¶ms.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, + name: &str, + from: &str, + to: &str, + ) -> AppResult { + 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 { + 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, + name: &str, + from: &str, + to: &str, + ) -> AppResult { + 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(()) + } +} diff --git a/src/domain/mod.rs b/src/domain/mod.rs index 70e887f..4a21816 100644 --- a/src/domain/mod.rs +++ b/src/domain/mod.rs @@ -1,4 +1,5 @@ pub mod agent; pub mod auth; +pub mod calendar; pub mod notifications; pub mod tasks; diff --git a/src/domain/notifications/push.rs b/src/domain/notifications/push.rs index 5aaa54a..478fe2a 100644 --- a/src/domain/notifications/push.rs +++ b/src/domain/notifications/push.rs @@ -1,5 +1,6 @@ use crate::error::AppResult; use serde::{Deserialize, Serialize}; +use uuid::Uuid; use web_push::*; pub struct PushSender { @@ -56,6 +57,8 @@ impl PushSender { subscription: &PushSubscription, title: &str, body: &str, + task_id: Option, + run_id: Option, ) -> AppResult<()> { let subscription_info = SubscriptionInfo::new( subscription.endpoint.clone(), @@ -67,8 +70,10 @@ impl PushSender { VapidSignatureBuilder::from_pem_no_sub(std::io::Cursor::new(&self.private_key)) .map_err(|e| crate::error::AppError::Internal(e.to_string()))?; + let mut builder = builder.add_sub_info(&subscription_info); + builder.add_claim("sub", "mailto:pavel@flegr.me"); + let vapid_signature = builder - .add_sub_info(&subscription_info) .build() .map_err(|e| crate::error::AppError::Internal(e.to_string()))?; @@ -79,6 +84,8 @@ impl PushSender { let payload = serde_json::to_vec(&serde_json::json!({ "title": title, "body": body, + "task_id": task_id, + "run_id": run_id, })) .map_err(|e| crate::error::AppError::Internal(e.to_string()))?; @@ -91,10 +98,15 @@ impl PushSender { let client = IsahcWebPushClient::new() .map_err(|e| crate::error::AppError::Internal(e.to_string()))?; - client - .send(message) - .await - .map_err(|e| crate::error::AppError::Internal(e.to_string()))?; + client.send(message).await.map_err(|e| { + tracing::error!("Failed to send push notification: {}", e); + crate::error::AppError::Internal(e.to_string()) + })?; + + tracing::info!( + "Push notification sent successfully to {}", + subscription.endpoint + ); Ok(()) } diff --git a/src/domain/tasks.rs b/src/domain/tasks.rs index 65d3090..bd518ef 100644 --- a/src/domain/tasks.rs +++ b/src/domain/tasks.rs @@ -56,11 +56,12 @@ pub async fn execute_agent_run( db: &DatabaseConnection, _scheduler: &Arc, config: &Arc, + calendar_client: Arc, task_id: Uuid, goal: String, ) -> AppResult { let run_id = Uuid::new_v4(); - tracing::info!(%task_id, %run_id, "Starting agent execution run"); + tracing::info!(%task_id, %run_id, "Starting background agent execution run"); let new_run = task_run::ActiveModel { id: Set(run_id), @@ -92,80 +93,92 @@ pub async fn execute_agent_run( db.clone(), config.zen_api_key.clone(), config.tavily_api_key.clone(), + calendar_client.clone(), + None, goal.clone(), )?; - let (logs, answer, status) = match agent.run(config).await { - Ok((logs, answer)) => { - tracing::info!(%task_id, %run_id, "Agent execution completed successfully"); - (logs, answer, "completed".to_string()) + let db_bg = db.clone(); + let config_bg = config.clone(); + let scheduler_bg = _scheduler.clone(); + let task_id_bg = task_id; + + tokio::spawn(async move { + let (logs, answer, status) = match agent.run(&config_bg).await { + Ok((logs, answer)) => { + tracing::info!(task_id = %task_id_bg, run_id = %run_id, "Agent execution completed successfully"); + (logs, answer, "completed".to_string()) + } + Err(e) => { + tracing::error!(task_id = %task_id_bg, run_id = %run_id, error = %e, "Agent execution failed"); + ( + format!("Execution failed: {}", e), + None, + "failed".to_string(), + ) + } + }; + + let run_update = task_run::ActiveModel { + id: Set(run_id), + logs: Set(logs), + answer: Set(answer), + status: Set(status.clone()), + ..Default::default() + }; + + if let Err(e) = run_update.update(&db_bg).await { + tracing::error!(task_id = %task_id_bg, run_id = %run_id, error = %e, "Failed to update run record"); } - Err(e) => { - tracing::error!(%task_id, %run_id, error = %e, "Agent execution failed"); - ( - format!("Execution failed: {}", e), - None, - "failed".to_string(), - ) + + if let Ok(task_response) = get_task_inner(task_id_bg, &db_bg).await { + let _ = scheduler_bg + .tx + .send(crate::server::notifications::WsEvent::RunFinished( + task_response.clone(), + )); + + // Send Push Notifications to subscribers + if let Ok(subscriptions) = task_subscription::Entity::find() + .filter(task_subscription::Column::TaskId.eq(task_id_bg)) + .all(&db_bg) + .await + { + for sub in subscriptions { + if let Ok(push_subs) = push_subscription::Entity::find() + .filter(push_subscription::Column::UserSub.eq(sub.user_sub.clone())) + .all(&db_bg) + .await + { + for push_sub in push_subs { + let sender = scheduler_bg.push_sender.clone(); + let goal = task_response.goal.clone(); + let status_bg = status.clone(); + let sub_data = crate::domain::notifications::push::PushSubscription { + endpoint: push_sub.endpoint, + p256dh: push_sub.p256dh, + auth: push_sub.auth, + }; + + tokio::spawn(async move { + let _ = sender + .send_notification( + &sub_data, + &format!("Task Completed: {}", status_bg), + &goal, + Some(task_id_bg), + Some(run_id), + ) + .await; + }); + } + } + } + } } - }; + }); - let run: task_run::ActiveModel = TaskRun::find_by_id(run_id) - .one(db) - .await - .map_err(crate::error::AppError::Database)? - .ok_or_else(|| crate::error::AppError::NotFound("Run not found after insert".into()))? - .into(); - - let mut run = run; - run.logs = Set(logs.clone()); - run.answer = Set(answer.clone()); - run.status = Set(status.clone()); - - run.update(db) - .await - .map_err(crate::error::AppError::Database)?; - - let task_response = get_task_inner(task_id, db).await?; - let _ = _scheduler - .tx - .send(crate::server::notifications::WsEvent::RunFinished( - task_response.clone(), - )); - - // Send Push Notifications to subscribers - let subscriptions = task_subscription::Entity::find() - .filter(task_subscription::Column::TaskId.eq(task_id)) - .all(db) - .await - .map_err(crate::error::AppError::Database)?; - - for sub in subscriptions { - let push_subs = push_subscription::Entity::find() - .filter(push_subscription::Column::UserSub.eq(sub.user_sub)) - .all(db) - .await - .map_err(crate::error::AppError::Database)?; - - for push_sub in push_subs { - let sender = _scheduler.push_sender.clone(); - let goal = task_response.goal.clone(); - let status = status.clone(); - let sub_data = crate::domain::notifications::push::PushSubscription { - endpoint: push_sub.endpoint, - p256dh: push_sub.p256dh, - auth: push_sub.auth, - }; - - tokio::spawn(async move { - let _ = sender - .send_notification(&sub_data, &format!("Task Completed: {}", status), &goal) - .await; - }); - } - } - - Ok(task_response) + get_task_inner(task_id, db).await } pub async fn get_task_inner(id: Uuid, db: &DatabaseConnection) -> AppResult { diff --git a/src/error.rs b/src/error.rs index 01943d5..81194ce 100644 --- a/src/error.rs +++ b/src/error.rs @@ -32,16 +32,20 @@ pub enum AppError { impl IntoResponse for AppError { fn into_response(self) -> Response { - let (status, error_message) = match self { + let (status, error_message) = match &self { AppError::Database(err) => (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()), - AppError::Config(err) => (StatusCode::INTERNAL_SERVER_ERROR, err), - AppError::NotFound(err) => (StatusCode::NOT_FOUND, err), - AppError::Unauthorized(err) => (StatusCode::UNAUTHORIZED, err), - AppError::Internal(err) => (StatusCode::INTERNAL_SERVER_ERROR, err), + AppError::Config(err) => (StatusCode::INTERNAL_SERVER_ERROR, err.clone()), + AppError::NotFound(err) => (StatusCode::NOT_FOUND, err.clone()), + AppError::Unauthorized(err) => (StatusCode::UNAUTHORIZED, err.clone()), + AppError::Internal(err) => (StatusCode::INTERNAL_SERVER_ERROR, err.clone()), AppError::Network(err) => (StatusCode::BAD_GATEWAY, err.to_string()), - AppError::InvalidRequest(err) => (StatusCode::BAD_REQUEST, err), + AppError::InvalidRequest(err) => (StatusCode::BAD_REQUEST, err.clone()), }; + if status.is_server_error() || status.is_client_error() { + tracing::error!(%status, error = %self, "AppError converted to response"); + } + let body = Json(json!({ "error": error_message, })); diff --git a/src/scheduler.rs b/src/scheduler.rs index 3b5b0a8..c98debd 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -13,6 +13,7 @@ pub struct Scheduler { db: DatabaseConnection, tasks_to_jobs: DashMap, config: Arc, + pub calendar_client: Arc, pub tx: tokio::sync::broadcast::Sender, pub push_sender: Arc, } @@ -21,6 +22,7 @@ impl Scheduler { pub async fn new( db: DatabaseConnection, config: Arc, + calendar_client: Arc, tx: tokio::sync::broadcast::Sender, push_sender: Arc, ) -> AppResult { @@ -36,6 +38,7 @@ impl Scheduler { db, tasks_to_jobs: DashMap::new(), config, + calendar_client, tx, push_sender, }) @@ -52,13 +55,17 @@ impl Scheduler { let tx = self.tx.clone(); let push_sender = self.push_sender.clone(); + let calendar_client = self.calendar_client.clone(); let job = Job::new_async(cron_expr, move |_uuid, _l| { let db = db.clone(); let config = config.clone(); let tx = tx.clone(); let push_sender = push_sender.clone(); + let calendar_client = calendar_client.clone(); Box::pin(async move { - if let Err(e) = Self::run_task(db, config, tx, push_sender, task_id).await { + if let Err(e) = + Self::run_task(db, config, calendar_client, tx, push_sender, task_id).await + { tracing::error!("Error in scheduled task {}: {}", task_id, e); } }) @@ -91,6 +98,7 @@ impl Scheduler { async fn run_task( db: DatabaseConnection, config: Arc, + calendar_client: Arc, tx: tokio::sync::broadcast::Sender, push_sender: Arc, task_id: Uuid, @@ -132,6 +140,8 @@ impl Scheduler { db.clone(), config.zen_api_key.clone(), config.tavily_api_key.clone(), + calendar_client.clone(), + None, task.goal.clone(), )?; @@ -201,6 +211,8 @@ impl Scheduler { &sub_data, &format!("Scheduled Task Completed: {}", status), &goal, + Some(task_id), + Some(run_id), ) .await; }); diff --git a/src/server/chat.rs b/src/server/chat.rs index 5163275..6e24d4b 100644 --- a/src/server/chat.rs +++ b/src/server/chat.rs @@ -17,6 +17,7 @@ pub struct ChatResult { pub async fn chat_handler( State(state): State>, + user: crate::server::auth::AuthenticatedUser, Json(payload): Json, ) -> Result, AppError> { let msg_count = payload.messages.len(); @@ -45,11 +46,21 @@ pub async fn chat_handler( state.db.clone(), state.config.zen_api_key.clone(), state.config.tavily_api_key.clone(), + state.calendar_client.clone(), + Some(user.0.sub), messages, )?; tracing::info!("Starting interactive agent turn"); - let assistant_message = agent.execute_turn().await?; + let assistant_message = + match tokio::time::timeout(std::time::Duration::from_secs(180), agent.execute_turn()).await + { + Ok(res) => res?, + Err(_) => { + tracing::error!("Interactive chat agent turn timed out after 180s"); + return Err(AppError::Internal("Agent turn timed out".into())); + } + }; Ok(Json(ChatResult { message: assistant_message, diff --git a/src/server/mod.rs b/src/server/mod.rs index 6a9d998..5160b4e 100644 --- a/src/server/mod.rs +++ b/src/server/mod.rs @@ -13,7 +13,6 @@ use sea_orm::{Database, DatabaseConnection, EntityTrait}; use std::sync::Arc; use tower_http::cors::{AllowOrigin, CorsLayer}; -use crate::entities::task::Entity as Task; use crate::scheduler::Scheduler; use crate::error::AppResult; @@ -25,6 +24,7 @@ pub struct AppState { pub config: Arc, pub verifier: Arc, pub authenticator: Arc, + pub calendar_client: Arc, pub tx: tokio::sync::broadcast::Sender, } @@ -39,13 +39,26 @@ pub async fn start(config: crate::config::Config) -> AppResult<()> { &config.vapid_private_key.clone(), )?); + let (verifier, authenticator) = setup_auth(&config).await?; + let calendar_client = Arc::new(crate::domain::calendar::CalendarClient::new( + config.calendar_api_url.clone(), + authenticator.clone(), + )); + let scheduler = Arc::new( - Scheduler::new(db.clone(), config.clone(), tx.clone(), push_sender.clone()) - .await - .map_err(|e| crate::error::AppError::Internal(e.to_string()))?, + Scheduler::new( + db.clone(), + config.clone(), + calendar_client.clone(), + tx.clone(), + push_sender.clone(), + ) + .await + .map_err(|e| crate::error::AppError::Internal(e.to_string()))?, ); // Load existing scheduled tasks + use crate::entities::task::Entity as Task; let existing_tasks = Task::find() .all(&db) .await @@ -56,14 +69,13 @@ pub async fn start(config: crate::config::Config) -> AppResult<()> { } } - let (verifier, authenticator) = setup_auth(&config).await?; - let state = Arc::new(AppState { db, scheduler, config: config.clone(), verifier, authenticator, + calendar_client, tx, }); @@ -136,6 +148,7 @@ fn build_app(state: Arc, config: &crate::config::Config) -> Router { .route("/api/notifications/vapid-key", get(notifications::push_handlers::get_vapid_key)) .route("/api/tasks/:id/subscription", get(notifications::push_handlers::get_subscription_status)) .route("/api/tasks/:id/subscribe", post(notifications::push_handlers::subscribe_task).delete(notifications::push_handlers::unsubscribe_task)) + .layer(axum::middleware::from_fn(log_error_responses)) .layer(cors) .layer(tower_http::set_header::SetResponseHeaderLayer::overriding( axum::http::header::CONTENT_SECURITY_POLICY, @@ -153,6 +166,22 @@ fn build_app(state: Arc, config: &crate::config::Config) -> Router { .with_state(state) } +async fn log_error_responses( + req: axum::extract::Request, + next: axum::middleware::Next, +) -> axum::response::Response { + let method = req.method().clone(); + let uri = req.uri().clone(); + let res = next.run(req).await; + let status = res.status(); + + if status.is_client_error() || status.is_server_error() { + tracing::error!(%method, %uri, %status, "Response error"); + } + + res +} + fn build_cors_layer(config: &crate::config::Config) -> CorsLayer { let allow_origin = if let Some(origins) = &config.cors_allowed_origins { let values: Vec = origins diff --git a/src/server/notifications/push_handlers.rs b/src/server/notifications/push_handlers.rs index 00860c4..cc38e19 100644 --- a/src/server/notifications/push_handlers.rs +++ b/src/server/notifications/push_handlers.rs @@ -26,6 +26,11 @@ pub async fn register_push( Json(payload): Json, ) -> AppResult> { let user_sub = user.0.sub; + tracing::info!( + "Registering push subscription for user: {} with endpoint: {}", + user_sub, + payload.endpoint + ); // Check if subscription exists let existing = push_subscription::Entity::find() diff --git a/src/server/tasks.rs b/src/server/tasks.rs index 5c5da2b..fddc842 100644 --- a/src/server/tasks.rs +++ b/src/server/tasks.rs @@ -109,6 +109,7 @@ pub async fn rerun_task( &state.db, &state.scheduler, &state.config, + state.calendar_client.clone(), task.id, task.goal, )