| 1 | from typing import Any |
| 2 | from typing import Dict |
| 3 | from typing import List |
| 4 | from typing import Optional |
| 5 | |
| 6 | from fastapi import HTTPException |
| 7 | from loguru import logger |
| 8 | |
| 9 | from app.agents.wazuh.syscollector.schema.packages import AgentPackagesResponse |
| 10 | from app.agents.wazuh.syscollector.schema.packages import IndexerPackageAgent |
| 11 | from app.agents.wazuh.syscollector.schema.packages import IndexerPackageDetail |
| 12 | from app.agents.wazuh.syscollector.schema.packages import IndexerPackageItem |
| 13 | from app.agents.wazuh.syscollector.schema.packages import IndexerPackagesResponse |
| 14 | from app.agents.wazuh.syscollector.schema.packages import PackageItem |
| 15 | from app.connectors.wazuh_indexer.utils.universal import ( |
| 16 | create_wazuh_indexer_client_async, |
| 17 | ) |
| 18 | from app.connectors.wazuh_manager.utils.universal import send_get_request |
| 19 | |
| 20 | PACKAGES_INDEX_PATTERN = "wazuh-states-inventory-packages-*" |
| 21 | |
| 22 | |
| 23 | async def collect_agent_packages( |
| 24 | agent_id: str, |
| 25 | limit: int = 500, |
| 26 | offset: int = 0, |
| 27 | sort: Optional[str] = None, |
| 28 | search: Optional[str] = None, |
| 29 | select: Optional[List[str]] = None, |
| 30 | vendor: Optional[str] = None, |
| 31 | name: Optional[str] = None, |
| 32 | architecture: Optional[str] = None, |
| 33 | format: Optional[str] = None, |
| 34 | version: Optional[str] = None, |
| 35 | q: Optional[str] = None, |
| 36 | ) -> AgentPackagesResponse: |
| 37 | """ |
| 38 | Fetch installed packages for a specific agent from the Wazuh Manager |
| 39 | syscollector API. |
| 40 | |
| 41 | Args: |
| 42 | agent_id: The Wazuh agent ID. |
| 43 | limit: Maximum number of packages to return (1-100000, default 500). |
| 44 | offset: First element to return (pagination). |
| 45 | sort: Sort field(s), prefixed with +/- for order. |
| 46 | search: Free-text search string. |
| 47 | select: List of fields to return. |
| 48 | vendor: Filter by vendor. |
| 49 | name: Filter by package name. |
| 50 | architecture: Filter by architecture. |
| 51 | format: Filter by package format (e.g. 'deb', 'rpm'). |
| 52 | version: Filter by package version. |
| 53 | q: Advanced query filter string. |
| 54 | |
| 55 | Returns: |
| 56 | AgentPackagesResponse with the list of packages. |
| 57 | """ |
| 58 | params: Dict[str, Any] = { |
| 59 | "limit": limit, |
| 60 | "offset": offset, |
| 61 | "wait_for_complete": True, |
| 62 | } |
| 63 | |
| 64 | if sort is not None: |
| 65 | params["sort"] = sort |
| 66 | if search is not None: |
| 67 | params["search"] = search |
| 68 | if select is not None: |
| 69 | params["select"] = ",".join(select) |
| 70 | if vendor is not None: |
| 71 | params["vendor"] = vendor |
| 72 | if name is not None: |
| 73 | params["name"] = name |
| 74 | if architecture is not None: |
| 75 | params["architecture"] = architecture |
| 76 | if format is not None: |
| 77 | params["format"] = format |
| 78 | if version is not None: |
| 79 | params["version"] = version |
| 80 | if q is not None: |
| 81 | params["q"] = q |
| 82 | |
| 83 | response = await send_get_request( |
| 84 | endpoint=f"/syscollector/{agent_id}/packages", |
| 85 | params=params, |
| 86 | ) |
| 87 | |
| 88 | if not response.get("success"): |
| 89 | raise HTTPException( |
| 90 | status_code=500, |
| 91 | detail=response.get("message", "Failed to fetch packages from Wazuh Manager"), |
| 92 | ) |
| 93 | |
| 94 | wazuh_data = response.get("data", {}).get("data", {}) |
| 95 | affected_items = wazuh_data.get("affected_items", []) |
| 96 | total_affected_items = wazuh_data.get("total_affected_items", len(affected_items)) |
| 97 | |
| 98 | packages = [PackageItem(**item) for item in affected_items] |
| 99 | |
| 100 | logger.info(f"Fetched {len(packages)} packages for agent {agent_id}") |
| 101 | |
| 102 | return AgentPackagesResponse( |
| 103 | packages=packages, |
| 104 | total_affected_items=total_affected_items, |
| 105 | success=True, |
| 106 | message=f"Successfully fetched {len(packages)} packages for agent {agent_id}", |
| 107 | ) |
| 108 | |
| 109 | |
| 110 | async def search_packages_in_indexer( |
| 111 | package_name: Optional[str] = None, |
| 112 | agent_name: Optional[str] = None, |
| 113 | agent_id: Optional[str] = None, |
| 114 | architecture: Optional[str] = None, |
| 115 | package_type: Optional[str] = None, |
| 116 | vendor: Optional[str] = None, |
| 117 | package_version: Optional[str] = None, |
| 118 | size: int = 500, |
| 119 | ) -> IndexerPackagesResponse: |
| 120 | """ |
| 121 | Search the Wazuh Indexer for package inventory data across all agents. |
| 122 | |
| 123 | Queries the ``wazuh-states-inventory-packages-*`` index pattern and returns |
| 124 | matching documents with optional filters. |
| 125 | |
| 126 | Args: |
| 127 | package_name: Filter by package name (wildcard match). |
| 128 | agent_name: Filter by agent name (wildcard match). |
| 129 | agent_id: Filter by agent ID (exact match). |
| 130 | architecture: Filter by architecture (exact match). |
| 131 | package_type: Filter by package type, e.g. ``deb``, ``rpm`` (exact match). |
| 132 | vendor: Filter by vendor (wildcard match). |
| 133 | package_version: Filter by package version (wildcard match). |
| 134 | size: Maximum number of documents to return (default 500). |
| 135 | |
| 136 | Returns: |
| 137 | IndexerPackagesResponse with the matching packages. |
| 138 | """ |
| 139 | must_clauses: List[Dict[str, Any]] = [] |
| 140 | |
| 141 | if package_name is not None: |
| 142 | must_clauses.append({"wildcard": {"package.name": {"value": f"*{package_name}*", "case_insensitive": True}}}) |
| 143 | if agent_name is not None: |
| 144 | must_clauses.append({"wildcard": {"agent.name": {"value": f"*{agent_name}*", "case_insensitive": True}}}) |
| 145 | if agent_id is not None: |
| 146 | must_clauses.append({"term": {"agent.id": agent_id}}) |
| 147 | if architecture is not None: |
| 148 | must_clauses.append({"term": {"package.architecture": architecture}}) |
| 149 | if package_type is not None: |
| 150 | must_clauses.append({"term": {"package.type": package_type}}) |
| 151 | if vendor is not None: |
| 152 | must_clauses.append({"wildcard": {"package.vendor": {"value": f"*{vendor}*", "case_insensitive": True}}}) |
| 153 | if package_version is not None: |
| 154 | must_clauses.append({"wildcard": {"package.version": {"value": f"*{package_version}*", "case_insensitive": True}}}) |
| 155 | |
| 156 | query: Dict[str, Any] = {"query": {"bool": {"must": must_clauses}}} if must_clauses else {"query": {"match_all": {}}} |
| 157 | |
| 158 | es_client = await create_wazuh_indexer_client_async("Wazuh-Indexer") |
| 159 | |
| 160 | try: |
| 161 | response = await es_client.search( |
| 162 | index=PACKAGES_INDEX_PATTERN, |
| 163 | body=query, |
| 164 | size=size, |
| 165 | ) |
| 166 | |
| 167 | hits = response.get("hits", {}) |
| 168 | total = hits.get("total", {}) |
| 169 | total_value = total.get("value", 0) if isinstance(total, dict) else total |
| 170 | |
| 171 | packages: List[IndexerPackageItem] = [] |
| 172 | for hit in hits.get("hits", []): |
| 173 | source = hit.get("_source", {}) |
| 174 | packages.append( |
| 175 | IndexerPackageItem( |
| 176 | _index=hit.get("_index"), |
| 177 | _id=hit.get("_id"), |
| 178 | agent=IndexerPackageAgent(**source.get("agent", {})) if source.get("agent") else None, |
| 179 | package=IndexerPackageDetail(**source.get("package", {})) if source.get("package") else None, |
| 180 | ), |
| 181 | ) |
| 182 | |
| 183 | logger.info(f"Indexer search returned {len(packages)} packages (total matched: {total_value})") |
| 184 | |
| 185 | return IndexerPackagesResponse( |
| 186 | packages=packages, |
| 187 | total=total_value, |
| 188 | success=True, |
| 189 | message=f"Successfully retrieved {len(packages)} packages from the indexer", |
| 190 | ) |
| 191 | except Exception as e: |
| 192 | logger.error(f"Error searching packages in Wazuh Indexer: {e}") |
| 193 | raise HTTPException(status_code=500, detail=f"Failed to search packages in Wazuh Indexer: {e}") |
| 194 | finally: |
| 195 | await es_client.close() |