diff --git a/backend/apps/agents/agent_manager.py b/backend/apps/agents/agent_manager.py index aa6b73f6..2e9389ff 100644 --- a/backend/apps/agents/agent_manager.py +++ b/backend/apps/agents/agent_manager.py @@ -23,6 +23,7 @@ from backend.apps.tools_lib.tools_lib import ( load_builtin_permissions, refresh_airtable_token, refresh_google_token, + refresh_hubspot_token, ) from backend.config.paths import SESSIONS_DIR from backend.apps.analytics.collector import record as _analytics @@ -153,6 +154,8 @@ class AgentManager: if tool.auth_type == "oauth2" and tool.auth_status == "connected": if tool.name.lower() == "airtable": refreshed = await refresh_airtable_token(tool) + elif tool.name.lower() == "hubspot": + refreshed = await refresh_hubspot_token(tool) else: refreshed = await refresh_google_token(tool) logger.info(f"[MCP-DEBUG] {tool.name} token refresh: {'OK' if refreshed else 'FAILED'}") @@ -1181,6 +1184,18 @@ class AgentManager: "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) + except BaseException as e: + # Catch BaseExceptionGroup from anyio task groups (e.g. concurrent + # CLI crash + pending approval cancellation) so it doesn't escape + # and kill the uvicorn process. + logger.exception(f"Agent {session_id} fatal error: {e}") + session.status = "error" + error_msg = Message(role="system", content=f"Error: {str(e)}", branch_id=session.active_branch_id) + session.messages.append(error_msg) + await ws_manager.send_to_session(session_id, "agent:message", { + "session_id": session_id, + "message": error_msg.model_dump(mode="json"), + }) finally: if session_id in self.sessions: await ws_manager.send_to_session(session_id, "agent:status", { diff --git a/backend/apps/tools_lib/tools_lib.py b/backend/apps/tools_lib/tools_lib.py index 59b7431f..c87e779b 100644 --- a/backend/apps/tools_lib/tools_lib.py +++ b/backend/apps/tools_lib/tools_lib.py @@ -64,6 +64,9 @@ AIRTABLE_SCOPES = [ "user.email:read", ] +HUBSPOT_AUTH_URL = "https://mcp-na2.hubspot.com/oauth/authorize/user" +HUBSPOT_TOKEN_URL = "https://api.hubapi.com/oauth/v1/token" + # Maps state -> {tool_id, code_verifier (for PKCE flows)} _pending_oauth: dict[str, dict] = {} @@ -178,6 +181,36 @@ async def oauth_callback(code: str = Query(...), state: str = Query("")): tool.auth_status = "connected" tool.connected_account_email = "Airtable account" + elif tool.name.lower() == "hubspot": + # HubSpot OAuth 2.1: PKCE flow + client_id = os.environ.get("HUBSPOT_OAUTH_CLIENT_ID", "") + client_secret = os.environ.get("HUBSPOT_OAUTH_CLIENT_SECRET", "") + async with httpx.AsyncClient(timeout=15.0) as client: + resp = await client.post(HUBSPOT_TOKEN_URL, data={ + "grant_type": "authorization_code", + "code": code, + "redirect_uri": redirect_uri, + "client_id": client_id, + "client_secret": client_secret, + "code_verifier": code_verifier or "", + }, headers={ + "Content-Type": "application/x-www-form-urlencoded", + }) + + if resp.status_code != 200: + logger.warning(f"HubSpot OAuth token exchange failed: {resp.text}") + return HTMLResponse(f"
{resp.text}", status_code=400)
+
+ tokens = resp.json()
+ tool.oauth_tokens = {
+ "access_token": tokens.get("access_token", ""),
+ "refresh_token": tokens.get("refresh_token", ""),
+ "token_expiry": time.time() + tokens.get("expires_in", 1800),
+ }
+ tool.auth_type = "oauth2"
+ tool.auth_status = "connected"
+ tool.connected_account_email = "HubSpot account"
+
elif tool.name.lower() == "notion":
# Notion OAuth: Basic auth with client_id:secret
notion_client_id = os.environ.get("NOTION_OAUTH_CLIENT_ID", "")
@@ -471,6 +504,8 @@ def derive_mcp_config(tool: ToolDefinition) -> Optional[dict]:
env["OAUTH_ACCESS_TOKEN"] = tool.oauth_tokens["access_token"]
if tool.name.lower() == "notion":
env["NOTION_TOKEN"] = tool.oauth_tokens["access_token"]
+ if tool.name.lower() == "hubspot":
+ env["PRIVATE_APP_ACCESS_TOKEN"] = tool.oauth_tokens["access_token"]
if tool.oauth_tokens.get("refresh_token"):
env["GOOGLE_WORKSPACE_REFRESH_TOKEN"] = tool.oauth_tokens["refresh_token"]
client_id = os.environ.get("GOOGLE_OAUTH_CLIENT_ID", "")
@@ -801,6 +836,8 @@ async def discover_tools(tool_id: str):
if tool.oauth_tokens.get("refresh_token"):
if tool.name.lower() == "airtable":
refreshed = await refresh_airtable_token(tool)
+ elif tool.name.lower() == "hubspot":
+ refreshed = await refresh_hubspot_token(tool)
else:
refreshed = await refresh_google_token(tool)
if not refreshed and tool.oauth_tokens.get("access_token"):
@@ -1099,6 +1136,23 @@ async def oauth_start(tool_id: str):
"code_challenge_method": "S256",
}
auth_url = f"{AIRTABLE_AUTH_URL}?{urlencode(params)}"
+ elif tool.name.lower() == "hubspot":
+ client_id = os.environ.get("HUBSPOT_OAUTH_CLIENT_ID", "")
+ if not client_id:
+ raise HTTPException(status_code=400, detail="HUBSPOT_OAUTH_CLIENT_ID not set in backend .env")
+ code_verifier = secrets.token_urlsafe(96)
+ code_challenge = hashlib.sha256(code_verifier.encode()).digest()
+ import base64
+ code_challenge_b64 = base64.urlsafe_b64encode(code_challenge).rstrip(b"=").decode()
+ _pending_oauth[state] = {"tool_id": tool_id, "code_verifier": code_verifier}
+ params = {
+ "client_id": client_id,
+ "redirect_uri": redirect_uri,
+ "code_challenge": code_challenge_b64,
+ "code_challenge_method": "S256",
+ "state": state,
+ }
+ auth_url = f"{HUBSPOT_AUTH_URL}?{urlencode(params)}"
elif tool.name.lower() == "notion":
_pending_oauth[state] = {"tool_id": tool_id}
client_id = os.environ.get("NOTION_OAUTH_CLIENT_ID", "")
@@ -1221,3 +1275,42 @@ async def refresh_airtable_token(tool: ToolDefinition) -> Optional[str]:
except Exception as e:
logger.warning(f"Airtable token refresh failed for tool {tool.id}: {e}")
return None
+
+
+async def refresh_hubspot_token(tool: ToolDefinition) -> Optional[str]:
+ """Refresh an expired HubSpot OAuth token. Returns the fresh access_token or None."""
+ if tool.auth_type != "oauth2":
+ return None
+ refresh_token = tool.oauth_tokens.get("refresh_token")
+ if not refresh_token:
+ return None
+ expiry = tool.oauth_tokens.get("token_expiry", 0)
+ if time.time() < expiry - 60:
+ return tool.oauth_tokens.get("access_token")
+
+ client_id = os.environ.get("HUBSPOT_OAUTH_CLIENT_ID", "")
+ client_secret = os.environ.get("HUBSPOT_OAUTH_CLIENT_SECRET", "")
+ if not client_id or not client_secret:
+ return None
+
+ try:
+ async with httpx.AsyncClient(timeout=15.0) as client:
+ resp = await client.post(HUBSPOT_TOKEN_URL, data={
+ "grant_type": "refresh_token",
+ "refresh_token": refresh_token,
+ "client_id": client_id,
+ "client_secret": client_secret,
+ }, headers={
+ "Content-Type": "application/x-www-form-urlencoded",
+ })
+ if resp.status_code == 200:
+ data = resp.json()
+ tool.oauth_tokens["access_token"] = data["access_token"]
+ tool.oauth_tokens["token_expiry"] = time.time() + data.get("expires_in", 1800)
+ if data.get("refresh_token"):
+ tool.oauth_tokens["refresh_token"] = data["refresh_token"]
+ _save(tool)
+ return data["access_token"]
+ except Exception as e:
+ logger.warning(f"HubSpot token refresh failed for tool {tool.id}: {e}")
+ return None
diff --git a/frontend/src/app/pages/Tools/Tools.tsx b/frontend/src/app/pages/Tools/Tools.tsx
index c176316a..497ff1ec 100644
--- a/frontend/src/app/pages/Tools/Tools.tsx
+++ b/frontend/src/app/pages/Tools/Tools.tsx
@@ -219,6 +219,20 @@ const INTEGRATIONS: Integration[] = [
),
authType: 'oauth2',
},
+ {
+ id: 'hubspot',
+ name: 'HubSpot',
+ description: 'CRM contacts, deals, companies, tickets, and more. Free CRM tier included.',
+ mcp_config: { type: 'stdio', command: 'npx', args: ['-y', '@hubspot/mcp-server'] },
+ color: '#FF7A59',
+ website: 'https://developers.hubspot.com/docs/guides/apps/developer-platform/build-apps/integrate-with-the-remote-hubspot-mcp-server',
+ icon: (
+
+ ),
+ authType: 'oauth2',
+ },
];
const CATEGORY_ORDER = ['filesystem', 'system', 'search', 'interaction', 'agents', 'planning', 'scheduling'];