diff --git a/frontend/src/pages/DataSourcesPage.tsx b/frontend/src/pages/DataSourcesPage.tsx index 5f61176d..b09d9d41 100644 --- a/frontend/src/pages/DataSourcesPage.tsx +++ b/frontend/src/pages/DataSourcesPage.tsx @@ -81,7 +81,7 @@ function InlineConnectForm({ borderRadius: 6, fontSize: 12, cursor: 'pointer', }} > - {loading ? 'Connecting...' : 'Connect'} + Connect ); @@ -327,20 +327,54 @@ function DataSourcesSection() { } }, [connectors, loadSyncStatuses]); + const [connectingId, setConnectingId] = useState(null); + const [connectStage, setConnectStage] = useState(''); + const [connectError, setConnectError] = useState(''); + const handleConnect = async (id: string, req: ConnectRequest) => { setLoading(true); + setConnectingId(id); + setConnectStage('Connecting...'); + setConnectError(''); try { await connectSource(id, req); - setExpandedId(null); - for (let i = 0; i < 30; i++) { - await new Promise((r) => setTimeout(r, 3000)); - await loadConnectors(); + setConnectStage('Connected! Starting sync...'); + + // Wait for connector to show as connected + for (let i = 0; i < 20; i++) { + await new Promise((r) => setTimeout(r, 2000)); const updated = await listConnectors(); const target = updated.find((c) => c.connector_id === id); - if (target?.connected) break; + if (target?.connected) { + setConnectors(updated.map((c) => ({ + connector_id: c.connector_id, + display_name: c.display_name, + connected: c.connected, + chunks: (c as any).chunks || 0, + }))); + break; + } + setConnectStage(i < 5 ? 'Authenticating...' : 'Waiting for connection...'); } - } catch { /* */ } finally { + + // Trigger sync + setConnectStage('Syncing data...'); + try { + await triggerSync(id); + } catch { /* sync may already be running */ } + + // Close form after a brief moment + await new Promise((r) => setTimeout(r, 1500)); + setExpandedId(null); + loadConnectors(); + loadSyncStatuses(); + } catch (err: any) { + setConnectError(err.message || 'Connection failed'); + setConnectStage(''); + } finally { setLoading(false); + setConnectingId(null); + setConnectStage(''); } }; @@ -524,10 +558,43 @@ function DataSourcesSection() { {meta?.inputFields && ( handleConnect(c.connector_id, req)} /> )} + {/* Connection progress */} + {connectingId === c.connector_id && connectStage && ( +
+
+
+ {connectStage} +
+
+
+
+
+ )} + {/* Connection error */} + {connectError && connectingId === null && expandedId === c.connector_id && ( +
+ {connectError} +
+ )}
)}
diff --git a/src/openjarvis/connectors/gcalendar.py b/src/openjarvis/connectors/gcalendar.py index 5b7a85c3..bc730d96 100644 --- a/src/openjarvis/connectors/gcalendar.py +++ b/src/openjarvis/connectors/gcalendar.py @@ -316,9 +316,12 @@ class GCalendarConnector(BaseConnector): page_token: Optional[str] = cursor while True: - events_resp = _gcal_api_events_list( - token, calendar_id, page_token=page_token - ) + try: + events_resp = _gcal_api_events_list( + token, calendar_id, page_token=page_token + ) + except httpx.HTTPStatusError: + break events: List[Dict[str, Any]] = events_resp.get("items", []) for event in events: diff --git a/src/openjarvis/connectors/obsidian.py b/src/openjarvis/connectors/obsidian.py index 1625d146..db8e6a9e 100644 --- a/src/openjarvis/connectors/obsidian.py +++ b/src/openjarvis/connectors/obsidian.py @@ -8,7 +8,7 @@ can be ingested by the knowledge pipeline. from __future__ import annotations import os -from datetime import datetime +from datetime import datetime, timezone from pathlib import Path from typing import Any, Dict, Iterator, List, Optional, Tuple from urllib.parse import quote @@ -161,7 +161,7 @@ class ObsidianConnector(BaseConnector): for fpath in collected_paths: # Apply since filter based on mtime - mtime = datetime.fromtimestamp(fpath.stat().st_mtime) + mtime = datetime.fromtimestamp(fpath.stat().st_mtime, tz=timezone.utc) if since is not None and mtime < since: continue diff --git a/src/openjarvis/server/agent_manager_routes.py b/src/openjarvis/server/agent_manager_routes.py index ae4b1e73..ef41b765 100644 --- a/src/openjarvis/server/agent_manager_routes.py +++ b/src/openjarvis/server/agent_manager_routes.py @@ -719,7 +719,7 @@ async def _stream_managed_agent( # Stream progress events and final content while True: try: - event = await asyncio.to_thread(progress_q.get, timeout=120) + event = await asyncio.to_thread(progress_q.get, timeout=600) except Exception: # Timeout yield _sse_chunk(chunk_id, model, "Agent timed out.")