Compare commits
16 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
705791d17f | ||
|
|
0ec9bbdc13 | ||
|
|
146c28ae27 | ||
|
|
fc61fb062e | ||
|
|
a46402fd08 | ||
|
|
47307a1045 | ||
|
|
b21e355ce1 | ||
|
|
0c69df107e | ||
|
|
4800807783 | ||
|
|
71ecc8f35b | ||
|
|
782c2aceff | ||
|
|
ff27fb157c | ||
|
|
bafbd494ed | ||
|
|
ab1f528a34 | ||
|
|
f319b89806 | ||
|
|
a1f6b828d5 |
@@ -39,6 +39,10 @@ SENDGRID_API_KEY="your-sendgrid-api-key"
|
|||||||
EMAIL_FROM="noreply@yourdomain.com"
|
EMAIL_FROM="noreply@yourdomain.com"
|
||||||
APP_URL="https://yourdomain.com"
|
APP_URL="https://yourdomain.com"
|
||||||
|
|
||||||
|
LANGFUSE_PUBLIC_KEY="your-langfuse-public-key"
|
||||||
|
LANGFUSE_SECRET_KEY="your-langfuse-secret-key"
|
||||||
|
OTEL_EXPORTER_OTLP_ENDPOINT="https://cloud.langfuse.com/api/public/otel"
|
||||||
|
|
||||||
# Server settings
|
# Server settings
|
||||||
HOST="0.0.0.0"
|
HOST="0.0.0.0"
|
||||||
PORT=8000
|
PORT=8000
|
||||||
|
|||||||
201
LICENSE
Normal file
201
LICENSE
Normal file
@@ -0,0 +1,201 @@
|
|||||||
|
Apache License
|
||||||
|
Version 2.0, January 2004
|
||||||
|
http://www.apache.org/licenses/
|
||||||
|
|
||||||
|
TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
|
||||||
|
|
||||||
|
1. Definitions.
|
||||||
|
|
||||||
|
"License" shall mean the terms and conditions for use, reproduction,
|
||||||
|
and distribution as defined by Sections 1 through 9 of this document.
|
||||||
|
|
||||||
|
"Licensor" shall mean the copyright owner or entity authorized by
|
||||||
|
the copyright owner that is granting the License.
|
||||||
|
|
||||||
|
"Legal Entity" shall mean the union of the acting entity and all
|
||||||
|
other entities that control, are controlled by, or are under common
|
||||||
|
control with that entity. For the purposes of this definition,
|
||||||
|
"control" means (i) the power, direct or indirect, to cause the
|
||||||
|
direction or management of such entity, whether by contract or
|
||||||
|
otherwise, or (ii) ownership of fifty percent (50%) or more of the
|
||||||
|
outstanding shares, or (iii) beneficial ownership of such entity.
|
||||||
|
|
||||||
|
"You" (or "Your") shall mean an individual or Legal Entity
|
||||||
|
exercising permissions granted by this License.
|
||||||
|
|
||||||
|
"Source" form shall mean the preferred form for making modifications,
|
||||||
|
including but not limited to software source code, documentation
|
||||||
|
source, and configuration files.
|
||||||
|
|
||||||
|
"Object" form shall mean any form resulting from mechanical
|
||||||
|
transformation or translation of a Source form, including but
|
||||||
|
not limited to compiled object code, generated documentation,
|
||||||
|
and conversions to other media types.
|
||||||
|
|
||||||
|
"Work" shall mean the work of authorship, whether in Source or
|
||||||
|
Object form, made available under the License, as indicated by a
|
||||||
|
copyright notice that is included in or attached to the work
|
||||||
|
(an example is provided in the Appendix below).
|
||||||
|
|
||||||
|
"Derivative Works" shall mean any work, whether in Source or Object
|
||||||
|
form, that is based on (or derived from) the Work and for which the
|
||||||
|
editorial revisions, annotations, elaborations, or other modifications
|
||||||
|
represent, as a whole, an original work of authorship. For the purposes
|
||||||
|
of this License, Derivative Works shall not include works that remain
|
||||||
|
separable from, or merely link (or bind by name) to the interfaces of,
|
||||||
|
the Work and Derivative Works thereof.
|
||||||
|
|
||||||
|
"Contribution" shall mean any work of authorship, including
|
||||||
|
the original version of the Work and any modifications or additions
|
||||||
|
to that Work or Derivative Works thereof, that is intentionally
|
||||||
|
submitted to Licensor for inclusion in the Work by the copyright owner
|
||||||
|
or by an individual or Legal Entity authorized to submit on behalf of
|
||||||
|
the copyright owner. For the purposes of this definition, "submitted"
|
||||||
|
means any form of electronic, verbal, or written communication sent
|
||||||
|
to the Licensor or its representatives, including but not limited to
|
||||||
|
communication on electronic mailing lists, source code control systems,
|
||||||
|
and issue tracking systems that are managed by, or on behalf of, the
|
||||||
|
Licensor for the purpose of discussing and improving the Work, but
|
||||||
|
excluding communication that is conspicuously marked or otherwise
|
||||||
|
designated in writing by the copyright owner as "Not a Contribution."
|
||||||
|
|
||||||
|
"Contributor" shall mean Licensor and any individual or Legal Entity
|
||||||
|
on behalf of whom a Contribution has been received by Licensor and
|
||||||
|
subsequently incorporated within the Work.
|
||||||
|
|
||||||
|
2. Grant of Copyright License. Subject to the terms and conditions of
|
||||||
|
this License, each Contributor hereby grants to You a perpetual,
|
||||||
|
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
|
||||||
|
copyright license to reproduce, prepare Derivative Works of,
|
||||||
|
publicly display, publicly perform, sublicense, and distribute the
|
||||||
|
Work and such Derivative Works in Source or Object form.
|
||||||
|
|
||||||
|
3. Grant of Patent License. Subject to the terms and conditions of
|
||||||
|
this License, each Contributor hereby grants to You a perpetual,
|
||||||
|
worldwide, non-exclusive, no-charge, royalty-free, irrevocable
|
||||||
|
(except as stated in this section) patent license to make, have made,
|
||||||
|
use, offer to sell, sell, import, and otherwise transfer the Work,
|
||||||
|
where such license applies only to those patent claims licensable
|
||||||
|
by such Contributor that are necessarily infringed by their
|
||||||
|
Contribution(s) alone or by combination of their Contribution(s)
|
||||||
|
with the Work to which such Contribution(s) was submitted. If You
|
||||||
|
institute patent litigation against any entity (including a
|
||||||
|
cross-claim or counterclaim in a lawsuit) alleging that the Work
|
||||||
|
or a Contribution incorporated within the Work constitutes direct
|
||||||
|
or contributory patent infringement, then any patent licenses
|
||||||
|
granted to You under this License for that Work shall terminate
|
||||||
|
as of the date such litigation is filed.
|
||||||
|
|
||||||
|
4. Redistribution. You may reproduce and distribute copies of the
|
||||||
|
Work or Derivative Works thereof in any medium, with or without
|
||||||
|
modifications, and in Source or Object form, provided that You
|
||||||
|
meet the following conditions:
|
||||||
|
|
||||||
|
(a) You must give any other recipients of the Work or
|
||||||
|
Derivative Works a copy of this License; and
|
||||||
|
|
||||||
|
(b) You must cause any modified files to carry prominent notices
|
||||||
|
stating that You changed the files; and
|
||||||
|
|
||||||
|
(c) You must retain, in the Source form of any Derivative Works
|
||||||
|
that You distribute, all copyright, patent, trademark, and
|
||||||
|
attribution notices from the Source form of the Work,
|
||||||
|
excluding those notices that do not pertain to any part of
|
||||||
|
the Derivative Works; and
|
||||||
|
|
||||||
|
(d) If the Work includes a "NOTICE" text file as part of its
|
||||||
|
distribution, then any Derivative Works that You distribute must
|
||||||
|
include a readable copy of the attribution notices contained
|
||||||
|
within such NOTICE file, excluding those notices that do not
|
||||||
|
pertain to any part of the Derivative Works, in at least one
|
||||||
|
of the following places: within a NOTICE text file distributed
|
||||||
|
as part of the Derivative Works; within the Source form or
|
||||||
|
documentation, if provided along with the Derivative Works; or,
|
||||||
|
within a display generated by the Derivative Works, if and
|
||||||
|
wherever such third-party notices normally appear. The contents
|
||||||
|
of the NOTICE file are for informational purposes only and
|
||||||
|
do not modify the License. You may add Your own attribution
|
||||||
|
notices within Derivative Works that You distribute, alongside
|
||||||
|
or as an addendum to the NOTICE text from the Work, provided
|
||||||
|
that such additional attribution notices cannot be construed
|
||||||
|
as modifying the License.
|
||||||
|
|
||||||
|
You may add Your own copyright statement to Your modifications and
|
||||||
|
may provide additional or different license terms and conditions
|
||||||
|
for use, reproduction, or distribution of Your modifications, or
|
||||||
|
for any such Derivative Works as a whole, provided Your use,
|
||||||
|
reproduction, and distribution of the Work otherwise complies with
|
||||||
|
the conditions stated in this License.
|
||||||
|
|
||||||
|
5. Submission of Contributions. Unless You explicitly state otherwise,
|
||||||
|
any Contribution intentionally submitted for inclusion in the Work
|
||||||
|
by You to the Licensor shall be under the terms and conditions of
|
||||||
|
this License, without any additional terms or conditions.
|
||||||
|
Notwithstanding the above, nothing herein shall supersede or modify
|
||||||
|
the terms of any separate license agreement you may have executed
|
||||||
|
with Licensor regarding such Contributions.
|
||||||
|
|
||||||
|
6. Trademarks. This License does not grant permission to use the trade
|
||||||
|
names, trademarks, service marks, or product names of the Licensor,
|
||||||
|
except as required for reasonable and customary use in describing the
|
||||||
|
origin of the Work and reproducing the content of the NOTICE file.
|
||||||
|
|
||||||
|
7. Disclaimer of Warranty. Unless required by applicable law or
|
||||||
|
agreed to in writing, Licensor provides the Work (and each
|
||||||
|
Contributor provides its Contributions) on an "AS IS" BASIS,
|
||||||
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
|
||||||
|
implied, including, without limitation, any warranties or conditions
|
||||||
|
of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
|
||||||
|
PARTICULAR PURPOSE. You are solely responsible for determining the
|
||||||
|
appropriateness of using or redistributing the Work and assume any
|
||||||
|
risks associated with Your exercise of permissions under this License.
|
||||||
|
|
||||||
|
8. Limitation of Liability. In no event and under no legal theory,
|
||||||
|
whether in tort (including negligence), contract, or otherwise,
|
||||||
|
unless required by applicable law (such as deliberate and grossly
|
||||||
|
negligent acts) or agreed to in writing, shall any Contributor be
|
||||||
|
liable to You for damages, including any direct, indirect, special,
|
||||||
|
incidental, or consequential damages of any character arising as a
|
||||||
|
result of this License or out of the use or inability to use the
|
||||||
|
Work (including but not limited to damages for loss of goodwill,
|
||||||
|
work stoppage, computer failure or malfunction, or any and all
|
||||||
|
other commercial damages or losses), even if such Contributor
|
||||||
|
has been advised of the possibility of such damages.
|
||||||
|
|
||||||
|
9. Accepting Warranty or Additional Liability. While redistributing
|
||||||
|
the Work or Derivative Works thereof, You may choose to offer,
|
||||||
|
and charge a fee for, acceptance of support, warranty, indemnity,
|
||||||
|
or other liability obligations and/or rights consistent with this
|
||||||
|
License. However, in accepting such obligations, You may act only
|
||||||
|
on Your own behalf and on Your sole responsibility, not on behalf
|
||||||
|
of any other Contributor, and only if You agree to indemnify,
|
||||||
|
defend, and hold each Contributor harmless for any liability
|
||||||
|
incurred by, or claims asserted against, such Contributor by reason
|
||||||
|
of your accepting any such warranty or additional liability.
|
||||||
|
|
||||||
|
END OF TERMS AND CONDITIONS
|
||||||
|
|
||||||
|
APPENDIX: How to apply the Apache License to your work.
|
||||||
|
|
||||||
|
To apply the Apache License to your work, attach the following
|
||||||
|
boilerplate notice, with the fields enclosed by brackets "[]"
|
||||||
|
replaced with your own identifying information. (Don't include
|
||||||
|
the brackets!) The text should be enclosed in the appropriate
|
||||||
|
comment syntax for the file format. We also recommend that a
|
||||||
|
file or class name and description of purpose be included on the
|
||||||
|
same "printed page" as the copyright notice for easier
|
||||||
|
identification within third-party archives.
|
||||||
|
|
||||||
|
Copyright 2025 Evolution API
|
||||||
|
|
||||||
|
Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
you may not use this file except in compliance with the License.
|
||||||
|
You may obtain a copy of the License at
|
||||||
|
|
||||||
|
http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
|
||||||
|
Unless required by applicable law or agreed to in writing, software
|
||||||
|
distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
See the License for the specific language governing permissions and
|
||||||
|
limitations under the License.
|
||||||
57
README.md
57
README.md
@@ -1,5 +1,7 @@
|
|||||||
# Evo AI - AI Agents Platform
|
# Evo AI - AI Agents Platform
|
||||||
|
|
||||||
|
> **License:** This project is licensed under the [Apache License 2.0](./LICENSE). The use of the "Evolution API" trademark is protected as described in section 6 (Trademarks) of the license.
|
||||||
|
|
||||||
Evo AI is an open-source platform for creating and managing AI agents, enabling integration with different AI models and services.
|
Evo AI is an open-source platform for creating and managing AI agents, enabling integration with different AI models and services.
|
||||||
|
|
||||||
## 🚀 Overview
|
## 🚀 Overview
|
||||||
@@ -301,6 +303,51 @@ Authorization: Bearer your-token-jwt
|
|||||||
- **LangGraph**: Framework for building stateful, multi-agent workflows
|
- **LangGraph**: Framework for building stateful, multi-agent workflows
|
||||||
- **ReactFlow**: Library for building node-based visual workflows
|
- **ReactFlow**: Library for building node-based visual workflows
|
||||||
|
|
||||||
|
## 📊 Langfuse Integration (Tracing & Observability)
|
||||||
|
|
||||||
|
Evo AI platform natively supports integration with [Langfuse](https://langfuse.com/) for detailed tracing of agent executions, prompts, model responses, and tool calls, using the OpenTelemetry (OTel) standard.
|
||||||
|
|
||||||
|
### Why use Langfuse?
|
||||||
|
|
||||||
|
- Visual dashboard for agent traces, prompts, and executions
|
||||||
|
- Detailed analytics for debugging and evaluating LLM apps
|
||||||
|
- Easy integration with Google ADK and other frameworks
|
||||||
|
|
||||||
|
### How it works
|
||||||
|
|
||||||
|
- Every agent execution (including streaming) is automatically traced via OpenTelemetry spans
|
||||||
|
- Data is sent to Langfuse, where it can be visualized and analyzed
|
||||||
|
|
||||||
|
### How to configure
|
||||||
|
|
||||||
|
1. **Set environment variables in your `.env`:**
|
||||||
|
|
||||||
|
```env
|
||||||
|
LANGFUSE_PUBLIC_KEY="pk-lf-..." # Your Langfuse public key
|
||||||
|
LANGFUSE_SECRET_KEY="sk-lf-..." # Your Langfuse secret key
|
||||||
|
OTEL_EXPORTER_OTLP_ENDPOINT="https://cloud.langfuse.com/api/public/otel" # (or us.cloud... for US region)
|
||||||
|
```
|
||||||
|
|
||||||
|
> **Attention:** Do not swap the keys! `pk-...` is public, `sk-...` is secret.
|
||||||
|
|
||||||
|
2. **Automatic initialization**
|
||||||
|
|
||||||
|
- Tracing is automatically initialized when the application starts (`src/main.py`).
|
||||||
|
- Agent execution functions are already instrumented with spans (`src/services/agent_runner.py`).
|
||||||
|
|
||||||
|
3. **View in the Langfuse dashboard**
|
||||||
|
- Access your Langfuse dashboard to see real-time traces.
|
||||||
|
|
||||||
|
### Troubleshooting
|
||||||
|
|
||||||
|
- **401 Error (Invalid credentials):**
|
||||||
|
- Check if the keys are correct and not swapped in your `.env`.
|
||||||
|
- Make sure the endpoint matches your region (EU or US).
|
||||||
|
- **Context error in async generator:**
|
||||||
|
- The code is already adjusted to avoid OpenTelemetry context issues in async generators.
|
||||||
|
- **Questions about integration:**
|
||||||
|
- See the [official Langfuse documentation - Google ADK](https://langfuse.com/docs/integrations/google-adk)
|
||||||
|
|
||||||
## 🤖 Agent 2 Agent (A2A) Protocol Support
|
## 🤖 Agent 2 Agent (A2A) Protocol Support
|
||||||
|
|
||||||
Evo AI implements the Google's Agent 2 Agent (A2A) protocol, enabling seamless communication and interoperability between AI agents. This implementation includes:
|
Evo AI implements the Google's Agent 2 Agent (A2A) protocol, enabling seamless communication and interoperability between AI agents. This implementation includes:
|
||||||
@@ -398,7 +445,7 @@ You'll also need the following accounts/API keys:
|
|||||||
1. Clone the repository:
|
1. Clone the repository:
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
git clone https://github.com/your-username/evo-ai.git
|
git clone https://github.com/EvolutionAPI/evo-ai.git
|
||||||
cd evo-ai
|
cd evo-ai
|
||||||
```
|
```
|
||||||
|
|
||||||
@@ -620,15 +667,17 @@ Please read our [Contributing Guidelines](CONTRIBUTING.md) for more details.
|
|||||||
|
|
||||||
## 📄 License
|
## 📄 License
|
||||||
|
|
||||||
This project is licensed under the MIT License - see the [LICENSE](LICENSE) file for details.
|
This project is licensed under the [Apache License 2.0](./LICENSE).
|
||||||
|
|
||||||
|
The use of the name, logo, or trademark "Evolution API" is protected and not automatically granted by the license. See section 6 (Trademarks) of the license for details about trademark usage.
|
||||||
|
|
||||||
## 📊 Stargazers
|
## 📊 Stargazers
|
||||||
|
|
||||||
[](https://github.com/your-username/evo-ai/stargazers)
|
[](https://github.com/EvolutionAPI/evo-ai/stargazers)
|
||||||
|
|
||||||
## 🔄 Forks
|
## 🔄 Forks
|
||||||
|
|
||||||
[](https://github.com/your-username/evo-ai/network/members)
|
[](https://github.com/EvolutionAPI/evo-ai/network/members)
|
||||||
|
|
||||||
## 🙏 Acknowledgments
|
## 🙏 Acknowledgments
|
||||||
|
|
||||||
|
|||||||
@@ -11,11 +11,11 @@ authors = [
|
|||||||
{name = "EvoAI Team", email = "admin@evoai.com"}
|
{name = "EvoAI Team", email = "admin@evoai.com"}
|
||||||
]
|
]
|
||||||
requires-python = ">=3.10"
|
requires-python = ">=3.10"
|
||||||
license = {text = "Proprietary"}
|
license = {text = "Apache-2.0"}
|
||||||
classifiers = [
|
classifiers = [
|
||||||
"Programming Language :: Python :: 3",
|
"Programming Language :: Python :: 3",
|
||||||
"Programming Language :: Python :: 3.10",
|
"Programming Language :: Python :: 3.10",
|
||||||
"License :: Other/Proprietary License",
|
"License :: OSI Approved :: Apache Software License",
|
||||||
"Operating System :: OS Independent",
|
"Operating System :: OS Independent",
|
||||||
]
|
]
|
||||||
|
|
||||||
@@ -49,6 +49,8 @@ dependencies = [
|
|||||||
"jwcrypto==1.5.6",
|
"jwcrypto==1.5.6",
|
||||||
"pyjwt[crypto]==2.9.0",
|
"pyjwt[crypto]==2.9.0",
|
||||||
"langgraph==0.4.1",
|
"langgraph==0.4.1",
|
||||||
|
"opentelemetry-sdk==1.33.0",
|
||||||
|
"opentelemetry-exporter-otlp==1.33.0",
|
||||||
]
|
]
|
||||||
|
|
||||||
[project.optional-dependencies]
|
[project.optional-dependencies]
|
||||||
|
|||||||
@@ -54,8 +54,7 @@ async def register_user(user_data: UserCreate, db: Session = Depends(get_db)):
|
|||||||
Raises:
|
Raises:
|
||||||
HTTPException: If there is an error in registration
|
HTTPException: If there is an error in registration
|
||||||
"""
|
"""
|
||||||
# TODO: remover o auto_verify temporariamente para teste
|
user, message = create_user(db, user_data, is_admin=False, auto_verify=False)
|
||||||
user, message = create_user(db, user_data, is_admin=False, auto_verify=True)
|
|
||||||
if not user:
|
if not user:
|
||||||
logger.error(f"Error registering user: {message}")
|
logger.error(f"Error registering user: {message}")
|
||||||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=message)
|
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=message)
|
||||||
|
|||||||
@@ -102,9 +102,8 @@ async def websocket_chat(
|
|||||||
memory_service=memory_service,
|
memory_service=memory_service,
|
||||||
db=db,
|
db=db,
|
||||||
):
|
):
|
||||||
# Send each chunk as a JSON message
|
|
||||||
await websocket.send_json(
|
await websocket.send_json(
|
||||||
{"message": chunk, "turn_complete": False}
|
{"message": json.loads(chunk), "turn_complete": False}
|
||||||
)
|
)
|
||||||
|
|
||||||
# Send signal of complete turn
|
# Send signal of complete turn
|
||||||
|
|||||||
@@ -84,6 +84,11 @@ class Settings(BaseSettings):
|
|||||||
DEMO_PASSWORD: str = os.getenv("DEMO_PASSWORD", "demo123")
|
DEMO_PASSWORD: str = os.getenv("DEMO_PASSWORD", "demo123")
|
||||||
DEMO_CLIENT_NAME: str = os.getenv("DEMO_CLIENT_NAME", "Demo Client")
|
DEMO_CLIENT_NAME: str = os.getenv("DEMO_CLIENT_NAME", "Demo Client")
|
||||||
|
|
||||||
|
# Langfuse / OpenTelemetry settings
|
||||||
|
LANGFUSE_PUBLIC_KEY: str = os.getenv("LANGFUSE_PUBLIC_KEY", "")
|
||||||
|
LANGFUSE_SECRET_KEY: str = os.getenv("LANGFUSE_SECRET_KEY", "")
|
||||||
|
OTEL_EXPORTER_OTLP_ENDPOINT: str = os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT", "")
|
||||||
|
|
||||||
class Config:
|
class Config:
|
||||||
env_file = ".env"
|
env_file = ".env"
|
||||||
env_file_encoding = "utf-8"
|
env_file_encoding = "utf-8"
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ from fastapi.staticfiles import StaticFiles
|
|||||||
from src.config.database import engine, Base
|
from src.config.database import engine, Base
|
||||||
from src.config.settings import settings
|
from src.config.settings import settings
|
||||||
from src.utils.logger import setup_logger
|
from src.utils.logger import setup_logger
|
||||||
|
from src.utils.otel import init_otel
|
||||||
|
|
||||||
# Necessary for other modules
|
# Necessary for other modules
|
||||||
from src.services.service_providers import session_service # noqa: F401
|
from src.services.service_providers import session_service # noqa: F401
|
||||||
@@ -85,6 +86,9 @@ app.include_router(session_router, prefix=API_PREFIX)
|
|||||||
app.include_router(agent_router, prefix=API_PREFIX)
|
app.include_router(agent_router, prefix=API_PREFIX)
|
||||||
app.include_router(a2a_router, prefix=API_PREFIX)
|
app.include_router(a2a_router, prefix=API_PREFIX)
|
||||||
|
|
||||||
|
# Inicializa o OpenTelemetry para Langfuse
|
||||||
|
init_otel()
|
||||||
|
|
||||||
|
|
||||||
@app.get("/")
|
@app.get("/")
|
||||||
def read_root():
|
def read_root():
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ import asyncio
|
|||||||
from collections.abc import AsyncIterable
|
from collections.abc import AsyncIterable
|
||||||
from typing import Dict, Optional
|
from typing import Dict, Optional
|
||||||
from uuid import UUID
|
from uuid import UUID
|
||||||
|
import json
|
||||||
|
|
||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
@@ -306,7 +307,7 @@ class A2ATaskManager:
|
|||||||
external_id = task_params.sessionId
|
external_id = task_params.sessionId
|
||||||
full_response = ""
|
full_response = ""
|
||||||
|
|
||||||
# We use the same streaming function used in the WebSocket
|
# Use the same pattern as chat_routes.py: deserialize each chunk
|
||||||
async for chunk in run_agent_stream(
|
async for chunk in run_agent_stream(
|
||||||
agent_id=str(agent.id),
|
agent_id=str(agent.id),
|
||||||
external_id=external_id,
|
external_id=external_id,
|
||||||
@@ -316,9 +317,14 @@ class A2ATaskManager:
|
|||||||
memory_service=memory_service,
|
memory_service=memory_service,
|
||||||
db=self.db,
|
db=self.db,
|
||||||
):
|
):
|
||||||
# Send incremental progress updates
|
try:
|
||||||
update_text_part = {"type": "text", "text": chunk}
|
chunk_data = json.loads(chunk)
|
||||||
update_message = Message(role="agent", parts=[update_text_part])
|
except Exception as e:
|
||||||
|
logger.warning(f"Invalid chunk received: {chunk} - {e}")
|
||||||
|
continue
|
||||||
|
|
||||||
|
# The chunk_data must be a dict with 'type' and 'text' (or other expected format)
|
||||||
|
update_message = Message(role="agent", parts=[chunk_data])
|
||||||
|
|
||||||
# Update the task with each intermediate message
|
# Update the task with each intermediate message
|
||||||
await self.update_store(
|
await self.update_store(
|
||||||
@@ -337,24 +343,24 @@ class A2ATaskManager:
|
|||||||
final=False,
|
final=False,
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
full_response += chunk
|
# If it's text, accumulate for the final response
|
||||||
|
if chunk_data.get("type") == "text":
|
||||||
|
full_response += chunk_data.get("text", "")
|
||||||
|
|
||||||
# Determine the task state
|
# Determine the final state of the task
|
||||||
task_state = (
|
task_state = (
|
||||||
TaskState.INPUT_REQUIRED
|
TaskState.INPUT_REQUIRED
|
||||||
if "MISSING_INFO:" in full_response
|
if "MISSING_INFO:" in full_response
|
||||||
else TaskState.COMPLETED
|
else TaskState.COMPLETED
|
||||||
)
|
)
|
||||||
|
|
||||||
# Create the final response part
|
# Create the final response
|
||||||
final_text_part = {"type": "text", "text": full_response}
|
final_text_part = {"type": "text", "text": full_response}
|
||||||
parts = [final_text_part]
|
parts = [final_text_part]
|
||||||
final_message = Message(role="agent", parts=parts)
|
final_message = Message(role="agent", parts=parts)
|
||||||
|
|
||||||
# Create the final artifact from the final response
|
|
||||||
final_artifact = Artifact(parts=parts, index=0)
|
final_artifact = Artifact(parts=parts, index=0)
|
||||||
|
|
||||||
# Update the task in the store with the final response
|
# Update the task with the final response
|
||||||
await self.update_store(
|
await self.update_store(
|
||||||
task_params.id,
|
task_params.id,
|
||||||
TaskStatus(state=task_state, message=final_message),
|
TaskStatus(state=task_state, message=final_message),
|
||||||
|
|||||||
@@ -10,6 +10,9 @@ from src.services.agent_builder import AgentBuilder
|
|||||||
from sqlalchemy.orm import Session
|
from sqlalchemy.orm import Session
|
||||||
from typing import Optional, AsyncGenerator
|
from typing import Optional, AsyncGenerator
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import json
|
||||||
|
from src.utils.otel import get_tracer
|
||||||
|
from opentelemetry import trace
|
||||||
|
|
||||||
logger = setup_logger(__name__)
|
logger = setup_logger(__name__)
|
||||||
|
|
||||||
@@ -25,159 +28,184 @@ async def run_agent(
|
|||||||
session_id: Optional[str] = None,
|
session_id: Optional[str] = None,
|
||||||
timeout: float = 60.0,
|
timeout: float = 60.0,
|
||||||
):
|
):
|
||||||
exit_stack = None
|
tracer = get_tracer()
|
||||||
try:
|
with tracer.start_as_current_span(
|
||||||
logger.info(
|
"run_agent",
|
||||||
f"Starting execution of agent {agent_id} for external_id {external_id}"
|
attributes={
|
||||||
)
|
"agent_id": agent_id,
|
||||||
logger.info(f"Received message: {message}")
|
"external_id": external_id,
|
||||||
|
"session_id": session_id or f"{external_id}_{agent_id}",
|
||||||
|
"message": message,
|
||||||
|
},
|
||||||
|
):
|
||||||
|
exit_stack = None
|
||||||
|
try:
|
||||||
|
logger.info(
|
||||||
|
f"Starting execution of agent {agent_id} for external_id {external_id}"
|
||||||
|
)
|
||||||
|
logger.info(f"Received message: {message}")
|
||||||
|
|
||||||
get_root_agent = get_agent(db, agent_id)
|
get_root_agent = get_agent(db, agent_id)
|
||||||
logger.info(
|
logger.info(
|
||||||
f"Root agent found: {get_root_agent.name} (type: {get_root_agent.type})"
|
f"Root agent found: {get_root_agent.name} (type: {get_root_agent.type})"
|
||||||
)
|
)
|
||||||
|
|
||||||
if get_root_agent is None:
|
if get_root_agent is None:
|
||||||
raise AgentNotFoundError(f"Agent with ID {agent_id} not found")
|
raise AgentNotFoundError(f"Agent with ID {agent_id} not found")
|
||||||
|
|
||||||
# Using the AgentBuilder to create the agent
|
# Using the AgentBuilder to create the agent
|
||||||
agent_builder = AgentBuilder(db)
|
agent_builder = AgentBuilder(db)
|
||||||
root_agent, exit_stack = await agent_builder.build_agent(get_root_agent)
|
root_agent, exit_stack = await agent_builder.build_agent(get_root_agent)
|
||||||
|
|
||||||
logger.info("Configuring Runner")
|
logger.info("Configuring Runner")
|
||||||
agent_runner = Runner(
|
agent_runner = Runner(
|
||||||
agent=root_agent,
|
agent=root_agent,
|
||||||
app_name=agent_id,
|
app_name=agent_id,
|
||||||
session_service=session_service,
|
session_service=session_service,
|
||||||
artifact_service=artifacts_service,
|
artifact_service=artifacts_service,
|
||||||
memory_service=memory_service,
|
memory_service=memory_service,
|
||||||
)
|
)
|
||||||
adk_session_id = external_id + "_" + agent_id
|
adk_session_id = external_id + "_" + agent_id
|
||||||
if session_id is None:
|
if session_id is None:
|
||||||
session_id = adk_session_id
|
session_id = adk_session_id
|
||||||
|
|
||||||
logger.info(f"Searching session for external_id {external_id}")
|
logger.info(f"Searching session for external_id {external_id}")
|
||||||
session = session_service.get_session(
|
session = session_service.get_session(
|
||||||
app_name=agent_id,
|
|
||||||
user_id=external_id,
|
|
||||||
session_id=adk_session_id,
|
|
||||||
)
|
|
||||||
|
|
||||||
if session is None:
|
|
||||||
logger.info(f"Creating new session for external_id {external_id}")
|
|
||||||
session = session_service.create_session(
|
|
||||||
app_name=agent_id,
|
app_name=agent_id,
|
||||||
user_id=external_id,
|
user_id=external_id,
|
||||||
session_id=adk_session_id,
|
session_id=adk_session_id,
|
||||||
)
|
)
|
||||||
|
|
||||||
content = Content(role="user", parts=[Part(text=message)])
|
if session is None:
|
||||||
logger.info("Starting agent execution")
|
logger.info(f"Creating new session for external_id {external_id}")
|
||||||
|
session = session_service.create_session(
|
||||||
|
app_name=agent_id,
|
||||||
|
user_id=external_id,
|
||||||
|
session_id=adk_session_id,
|
||||||
|
)
|
||||||
|
|
||||||
final_response_text = "No final response captured."
|
content = Content(role="user", parts=[Part(text=message)])
|
||||||
try:
|
logger.info("Starting agent execution")
|
||||||
response_queue = asyncio.Queue()
|
|
||||||
execution_completed = asyncio.Event()
|
|
||||||
|
|
||||||
async def process_events():
|
|
||||||
try:
|
|
||||||
events_async = agent_runner.run_async(
|
|
||||||
user_id=external_id,
|
|
||||||
session_id=adk_session_id,
|
|
||||||
new_message=content,
|
|
||||||
)
|
|
||||||
|
|
||||||
last_response = None
|
|
||||||
all_responses = []
|
|
||||||
|
|
||||||
async for event in events_async:
|
|
||||||
if (
|
|
||||||
event.content
|
|
||||||
and event.content.parts
|
|
||||||
and event.content.parts[0].text
|
|
||||||
):
|
|
||||||
current_text = event.content.parts[0].text
|
|
||||||
last_response = current_text
|
|
||||||
all_responses.append(current_text)
|
|
||||||
|
|
||||||
if event.actions and event.actions.escalate:
|
|
||||||
escalate_text = f"Agent escalated: {event.error_message or 'No specific message.'}"
|
|
||||||
await response_queue.put(escalate_text)
|
|
||||||
execution_completed.set()
|
|
||||||
return
|
|
||||||
|
|
||||||
if last_response:
|
|
||||||
await response_queue.put(last_response)
|
|
||||||
else:
|
|
||||||
await response_queue.put("Finished without specific response")
|
|
||||||
|
|
||||||
execution_completed.set()
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Error in process_events: {str(e)}")
|
|
||||||
await response_queue.put(f"Error: {str(e)}")
|
|
||||||
execution_completed.set()
|
|
||||||
|
|
||||||
task = asyncio.create_task(process_events())
|
|
||||||
|
|
||||||
|
final_response_text = "No final response captured."
|
||||||
try:
|
try:
|
||||||
wait_task = asyncio.create_task(execution_completed.wait())
|
response_queue = asyncio.Queue()
|
||||||
done, pending = await asyncio.wait({wait_task}, timeout=timeout)
|
execution_completed = asyncio.Event()
|
||||||
|
|
||||||
for p in pending:
|
async def process_events():
|
||||||
p.cancel()
|
try:
|
||||||
|
events_async = agent_runner.run_async(
|
||||||
|
user_id=external_id,
|
||||||
|
session_id=adk_session_id,
|
||||||
|
new_message=content,
|
||||||
|
)
|
||||||
|
|
||||||
if not execution_completed.is_set():
|
last_response = None
|
||||||
logger.warning(f"Agent execution timed out after {timeout} seconds")
|
all_responses = []
|
||||||
await response_queue.put(
|
|
||||||
"The response took too long and was interrupted."
|
|
||||||
)
|
|
||||||
|
|
||||||
final_response_text = await response_queue.get()
|
async for event in events_async:
|
||||||
|
if (
|
||||||
|
event.content
|
||||||
|
and event.content.parts
|
||||||
|
and event.content.parts[0].text
|
||||||
|
):
|
||||||
|
current_text = event.content.parts[0].text
|
||||||
|
last_response = current_text
|
||||||
|
all_responses.append(current_text)
|
||||||
|
|
||||||
|
if event.actions and event.actions.escalate:
|
||||||
|
escalate_text = f"Agent escalated: {event.error_message or 'No specific message.'}"
|
||||||
|
await response_queue.put(escalate_text)
|
||||||
|
execution_completed.set()
|
||||||
|
return
|
||||||
|
|
||||||
|
if last_response:
|
||||||
|
await response_queue.put(last_response)
|
||||||
|
else:
|
||||||
|
await response_queue.put(
|
||||||
|
"Finished without specific response"
|
||||||
|
)
|
||||||
|
|
||||||
|
execution_completed.set()
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Error in process_events: {str(e)}")
|
||||||
|
await response_queue.put(f"Error: {str(e)}")
|
||||||
|
execution_completed.set()
|
||||||
|
|
||||||
|
task = asyncio.create_task(process_events())
|
||||||
|
|
||||||
|
try:
|
||||||
|
wait_task = asyncio.create_task(execution_completed.wait())
|
||||||
|
done, pending = await asyncio.wait({wait_task}, timeout=timeout)
|
||||||
|
|
||||||
|
for p in pending:
|
||||||
|
p.cancel()
|
||||||
|
|
||||||
|
if not execution_completed.is_set():
|
||||||
|
logger.warning(
|
||||||
|
f"Agent execution timed out after {timeout} seconds"
|
||||||
|
)
|
||||||
|
await response_queue.put(
|
||||||
|
"The response took too long and was interrupted."
|
||||||
|
)
|
||||||
|
|
||||||
|
final_response_text = await response_queue.get()
|
||||||
|
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Error waiting for response: {str(e)}")
|
||||||
|
final_response_text = f"Error processing response: {str(e)}"
|
||||||
|
|
||||||
|
# Add the session to memory after completion
|
||||||
|
completed_session = session_service.get_session(
|
||||||
|
app_name=agent_id,
|
||||||
|
user_id=external_id,
|
||||||
|
session_id=adk_session_id,
|
||||||
|
)
|
||||||
|
|
||||||
|
memory_service.add_session_to_memory(completed_session)
|
||||||
|
|
||||||
|
# Cancel the processing task if it is still running
|
||||||
|
if not task.done():
|
||||||
|
task.cancel()
|
||||||
|
try:
|
||||||
|
await task
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
logger.info("Task cancelled successfully")
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Error cancelling task: {str(e)}")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error waiting for response: {str(e)}")
|
logger.error(f"Error processing request: {str(e)}")
|
||||||
final_response_text = f"Error processing response: {str(e)}"
|
raise e
|
||||||
|
|
||||||
# Add the session to memory after completion
|
logger.info("Agent execution completed successfully")
|
||||||
completed_session = session_service.get_session(
|
return final_response_text
|
||||||
app_name=agent_id,
|
except AgentNotFoundError as e:
|
||||||
user_id=external_id,
|
|
||||||
session_id=adk_session_id,
|
|
||||||
)
|
|
||||||
|
|
||||||
memory_service.add_session_to_memory(completed_session)
|
|
||||||
|
|
||||||
# Cancel the processing task if it is still running
|
|
||||||
if not task.done():
|
|
||||||
task.cancel()
|
|
||||||
try:
|
|
||||||
await task
|
|
||||||
except asyncio.CancelledError:
|
|
||||||
logger.info("Task cancelled successfully")
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Error cancelling task: {str(e)}")
|
|
||||||
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Error processing request: {str(e)}")
|
logger.error(f"Error processing request: {str(e)}")
|
||||||
raise e
|
raise e
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Internal error processing request: {str(e)}", exc_info=True)
|
||||||
|
raise InternalServerError(str(e))
|
||||||
|
finally:
|
||||||
|
# Clean up MCP connection - MUST be executed in the same task
|
||||||
|
if exit_stack:
|
||||||
|
logger.info("Closing MCP server connection...")
|
||||||
|
try:
|
||||||
|
await exit_stack.aclose()
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Error closing MCP connection: {e}")
|
||||||
|
# Do not raise the exception to not obscure the original error
|
||||||
|
|
||||||
logger.info("Agent execution completed successfully")
|
|
||||||
return final_response_text
|
def convert_sets(obj):
|
||||||
except AgentNotFoundError as e:
|
if isinstance(obj, set):
|
||||||
logger.error(f"Error processing request: {str(e)}")
|
return list(obj)
|
||||||
raise e
|
elif isinstance(obj, dict):
|
||||||
except Exception as e:
|
return {k: convert_sets(v) for k, v in obj.items()}
|
||||||
logger.error(f"Internal error processing request: {str(e)}", exc_info=True)
|
elif isinstance(obj, list):
|
||||||
raise InternalServerError(str(e))
|
return [convert_sets(i) for i in obj]
|
||||||
finally:
|
else:
|
||||||
# Clean up MCP connection - MUST be executed in the same task
|
return obj
|
||||||
if exit_stack:
|
|
||||||
logger.info("Closing MCP server connection...")
|
|
||||||
try:
|
|
||||||
await exit_stack.aclose()
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Error closing MCP connection: {e}")
|
|
||||||
# Do not raise the exception to not obscure the original error
|
|
||||||
|
|
||||||
|
|
||||||
async def run_agent_stream(
|
async def run_agent_stream(
|
||||||
@@ -190,91 +218,105 @@ async def run_agent_stream(
|
|||||||
db: Session,
|
db: Session,
|
||||||
session_id: Optional[str] = None,
|
session_id: Optional[str] = None,
|
||||||
) -> AsyncGenerator[str, None]:
|
) -> AsyncGenerator[str, None]:
|
||||||
|
tracer = get_tracer()
|
||||||
|
span = tracer.start_span(
|
||||||
|
"run_agent_stream",
|
||||||
|
attributes={
|
||||||
|
"agent_id": agent_id,
|
||||||
|
"external_id": external_id,
|
||||||
|
"session_id": session_id or f"{external_id}_{agent_id}",
|
||||||
|
"message": message,
|
||||||
|
},
|
||||||
|
)
|
||||||
try:
|
try:
|
||||||
logger.info(
|
with trace.use_span(span, end_on_exit=True):
|
||||||
f"Starting streaming execution of agent {agent_id} for external_id {external_id}"
|
try:
|
||||||
)
|
logger.info(
|
||||||
logger.info(f"Received message: {message}")
|
f"Starting streaming execution of agent {agent_id} for external_id {external_id}"
|
||||||
|
)
|
||||||
|
logger.info(f"Received message: {message}")
|
||||||
|
|
||||||
get_root_agent = get_agent(db, agent_id)
|
get_root_agent = get_agent(db, agent_id)
|
||||||
logger.info(
|
logger.info(
|
||||||
f"Root agent found: {get_root_agent.name} (type: {get_root_agent.type})"
|
f"Root agent found: {get_root_agent.name} (type: {get_root_agent.type})"
|
||||||
)
|
)
|
||||||
|
|
||||||
if get_root_agent is None:
|
if get_root_agent is None:
|
||||||
raise AgentNotFoundError(f"Agent with ID {agent_id} not found")
|
raise AgentNotFoundError(f"Agent with ID {agent_id} not found")
|
||||||
|
|
||||||
# Using the AgentBuilder to create the agent
|
# Using the AgentBuilder to create the agent
|
||||||
agent_builder = AgentBuilder(db)
|
agent_builder = AgentBuilder(db)
|
||||||
root_agent, exit_stack = await agent_builder.build_agent(get_root_agent)
|
root_agent, exit_stack = await agent_builder.build_agent(get_root_agent)
|
||||||
|
|
||||||
logger.info("Configuring Runner")
|
logger.info("Configuring Runner")
|
||||||
agent_runner = Runner(
|
agent_runner = Runner(
|
||||||
agent=root_agent,
|
agent=root_agent,
|
||||||
app_name=agent_id,
|
app_name=agent_id,
|
||||||
session_service=session_service,
|
session_service=session_service,
|
||||||
artifact_service=artifacts_service,
|
artifact_service=artifacts_service,
|
||||||
memory_service=memory_service,
|
memory_service=memory_service,
|
||||||
)
|
)
|
||||||
adk_session_id = external_id + "_" + agent_id
|
adk_session_id = external_id + "_" + agent_id
|
||||||
if session_id is None:
|
if session_id is None:
|
||||||
session_id = adk_session_id
|
session_id = adk_session_id
|
||||||
|
|
||||||
logger.info(f"Searching session for external_id {external_id}")
|
logger.info(f"Searching session for external_id {external_id}")
|
||||||
session = session_service.get_session(
|
session = session_service.get_session(
|
||||||
app_name=agent_id,
|
app_name=agent_id,
|
||||||
user_id=external_id,
|
user_id=external_id,
|
||||||
session_id=adk_session_id,
|
session_id=adk_session_id,
|
||||||
)
|
)
|
||||||
|
|
||||||
if session is None:
|
if session is None:
|
||||||
logger.info(f"Creating new session for external_id {external_id}")
|
logger.info(f"Creating new session for external_id {external_id}")
|
||||||
session = session_service.create_session(
|
session = session_service.create_session(
|
||||||
app_name=agent_id,
|
app_name=agent_id,
|
||||||
user_id=external_id,
|
user_id=external_id,
|
||||||
session_id=adk_session_id,
|
session_id=adk_session_id,
|
||||||
)
|
)
|
||||||
|
|
||||||
content = Content(role="user", parts=[Part(text=message)])
|
content = Content(role="user", parts=[Part(text=message)])
|
||||||
logger.info("Starting agent streaming execution")
|
logger.info("Starting agent streaming execution")
|
||||||
|
|
||||||
try:
|
|
||||||
events_async = agent_runner.run_async(
|
|
||||||
user_id=external_id,
|
|
||||||
session_id=adk_session_id,
|
|
||||||
new_message=content,
|
|
||||||
)
|
|
||||||
|
|
||||||
async for event in events_async:
|
|
||||||
if event.content and event.content.parts:
|
|
||||||
text = event.content.parts[0].text
|
|
||||||
if text:
|
|
||||||
yield text
|
|
||||||
await asyncio.sleep(0) # Allow other tasks to run
|
|
||||||
|
|
||||||
completed_session = session_service.get_session(
|
|
||||||
app_name=agent_id,
|
|
||||||
user_id=external_id,
|
|
||||||
session_id=adk_session_id,
|
|
||||||
)
|
|
||||||
|
|
||||||
memory_service.add_session_to_memory(completed_session)
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Error processing request: {str(e)}")
|
|
||||||
raise e
|
|
||||||
finally:
|
|
||||||
# Clean up MCP connection
|
|
||||||
if exit_stack:
|
|
||||||
logger.info("Closing MCP server connection...")
|
|
||||||
try:
|
try:
|
||||||
await exit_stack.aclose()
|
events_async = agent_runner.run_async(
|
||||||
except Exception as e:
|
user_id=external_id,
|
||||||
logger.error(f"Error closing MCP connection: {e}")
|
session_id=adk_session_id,
|
||||||
|
new_message=content,
|
||||||
|
)
|
||||||
|
|
||||||
logger.info("Agent streaming execution completed successfully")
|
async for event in events_async:
|
||||||
except AgentNotFoundError as e:
|
event_dict = event.dict()
|
||||||
logger.error(f"Error processing request: {str(e)}")
|
event_dict = convert_sets(event_dict)
|
||||||
raise e
|
yield json.dumps(event_dict)
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"Internal error processing request: {str(e)}", exc_info=True)
|
completed_session = session_service.get_session(
|
||||||
raise InternalServerError(str(e))
|
app_name=agent_id,
|
||||||
|
user_id=external_id,
|
||||||
|
session_id=adk_session_id,
|
||||||
|
)
|
||||||
|
|
||||||
|
memory_service.add_session_to_memory(completed_session)
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Error processing request: {str(e)}")
|
||||||
|
raise e
|
||||||
|
finally:
|
||||||
|
# Clean up MCP connection
|
||||||
|
if exit_stack:
|
||||||
|
logger.info("Closing MCP server connection...")
|
||||||
|
try:
|
||||||
|
await exit_stack.aclose()
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Error closing MCP connection: {e}")
|
||||||
|
|
||||||
|
logger.info("Agent streaming execution completed successfully")
|
||||||
|
except AgentNotFoundError as e:
|
||||||
|
logger.error(f"Error processing request: {str(e)}")
|
||||||
|
raise e
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(
|
||||||
|
f"Internal error processing request: {str(e)}", exc_info=True
|
||||||
|
)
|
||||||
|
raise InternalServerError(str(e))
|
||||||
|
finally:
|
||||||
|
span.end()
|
||||||
|
|||||||
@@ -1,6 +1,5 @@
|
|||||||
import sendgrid
|
import sendgrid
|
||||||
from sendgrid.helpers.mail import Mail, Email, To, Content
|
from sendgrid.helpers.mail import Mail, Email, To, Content
|
||||||
from src.config.settings import settings
|
|
||||||
import logging
|
import logging
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from jinja2 import Environment, FileSystemLoader, select_autoescape
|
from jinja2 import Environment, FileSystemLoader, select_autoescape
|
||||||
@@ -51,12 +50,12 @@ def send_verification_email(email: str, token: str) -> bool:
|
|||||||
bool: True if the email was sent successfully, False otherwise
|
bool: True if the email was sent successfully, False otherwise
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
sg = sendgrid.SendGridAPIClient(api_key=settings.SENDGRID_API_KEY)
|
sg = sendgrid.SendGridAPIClient(api_key=os.getenv("SENDGRID_API_KEY"))
|
||||||
from_email = Email(settings.EMAIL_FROM)
|
from_email = Email(os.getenv("EMAIL_FROM"))
|
||||||
to_email = To(email)
|
to_email = To(email)
|
||||||
subject = "Email Verification - Evo AI"
|
subject = "Email Verification - Evo AI"
|
||||||
|
|
||||||
verification_link = f"{settings.APP_URL}/api/v1/auth/verify-email/{token}"
|
verification_link = f"{os.getenv('APP_URL')}/security/verify-email?code={token}"
|
||||||
|
|
||||||
html_content = _render_template(
|
html_content = _render_template(
|
||||||
"verification_email",
|
"verification_email",
|
||||||
@@ -100,12 +99,12 @@ def send_password_reset_email(email: str, token: str) -> bool:
|
|||||||
bool: True if the email was sent successfully, False otherwise
|
bool: True if the email was sent successfully, False otherwise
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
sg = sendgrid.SendGridAPIClient(api_key=settings.SENDGRID_API_KEY)
|
sg = sendgrid.SendGridAPIClient(api_key=os.getenv("SENDGRID_API_KEY"))
|
||||||
from_email = Email(settings.EMAIL_FROM)
|
from_email = Email(os.getenv("EMAIL_FROM"))
|
||||||
to_email = To(email)
|
to_email = To(email)
|
||||||
subject = "Password Reset - Evo AI"
|
subject = "Password Reset - Evo AI"
|
||||||
|
|
||||||
reset_link = f"{settings.APP_URL}/reset-password?token={token}"
|
reset_link = f"{os.getenv('APP_URL')}/security/reset-password?token={token}"
|
||||||
|
|
||||||
html_content = _render_template(
|
html_content = _render_template(
|
||||||
"password_reset",
|
"password_reset",
|
||||||
@@ -149,12 +148,12 @@ def send_welcome_email(email: str, user_name: str = None) -> bool:
|
|||||||
bool: True if the email was sent successfully, False otherwise
|
bool: True if the email was sent successfully, False otherwise
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
sg = sendgrid.SendGridAPIClient(api_key=settings.SENDGRID_API_KEY)
|
sg = sendgrid.SendGridAPIClient(api_key=os.getenv("SENDGRID_API_KEY"))
|
||||||
from_email = Email(settings.EMAIL_FROM)
|
from_email = Email(os.getenv("EMAIL_FROM"))
|
||||||
to_email = To(email)
|
to_email = To(email)
|
||||||
subject = "Welcome to Evo AI"
|
subject = "Welcome to Evo AI"
|
||||||
|
|
||||||
dashboard_link = f"{settings.APP_URL}/dashboard"
|
dashboard_link = f"{os.getenv('APP_URL')}/dashboard"
|
||||||
|
|
||||||
html_content = _render_template(
|
html_content = _render_template(
|
||||||
"welcome_email",
|
"welcome_email",
|
||||||
@@ -200,12 +199,14 @@ def send_account_locked_email(
|
|||||||
bool: True if the email was sent successfully, False otherwise
|
bool: True if the email was sent successfully, False otherwise
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
sg = sendgrid.SendGridAPIClient(api_key=settings.SENDGRID_API_KEY)
|
sg = sendgrid.SendGridAPIClient(api_key=os.getenv("SENDGRID_API_KEY"))
|
||||||
from_email = Email(settings.EMAIL_FROM)
|
from_email = Email(os.getenv("EMAIL_FROM"))
|
||||||
to_email = To(email)
|
to_email = To(email)
|
||||||
subject = "Security Alert - Account Locked"
|
subject = "Security Alert - Account Locked"
|
||||||
|
|
||||||
reset_link = f"{settings.APP_URL}/reset-password?token={reset_token}"
|
reset_link = (
|
||||||
|
f"{os.getenv('APP_URL')}/security/reset-password?token={reset_token}"
|
||||||
|
)
|
||||||
|
|
||||||
html_content = _render_template(
|
html_content = _render_template(
|
||||||
"account_locked",
|
"account_locked",
|
||||||
|
|||||||
@@ -7,7 +7,7 @@ from src.services.email_service import (
|
|||||||
send_verification_email,
|
send_verification_email,
|
||||||
send_password_reset_email,
|
send_password_reset_email,
|
||||||
)
|
)
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta, timezone
|
||||||
import uuid
|
import uuid
|
||||||
import logging
|
import logging
|
||||||
from typing import Optional, Tuple
|
from typing import Optional, Tuple
|
||||||
@@ -295,7 +295,11 @@ def reset_password(db: Session, token: str, new_password: str) -> Tuple[bool, st
|
|||||||
return False, "Invalid password reset token"
|
return False, "Invalid password reset token"
|
||||||
|
|
||||||
# Check if the token has expired
|
# Check if the token has expired
|
||||||
if user.password_reset_expiry < datetime.utcnow():
|
now = datetime.now(timezone.utc)
|
||||||
|
expiry = user.password_reset_expiry
|
||||||
|
if expiry is not None and expiry.tzinfo is None:
|
||||||
|
expiry = expiry.replace(tzinfo=timezone.utc)
|
||||||
|
if expiry is None or expiry < now:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
f"Attempt to reset password with expired token for user: {user.email}"
|
f"Attempt to reset password with expired token for user: {user.email}"
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -304,7 +304,6 @@ class WorkflowAgent(BaseAgent):
|
|||||||
"session_id": session_id,
|
"session_id": session_id,
|
||||||
}
|
}
|
||||||
|
|
||||||
# Função para message-node
|
|
||||||
async def message_node_function(
|
async def message_node_function(
|
||||||
state: State, node_id: str, node_data: Dict[str, Any]
|
state: State, node_id: str, node_data: Dict[str, Any]
|
||||||
) -> AsyncGenerator[State, None]:
|
) -> AsyncGenerator[State, None]:
|
||||||
@@ -318,7 +317,6 @@ class WorkflowAgent(BaseAgent):
|
|||||||
session_id = state.get("session_id", "")
|
session_id = state.get("session_id", "")
|
||||||
conversation_history = state.get("conversation_history", [])
|
conversation_history = state.get("conversation_history", [])
|
||||||
|
|
||||||
# Adiciona a mensagem como um novo Event do tipo agent
|
|
||||||
new_event = Event(
|
new_event = Event(
|
||||||
author="agent",
|
author="agent",
|
||||||
content=Content(parts=[Part(text=message_content)]),
|
content=Content(parts=[Part(text=message_content)]),
|
||||||
@@ -750,7 +748,7 @@ class WorkflowAgent(BaseAgent):
|
|||||||
content=Content(parts=[Part(text=user_message)]),
|
content=Content(parts=[Part(text=user_message)]),
|
||||||
)
|
)
|
||||||
|
|
||||||
# Se o histórico estiver vazio, adiciona a mensagem do usuário
|
# If the conversation history is empty, add the user message
|
||||||
conversation_history = ctx.session.events or []
|
conversation_history = ctx.session.events or []
|
||||||
if not conversation_history or (len(conversation_history) == 0):
|
if not conversation_history or (len(conversation_history) == 0):
|
||||||
conversation_history = [user_event]
|
conversation_history = [user_event]
|
||||||
@@ -768,16 +766,17 @@ class WorkflowAgent(BaseAgent):
|
|||||||
print("\n🚀 Starting workflow execution:")
|
print("\n🚀 Starting workflow execution:")
|
||||||
print(f"Initial content: {user_message[:100]}...")
|
print(f"Initial content: {user_message[:100]}...")
|
||||||
|
|
||||||
# Execute the graph with a recursion limit to avoid infinite loops
|
sent_events = 0 # Count of events already sent
|
||||||
result = await graph.ainvoke(initial_state, {"recursion_limit": 20})
|
|
||||||
|
|
||||||
# 6. Process and return the result
|
async for state in graph.astream(initial_state, {"recursion_limit": 20}):
|
||||||
final_content = result.get("content", [])
|
# The state can be a dict with the node name as a key
|
||||||
print(f"\n✅ FINAL RESULT: {final_content[:100]}...")
|
for node_state in state.values():
|
||||||
|
content = node_state.get("content", [])
|
||||||
for content in final_content:
|
# Only send new events
|
||||||
if content.author != "user":
|
for event in content[sent_events:]:
|
||||||
yield content
|
if event.author != "user":
|
||||||
|
yield event
|
||||||
|
sent_events = len(content)
|
||||||
|
|
||||||
# Execute sub-agents
|
# Execute sub-agents
|
||||||
for sub_agent in self.sub_agents:
|
for sub_agent in self.sub_agents:
|
||||||
|
|||||||
@@ -1,83 +1,103 @@
|
|||||||
<!DOCTYPE html>
|
<!DOCTYPE html>
|
||||||
<html>
|
<html>
|
||||||
<head>
|
<head>
|
||||||
<meta charset="UTF-8">
|
<meta charset="UTF-8" />
|
||||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
<meta name="viewport" content="width=device-width, initial-scale=1.0" />
|
||||||
<title>{% block title %}Evo AI{% endblock %}</title>
|
<title>{% block title %}Evo AI{% endblock %}</title>
|
||||||
<style>
|
<style>
|
||||||
body {
|
body {
|
||||||
font-family: Arial, sans-serif;
|
font-family: "Segoe UI", Arial, sans-serif;
|
||||||
line-height: 1.6;
|
background-color: #f7f7f7;
|
||||||
margin: 0;
|
color: #222;
|
||||||
padding: 0;
|
margin: 0;
|
||||||
background-color: #f7f7f7;
|
padding: 0;
|
||||||
|
}
|
||||||
|
.container {
|
||||||
|
max-width: 480px;
|
||||||
|
margin: 32px auto;
|
||||||
|
background: #fff;
|
||||||
|
border-radius: 12px;
|
||||||
|
box-shadow: 0 4px 24px rgba(0, 0, 0, 0.08);
|
||||||
|
padding: 0;
|
||||||
|
overflow: hidden;
|
||||||
|
}
|
||||||
|
.header {
|
||||||
|
background: #f7f7f7;
|
||||||
|
border-bottom: 1px solid #e5e5e5;
|
||||||
|
padding: 32px 0 16px 0;
|
||||||
|
text-align: center;
|
||||||
|
}
|
||||||
|
.header h1 {
|
||||||
|
color: #155a2c;
|
||||||
|
font-size: 2rem;
|
||||||
|
margin: 0;
|
||||||
|
letter-spacing: 1px;
|
||||||
|
}
|
||||||
|
.content {
|
||||||
|
padding: 32px 24px 24px 24px;
|
||||||
|
}
|
||||||
|
.button {
|
||||||
|
background: linear-gradient(90deg, #155a2c 0%, #1f7a3dff 100%);
|
||||||
|
color: #fff !important;
|
||||||
|
padding: 14px 0;
|
||||||
|
border-radius: 6px;
|
||||||
|
text-decoration: none;
|
||||||
|
display: block;
|
||||||
|
font-weight: bold;
|
||||||
|
text-align: center;
|
||||||
|
font-size: 1.1rem;
|
||||||
|
margin: 32px 0 0 0;
|
||||||
|
transition: filter 0.2s;
|
||||||
|
box-shadow: 0 2px 8px rgba(37, 99, 235, 0.08);
|
||||||
|
}
|
||||||
|
.button:hover {
|
||||||
|
filter: brightness(1.08);
|
||||||
|
}
|
||||||
|
.footer {
|
||||||
|
font-size: 12px;
|
||||||
|
text-align: center;
|
||||||
|
color: #888;
|
||||||
|
background: #f7f7f7;
|
||||||
|
border-top: 1px solid #e5e5e5;
|
||||||
|
padding: 20px 0 10px 0;
|
||||||
|
}
|
||||||
|
.link {
|
||||||
|
color: #155a2c;
|
||||||
|
text-decoration: underline;
|
||||||
|
word-break: break-all;
|
||||||
|
}
|
||||||
|
.warning {
|
||||||
|
color: #b91c1c;
|
||||||
|
background: #fee2e2;
|
||||||
|
border-radius: 4px;
|
||||||
|
padding: 12px;
|
||||||
|
margin-top: 20px;
|
||||||
|
}
|
||||||
|
@media (max-width: 600px) {
|
||||||
|
.container {
|
||||||
|
max-width: 98vw;
|
||||||
|
margin: 8px;
|
||||||
}
|
}
|
||||||
.container {
|
.content {
|
||||||
max-width: 600px;
|
padding: 18px 8px 16px 8px;
|
||||||
margin: 0 auto;
|
|
||||||
padding: 20px;
|
|
||||||
background-color: #ffffff;
|
|
||||||
border-radius: 8px;
|
|
||||||
box-shadow: 0 2px 4px rgba(0,0,0,0.1);
|
|
||||||
}
|
|
||||||
.header {
|
|
||||||
background-color: #4A90E2;
|
|
||||||
color: white;
|
|
||||||
padding: 15px;
|
|
||||||
text-align: center;
|
|
||||||
border-radius: 6px 6px 0 0;
|
|
||||||
}
|
|
||||||
.content {
|
|
||||||
padding: 20px;
|
|
||||||
}
|
|
||||||
.button {
|
|
||||||
background-color: #4A90E2;
|
|
||||||
color: white !important;
|
|
||||||
padding: 12px 24px;
|
|
||||||
text-decoration: none;
|
|
||||||
border-radius: 4px;
|
|
||||||
display: inline-block;
|
|
||||||
font-weight: bold;
|
|
||||||
text-align: center;
|
|
||||||
transition: background-color 0.3s;
|
|
||||||
}
|
|
||||||
.button:hover {
|
|
||||||
background-color: #3a7bc8;
|
|
||||||
}
|
|
||||||
.footer {
|
|
||||||
font-size: 12px;
|
|
||||||
text-align: center;
|
|
||||||
margin-top: 30px;
|
|
||||||
color: #888;
|
|
||||||
border-top: 1px solid #eee;
|
|
||||||
padding-top: 20px;
|
|
||||||
}
|
|
||||||
.link {
|
|
||||||
word-break: break-all;
|
|
||||||
color: #4A90E2;
|
|
||||||
}
|
|
||||||
.warning {
|
|
||||||
color: #E74C3C;
|
|
||||||
padding: 10px;
|
|
||||||
background-color: #FADBD8;
|
|
||||||
border-radius: 4px;
|
|
||||||
margin-top: 20px;
|
|
||||||
}
|
}
|
||||||
|
}
|
||||||
</style>
|
</style>
|
||||||
{% block additional_styles %}{% endblock %}
|
{% block additional_styles %}{% endblock %}
|
||||||
</head>
|
</head>
|
||||||
<body>
|
<body>
|
||||||
<div class="container">
|
<div class="container">
|
||||||
<div class="header">
|
<div class="header">
|
||||||
<h1>{% block header %}Evo AI{% endblock %}</h1>
|
<h1>{% block header %}Evo AI{% endblock %}</h1>
|
||||||
</div>
|
</div>
|
||||||
<div class="content">
|
<div class="content">{% block content %}{% endblock %}</div>
|
||||||
{% block content %}{% endblock %}
|
<div class="footer">
|
||||||
</div>
|
<p>
|
||||||
<div class="footer">
|
{% block footer_message %}This is an automated email, please do not
|
||||||
<p>{% block footer_message %}This is an automated email, please do not reply.{% endblock %}</p>
|
reply.{% endblock %}
|
||||||
<p>© {{ current_year }} Evo AI. All rights reserved.</p>
|
</p>
|
||||||
</div>
|
<p>© {{ current_year }} Evo AI. All rights reserved.</p>
|
||||||
|
</div>
|
||||||
</div>
|
</div>
|
||||||
</body>
|
</body>
|
||||||
</html>
|
</html>
|
||||||
|
|||||||
41
src/utils/otel.py
Normal file
41
src/utils/otel.py
Normal file
@@ -0,0 +1,41 @@
|
|||||||
|
import os
|
||||||
|
import base64
|
||||||
|
from src.config.settings import settings
|
||||||
|
|
||||||
|
from opentelemetry import trace
|
||||||
|
from opentelemetry.sdk.resources import Resource
|
||||||
|
from opentelemetry.sdk.trace import TracerProvider
|
||||||
|
from opentelemetry.sdk.trace.export import BatchSpanProcessor
|
||||||
|
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
|
||||||
|
|
||||||
|
_otlp_initialized = False
|
||||||
|
|
||||||
|
|
||||||
|
def init_otel():
|
||||||
|
global _otlp_initialized
|
||||||
|
if _otlp_initialized:
|
||||||
|
return
|
||||||
|
if not (
|
||||||
|
settings.LANGFUSE_PUBLIC_KEY
|
||||||
|
and settings.LANGFUSE_SECRET_KEY
|
||||||
|
and settings.OTEL_EXPORTER_OTLP_ENDPOINT
|
||||||
|
):
|
||||||
|
return
|
||||||
|
|
||||||
|
langfuse_auth = base64.b64encode(
|
||||||
|
f"{settings.LANGFUSE_PUBLIC_KEY}:{settings.LANGFUSE_SECRET_KEY}".encode()
|
||||||
|
).decode()
|
||||||
|
os.environ["OTEL_EXPORTER_OTLP_ENDPOINT"] = settings.OTEL_EXPORTER_OTLP_ENDPOINT
|
||||||
|
os.environ["OTEL_EXPORTER_OTLP_HEADERS"] = f"Authorization=Basic {langfuse_auth}"
|
||||||
|
|
||||||
|
provider = TracerProvider(
|
||||||
|
resource=Resource.create({"service.name": "evo_ai_agent"})
|
||||||
|
)
|
||||||
|
exporter = OTLPSpanExporter()
|
||||||
|
provider.add_span_processor(BatchSpanProcessor(exporter))
|
||||||
|
trace.set_tracer_provider(provider)
|
||||||
|
_otlp_initialized = True
|
||||||
|
|
||||||
|
|
||||||
|
def get_tracer(name: str = "evo_ai_agent"):
|
||||||
|
return trace.get_tracer(name)
|
||||||
Reference in New Issue
Block a user