Compare commits
4 commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f9e976c224 | |||
| 7384bcaf40 | |||
| 9aca4cadf0 | |||
| d3cdcea0df |
16 changed files with 1258 additions and 124 deletions
|
|
@ -1,4 +1,4 @@
|
||||||
const CACHE_NAME = 'agency-cache-v3';
|
const CACHE_NAME = 'agency-cache-v4';
|
||||||
const ASSETS = [
|
const ASSETS = [
|
||||||
'/',
|
'/',
|
||||||
'/index.html',
|
'/index.html',
|
||||||
|
|
@ -7,12 +7,29 @@ const ASSETS = [
|
||||||
'/icon-512.png'
|
'/icon-512.png'
|
||||||
];
|
];
|
||||||
|
|
||||||
|
// Force immediate update to the latest SW
|
||||||
self.addEventListener('install', (event) => {
|
self.addEventListener('install', (event) => {
|
||||||
event.waitUntil(
|
event.waitUntil(
|
||||||
caches.open(CACHE_NAME).then((cache) => {
|
caches.open(CACHE_NAME).then((cache) => {
|
||||||
return cache.addAll(ASSETS);
|
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) => {
|
self.addEventListener('fetch', (event) => {
|
||||||
|
|
@ -20,16 +37,27 @@ self.addEventListener('fetch', (event) => {
|
||||||
if (!event.request.url.startsWith('http')) return;
|
if (!event.request.url.startsWith('http')) return;
|
||||||
|
|
||||||
event.respondWith(
|
event.respondWith(
|
||||||
caches.match(event.request).then((response) => {
|
caches.match(event.request).then((cachedResponse) => {
|
||||||
// Return cached response if found, otherwise fetch from network
|
if (cachedResponse) {
|
||||||
return response || fetch(event.request).catch(error => {
|
return cachedResponse;
|
||||||
|
}
|
||||||
|
|
||||||
|
return fetch(event.request).catch((error) => {
|
||||||
// If network fetch fails and it's a navigation request, return index.html
|
// If network fetch fails and it's a navigation request, return index.html
|
||||||
if (event.request.mode === 'navigate') {
|
if (event.request.mode === 'navigate') {
|
||||||
return caches.match('/index.html');
|
return caches.match('/index.html');
|
||||||
}
|
}
|
||||||
// For assets, let the browser handle the failure normally
|
|
||||||
console.error('Fetch failed:', event.request.url, error);
|
// For assets, return a failure response instead of throwing.
|
||||||
throw error;
|
// 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' })
|
||||||
|
});
|
||||||
});
|
});
|
||||||
})
|
})
|
||||||
);
|
);
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,7 @@ export default defineConfig({
|
||||||
server: {
|
server: {
|
||||||
proxy: {
|
proxy: {
|
||||||
'/api': {
|
'/api': {
|
||||||
target: 'http://localhost:3000',
|
target: 'http://localhost:3001',
|
||||||
changeOrigin: true,
|
changeOrigin: true,
|
||||||
ws: true,
|
ws: true,
|
||||||
}
|
}
|
||||||
|
|
|
||||||
593
openapi/calendar/openapi.json
Normal file
593
openapi/calendar/openapi.json
Normal file
|
|
@ -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"
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
|
@ -15,6 +15,7 @@ pub struct Config {
|
||||||
pub agent_max_turns: u32,
|
pub agent_max_turns: u32,
|
||||||
pub agent_max_duration_secs: u64,
|
pub agent_max_duration_secs: u64,
|
||||||
pub vapid_private_key: String,
|
pub vapid_private_key: String,
|
||||||
|
pub calendar_api_url: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Config {
|
impl Config {
|
||||||
|
|
@ -58,6 +59,9 @@ impl Config {
|
||||||
let vapid_private_key = env::var("VAPID_PRIVATE_KEY")
|
let vapid_private_key = env::var("VAPID_PRIVATE_KEY")
|
||||||
.map_err(|_| AppError::Config("VAPID_PRIVATE_KEY must be set".into()))?;
|
.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 {
|
Ok(Config {
|
||||||
database_url,
|
database_url,
|
||||||
port,
|
port,
|
||||||
|
|
@ -71,6 +75,7 @@ impl Config {
|
||||||
agent_max_turns,
|
agent_max_turns,
|
||||||
agent_max_duration_secs,
|
agent_max_duration_secs,
|
||||||
vapid_private_key,
|
vapid_private_key,
|
||||||
|
calendar_api_url,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -75,6 +75,8 @@ pub async fn perform_search(
|
||||||
tracing::info!(query = %query, "Performing Tavily web search");
|
tracing::info!(query = %query, "Performing Tavily web search");
|
||||||
let client = reqwest::Client::builder()
|
let client = reqwest::Client::builder()
|
||||||
.timeout(std::time::Duration::from_secs(30))
|
.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()?;
|
.build()?;
|
||||||
let response = client
|
let response = client
|
||||||
.post("https://api.tavily.com/search")
|
.post("https://api.tavily.com/search")
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,7 @@ pub mod tools;
|
||||||
|
|
||||||
use chrono::Utc;
|
use chrono::Utc;
|
||||||
use sea_orm::DatabaseConnection;
|
use sea_orm::DatabaseConnection;
|
||||||
|
use std::sync::Arc;
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
use self::api::{ChatRequest, ChatResponse, Message, Tool};
|
use self::api::{ChatRequest, ChatResponse, Message, Tool};
|
||||||
|
|
@ -13,6 +14,8 @@ pub struct Agent {
|
||||||
url: String,
|
url: String,
|
||||||
zen_api_key: Option<String>,
|
zen_api_key: Option<String>,
|
||||||
tavily_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>,
|
pub messages: Vec<Message>,
|
||||||
tools: Option<Vec<Tool>>,
|
tools: Option<Vec<Tool>>,
|
||||||
logs: String,
|
logs: String,
|
||||||
|
|
@ -26,6 +29,8 @@ impl Agent {
|
||||||
db: DatabaseConnection,
|
db: DatabaseConnection,
|
||||||
zen_api_key: Option<String>,
|
zen_api_key: Option<String>,
|
||||||
tavily_api_key: Option<String>,
|
tavily_api_key: Option<String>,
|
||||||
|
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
|
||||||
|
user_sub: Option<String>,
|
||||||
initial_message: String,
|
initial_message: String,
|
||||||
) -> AppResult<Self> {
|
) -> AppResult<Self> {
|
||||||
let intro = format!(
|
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(
|
pub fn with_messages(
|
||||||
db: DatabaseConnection,
|
db: DatabaseConnection,
|
||||||
zen_api_key: Option<String>,
|
zen_api_key: Option<String>,
|
||||||
tavily_api_key: Option<String>,
|
tavily_api_key: Option<String>,
|
||||||
|
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
|
||||||
|
user_sub: Option<String>,
|
||||||
messages: Vec<Message>,
|
messages: Vec<Message>,
|
||||||
) -> AppResult<Self> {
|
) -> AppResult<Self> {
|
||||||
let tools = Some(tools::get_tools());
|
let tools = Some(tools::get_tools());
|
||||||
|
|
||||||
let client = reqwest::Client::builder()
|
let client = reqwest::Client::builder()
|
||||||
.timeout(std::time::Duration::from_secs(120))
|
.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()
|
.build()
|
||||||
.map_err(|e| AppError::Internal(format!("Failed to build HTTP client: {}", e)))?;
|
.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(),
|
url: "https://opencode.ai/zen/v1/chat/completions".to_string(),
|
||||||
zen_api_key,
|
zen_api_key,
|
||||||
tavily_api_key,
|
tavily_api_key,
|
||||||
|
calendar_client,
|
||||||
|
user_sub,
|
||||||
messages,
|
messages,
|
||||||
tools,
|
tools,
|
||||||
logs: String::new(),
|
logs: String::new(),
|
||||||
|
|
@ -103,6 +122,8 @@ impl Agent {
|
||||||
return Err(AppError::Internal("Agent run exceeded max turns".into()));
|
return Err(AppError::Internal("Agent run exceeded max turns".into()));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
tracing::info!("Turn {}", turns);
|
||||||
|
|
||||||
turns += 1;
|
turns += 1;
|
||||||
let current_role = self
|
let current_role = self
|
||||||
.messages
|
.messages
|
||||||
|
|
@ -114,17 +135,39 @@ impl Agent {
|
||||||
turns, current_role
|
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 {
|
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") {
|
if tool_calls.iter().any(|tc| tc.function.name == "answer") {
|
||||||
finished = true;
|
finished = true;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if self.answer.is_some() {
|
if self.answer.is_some() {
|
||||||
|
tracing::info!("Answer: {}", self.answer.as_ref().unwrap());
|
||||||
finished = true;
|
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 ---");
|
self.log("\n--- Execution Finished ---");
|
||||||
|
|
@ -132,10 +175,12 @@ impl Agent {
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn execute_turn(&mut self) -> AppResult<Message> {
|
pub async fn execute_turn(&mut self) -> AppResult<Message> {
|
||||||
let max_sub_turns = 10;
|
let max_sub_turns = 20;
|
||||||
let mut sub_turns = 0;
|
let mut sub_turns = 0;
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
|
tracing::info!("Sub turn {}", sub_turns);
|
||||||
|
|
||||||
sub_turns += 1;
|
sub_turns += 1;
|
||||||
if sub_turns > max_sub_turns {
|
if sub_turns > max_sub_turns {
|
||||||
return Err(AppError::Internal(
|
return Err(AppError::Internal(
|
||||||
|
|
@ -151,27 +196,34 @@ impl Agent {
|
||||||
.message
|
.message
|
||||||
.clone();
|
.clone();
|
||||||
|
|
||||||
|
tracing::info!("Assistant message: {:#?}", assistant_message);
|
||||||
|
|
||||||
self.messages.push(assistant_message.clone());
|
self.messages.push(assistant_message.clone());
|
||||||
|
|
||||||
if let Some(content) = &assistant_message.content {
|
if let Some(content) = &assistant_message.content {
|
||||||
|
tracing::info!("Assistant content: {}", content);
|
||||||
if !content.is_empty() {
|
if !content.is_empty() {
|
||||||
self.log(&format!("\nAssistant: {}", content));
|
self.log(&format!("\nAssistant: {}", content));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if let Some(tool_calls) = &assistant_message.tool_calls {
|
if let Some(tool_calls) = &assistant_message.tool_calls {
|
||||||
|
tracing::info!("Assistant tool calls: {:#?}", tool_calls);
|
||||||
let mut is_final_cycle = false;
|
let mut is_final_cycle = false;
|
||||||
let mut final_answer = None;
|
let mut final_answer = None;
|
||||||
|
|
||||||
for tool_call in tool_calls {
|
for tool_call in tool_calls {
|
||||||
self.log(&format!("Calling tool: {}", tool_call.function.name));
|
self.log(&format!("Calling tool: {}", tool_call.function.name));
|
||||||
|
|
||||||
let (tool_message, is_final, tool_answer) =
|
let (tool_message, is_final, tool_answer) = tools::handle_tool_call(
|
||||||
tools::handle_tool_call(tool_call, &self.tavily_api_key, &self.db)
|
tool_call,
|
||||||
|
&self.tavily_api_key,
|
||||||
|
&self.db,
|
||||||
|
&self.calendar_client,
|
||||||
|
self.user_sub.as_deref(),
|
||||||
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| {
|
.map_err(|e| AppError::Internal(format!("Tool execution failed: {}", e)))?;
|
||||||
AppError::Internal(format!("Tool execution failed: {}", e))
|
|
||||||
})?;
|
|
||||||
|
|
||||||
if let Some(ans) = tool_answer {
|
if let Some(ans) = tool_answer {
|
||||||
self.answer = Some(ans.clone());
|
self.answer = Some(ans.clone());
|
||||||
|
|
@ -190,6 +242,7 @@ impl Agent {
|
||||||
}
|
}
|
||||||
|
|
||||||
if is_final_cycle {
|
if is_final_cycle {
|
||||||
|
tracing::info!("Final answer: {}", final_answer.as_ref().unwrap());
|
||||||
return Ok(Message {
|
return Ok(Message {
|
||||||
role: "assistant".to_string(),
|
role: "assistant".to_string(),
|
||||||
content: final_answer.or(assistant_message.content),
|
content: final_answer.or(assistant_message.content),
|
||||||
|
|
@ -212,7 +265,11 @@ impl Agent {
|
||||||
tools: self.tools.clone(),
|
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 {
|
if let Some(key) = &self.zen_api_key {
|
||||||
request_builder = request_builder.header("Authorization", format!("Bearer {}", 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 response = request_builder.send().await.map_err(|e| {
|
||||||
let duration = start.elapsed();
|
let duration = start.elapsed();
|
||||||
let is_timeout = e.is_timeout();
|
let is_timeout = e.is_timeout();
|
||||||
|
let is_connect = e.is_connect();
|
||||||
tracing::error!(
|
tracing::error!(
|
||||||
"Network error after {:?} during LLM call (Timeout: {}): {:?}",
|
"Network error after {:?} during LLM call (Timeout: {}, Connect: {}): {:?}",
|
||||||
duration,
|
duration,
|
||||||
is_timeout,
|
is_timeout,
|
||||||
|
is_connect,
|
||||||
e
|
e
|
||||||
);
|
);
|
||||||
AppError::Network(e)
|
AppError::Network(e)
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,7 @@ use sea_orm::{
|
||||||
ColumnTrait, Condition, DatabaseConnection, EntityTrait, QueryFilter, QueryOrder, QuerySelect,
|
ColumnTrait, Condition, DatabaseConnection, EntityTrait, QueryFilter, QueryOrder, QuerySelect,
|
||||||
};
|
};
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
|
use std::sync::Arc;
|
||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
#[derive(Deserialize)]
|
#[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,
|
tool_call: &ToolCall,
|
||||||
tavily_api_key: &Option<String>,
|
tavily_api_key: &Option<String>,
|
||||||
db: &DatabaseConnection,
|
db: &DatabaseConnection,
|
||||||
|
calendar: &Arc<crate::domain::calendar::CalendarClient>,
|
||||||
|
user_sub: Option<&str>,
|
||||||
) -> Result<(Message, bool, Option<String>), Box<dyn std::error::Error>> {
|
) -> Result<(Message, bool, Option<String>), Box<dyn std::error::Error>> {
|
||||||
let mut answer = None;
|
let mut answer = None;
|
||||||
let name = &tool_call.function.name;
|
let name = &tool_call.function.name;
|
||||||
|
|
@ -198,6 +242,34 @@ pub async fn handle_tool_call(
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
(out, false)
|
(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 {
|
} else {
|
||||||
(format!("Error: Unknown tool {}", name), false)
|
(format!("Error: Unknown tool {}", name), false)
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -46,7 +46,10 @@ pub struct JwksVerifier {
|
||||||
|
|
||||||
impl JwksVerifier {
|
impl JwksVerifier {
|
||||||
pub async fn new(issuer: String, audience: String) -> Result<Self, Box<dyn std::error::Error>> {
|
pub async fn new(issuer: String, audience: String) -> Result<Self, Box<dyn std::error::Error>> {
|
||||||
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!(
|
let discovery_url = format!(
|
||||||
"{}/.well-known/openid-configuration",
|
"{}/.well-known/openid-configuration",
|
||||||
issuer.trim_end_matches('/')
|
issuer.trim_end_matches('/')
|
||||||
|
|
@ -123,7 +126,10 @@ impl Authenticator {
|
||||||
client_id: String,
|
client_id: String,
|
||||||
client_secret: String,
|
client_secret: String,
|
||||||
) -> Result<Self, Box<dyn std::error::Error>> {
|
) -> Result<Self, Box<dyn std::error::Error>> {
|
||||||
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!(
|
let discovery_url = format!(
|
||||||
"{}/.well-known/openid-configuration",
|
"{}/.well-known/openid-configuration",
|
||||||
issuer.trim_end_matches('/')
|
issuer.trim_end_matches('/')
|
||||||
|
|
@ -190,4 +196,27 @@ impl Authenticator {
|
||||||
|
|
||||||
Ok(res)
|
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(¶ms)
|
||||||
|
.send()
|
||||||
|
.await?
|
||||||
|
.json()
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
Ok(res)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
298
src/domain/calendar/mod.rs
Normal file
298
src/domain/calendar/mod.rs
Normal file
|
|
@ -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<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 {
|
||||||
|
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<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(¶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<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(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -1,4 +1,5 @@
|
||||||
pub mod agent;
|
pub mod agent;
|
||||||
pub mod auth;
|
pub mod auth;
|
||||||
|
pub mod calendar;
|
||||||
pub mod notifications;
|
pub mod notifications;
|
||||||
pub mod tasks;
|
pub mod tasks;
|
||||||
|
|
|
||||||
|
|
@ -56,11 +56,12 @@ pub async fn execute_agent_run(
|
||||||
db: &DatabaseConnection,
|
db: &DatabaseConnection,
|
||||||
_scheduler: &Arc<Scheduler>,
|
_scheduler: &Arc<Scheduler>,
|
||||||
config: &Arc<Config>,
|
config: &Arc<Config>,
|
||||||
|
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
|
||||||
task_id: Uuid,
|
task_id: Uuid,
|
||||||
goal: String,
|
goal: String,
|
||||||
) -> AppResult<TaskResponse> {
|
) -> AppResult<TaskResponse> {
|
||||||
let run_id = Uuid::new_v4();
|
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 {
|
let new_run = task_run::ActiveModel {
|
||||||
id: Set(run_id),
|
id: Set(run_id),
|
||||||
|
|
@ -92,16 +93,24 @@ pub async fn execute_agent_run(
|
||||||
db.clone(),
|
db.clone(),
|
||||||
config.zen_api_key.clone(),
|
config.zen_api_key.clone(),
|
||||||
config.tavily_api_key.clone(),
|
config.tavily_api_key.clone(),
|
||||||
|
calendar_client.clone(),
|
||||||
|
None,
|
||||||
goal.clone(),
|
goal.clone(),
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
let (logs, answer, status) = match agent.run(config).await {
|
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)) => {
|
Ok((logs, answer)) => {
|
||||||
tracing::info!(%task_id, %run_id, "Agent execution completed successfully");
|
tracing::info!(task_id = %task_id_bg, run_id = %run_id, "Agent execution completed successfully");
|
||||||
(logs, answer, "completed".to_string())
|
(logs, answer, "completed".to_string())
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
tracing::error!(%task_id, %run_id, error = %e, "Agent execution failed");
|
tracing::error!(task_id = %task_id_bg, run_id = %run_id, error = %e, "Agent execution failed");
|
||||||
(
|
(
|
||||||
format!("Execution failed: {}", e),
|
format!("Execution failed: {}", e),
|
||||||
None,
|
None,
|
||||||
|
|
@ -110,59 +119,41 @@ pub async fn execute_agent_run(
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
let run: task_run::ActiveModel = TaskRun::find_by_id(run_id)
|
let run_update = task_run::ActiveModel {
|
||||||
.one(db)
|
id: Set(run_id),
|
||||||
.await
|
logs: Set(logs),
|
||||||
.map_err(crate::error::AppError::Database)?
|
answer: Set(answer),
|
||||||
.ok_or_else(|| crate::error::AppError::NotFound("Run not found after insert".into()))?
|
status: Set(status.clone()),
|
||||||
.into();
|
..Default::default()
|
||||||
|
};
|
||||||
|
|
||||||
let mut run = run;
|
if let Err(e) = run_update.update(&db_bg).await {
|
||||||
run.logs = Set(logs.clone());
|
tracing::error!(task_id = %task_id_bg, run_id = %run_id, error = %e, "Failed to update run record");
|
||||||
run.answer = Set(answer.clone());
|
}
|
||||||
run.status = Set(status.clone());
|
|
||||||
|
|
||||||
run.update(db)
|
if let Ok(task_response) = get_task_inner(task_id_bg, &db_bg).await {
|
||||||
.await
|
let _ = scheduler_bg
|
||||||
.map_err(crate::error::AppError::Database)?;
|
|
||||||
|
|
||||||
let task_response = get_task_inner(task_id, db).await?;
|
|
||||||
let _ = _scheduler
|
|
||||||
.tx
|
.tx
|
||||||
.send(crate::server::notifications::WsEvent::RunFinished(
|
.send(crate::server::notifications::WsEvent::RunFinished(
|
||||||
task_response.clone(),
|
task_response.clone(),
|
||||||
));
|
));
|
||||||
|
|
||||||
// Send Push Notifications to subscribers
|
// Send Push Notifications to subscribers
|
||||||
let subscriptions = task_subscription::Entity::find()
|
if let Ok(subscriptions) = task_subscription::Entity::find()
|
||||||
.filter(task_subscription::Column::TaskId.eq(task_id))
|
.filter(task_subscription::Column::TaskId.eq(task_id_bg))
|
||||||
.all(db)
|
.all(&db_bg)
|
||||||
.await
|
.await
|
||||||
.map_err(crate::error::AppError::Database)?;
|
{
|
||||||
|
|
||||||
tracing::info!(
|
|
||||||
"Found {} task subscriptions for task {}",
|
|
||||||
subscriptions.len(),
|
|
||||||
task_id
|
|
||||||
);
|
|
||||||
|
|
||||||
for sub in subscriptions {
|
for sub in subscriptions {
|
||||||
let push_subs = push_subscription::Entity::find()
|
if let Ok(push_subs) = push_subscription::Entity::find()
|
||||||
.filter(push_subscription::Column::UserSub.eq(sub.user_sub.clone()))
|
.filter(push_subscription::Column::UserSub.eq(sub.user_sub.clone()))
|
||||||
.all(db)
|
.all(&db_bg)
|
||||||
.await
|
.await
|
||||||
.map_err(crate::error::AppError::Database)?;
|
{
|
||||||
|
|
||||||
tracing::info!(
|
|
||||||
"Found {} push subscriptions for user {}",
|
|
||||||
push_subs.len(),
|
|
||||||
sub.user_sub
|
|
||||||
);
|
|
||||||
|
|
||||||
for push_sub in push_subs {
|
for push_sub in push_subs {
|
||||||
let sender = _scheduler.push_sender.clone();
|
let sender = scheduler_bg.push_sender.clone();
|
||||||
let goal = task_response.goal.clone();
|
let goal = task_response.goal.clone();
|
||||||
let status = status.clone();
|
let status_bg = status.clone();
|
||||||
let sub_data = crate::domain::notifications::push::PushSubscription {
|
let sub_data = crate::domain::notifications::push::PushSubscription {
|
||||||
endpoint: push_sub.endpoint,
|
endpoint: push_sub.endpoint,
|
||||||
p256dh: push_sub.p256dh,
|
p256dh: push_sub.p256dh,
|
||||||
|
|
@ -170,23 +161,24 @@ pub async fn execute_agent_run(
|
||||||
};
|
};
|
||||||
|
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
if let Err(e) = sender
|
let _ = sender
|
||||||
.send_notification(
|
.send_notification(
|
||||||
&sub_data,
|
&sub_data,
|
||||||
&format!("Task Completed: {}", status),
|
&format!("Task Completed: {}", status_bg),
|
||||||
&goal,
|
&goal,
|
||||||
Some(task_id),
|
Some(task_id_bg),
|
||||||
Some(run_id),
|
Some(run_id),
|
||||||
)
|
)
|
||||||
.await
|
.await;
|
||||||
{
|
|
||||||
tracing::error!("Failed to send notification in background task: {}", e);
|
|
||||||
}
|
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
Ok(task_response)
|
get_task_inner(task_id, db).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn get_task_inner(id: Uuid, db: &DatabaseConnection) -> AppResult<TaskResponse> {
|
pub async fn get_task_inner(id: Uuid, db: &DatabaseConnection) -> AppResult<TaskResponse> {
|
||||||
|
|
|
||||||
16
src/error.rs
16
src/error.rs
|
|
@ -32,16 +32,20 @@ pub enum AppError {
|
||||||
|
|
||||||
impl IntoResponse for AppError {
|
impl IntoResponse for AppError {
|
||||||
fn into_response(self) -> Response {
|
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::Database(err) => (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()),
|
||||||
AppError::Config(err) => (StatusCode::INTERNAL_SERVER_ERROR, err),
|
AppError::Config(err) => (StatusCode::INTERNAL_SERVER_ERROR, err.clone()),
|
||||||
AppError::NotFound(err) => (StatusCode::NOT_FOUND, err),
|
AppError::NotFound(err) => (StatusCode::NOT_FOUND, err.clone()),
|
||||||
AppError::Unauthorized(err) => (StatusCode::UNAUTHORIZED, err),
|
AppError::Unauthorized(err) => (StatusCode::UNAUTHORIZED, err.clone()),
|
||||||
AppError::Internal(err) => (StatusCode::INTERNAL_SERVER_ERROR, err),
|
AppError::Internal(err) => (StatusCode::INTERNAL_SERVER_ERROR, err.clone()),
|
||||||
AppError::Network(err) => (StatusCode::BAD_GATEWAY, err.to_string()),
|
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!({
|
let body = Json(json!({
|
||||||
"error": error_message,
|
"error": error_message,
|
||||||
}));
|
}));
|
||||||
|
|
|
||||||
|
|
@ -13,6 +13,7 @@ pub struct Scheduler {
|
||||||
db: DatabaseConnection,
|
db: DatabaseConnection,
|
||||||
tasks_to_jobs: DashMap<Uuid, Uuid>,
|
tasks_to_jobs: DashMap<Uuid, Uuid>,
|
||||||
config: Arc<crate::config::Config>,
|
config: Arc<crate::config::Config>,
|
||||||
|
pub calendar_client: Arc<crate::domain::calendar::CalendarClient>,
|
||||||
pub tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
pub tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
||||||
pub push_sender: Arc<crate::domain::notifications::push::PushSender>,
|
pub push_sender: Arc<crate::domain::notifications::push::PushSender>,
|
||||||
}
|
}
|
||||||
|
|
@ -21,6 +22,7 @@ impl Scheduler {
|
||||||
pub async fn new(
|
pub async fn new(
|
||||||
db: DatabaseConnection,
|
db: DatabaseConnection,
|
||||||
config: Arc<crate::config::Config>,
|
config: Arc<crate::config::Config>,
|
||||||
|
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
|
||||||
tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
||||||
push_sender: Arc<crate::domain::notifications::push::PushSender>,
|
push_sender: Arc<crate::domain::notifications::push::PushSender>,
|
||||||
) -> AppResult<Self> {
|
) -> AppResult<Self> {
|
||||||
|
|
@ -36,6 +38,7 @@ impl Scheduler {
|
||||||
db,
|
db,
|
||||||
tasks_to_jobs: DashMap::new(),
|
tasks_to_jobs: DashMap::new(),
|
||||||
config,
|
config,
|
||||||
|
calendar_client,
|
||||||
tx,
|
tx,
|
||||||
push_sender,
|
push_sender,
|
||||||
})
|
})
|
||||||
|
|
@ -52,13 +55,17 @@ impl Scheduler {
|
||||||
let tx = self.tx.clone();
|
let tx = self.tx.clone();
|
||||||
let push_sender = self.push_sender.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 job = Job::new_async(cron_expr, move |_uuid, _l| {
|
||||||
let db = db.clone();
|
let db = db.clone();
|
||||||
let config = config.clone();
|
let config = config.clone();
|
||||||
let tx = tx.clone();
|
let tx = tx.clone();
|
||||||
let push_sender = push_sender.clone();
|
let push_sender = push_sender.clone();
|
||||||
|
let calendar_client = calendar_client.clone();
|
||||||
Box::pin(async move {
|
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);
|
tracing::error!("Error in scheduled task {}: {}", task_id, e);
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
@ -91,6 +98,7 @@ impl Scheduler {
|
||||||
async fn run_task(
|
async fn run_task(
|
||||||
db: DatabaseConnection,
|
db: DatabaseConnection,
|
||||||
config: Arc<crate::config::Config>,
|
config: Arc<crate::config::Config>,
|
||||||
|
calendar_client: Arc<crate::domain::calendar::CalendarClient>,
|
||||||
tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
||||||
push_sender: Arc<crate::domain::notifications::push::PushSender>,
|
push_sender: Arc<crate::domain::notifications::push::PushSender>,
|
||||||
task_id: Uuid,
|
task_id: Uuid,
|
||||||
|
|
@ -132,6 +140,8 @@ impl Scheduler {
|
||||||
db.clone(),
|
db.clone(),
|
||||||
config.zen_api_key.clone(),
|
config.zen_api_key.clone(),
|
||||||
config.tavily_api_key.clone(),
|
config.tavily_api_key.clone(),
|
||||||
|
calendar_client.clone(),
|
||||||
|
None,
|
||||||
task.goal.clone(),
|
task.goal.clone(),
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@ pub struct ChatResult {
|
||||||
|
|
||||||
pub async fn chat_handler(
|
pub async fn chat_handler(
|
||||||
State(state): State<Arc<AppState>>,
|
State(state): State<Arc<AppState>>,
|
||||||
|
user: crate::server::auth::AuthenticatedUser,
|
||||||
Json(payload): Json<ChatPayload>,
|
Json(payload): Json<ChatPayload>,
|
||||||
) -> Result<Json<ChatResult>, AppError> {
|
) -> Result<Json<ChatResult>, AppError> {
|
||||||
let msg_count = payload.messages.len();
|
let msg_count = payload.messages.len();
|
||||||
|
|
@ -45,11 +46,21 @@ pub async fn chat_handler(
|
||||||
state.db.clone(),
|
state.db.clone(),
|
||||||
state.config.zen_api_key.clone(),
|
state.config.zen_api_key.clone(),
|
||||||
state.config.tavily_api_key.clone(),
|
state.config.tavily_api_key.clone(),
|
||||||
|
state.calendar_client.clone(),
|
||||||
|
Some(user.0.sub),
|
||||||
messages,
|
messages,
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
tracing::info!("Starting interactive agent turn");
|
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 {
|
Ok(Json(ChatResult {
|
||||||
message: assistant_message,
|
message: assistant_message,
|
||||||
|
|
|
||||||
|
|
@ -13,7 +13,6 @@ use sea_orm::{Database, DatabaseConnection, EntityTrait};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tower_http::cors::{AllowOrigin, CorsLayer};
|
use tower_http::cors::{AllowOrigin, CorsLayer};
|
||||||
|
|
||||||
use crate::entities::task::Entity as Task;
|
|
||||||
use crate::scheduler::Scheduler;
|
use crate::scheduler::Scheduler;
|
||||||
|
|
||||||
use crate::error::AppResult;
|
use crate::error::AppResult;
|
||||||
|
|
@ -25,6 +24,7 @@ pub struct AppState {
|
||||||
pub config: Arc<crate::config::Config>,
|
pub config: Arc<crate::config::Config>,
|
||||||
pub verifier: Arc<crate::domain::auth::JwksVerifier>,
|
pub verifier: Arc<crate::domain::auth::JwksVerifier>,
|
||||||
pub authenticator: Arc<crate::domain::auth::Authenticator>,
|
pub authenticator: Arc<crate::domain::auth::Authenticator>,
|
||||||
|
pub calendar_client: Arc<crate::domain::calendar::CalendarClient>,
|
||||||
pub tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
pub tx: tokio::sync::broadcast::Sender<crate::server::notifications::WsEvent>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -39,13 +39,26 @@ pub async fn start(config: crate::config::Config) -> AppResult<()> {
|
||||||
&config.vapid_private_key.clone(),
|
&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(
|
let scheduler = Arc::new(
|
||||||
Scheduler::new(db.clone(), config.clone(), tx.clone(), push_sender.clone())
|
Scheduler::new(
|
||||||
|
db.clone(),
|
||||||
|
config.clone(),
|
||||||
|
calendar_client.clone(),
|
||||||
|
tx.clone(),
|
||||||
|
push_sender.clone(),
|
||||||
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| crate::error::AppError::Internal(e.to_string()))?,
|
.map_err(|e| crate::error::AppError::Internal(e.to_string()))?,
|
||||||
);
|
);
|
||||||
|
|
||||||
// Load existing scheduled tasks
|
// Load existing scheduled tasks
|
||||||
|
use crate::entities::task::Entity as Task;
|
||||||
let existing_tasks = Task::find()
|
let existing_tasks = Task::find()
|
||||||
.all(&db)
|
.all(&db)
|
||||||
.await
|
.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 {
|
let state = Arc::new(AppState {
|
||||||
db,
|
db,
|
||||||
scheduler,
|
scheduler,
|
||||||
config: config.clone(),
|
config: config.clone(),
|
||||||
verifier,
|
verifier,
|
||||||
authenticator,
|
authenticator,
|
||||||
|
calendar_client,
|
||||||
tx,
|
tx,
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
@ -136,6 +148,7 @@ fn build_app(state: Arc<AppState>, config: &crate::config::Config) -> Router {
|
||||||
.route("/api/notifications/vapid-key", get(notifications::push_handlers::get_vapid_key))
|
.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/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))
|
.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(cors)
|
||||||
.layer(tower_http::set_header::SetResponseHeaderLayer::overriding(
|
.layer(tower_http::set_header::SetResponseHeaderLayer::overriding(
|
||||||
axum::http::header::CONTENT_SECURITY_POLICY,
|
axum::http::header::CONTENT_SECURITY_POLICY,
|
||||||
|
|
@ -153,6 +166,22 @@ fn build_app(state: Arc<AppState>, config: &crate::config::Config) -> Router {
|
||||||
.with_state(state)
|
.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 {
|
fn build_cors_layer(config: &crate::config::Config) -> CorsLayer {
|
||||||
let allow_origin = if let Some(origins) = &config.cors_allowed_origins {
|
let allow_origin = if let Some(origins) = &config.cors_allowed_origins {
|
||||||
let values: Vec<HeaderValue> = origins
|
let values: Vec<HeaderValue> = origins
|
||||||
|
|
|
||||||
|
|
@ -109,6 +109,7 @@ pub async fn rerun_task(
|
||||||
&state.db,
|
&state.db,
|
||||||
&state.scheduler,
|
&state.scheduler,
|
||||||
&state.config,
|
&state.config,
|
||||||
|
state.calendar_client.clone(),
|
||||||
task.id,
|
task.id,
|
||||||
task.goal,
|
task.goal,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue