From 12c91b7c2502a1f64daf9dcb58771c37bc4cf200 Mon Sep 17 00:00:00 2001 From: A8072846 Date: Tue, 3 Mar 2026 22:24:31 +0000 Subject: [PATCH 1/5] Add notification service using Google ADK --- config.yaml | 4 + src/va_agent/agent.py | 36 ++++- src/va_agent/config.py | 6 + src/va_agent/dynamic_instruction.py | 120 ++++++++++++++ src/va_agent/notifications.py | 232 ++++++++++++++++++++++++++++ 5 files changed, 392 insertions(+), 6 deletions(-) create mode 100644 src/va_agent/dynamic_instruction.py create mode 100644 src/va_agent/notifications.py diff --git a/config.yaml b/config.yaml index f8be685..2af4060 100644 --- a/config.yaml +++ b/config.yaml @@ -3,6 +3,10 @@ google_cloud_location: us-central1 firestore_db: bnt-orquestador-cognitivo-firestore-bdo-dev +# Notifications configuration +notifications_collection_path: "artifacts/bnt-orquestador-cognitivo-dev/notifications" +notifications_max_to_notify: 5 + mcp_remote_url: "https://ap01194-orq-cog-rag-connector-1007577023101.us-central1.run.app/mcp" # audience sin la ruta, para emitir el ID Token: mcp_audience: "https://ap01194-orq-cog-rag-connector-1007577023101.us-central1.run.app" diff --git a/src/va_agent/agent.py b/src/va_agent/agent.py index cbccff9..72085d6 100644 --- a/src/va_agent/agent.py +++ b/src/va_agent/agent.py @@ -1,37 +1,61 @@ """ADK agent with vector search RAG tool.""" +from functools import partial + from google import genai from google.adk.agents.llm_agent import Agent from google.adk.runners import Runner from google.adk.tools.mcp_tool import McpToolset from google.adk.tools.mcp_tool.mcp_session_manager import StreamableHTTPConnectionParams from google.cloud.firestore_v1.async_client import AsyncClient +from google.genai.types import Content, Part from va_agent.auth import auth_headers_provider from va_agent.config import settings +from va_agent.dynamic_instruction import provide_dynamic_instruction +from va_agent.notifications import NotificationService from va_agent.session import FirestoreSessionService from va_agent.governance import GovernancePlugin +# MCP Toolset for RAG knowledge search toolset = McpToolset( connection_params=StreamableHTTPConnectionParams(url=settings.mcp_remote_url), header_provider=auth_headers_provider, ) + +# Shared Firestore client for session service and notifications +firestore_db = AsyncClient(database=settings.firestore_db) + +# Session service with compaction +session_service = FirestoreSessionService( + db=firestore_db, + compaction_token_threshold=10_000, + genai_client=genai.Client(), +) + +# Notification service +notification_service = NotificationService( + db=firestore_db, + collection_path=settings.notifications_collection_path, + max_to_notify=settings.notifications_max_to_notify, +) + +# Agent with static and dynamic instructions governance = GovernancePlugin() agent = Agent( model=settings.agent_model, name=settings.agent_name, + static_instruction=Content( + role="user", + parts=[Part(text=settings.agent_instructions)], + ), instruction=settings.agent_instructions, tools=[toolset], after_model_callback=governance.after_model_callback, ) -session_service = FirestoreSessionService( - db=AsyncClient(database=settings.firestore_db), - compaction_token_threshold=10_000, - genai_client=genai.Client(), -) - +# Runner runner = Runner( app_name="va_agent", agent=agent, diff --git a/src/va_agent/config.py b/src/va_agent/config.py index 5e112b1..f3502b5 100644 --- a/src/va_agent/config.py +++ b/src/va_agent/config.py @@ -26,6 +26,12 @@ class AgentSettings(BaseSettings): # Firestore configuration firestore_db: str + # Notifications configuration + notifications_collection_path: str = ( + "artifacts/bnt-orquestador-cognitivo-dev/notifications" + ) + notifications_max_to_notify: int = 5 + # MCP configuration mcp_audience: str mcp_remote_url: str diff --git a/src/va_agent/dynamic_instruction.py b/src/va_agent/dynamic_instruction.py new file mode 100644 index 0000000..86f190d --- /dev/null +++ b/src/va_agent/dynamic_instruction.py @@ -0,0 +1,120 @@ +"""Dynamic instruction provider for VAia agent.""" + +from __future__ import annotations + +import logging +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from google.adk.agents.readonly_context import ReadonlyContext + + from va_agent.notifications import NotificationService + +logger = logging.getLogger(__name__) + + +async def provide_dynamic_instruction( + notification_service: NotificationService, + ctx: ReadonlyContext | None = None, +) -> str: + """Provide dynamic instructions based on pending notifications. + + This function is called by the ADK agent on each message. It: + 1. Checks if this is the first message in the session (< 2 events) + 2. Queries Firestore for pending notifications + 3. Marks them as notified + 4. Returns a dynamic instruction for the agent to mention them + + Args: + notification_service: Service for fetching/marking notifications + ctx: Agent context containing session information + + Returns: + Dynamic instruction string (empty if no notifications or not first message) + + """ + # Only check notifications on the first message + if not ctx or not ctx._invocation_context: + logger.debug("No context available for dynamic instruction") + return "" + + session = ctx._invocation_context.session + if not session: + logger.debug("No session available for dynamic instruction") + return "" + + # FOR TESTING: Always check for notifications (comment out to enable first-message-only) + # Only check on first message (when events list is empty or has only 1-2 events) + # Events include both user and agent messages, so < 2 means first interaction + # event_count = len(session.events) if session.events else 0 + # + # if event_count >= 2: + # logger.debug( + # "Skipping notification check: not first message (event_count=%d)", + # event_count, + # ) + # return "" + + # Extract phone number from user_id (they are the same in this implementation) + phone_number = session.user_id + + logger.info( + "First message detected for user %s, checking for pending notifications", + phone_number, + ) + + try: + # Fetch pending notifications + pending_notifications = await notification_service.get_pending_notifications( + phone_number + ) + + if not pending_notifications: + logger.info("No pending notifications for user %s", phone_number) + return "" + + # Build dynamic instruction with notification details + notification_ids = [n.get("id_notificacion") for n in pending_notifications] + count = len(pending_notifications) + + # Format notification details for the agent + notification_details = [] + for notif in pending_notifications: + evento = notif.get("nombre_evento_dialogflow", "notificacion") + texto = notif.get("texto", "Sin texto") + notification_details.append(f" - Evento: {evento} | Texto: {texto}") + + details_text = "\n".join(notification_details) + + instruction = f""" +IMPORTANTE - NOTIFICACIONES PENDIENTES: + +El usuario tiene {count} notificación(es) sin leer: + +{details_text} + +INSTRUCCIONES: +- Menciona estas notificaciones de forma natural en tu respuesta inicial +- No necesitas leerlas todas literalmente, solo hazle saber que las tiene +- Sé breve y directo según tu personalidad (directo y cálido) +- Si el usuario pregunta algo específico, prioriza responder eso primero y luego menciona las notificaciones + +Ejemplo: "¡Hola! 👋 Antes de empezar, veo que tienes {count} notificación(es) pendiente(s) en tu cuenta. ¿Te gustaría revisarlas o prefieres que te ayude con algo más?" +""" + + # Mark notifications as notified in Firestore + await notification_service.mark_as_notified(phone_number, notification_ids) + + logger.info( + "Returning dynamic instruction with %d notification(s) for user %s", + count, + phone_number, + ) + + return instruction + + except Exception: + logger.exception( + "Error building dynamic instruction for user %s", phone_number + ) + return "" diff --git a/src/va_agent/notifications.py b/src/va_agent/notifications.py new file mode 100644 index 0000000..e7cc57a --- /dev/null +++ b/src/va_agent/notifications.py @@ -0,0 +1,232 @@ +"""Notification management for VAia agent.""" + +from __future__ import annotations + +import logging +import time +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from google.cloud.firestore_v1.async_client import AsyncClient + +logger = logging.getLogger(__name__) + + +class NotificationService: + """Service for fetching and managing user notifications from Firestore.""" + + def __init__( + self, + *, + db: AsyncClient, + collection_path: str, + max_to_notify: int = 5, + ) -> None: + """Initialize NotificationService. + + Args: + db: Firestore async client + collection_path: Path to notifications collection + max_to_notify: Maximum number of notifications to return + + """ + self._db = db + self._collection_path = collection_path + self._max_to_notify = max_to_notify + + async def get_pending_notifications( + self, phone_number: str + ) -> list[dict[str, Any]]: + """Get pending notifications for a user. + + Retrieves notifications that have not been notified by the agent yet, + ordered by timestamp (most recent first), limited to max_to_notify. + + Args: + phone_number: User's phone number (used as document ID) + + Returns: + List of notification dictionaries with structure: + { + "id_notificacion": str, + "texto": str, + "status": str, + "timestamp_creacion": timestamp, + "parametros": {...} + } + + """ + try: + # Query Firestore document by phone number + doc_ref = self._db.collection(self._collection_path).document( + phone_number + ) + doc = await doc_ref.get() + + if not doc.exists: + logger.info( + "No notification document found for phone: %s", phone_number + ) + return [] + + data = doc.to_dict() or {} + all_notifications = data.get("notificaciones", []) + + if not all_notifications: + logger.info("No notifications in array for phone: %s", phone_number) + return [] + + # Filter notifications that have NOT been notified by the agent + pending = [ + n + for n in all_notifications + if not n.get("notified_by_agent", False) + ] + + if not pending: + logger.info( + "All notifications already notified for phone: %s", phone_number + ) + return [] + + # Sort by timestamp_creacion (most recent first) + pending.sort( + key=lambda n: n.get("timestamp_creacion", 0), reverse=True + ) + + # Return top N most recent + result = pending[: self._max_to_notify] + + logger.info( + "Found %d pending notifications for phone: %s (returning top %d)", + len(pending), + phone_number, + len(result), + ) + + return result + + except Exception: + logger.exception( + "Failed to fetch notifications for phone: %s", phone_number + ) + return [] + + async def mark_as_notified( + self, phone_number: str, notification_ids: list[str] + ) -> bool: + """Mark notifications as notified by the agent. + + Updates the notifications in Firestore by adding: + - notified_by_agent: true + - notified_at: current timestamp + + Args: + phone_number: User's phone number (document ID) + notification_ids: List of id_notificacion values to mark + + Returns: + True if update was successful, False otherwise + + """ + if not notification_ids: + return True + + try: + doc_ref = self._db.collection(self._collection_path).document( + phone_number + ) + doc = await doc_ref.get() + + if not doc.exists: + logger.warning( + "Cannot mark notifications as notified: document not found for %s", + phone_number, + ) + return False + + data = doc.to_dict() or {} + notificaciones = data.get("notificaciones", []) + + if not notificaciones: + logger.warning( + "Cannot mark notifications: empty array for %s", phone_number + ) + return False + + # Update matching notifications + now = time.time() + updated_count = 0 + + for notif in notificaciones: + if notif.get("id_notificacion") in notification_ids: + notif["notified_by_agent"] = True + notif["notified_at"] = now + updated_count += 1 + + if updated_count == 0: + logger.warning( + "No notifications matched IDs for phone: %s", phone_number + ) + return False + + # Save back to Firestore + await doc_ref.update( + { + "notificaciones": notificaciones, + "ultima_actualizacion": now, + } + ) + + logger.info( + "Marked %d notification(s) as notified for phone: %s", + updated_count, + phone_number, + ) + + return True + + except Exception: + logger.exception( + "Failed to mark notifications as notified for phone: %s", + phone_number, + ) + return False + + def format_notification_summary( + self, notifications: list[dict[str, Any]] + ) -> str: + """Format notifications into a human-readable summary. + + Args: + notifications: List of notification dictionaries + + Returns: + Formatted string summarizing the notifications + + """ + if not notifications: + return "" + + count = len(notifications) + summary_lines = [ + f"El usuario tiene {count} notificación(es) pendiente(s):" + ] + + for i, notif in enumerate(notifications, 1): + texto = notif.get("texto", "Sin texto") + params = notif.get("parametros", {}) + + # Extract key parameters if available + amount = params.get("notification_po_amount") + tx_id = params.get("notification_po_transaction_id") + + line = f"{i}. {texto}" + if amount: + line += f" (monto: ${amount})" + if tx_id: + line += f" [ID: {tx_id}]" + + summary_lines.append(line) + + return "\n".join(summary_lines) From 5941c41296b5712c7212d6191173fd8f79966ce4 Mon Sep 17 00:00:00 2001 From: Anibal Angulo Date: Thu, 5 Mar 2026 05:55:09 +0000 Subject: [PATCH 2/5] Remove firestore emulator from test dependencies --- tests/conftest.py | 8 +- tests/fake_firestore.py | 284 ++++++++++++++++++++++++++++++++++++++++ 2 files changed, 287 insertions(+), 5 deletions(-) create mode 100644 tests/fake_firestore.py diff --git a/tests/conftest.py b/tests/conftest.py index 959b677..2c1c3ea 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -2,25 +2,23 @@ from __future__ import annotations -import os import uuid import pytest import pytest_asyncio -from google.cloud.firestore_v1.async_client import AsyncClient from va_agent.session import FirestoreSessionService -os.environ.setdefault("FIRESTORE_EMULATOR_HOST", "localhost:8602") +from .fake_firestore import FakeAsyncClient @pytest_asyncio.fixture async def db(): - return AsyncClient(project="test-project") + return FakeAsyncClient() @pytest_asyncio.fixture -async def service(db: AsyncClient): +async def service(db): prefix = f"test_{uuid.uuid4().hex[:8]}" return FirestoreSessionService(db=db, collection_prefix=prefix) diff --git a/tests/fake_firestore.py b/tests/fake_firestore.py new file mode 100644 index 0000000..f7fce9b --- /dev/null +++ b/tests/fake_firestore.py @@ -0,0 +1,284 @@ +"""In-memory fake of the Firestore async surface used by this project. + +Covers: AsyncClient, DocumentReference, CollectionReference, Query, +DocumentSnapshot, WriteBatch, and basic transaction support (enough for +``@async_transactional``). +""" + +from __future__ import annotations + +import copy +from typing import Any + + +# ------------------------------------------------------------------ # +# DocumentSnapshot +# ------------------------------------------------------------------ # + +class FakeDocumentSnapshot: + def __init__(self, *, exists: bool, data: dict[str, Any] | None, reference: FakeDocumentReference) -> None: + self._exists = exists + self._data = data + self._reference = reference + + @property + def exists(self) -> bool: + return self._exists + + @property + def reference(self) -> FakeDocumentReference: + return self._reference + + def to_dict(self) -> dict[str, Any] | None: + if not self._exists: + return None + return copy.deepcopy(self._data) + + +# ------------------------------------------------------------------ # +# DocumentReference +# ------------------------------------------------------------------ # + +class FakeDocumentReference: + def __init__(self, store: FakeStore, path: str) -> None: + self._store = store + self._path = path + + @property + def path(self) -> str: + return self._path + + # --- read --- + + async def get(self, *, transaction: FakeTransaction | None = None) -> FakeDocumentSnapshot: + data = self._store.get_doc(self._path) + if data is None: + return FakeDocumentSnapshot(exists=False, data=None, reference=self) + return FakeDocumentSnapshot(exists=True, data=copy.deepcopy(data), reference=self) + + # --- write --- + + async def set(self, document_data: dict[str, Any], merge: bool = False) -> None: + if merge: + existing = self._store.get_doc(self._path) or {} + existing.update(document_data) + self._store.set_doc(self._path, existing) + else: + self._store.set_doc(self._path, copy.deepcopy(document_data)) + + async def update(self, field_updates: dict[str, Any]) -> None: + data = self._store.get_doc(self._path) + if data is None: + msg = f"Document {self._path} does not exist" + raise ValueError(msg) + for key, value in field_updates.items(): + _nested_set(data, key, value) + self._store.set_doc(self._path, data) + + # --- subcollection --- + + def collection(self, subcollection_name: str) -> FakeCollectionReference: + return FakeCollectionReference(self._store, f"{self._path}/{subcollection_name}") + + +# ------------------------------------------------------------------ # +# Helpers for nested field-path updates ("state.counter" → data["state"]["counter"]) +# ------------------------------------------------------------------ # + +def _nested_set(data: dict[str, Any], dotted_key: str, value: Any) -> None: + parts = dotted_key.split(".") + for part in parts[:-1]: + # Backtick-quoted segments (Firestore FieldPath encoding) + part = part.strip("`") + data = data.setdefault(part, {}) + final = parts[-1].strip("`") + data[final] = value + + +# ------------------------------------------------------------------ # +# Query +# ------------------------------------------------------------------ # + +class FakeQuery: + """Supports chained .where() / .order_by() / .get().""" + + def __init__(self, store: FakeStore, collection_path: str) -> None: + self._store = store + self._collection_path = collection_path + self._filters: list[tuple[str, str, Any]] = [] + self._order_by_field: str | None = None + + def where(self, *, filter: Any) -> FakeQuery: # noqa: A002 + clone = FakeQuery(self._store, self._collection_path) + clone._filters = [*self._filters, (filter.field_path, filter.op_string, filter.value)] + clone._order_by_field = self._order_by_field + return clone + + def order_by(self, field_path: str) -> FakeQuery: + clone = FakeQuery(self._store, self._collection_path) + clone._filters = list(self._filters) + clone._order_by_field = field_path + return clone + + async def get(self) -> list[FakeDocumentSnapshot]: + docs = self._store.list_collection(self._collection_path) + results: list[tuple[str, dict[str, Any]]] = [] + + for doc_path, data in docs: + if all(_match(data, field, op, val) for field, op, val in self._filters): + results.append((doc_path, data)) + + if self._order_by_field: + field = self._order_by_field + results.sort(key=lambda item: item[1].get(field, 0)) + + return [ + FakeDocumentSnapshot( + exists=True, + data=copy.deepcopy(data), + reference=FakeDocumentReference(self._store, path), + ) + for path, data in results + ] + + +def _match(data: dict[str, Any], field: str, op: str, value: Any) -> bool: + doc_val = data.get(field) + if op == "==": + return doc_val == value + if op == ">=": + return doc_val is not None and doc_val >= value + return False + + +# ------------------------------------------------------------------ # +# CollectionReference (extends Query behaviour) +# ------------------------------------------------------------------ # + +class FakeCollectionReference(FakeQuery): + def document(self, document_id: str) -> FakeDocumentReference: + return FakeDocumentReference(self._store, f"{self._collection_path}/{document_id}") + + +# ------------------------------------------------------------------ # +# WriteBatch +# ------------------------------------------------------------------ # + +class FakeWriteBatch: + def __init__(self, store: FakeStore) -> None: + self._store = store + self._deletes: list[str] = [] + + def delete(self, doc_ref: FakeDocumentReference) -> None: + self._deletes.append(doc_ref.path) + + async def commit(self) -> None: + for path in self._deletes: + self._store.delete_doc(path) + + +# ------------------------------------------------------------------ # +# Transaction (minimal, supports @async_transactional) +# ------------------------------------------------------------------ # + +class FakeTransaction: + """Minimal transaction compatible with ``@async_transactional``. + + The decorator calls ``_clean_up()``, ``_begin()``, the wrapped function, + then ``_commit()``. On error it calls ``_rollback()``. + ``in_progress`` is a property that checks ``_id is not None``. + """ + + def __init__(self, store: FakeStore) -> None: + self._store = store + self._staged_updates: list[tuple[str, dict[str, Any]]] = [] + self._id: bytes | None = None + self._max_attempts = 1 + self._read_only = False + + @property + def in_progress(self) -> bool: + return self._id is not None + + def _clean_up(self) -> None: + self._id = None + + async def _begin(self, retry_id: bytes | None = None) -> None: + self._id = b"fake-txn" + + async def _commit(self) -> list: + for path, updates in self._staged_updates: + data = self._store.get_doc(path) + if data is not None: + for key, value in updates.items(): + _nested_set(data, key, value) + self._store.set_doc(path, data) + self._staged_updates.clear() + self._clean_up() + return [] + + async def _rollback(self) -> None: + self._staged_updates.clear() + self._clean_up() + + def update(self, doc_ref: FakeDocumentReference, field_updates: dict[str, Any]) -> None: + self._staged_updates.append((doc_ref.path, field_updates)) + + +# ------------------------------------------------------------------ # +# Document store (flat dict keyed by path) +# ------------------------------------------------------------------ # + +class FakeStore: + def __init__(self) -> None: + self._docs: dict[str, dict[str, Any]] = {} + + def get_doc(self, path: str) -> dict[str, Any] | None: + data = self._docs.get(path) + return data # returns reference, callers deepcopy where needed + + def set_doc(self, path: str, data: dict[str, Any]) -> None: + self._docs[path] = data + + def delete_doc(self, path: str) -> None: + self._docs.pop(path, None) + + def list_collection(self, collection_path: str) -> list[tuple[str, dict[str, Any]]]: + """Return (path, data) for every direct child doc of *collection_path*.""" + prefix = collection_path + "/" + results: list[tuple[str, dict[str, Any]]] = [] + for doc_path, data in self._docs.items(): + if not doc_path.startswith(prefix): + continue + # Must be a direct child (no further '/' after the prefix, except maybe subcollection paths) + remainder = doc_path[len(prefix):] + if "/" not in remainder: + results.append((doc_path, data)) + return results + + def recursive_delete(self, path: str) -> None: + """Delete a document and everything nested under it.""" + to_delete = [p for p in self._docs if p == path or p.startswith(path + "/")] + for p in to_delete: + del self._docs[p] + + +# ------------------------------------------------------------------ # +# FakeAsyncClient (drop-in for AsyncClient) +# ------------------------------------------------------------------ # + +class FakeAsyncClient: + def __init__(self, **_kwargs: Any) -> None: + self._store = FakeStore() + + def collection(self, collection_path: str) -> FakeCollectionReference: + return FakeCollectionReference(self._store, collection_path) + + def batch(self) -> FakeWriteBatch: + return FakeWriteBatch(self._store) + + def transaction(self, **kwargs: Any) -> FakeTransaction: + return FakeTransaction(self._store) + + async def recursive_delete(self, doc_ref: FakeDocumentReference) -> None: + self._store.recursive_delete(doc_ref.path) From db879cee9f4dc19cb66fb2a930eadfd8725d806e Mon Sep 17 00:00:00 2001 From: Anibal Angulo Date: Thu, 5 Mar 2026 06:06:01 +0000 Subject: [PATCH 3/5] Format/Lint --- AGENTS.md | 3 +- src/va_agent/agent.py | 5 +- src/va_agent/config.py | 2 +- src/va_agent/dynamic_instruction.py | 29 ++++++---- src/va_agent/governance.py | 88 ++++++++++++++++++++++------- src/va_agent/notifications.py | 32 ++++------- 6 files changed, 102 insertions(+), 57 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 3434537..a4792cf 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,3 +1,4 @@ Use `uv` for project management. -Use `uv run ruff check` for linting, and `uv run ty check` for type checking +Use `uv run ruff check` for linting +Use `uv run ty check` for type checking Use `uv run pytest` for testing. diff --git a/src/va_agent/agent.py b/src/va_agent/agent.py index 72085d6..11b458b 100644 --- a/src/va_agent/agent.py +++ b/src/va_agent/agent.py @@ -1,7 +1,5 @@ """ADK agent with vector search RAG tool.""" -from functools import partial - from google import genai from google.adk.agents.llm_agent import Agent from google.adk.runners import Runner @@ -12,10 +10,9 @@ from google.genai.types import Content, Part from va_agent.auth import auth_headers_provider from va_agent.config import settings -from va_agent.dynamic_instruction import provide_dynamic_instruction +from va_agent.governance import GovernancePlugin from va_agent.notifications import NotificationService from va_agent.session import FirestoreSessionService -from va_agent.governance import GovernancePlugin # MCP Toolset for RAG knowledge search toolset = McpToolset( diff --git a/src/va_agent/config.py b/src/va_agent/config.py index f3502b5..ae33d2d 100644 --- a/src/va_agent/config.py +++ b/src/va_agent/config.py @@ -39,7 +39,7 @@ class AgentSettings(BaseSettings): model_config = SettingsConfigDict( yaml_file=CONFIG_FILE_PATH, extra="ignore", # Ignore extra fields from config.yaml - env_file=".env" + env_file=".env", ) @classmethod diff --git a/src/va_agent/dynamic_instruction.py b/src/va_agent/dynamic_instruction.py index 86f190d..cc64fb9 100644 --- a/src/va_agent/dynamic_instruction.py +++ b/src/va_agent/dynamic_instruction.py @@ -34,17 +34,19 @@ async def provide_dynamic_instruction( """ # Only check notifications on the first message - if not ctx or not ctx._invocation_context: + if not ctx: logger.debug("No context available for dynamic instruction") return "" - session = ctx._invocation_context.session + session = ctx.session if not session: logger.debug("No session available for dynamic instruction") return "" - # FOR TESTING: Always check for notifications (comment out to enable first-message-only) - # Only check on first message (when events list is empty or has only 1-2 events) + # FOR TESTING: Always check for notifications + # (comment out to enable first-message-only) + # Only check on first message (when events list is empty + # or has only 1-2 events) # Events include both user and agent messages, so < 2 means first interaction # event_count = len(session.events) if session.events else 0 # @@ -74,7 +76,11 @@ async def provide_dynamic_instruction( return "" # Build dynamic instruction with notification details - notification_ids = [n.get("id_notificacion") for n in pending_notifications] + notification_ids = [ + nid + for n in pending_notifications + if (nid := n.get("id_notificacion")) is not None + ] count = len(pending_notifications) # Format notification details for the agent @@ -97,9 +103,11 @@ INSTRUCCIONES: - Menciona estas notificaciones de forma natural en tu respuesta inicial - No necesitas leerlas todas literalmente, solo hazle saber que las tiene - Sé breve y directo según tu personalidad (directo y cálido) -- Si el usuario pregunta algo específico, prioriza responder eso primero y luego menciona las notificaciones +- Si el usuario pregunta algo específico, prioriza responder eso primero\ + y luego menciona las notificaciones -Ejemplo: "¡Hola! 👋 Antes de empezar, veo que tienes {count} notificación(es) pendiente(s) en tu cuenta. ¿Te gustaría revisarlas o prefieres que te ayude con algo más?" +Ejemplo: "¡Hola! 👋 Tienes {count} notificación(es)\ + pendiente(s). ¿Te gustaría revisarlas?" """ # Mark notifications as notified in Firestore @@ -111,10 +119,11 @@ Ejemplo: "¡Hola! 👋 Antes de empezar, veo que tienes {count} notificación(es phone_number, ) - return instruction - except Exception: logger.exception( - "Error building dynamic instruction for user %s", phone_number + "Error building dynamic instruction for user %s", + phone_number, ) return "" + else: + return instruction diff --git a/src/va_agent/governance.py b/src/va_agent/governance.py index 936c668..a65d5a3 100644 --- a/src/va_agent/governance.py +++ b/src/va_agent/governance.py @@ -1,4 +1,5 @@ """GovernancePlugin: Guardrails for VAia, the virtual assistant for VA.""" + import logging import re @@ -9,10 +10,57 @@ logger = logging.getLogger(__name__) FORBIDDEN_EMOJIS = [ - "🥵","🔪","🎰","🎲","🃏","😤","🤬","😡","😠","🩸","🧨","🪓","☠️","💀", - "💣","🔫","👗","💦","🍑","🍆","👄","👅","🫦","💩","⚖️","⚔️","✝️","🕍", - "🕌","⛪","🍻","🍸","🥃","🍷","🍺","🚬","👹","👺","👿","😈","🤡","🧙", - "🧙‍♀️", "🧙‍♂️", "🧛", "🧛‍♀️", "🧛‍♂️", "🔞","🧿","💊", "💏" + "🥵", + "🔪", + "🎰", + "🎲", + "🃏", + "😤", + "🤬", + "😡", + "😠", + "🩸", + "🧨", + "🪓", + "☠️", + "💀", + "💣", + "🔫", + "👗", + "💦", + "🍑", + "🍆", + "👄", + "👅", + "🫦", + "💩", + "⚖️", + "⚔️", + "✝️", + "🕍", + "🕌", + "⛪", + "🍻", + "🍸", + "🥃", + "🍷", + "🍺", + "🚬", + "👹", + "👺", + "👿", + "😈", + "🤡", + "🧙", + "🧙‍♀️", + "🧙‍♂️", + "🧛", + "🧛‍♀️", + "🧛‍♂️", + "🔞", + "🧿", + "💊", + "💏", ] @@ -20,29 +68,31 @@ class GovernancePlugin: """Guardrail executor for VAia requests as a Agent engine callbacks.""" def __init__(self) -> None: - """Initialize guardrail model (structured output), prompt and emojis patterns.""" + """Initialize guardrail model, prompt and emojis patterns.""" self._combined_pattern = self._get_combined_pattern() - def _get_combined_pattern(self): - person_pattern = r"(?:🧑|👩|👨)" - tone_pattern = r"[\U0001F3FB-\U0001F3FF]?" - - # Unique pattern that combines all forbidden emojis, including complex ones with skin tones - combined_pattern = re.compile( - rf"{person_pattern}{tone_pattern}\u200d❤️?\u200d💋\u200d{person_pattern}{tone_pattern}" # kiss - rf"|{person_pattern}{tone_pattern}\u200d❤️?\u200d{person_pattern}{tone_pattern}" # lovers - rf"|🖕{tone_pattern}" # middle finger with all skin tone variations - rf"|{'|'.join(map(re.escape, sorted(FORBIDDEN_EMOJIS, key=len, reverse=True)))}" # simple emojis - rf"|\u200d|\uFE0F" # residual ZWJ and variation selectors + def _get_combined_pattern(self) -> re.Pattern[str]: + person = r"(?:🧑|👩|👨)" + tone = r"[\U0001F3FB-\U0001F3FF]?" + simple = "|".join( + map(re.escape, sorted(FORBIDDEN_EMOJIS, key=len, reverse=True)) ) - return combined_pattern - + + # Combines all forbidden emojis, including complex + # ones with skin tones + return re.compile( + rf"{person}{tone}\u200d❤️?\u200d💋\u200d{person}{tone}" + rf"|{person}{tone}\u200d❤️?\u200d{person}{tone}" + rf"|🖕{tone}" + rf"|{simple}" + rf"|\u200d|\uFE0F" + ) + def _remove_emojis(self, text: str) -> tuple[str, list[str]]: removed = self._combined_pattern.findall(text) text = self._combined_pattern.sub("", text) return text.strip(), removed - def after_model_callback( self, callback_context: CallbackContext | None = None, diff --git a/src/va_agent/notifications.py b/src/va_agent/notifications.py index e7cc57a..8536fb2 100644 --- a/src/va_agent/notifications.py +++ b/src/va_agent/notifications.py @@ -58,9 +58,7 @@ class NotificationService: """ try: # Query Firestore document by phone number - doc_ref = self._db.collection(self._collection_path).document( - phone_number - ) + doc_ref = self._db.collection(self._collection_path).document(phone_number) doc = await doc_ref.get() if not doc.exists: @@ -78,9 +76,7 @@ class NotificationService: # Filter notifications that have NOT been notified by the agent pending = [ - n - for n in all_notifications - if not n.get("notified_by_agent", False) + n for n in all_notifications if not n.get("notified_by_agent", False) ] if not pending: @@ -90,9 +86,7 @@ class NotificationService: return [] # Sort by timestamp_creacion (most recent first) - pending.sort( - key=lambda n: n.get("timestamp_creacion", 0), reverse=True - ) + pending.sort(key=lambda n: n.get("timestamp_creacion", 0), reverse=True) # Return top N most recent result = pending[: self._max_to_notify] @@ -104,13 +98,13 @@ class NotificationService: len(result), ) - return result - except Exception: logger.exception( "Failed to fetch notifications for phone: %s", phone_number ) return [] + else: + return result async def mark_as_notified( self, phone_number: str, notification_ids: list[str] @@ -133,9 +127,7 @@ class NotificationService: return True try: - doc_ref = self._db.collection(self._collection_path).document( - phone_number - ) + doc_ref = self._db.collection(self._collection_path).document(phone_number) doc = await doc_ref.get() if not doc.exists: @@ -184,18 +176,16 @@ class NotificationService: phone_number, ) - return True - except Exception: logger.exception( "Failed to mark notifications as notified for phone: %s", phone_number, ) return False + else: + return True - def format_notification_summary( - self, notifications: list[dict[str, Any]] - ) -> str: + def format_notification_summary(self, notifications: list[dict[str, Any]]) -> str: """Format notifications into a human-readable summary. Args: @@ -209,9 +199,7 @@ class NotificationService: return "" count = len(notifications) - summary_lines = [ - f"El usuario tiene {count} notificación(es) pendiente(s):" - ] + summary_lines = [f"El usuario tiene {count} notificación(es) pendiente(s):"] for i, notif in enumerate(notifications, 1): texto = notif.get("texto", "Sin texto") From 670c00b1da333baa6c8c5d8d4e2c52f3c0cb4961 Mon Sep 17 00:00:00 2001 From: Anibal Angulo Date: Thu, 5 Mar 2026 06:14:32 +0000 Subject: [PATCH 4/5] Add CI --- .github/workflows/ci.yml | 33 +++++++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) create mode 100644 .github/workflows/ci.yml diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..b98199c --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,33 @@ +name: CI + +on: + push: + branches: [main] + pull_request: + branches: [main] + +jobs: + ci: + runs-on: ubuntu-latest + + steps: + - uses: actions/checkout@v4 + + - uses: astral-sh/setup-uv@v6 + with: + enable-cache: true + + - name: Install dependencies + run: uv sync --frozen + + - name: Format check + run: uv run ruff format --check + + - name: Lint + run: uv run ruff check + + - name: Type check + run: uv run ty check + + - name: Test + run: uv run pytest From 1803d011d0380822d3cb26e99320e805e1bf1b33 Mon Sep 17 00:00:00 2001 From: Anibal Angulo Date: Mon, 9 Mar 2026 07:36:47 +0000 Subject: [PATCH 5/5] Add Notification Backend Protocol (#24) Reviewed-on: https://gitea.ia-innovacion.work/va/agent/pulls/24 --- pyproject.toml | 1 + src/va_agent/agent.py | 10 +- src/va_agent/config.py | 1 + src/va_agent/dynamic_instruction.py | 93 ++++----- src/va_agent/notifications.py | 234 ++++++++++++----------- src/va_agent/server.py | 32 +--- utils/check_notifications.py | 108 +++++++++++ utils/check_notifications_firestore.py | 107 +++++++++++ utils/register_notification.py | 159 +++++++++++++++ utils/register_notification_firestore.py | 108 +++++++++++ uv.lock | 12 ++ 11 files changed, 676 insertions(+), 189 deletions(-) create mode 100644 utils/check_notifications.py create mode 100644 utils/check_notifications_firestore.py create mode 100644 utils/register_notification.py create mode 100644 utils/register_notification_firestore.py diff --git a/pyproject.toml b/pyproject.toml index f66b064..b454aba 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -14,6 +14,7 @@ dependencies = [ "pydantic-settings[yaml]>=2.13.1", "google-auth>=2.34.0", "google-genai>=1.64.0", + "redis>=5.0", ] [build-system] diff --git a/src/va_agent/agent.py b/src/va_agent/agent.py index 11b458b..6369737 100644 --- a/src/va_agent/agent.py +++ b/src/va_agent/agent.py @@ -1,5 +1,7 @@ """ADK agent with vector search RAG tool.""" +from functools import partial + from google import genai from google.adk.agents.llm_agent import Agent from google.adk.runners import Runner @@ -10,8 +12,9 @@ from google.genai.types import Content, Part from va_agent.auth import auth_headers_provider from va_agent.config import settings +from va_agent.dynamic_instruction import provide_dynamic_instruction from va_agent.governance import GovernancePlugin -from va_agent.notifications import NotificationService +from va_agent.notifications import FirestoreNotificationBackend from va_agent.session import FirestoreSessionService # MCP Toolset for RAG knowledge search @@ -32,10 +35,11 @@ session_service = FirestoreSessionService( ) # Notification service -notification_service = NotificationService( +notification_service = FirestoreNotificationBackend( db=firestore_db, collection_path=settings.notifications_collection_path, max_to_notify=settings.notifications_max_to_notify, + window_hours=settings.notifications_window_hours, ) # Agent with static and dynamic instructions @@ -43,11 +47,11 @@ governance = GovernancePlugin() agent = Agent( model=settings.agent_model, name=settings.agent_name, + instruction=partial(provide_dynamic_instruction, notification_service), static_instruction=Content( role="user", parts=[Part(text=settings.agent_instructions)], ), - instruction=settings.agent_instructions, tools=[toolset], after_model_callback=governance.after_model_callback, ) diff --git a/src/va_agent/config.py b/src/va_agent/config.py index ae33d2d..a4e202e 100644 --- a/src/va_agent/config.py +++ b/src/va_agent/config.py @@ -31,6 +31,7 @@ class AgentSettings(BaseSettings): "artifacts/bnt-orquestador-cognitivo-dev/notifications" ) notifications_max_to_notify: int = 5 + notifications_window_hours: float = 48 # MCP configuration mcp_audience: str diff --git a/src/va_agent/dynamic_instruction.py b/src/va_agent/dynamic_instruction.py index cc64fb9..664bad4 100644 --- a/src/va_agent/dynamic_instruction.py +++ b/src/va_agent/dynamic_instruction.py @@ -3,27 +3,49 @@ from __future__ import annotations import logging +import time from typing import TYPE_CHECKING if TYPE_CHECKING: from google.adk.agents.readonly_context import ReadonlyContext - from va_agent.notifications import NotificationService + from va_agent.notifications import NotificationBackend logger = logging.getLogger(__name__) +_SECONDS_PER_MINUTE = 60 +_SECONDS_PER_HOUR = 3600 +_MINUTES_PER_HOUR = 60 +_HOURS_PER_DAY = 24 + + +def _format_time_ago(now: float, ts: float) -> str: + """Return a human-readable Spanish label like 'hace 3 horas'.""" + diff = max(now - ts, 0) + minutes = int(diff // _SECONDS_PER_MINUTE) + hours = int(diff // _SECONDS_PER_HOUR) + + if minutes < 1: + return "justo ahora" + if minutes < _MINUTES_PER_HOUR: + return f"hace {minutes} min" + if hours < _HOURS_PER_DAY: + return f"hace {hours}h" + days = hours // _HOURS_PER_DAY + return f"hace {days}d" + + async def provide_dynamic_instruction( - notification_service: NotificationService, + notification_service: NotificationBackend, ctx: ReadonlyContext | None = None, ) -> str: - """Provide dynamic instructions based on pending notifications. + """Provide dynamic instructions based on recent notifications. This function is called by the ADK agent on each message. It: - 1. Checks if this is the first message in the session (< 2 events) - 2. Queries Firestore for pending notifications - 3. Marks them as notified - 4. Returns a dynamic instruction for the agent to mention them + 1. Queries Firestore for recent notifications + 2. Marks them as notified + 3. Returns a dynamic instruction for the agent to mention them Args: notification_service: Service for fetching/marking notifications @@ -43,71 +65,54 @@ async def provide_dynamic_instruction( logger.debug("No session available for dynamic instruction") return "" - # FOR TESTING: Always check for notifications - # (comment out to enable first-message-only) - # Only check on first message (when events list is empty - # or has only 1-2 events) - # Events include both user and agent messages, so < 2 means first interaction - # event_count = len(session.events) if session.events else 0 - # - # if event_count >= 2: - # logger.debug( - # "Skipping notification check: not first message (event_count=%d)", - # event_count, - # ) - # return "" - # Extract phone number from user_id (they are the same in this implementation) phone_number = session.user_id logger.info( - "First message detected for user %s, checking for pending notifications", + "Checking recent notifications for user %s", phone_number, ) try: - # Fetch pending notifications - pending_notifications = await notification_service.get_pending_notifications( + # Fetch recent notifications + recent_notifications = await notification_service.get_recent_notifications( phone_number ) - if not pending_notifications: - logger.info("No pending notifications for user %s", phone_number) + if not recent_notifications: + logger.info("No recent notifications for user %s", phone_number) return "" # Build dynamic instruction with notification details notification_ids = [ nid - for n in pending_notifications + for n in recent_notifications if (nid := n.get("id_notificacion")) is not None ] - count = len(pending_notifications) + count = len(recent_notifications) - # Format notification details for the agent + # Format notification details for the agent (most recent first) + now = time.time() notification_details = [] - for notif in pending_notifications: + for i, notif in enumerate(recent_notifications, 1): evento = notif.get("nombre_evento_dialogflow", "notificacion") texto = notif.get("texto", "Sin texto") - notification_details.append(f" - Evento: {evento} | Texto: {texto}") + ts = notif.get("timestamp_creacion", notif.get("timestampCreacion", 0)) + ago = _format_time_ago(now, ts) + notification_details.append( + f" {i}. [{ago}] Evento: {evento} | Texto: {texto}" + ) details_text = "\n".join(notification_details) + header = ( + f"Estas son {count} notificación(es) reciente(s)" + " de las cuales el usuario podría preguntar más:" + ) instruction = f""" -IMPORTANTE - NOTIFICACIONES PENDIENTES: - -El usuario tiene {count} notificación(es) sin leer: +{header} {details_text} - -INSTRUCCIONES: -- Menciona estas notificaciones de forma natural en tu respuesta inicial -- No necesitas leerlas todas literalmente, solo hazle saber que las tiene -- Sé breve y directo según tu personalidad (directo y cálido) -- Si el usuario pregunta algo específico, prioriza responder eso primero\ - y luego menciona las notificaciones - -Ejemplo: "¡Hola! 👋 Tienes {count} notificación(es)\ - pendiente(s). ¿Te gustaría revisarlas?" """ # Mark notifications as notified in Firestore diff --git a/src/va_agent/notifications.py b/src/va_agent/notifications.py index 8536fb2..0144608 100644 --- a/src/va_agent/notifications.py +++ b/src/va_agent/notifications.py @@ -4,7 +4,7 @@ from __future__ import annotations import logging import time -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING, Any, Protocol, runtime_checkable if TYPE_CHECKING: from google.cloud.firestore_v1.async_client import AsyncClient @@ -12,8 +12,28 @@ if TYPE_CHECKING: logger = logging.getLogger(__name__) -class NotificationService: - """Service for fetching and managing user notifications from Firestore.""" +@runtime_checkable +class NotificationBackend(Protocol): + """Backend-agnostic interface for notification storage.""" + + async def get_recent_notifications(self, phone_number: str) -> list[dict[str, Any]]: + """Return recent notifications for *phone_number*.""" + ... + + async def mark_as_notified( + self, phone_number: str, notification_ids: list[str] + ) -> bool: + """Mark the given notification IDs as notified. Return success.""" + ... + + +class FirestoreNotificationBackend: + """Firestore-backed notification backend (read-only). + + Reads notifications from a Firestore document keyed by phone number. + Filters by a configurable time window instead of tracking read/unread + state — the agent is awareness-only; delivery happens in the app. + """ def __init__( self, @@ -21,25 +41,18 @@ class NotificationService: db: AsyncClient, collection_path: str, max_to_notify: int = 5, + window_hours: float = 48, ) -> None: - """Initialize NotificationService. - - Args: - db: Firestore async client - collection_path: Path to notifications collection - max_to_notify: Maximum number of notifications to return - - """ + """Initialize with Firestore client and collection path.""" self._db = db self._collection_path = collection_path self._max_to_notify = max_to_notify + self._window_hours = window_hours - async def get_pending_notifications( - self, phone_number: str - ) -> list[dict[str, Any]]: - """Get pending notifications for a user. + async def get_recent_notifications(self, phone_number: str) -> list[dict[str, Any]]: + """Get recent notifications for a user. - Retrieves notifications that have not been notified by the agent yet, + Retrieves notifications created within the configured time window, ordered by timestamp (most recent first), limited to max_to_notify. Args: @@ -57,7 +70,6 @@ class NotificationService: """ try: - # Query Firestore document by phone number doc_ref = self._db.collection(self._collection_path).document(phone_number) doc = await doc_ref.get() @@ -74,26 +86,31 @@ class NotificationService: logger.info("No notifications in array for phone: %s", phone_number) return [] - # Filter notifications that have NOT been notified by the agent - pending = [ - n for n in all_notifications if not n.get("notified_by_agent", False) - ] + cutoff = time.time() - (self._window_hours * 3600) - if not pending: + def _ts(n: dict[str, Any]) -> Any: + return n.get( + "timestamp_creacion", + n.get("timestampCreacion", 0), + ) + + recent = [n for n in all_notifications if _ts(n) >= cutoff] + + if not recent: logger.info( - "All notifications already notified for phone: %s", phone_number + "No notifications within the last %.0fh for phone: %s", + self._window_hours, + phone_number, ) return [] - # Sort by timestamp_creacion (most recent first) - pending.sort(key=lambda n: n.get("timestamp_creacion", 0), reverse=True) + recent.sort(key=_ts, reverse=True) - # Return top N most recent - result = pending[: self._max_to_notify] + result = recent[: self._max_to_notify] logger.info( - "Found %d pending notifications for phone: %s (returning top %d)", - len(pending), + "Found %d recent notifications for phone: %s (returning top %d)", + len(recent), phone_number, len(result), ) @@ -107,114 +124,109 @@ class NotificationService: return result async def mark_as_notified( - self, phone_number: str, notification_ids: list[str] + self, + phone_number: str, # noqa: ARG002 + notification_ids: list[str], # noqa: ARG002 ) -> bool: - """Mark notifications as notified by the agent. + """No-op — the agent is not the delivery mechanism.""" + return True - Updates the notifications in Firestore by adding: - - notified_by_agent: true - - notified_at: current timestamp - Args: - phone_number: User's phone number (document ID) - notification_ids: List of id_notificacion values to mark +class RedisNotificationBackend: + """Redis-backed notification backend (read-only).""" - Returns: - True if update was successful, False otherwise + def __init__( + self, + *, + host: str = "127.0.0.1", + port: int = 6379, + max_to_notify: int = 5, + window_hours: float = 48, + ) -> None: + """Initialize with Redis connection parameters.""" + import redis.asyncio as aioredis # noqa: PLC0415 + self._client = aioredis.Redis( + host=host, + port=port, + decode_responses=True, + socket_connect_timeout=5, + ) + self._max_to_notify = max_to_notify + self._window_hours = window_hours + + async def get_recent_notifications(self, phone_number: str) -> list[dict[str, Any]]: + """Get recent notifications for a user from Redis. + + Reads from the ``notification:{phone}`` key, parses the JSON + payload, and returns notifications created within the configured + time window, sorted by creation timestamp (most recent first), + limited to *max_to_notify*. """ - if not notification_ids: - return True + import json # noqa: PLC0415 try: - doc_ref = self._db.collection(self._collection_path).document(phone_number) - doc = await doc_ref.get() + raw = await self._client.get(f"notification:{phone_number}") - if not doc.exists: - logger.warning( - "Cannot mark notifications as notified: document not found for %s", + if not raw: + logger.info( + "No notification data in Redis for phone: %s", phone_number, ) - return False + return [] - data = doc.to_dict() or {} - notificaciones = data.get("notificaciones", []) + data = json.loads(raw) + all_notifications: list[dict[str, Any]] = data.get("notificaciones", []) - if not notificaciones: - logger.warning( - "Cannot mark notifications: empty array for %s", phone_number + if not all_notifications: + logger.info( + "No notifications in array for phone: %s", + phone_number, ) - return False + return [] - # Update matching notifications - now = time.time() - updated_count = 0 + cutoff = time.time() - (self._window_hours * 3600) - for notif in notificaciones: - if notif.get("id_notificacion") in notification_ids: - notif["notified_by_agent"] = True - notif["notified_at"] = now - updated_count += 1 - - if updated_count == 0: - logger.warning( - "No notifications matched IDs for phone: %s", phone_number + def _ts(n: dict[str, Any]) -> Any: + return n.get( + "timestamp_creacion", + n.get("timestampCreacion", 0), ) - return False - # Save back to Firestore - await doc_ref.update( - { - "notificaciones": notificaciones, - "ultima_actualizacion": now, - } - ) + recent = [n for n in all_notifications if _ts(n) >= cutoff] + + if not recent: + logger.info( + "No notifications within the last %.0fh for phone: %s", + self._window_hours, + phone_number, + ) + return [] + + recent.sort(key=_ts, reverse=True) + + result = recent[: self._max_to_notify] logger.info( - "Marked %d notification(s) as notified for phone: %s", - updated_count, + "Found %d recent notifications for phone: %s (returning top %d)", + len(recent), phone_number, + len(result), ) except Exception: logger.exception( - "Failed to mark notifications as notified for phone: %s", + "Failed to fetch notifications from Redis for phone: %s", phone_number, ) - return False + return [] else: - return True + return result - def format_notification_summary(self, notifications: list[dict[str, Any]]) -> str: - """Format notifications into a human-readable summary. - - Args: - notifications: List of notification dictionaries - - Returns: - Formatted string summarizing the notifications - - """ - if not notifications: - return "" - - count = len(notifications) - summary_lines = [f"El usuario tiene {count} notificación(es) pendiente(s):"] - - for i, notif in enumerate(notifications, 1): - texto = notif.get("texto", "Sin texto") - params = notif.get("parametros", {}) - - # Extract key parameters if available - amount = params.get("notification_po_amount") - tx_id = params.get("notification_po_transaction_id") - - line = f"{i}. {texto}" - if amount: - line += f" (monto: ${amount})" - if tx_id: - line += f" [ID: {tx_id}]" - - summary_lines.append(line) - - return "\n".join(summary_lines) + async def mark_as_notified( + self, + phone_number: str, # noqa: ARG002 + notification_ids: list[str], # noqa: ARG002 + ) -> bool: + """No-op — the agent is not the delivery mechanism.""" + return True diff --git a/src/va_agent/server.py b/src/va_agent/server.py index b511653..bfa2597 100644 --- a/src/va_agent/server.py +++ b/src/va_agent/server.py @@ -22,20 +22,11 @@ app = FastAPI(title="Vaia Agent") # --------------------------------------------------------------------------- -class NotificationPayload(BaseModel): - """Notification context sent alongside a user query.""" - - text: str | None = None - parameters: dict[str, Any] = Field(default_factory=dict) - - class QueryRequest(BaseModel): """Incoming query request from the integration layer.""" phone_number: str text: str - type: str = "conversation" - notification: NotificationPayload | None = None language_code: str = "es" @@ -56,26 +47,6 @@ class ErrorResponse(BaseModel): status: int -# --------------------------------------------------------------------------- -# Helpers -# --------------------------------------------------------------------------- - - -def _build_user_message(request: QueryRequest) -> str: - """Compose the text sent to the agent, including notification context.""" - if request.type == "notification" and request.notification: - parts = [request.text] - if request.notification.text: - parts.append(f"\n[Notificación recibida]: {request.notification.text}") - if request.notification.parameters: - formatted = ", ".join( - f"{k}: {v}" for k, v in request.notification.parameters.items() - ) - parts.append(f"[Parámetros de notificación]: {formatted}") - return "\n".join(parts) - return request.text - - # --------------------------------------------------------------------------- # Endpoints # --------------------------------------------------------------------------- @@ -92,13 +63,12 @@ def _build_user_message(request: QueryRequest) -> str: ) async def query(request: QueryRequest) -> QueryResponse: """Process a user message and return a generated response.""" - user_message = _build_user_message(request) session_id = request.phone_number user_id = request.phone_number new_message = Content( role="user", - parts=[Part(text=user_message)], + parts=[Part(text=request.text)], ) try: diff --git a/utils/check_notifications.py b/utils/check_notifications.py new file mode 100644 index 0000000..b2d0629 --- /dev/null +++ b/utils/check_notifications.py @@ -0,0 +1,108 @@ +# /// script +# requires-python = ">=3.12" +# dependencies = ["redis>=5.0", "pydantic>=2.0"] +# /// +"""Check pending notifications for a phone number. + +Usage: + REDIS_HOST=10.33.22.4 uv run utils/check_notifications.py + REDIS_HOST=10.33.22.4 uv run utils/check_notifications.py --since 2026-01-01 +""" + +import json +import os +import sys +from datetime import UTC, datetime + +import redis +from pydantic import AliasChoices, BaseModel, Field, ValidationError + + +class Notification(BaseModel): + id_notificacion: str = Field( + validation_alias=AliasChoices("id_notificacion", "idNotificacion"), + ) + telefono: str + timestamp_creacion: datetime = Field( + validation_alias=AliasChoices("timestamp_creacion", "timestampCreacion"), + ) + texto: str + nombre_evento_dialogflow: str = Field( + validation_alias=AliasChoices( + "nombre_evento_dialogflow", "nombreEventoDialogflow" + ), + ) + codigo_idioma_dialogflow: str = Field( + default="es", + validation_alias=AliasChoices( + "codigo_idioma_dialogflow", "codigoIdiomaDialogflow" + ), + ) + parametros: dict = Field(default_factory=dict) + status: str + + +class NotificationSession(BaseModel): + session_id: str = Field( + validation_alias=AliasChoices("session_id", "sessionId"), + ) + telefono: str + fecha_creacion: datetime = Field( + validation_alias=AliasChoices("fecha_creacion", "fechaCreacion"), + ) + ultima_actualizacion: datetime = Field( + validation_alias=AliasChoices("ultima_actualizacion", "ultimaActualizacion"), + ) + notificaciones: list[Notification] + + +HOST = os.environ.get("REDIS_HOST", "127.0.0.1") +PORT = int(os.environ.get("REDIS_PORT", "6379")) + + +def main() -> None: + if len(sys.argv) < 2: + print(f"Usage: {sys.argv[0]} [--since YYYY-MM-DD]") + sys.exit(1) + + phone = sys.argv[1] + since = None + if "--since" in sys.argv: + idx = sys.argv.index("--since") + since = datetime.fromisoformat(sys.argv[idx + 1]).replace(tzinfo=UTC) + + r = redis.Redis(host=HOST, port=PORT, decode_responses=True, socket_connect_timeout=5) + raw = r.get(f"notification:{phone}") + + if not raw: + print(f"📭 No notifications found for {phone}") + sys.exit(0) + + try: + session = NotificationSession.model_validate(json.loads(raw)) + except ValidationError as e: + print(f"❌ Invalid notification data for {phone}:\n{e}") + sys.exit(1) + + active = [n for n in session.notificaciones if n.status == "active"] + + if since: + active = [n for n in active if n.timestamp_creacion >= since] + + if not active: + print(f"📭 No {'new ' if since else ''}active notifications for {phone}") + sys.exit(0) + + print(f"🔔 {len(active)} active notification(s) for {phone}\n") + for i, n in enumerate(active, 1): + categoria = n.parametros.get("notification_po_Categoria", "") + print(f" [{i}] {n.timestamp_creacion.isoformat()}") + print(f" ID: {n.id_notificacion}") + if categoria: + print(f" Category: {categoria}") + print(f" {n.texto[:120]}{'…' if len(n.texto) > 120 else ''}") + print() + + +if __name__ == "__main__": + main() diff --git a/utils/check_notifications_firestore.py b/utils/check_notifications_firestore.py new file mode 100644 index 0000000..1408943 --- /dev/null +++ b/utils/check_notifications_firestore.py @@ -0,0 +1,107 @@ +# /// script +# requires-python = ">=3.12" +# dependencies = ["google-cloud-firestore>=2.0", "pyyaml>=6.0"] +# /// +"""Check recent notifications in Firestore for a phone number. + +Usage: + uv run utils/check_notifications_firestore.py + uv run utils/check_notifications_firestore.py --hours 24 +""" + +import sys +import time + +import yaml +from google.cloud.firestore import Client + +_SECONDS_PER_HOUR = 3600 +_DEFAULT_WINDOW_HOURS = 48 + + +def main() -> None: + if len(sys.argv) < 2: + print(f"Usage: {sys.argv[0]} [--hours N]") + sys.exit(1) + + phone = sys.argv[1] + window_hours = _DEFAULT_WINDOW_HOURS + if "--hours" in sys.argv: + idx = sys.argv.index("--hours") + window_hours = float(sys.argv[idx + 1]) + + with open("config.yaml") as f: + cfg = yaml.safe_load(f) + + db = Client( + project=cfg["google_cloud_project"], + database=cfg["firestore_db"], + ) + + collection_path = cfg["notifications_collection_path"] + doc_ref = db.collection(collection_path).document(phone) + doc = doc_ref.get() + + if not doc.exists: + print(f"📭 No notifications found for {phone}") + sys.exit(0) + + data = doc.to_dict() or {} + all_notifications = data.get("notificaciones", []) + + if not all_notifications: + print(f"📭 No notifications found for {phone}") + sys.exit(0) + + cutoff = time.time() - (window_hours * _SECONDS_PER_HOUR) + + def _ts(n: dict) -> float: + return n.get("timestamp_creacion", n.get("timestampCreacion", 0)) + + recent = [n for n in all_notifications if _ts(n) >= cutoff] + recent.sort(key=_ts, reverse=True) + + if not recent: + print( + f"📭 No notifications within the last" + f" {window_hours:.0f}h for {phone}" + ) + sys.exit(0) + + print( + f"🔔 {len(recent)} notification(s) for {phone}" + f" (last {window_hours:.0f}h)\n" + ) + now = time.time() + for i, n in enumerate(recent, 1): + ts = _ts(n) + ago = _format_time_ago(now, ts) + categoria = n.get("parametros", {}).get( + "notification_po_Categoria", "" + ) + texto = n.get("texto", "") + print(f" [{i}] {ago}") + print(f" ID: {n.get('id_notificacion', '?')}") + if categoria: + print(f" Category: {categoria}") + print(f" {texto[:120]}{'…' if len(texto) > 120 else ''}") + print() + + +def _format_time_ago(now: float, ts: float) -> str: + diff = max(now - ts, 0) + minutes = int(diff // 60) + hours = int(diff // _SECONDS_PER_HOUR) + + if minutes < 1: + return "justo ahora" + if minutes < 60: + return f"hace {minutes} min" + if hours < 24: + return f"hace {hours}h" + days = hours // 24 + return f"hace {days}d" + + +if __name__ == "__main__": + main() diff --git a/utils/register_notification.py b/utils/register_notification.py new file mode 100644 index 0000000..d5df458 --- /dev/null +++ b/utils/register_notification.py @@ -0,0 +1,159 @@ +# /// script +# requires-python = ">=3.12" +# dependencies = ["redis>=5.0"] +# /// +"""Register a new notification in Redis for a given phone number. + +Usage: + REDIS_HOST=10.33.22.4 uv run utils/register_notification.py + +The notification content is randomly picked from a predefined set based on +existing entries in Memorystore. +""" + +import json +import os +import random +import sys +import uuid +from datetime import UTC, datetime + +import redis + +HOST = os.environ.get("REDIS_HOST", "127.0.0.1") +PORT = int(os.environ.get("REDIS_PORT", "6379")) +TTL_SECONDS = 18 * 24 * 3600 # ~18 days, matching existing keys + +NOTIFICATION_TEMPLATES = [ + { + "texto": ( + "Se detectó un cargo de $1,500 en tu cuenta" + ), + "parametros": { + "notification_po_transaction_id": "TXN15367", + "notification_po_amount": 5814, + }, + }, + { + "texto": ( + "💡 Recuerda que puedes obtener tu Adelanto de Nómina en cualquier" + " momento, sólo tienes que seleccionar Solicitud adelanto de Nómina" + " en tu app." + ), + "parametros": { + "notification_po_Categoria": "Adelanto de Nómina solicitud", + "notification_po_caption": "Adelanto de Nómina", + "notification_po_CTA": "Realiza la solicitud desde tu app", + "notification_po_Descripcion": ( + "Notificación para incentivar la solicitud de Adelanto de" + " Nómina desde la APP" + ), + "notification_po_link": ( + "https://public-media.yalochat.com/banorte/" + "1764025754-10e06fb8-b4e6-484c-ad0b-7f677429380e-03-ADN-Toque-1.jpg" + ), + "notification_po_Beneficios": ( + "Tasa de interés de 0%: Solicita tu Adelanto sin preocuparte" + " por los intereses, así de fácil. No requiere garantías o aval." + ), + "notification_po_Requisitos": ( + "Tener Cuenta Digital o Cuenta Digital Ilimitada con dispersión" + " de Nómina No tener otro Adelanto vigente Ingreso neto mensual" + " mayor a $2,000" + ), + }, + }, + { + "texto": ( + "Estás a un clic de Programa de Lealtad, entra a tu app y finaliza" + " Tu contratación en instantes. ⏱ 🤳" + ), + "parametros": { + "notification_po_Categoria": "Tarjeta de Crédito Contratación", + "notification_po_caption": "Tarjeta de Crédito", + "notification_po_CTA": "Entra a tu app y contrata en instantes", + "notification_po_Descripcion": ( + "Notificación para terminar el proceso de contratación de la" + " Tarjeta de Crédito, desde la app" + ), + "notification_po_link": ( + "https://public-media.yalochat.com/banorte/" + "1764363798-05dadc23-6e47-447c-8e38-0346f25e31c0-15-TDC-Toque-1.jpg" + ), + "notification_po_Beneficios": ( + "Acceso al Programa de Lealtad: Cada compra suma, gana" + " experiencias exclusivas" + ), + "notification_po_Requisitos": ( + "Ser persona física o física con actividad empresarial." + " Ingresos mínimos de $2,000 pesos mensuales. Sin historial de" + " crédito o con buró positivo" + ), + }, + }, + { + "texto": ( + "🚀 ¿Listo para obtener tu Cápsula Plus? Continúa en tu app y" + " termina al instante. Conoce más en: va.app" + ), + "parametros": {}, + }, + { + "texto": ( + "🚀 ¿Listo para obtener tu Cuenta Digital ilimitada? Continúa en" + " tu app y termina al instante. Conoce más en: va.app" + ), + "parametros": {}, + }, +] + + +def main() -> None: + if len(sys.argv) < 2: + print(f"Usage: {sys.argv[0]} ") + sys.exit(1) + + phone = sys.argv[1] + r = redis.Redis(host=HOST, port=PORT, decode_responses=True, socket_connect_timeout=5) + + now = datetime.now(UTC).isoformat() + template = random.choice(NOTIFICATION_TEMPLATES) + notification = { + "id_notificacion": str(uuid.uuid4()), + "telefono": phone, + "timestamp_creacion": now, + "texto": template["texto"], + "nombre_evento_dialogflow": "notificacion", + "codigo_idioma_dialogflow": "es", + "parametros": template["parametros"], + "status": "active", + } + + session_key = f"notification:{phone}" + existing = r.get(session_key) + + if existing: + session = json.loads(existing) + session["ultima_actualizacion"] = now + session["notificaciones"].append(notification) + else: + session = { + "session_id": phone, + "telefono": phone, + "fecha_creacion": now, + "ultima_actualizacion": now, + "notificaciones": [notification], + } + + r.set(session_key, json.dumps(session, ensure_ascii=False), ex=TTL_SECONDS) + r.set(f"notification:phone_to_notification:{phone}", phone, ex=TTL_SECONDS) + + total = len(session["notificaciones"]) + print(f"✅ Registered notification for {phone}") + print(f" ID: {notification['id_notificacion']}") + print(f" Text: {template['texto'][:80]}...") + print(f" Total notifications for this phone: {total}") + + +if __name__ == "__main__": + main() diff --git a/utils/register_notification_firestore.py b/utils/register_notification_firestore.py new file mode 100644 index 0000000..c5f01dc --- /dev/null +++ b/utils/register_notification_firestore.py @@ -0,0 +1,108 @@ +# /// script +# requires-python = ">=3.12" +# dependencies = ["google-cloud-firestore>=2.0", "pyyaml>=6.0"] +# /// +"""Register a new notification in Firestore for a given phone number. + +Usage: + uv run utils/register_notification_firestore.py + +Reads project/database/collection settings from config.yaml. +""" + +import random +import sys +import time +import uuid + +import yaml +from google.cloud.firestore import Client + +NOTIFICATION_TEMPLATES = [ + { + "texto": "Se detectó un cargo de $1,500 en tu cuenta", + "parametros": { + "notification_po_transaction_id": "TXN15367", + "notification_po_amount": 5814, + }, + }, + { + "texto": ( + "💡 Recuerda que puedes obtener tu Adelanto de Nómina en" + " cualquier momento, sólo tienes que seleccionar Solicitud" + " adelanto de Nómina en tu app." + ), + "parametros": { + "notification_po_Categoria": "Adelanto de Nómina solicitud", + "notification_po_caption": "Adelanto de Nómina", + }, + }, + { + "texto": ( + "Estás a un clic de Programa de Lealtad, entra a tu app y" + " finaliza Tu contratación en instantes. ⏱ 🤳" + ), + "parametros": { + "notification_po_Categoria": "Tarjeta de Crédito Contratación", + "notification_po_caption": "Tarjeta de Crédito", + }, + }, + { + "texto": ( + "🚀 ¿Listo para obtener tu Cápsula Plus? Continúa en tu app" + " y termina al instante. Conoce más en: va.app" + ), + "parametros": {}, + }, +] + + +def main() -> None: + if len(sys.argv) < 2: + print(f"Usage: {sys.argv[0]} ") + sys.exit(1) + + phone = sys.argv[1] + + with open("config.yaml") as f: + cfg = yaml.safe_load(f) + + db = Client( + project=cfg["google_cloud_project"], + database=cfg["firestore_db"], + ) + + collection_path = cfg["notifications_collection_path"] + doc_ref = db.collection(collection_path).document(phone) + + template = random.choice(NOTIFICATION_TEMPLATES) + notification = { + "id_notificacion": str(uuid.uuid4()), + "telefono": phone, + "timestamp_creacion": time.time(), + "texto": template["texto"], + "nombre_evento_dialogflow": "notificacion", + "codigo_idioma_dialogflow": "es", + "parametros": template["parametros"], + "status": "active", + } + + doc = doc_ref.get() + if doc.exists: + data = doc.to_dict() or {} + notifications = data.get("notificaciones", []) + notifications.append(notification) + doc_ref.update({"notificaciones": notifications}) + else: + doc_ref.set({"notificaciones": [notification]}) + + total = len(doc_ref.get().to_dict().get("notificaciones", [])) + print(f"✅ Registered notification for {phone}") + print(f" ID: {notification['id_notificacion']}") + print(f" Text: {template['texto'][:80]}...") + print(f" Collection: {collection_path}") + print(f" Total notifications for this phone: {total}") + + +if __name__ == "__main__": + main() diff --git a/uv.lock b/uv.lock index 5598062..952f9fa 100644 --- a/uv.lock +++ b/uv.lock @@ -871,6 +871,7 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/ea/ab/1608e5a7578e62113506740b88066bf09888322a311cff602105e619bd87/greenlet-3.3.2-cp312-cp312-macosx_11_0_universal2.whl", hash = "sha256:ac8d61d4343b799d1e526db579833d72f23759c71e07181c2d2944e429eb09cd", size = 280358, upload-time = "2026-02-20T20:17:43.971Z" }, { url = "https://files.pythonhosted.org/packages/a5/23/0eae412a4ade4e6623ff7626e38998cb9b11e9ff1ebacaa021e4e108ec15/greenlet-3.3.2-cp312-cp312-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:3ceec72030dae6ac0c8ed7591b96b70410a8be370b6a477b1dbc072856ad02bd", size = 601217, upload-time = "2026-02-20T20:47:31.462Z" }, { url = "https://files.pythonhosted.org/packages/f8/16/5b1678a9c07098ecb9ab2dd159fafaf12e963293e61ee8d10ecb55273e5e/greenlet-3.3.2-cp312-cp312-manylinux_2_24_ppc64le.manylinux_2_28_ppc64le.whl", hash = "sha256:a2a5be83a45ce6188c045bcc44b0ee037d6a518978de9a5d97438548b953a1ac", size = 611792, upload-time = "2026-02-20T20:55:58.423Z" }, + { url = "https://files.pythonhosted.org/packages/5c/c5/cc09412a29e43406eba18d61c70baa936e299bc27e074e2be3806ed29098/greenlet-3.3.2-cp312-cp312-manylinux_2_24_s390x.manylinux_2_28_s390x.whl", hash = "sha256:ae9e21c84035c490506c17002f5c8ab25f980205c3e61ddb3a2a2a2e6c411fcb", size = 626250, upload-time = "2026-02-20T21:02:46.596Z" }, { url = "https://files.pythonhosted.org/packages/50/1f/5155f55bd71cabd03765a4aac9ac446be129895271f73872c36ebd4b04b6/greenlet-3.3.2-cp312-cp312-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:43e99d1749147ac21dde49b99c9abffcbc1e2d55c67501465ef0930d6e78e070", size = 613875, upload-time = "2026-02-20T20:21:01.102Z" }, { url = "https://files.pythonhosted.org/packages/fc/dd/845f249c3fcd69e32df80cdab059b4be8b766ef5830a3d0aa9d6cad55beb/greenlet-3.3.2-cp312-cp312-musllinux_1_2_aarch64.whl", hash = "sha256:4c956a19350e2c37f2c48b336a3afb4bff120b36076d9d7fb68cb44e05d95b79", size = 1571467, upload-time = "2026-02-20T20:49:33.495Z" }, { url = "https://files.pythonhosted.org/packages/2a/50/2649fe21fcc2b56659a452868e695634722a6655ba245d9f77f5656010bf/greenlet-3.3.2-cp312-cp312-musllinux_1_2_x86_64.whl", hash = "sha256:6c6f8ba97d17a1e7d664151284cb3315fc5f8353e75221ed4324f84eb162b395", size = 1640001, upload-time = "2026-02-20T20:21:09.154Z" }, @@ -1625,6 +1626,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/1a/08/67bd04656199bbb51dbed1439b7f27601dfb576fb864099c7ef0c3e55531/pyyaml-6.0.3-cp312-cp312-win_arm64.whl", hash = "sha256:64386e5e707d03a7e172c0701abfb7e10f0fb753ee1d773128192742712a98fd", size = 140344, upload-time = "2025-09-25T21:32:22.617Z" }, ] +[[package]] +name = "redis" +version = "7.2.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/e9/31/1476f206482dd9bc53fdbbe9f6fbd5e05d153f18e54667ce839df331f2e6/redis-7.2.1.tar.gz", hash = "sha256:6163c1a47ee2d9d01221d8456bc1c75ab953cbda18cfbc15e7140e9ba16ca3a5", size = 4906735, upload-time = "2026-02-25T20:05:18.171Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/ca/98/1dd1a5c060916cf21d15e67b7d6a7078e26e2605d5c37cbc9f4f5454c478/redis-7.2.1-py3-none-any.whl", hash = "sha256:49e231fbc8df2001436ae5252b3f0f3dc930430239bfeb6da4c7ee92b16e5d33", size = 396057, upload-time = "2026-02-25T20:05:16.533Z" }, +] + [[package]] name = "referencing" version = "0.37.0" @@ -1926,6 +1936,7 @@ dependencies = [ { name = "google-cloud-firestore" }, { name = "google-genai" }, { name = "pydantic-settings", extra = ["yaml"] }, + { name = "redis" }, ] [package.dev-dependencies] @@ -1944,6 +1955,7 @@ requires-dist = [ { name = "google-cloud-firestore", specifier = ">=2.23.0" }, { name = "google-genai", specifier = ">=1.64.0" }, { name = "pydantic-settings", extras = ["yaml"], specifier = ">=2.13.1" }, + { name = "redis", specifier = ">=5.0" }, ] [package.metadata.requires-dev]