main
py 195 lines 7.21 KB
Raw
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()