Skip to content

Commit 4e4d2bc

Browse files
committed
Feat: Add streaming bulk methods
- Add bulk_export_stream() generator for memory-efficient streaming - Add bulk_service_stream() generator for memory-efficient streaming
1 parent d1d237b commit 4e4d2bc

1 file changed

Lines changed: 62 additions & 0 deletions

File tree

leakix/client.py

Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -283,3 +283,65 @@ def search(self, query: str, scope: str = "leak", page: int = 0) -> AbstractResp
283283
if scope == "service":
284284
return self.get_service(queries=queries, page=page)
285285
return self.get_leak(queries=queries, page=page)
286+
287+
def bulk_export_stream(self, queries: list[Query] | None = None):
288+
"""
289+
Streaming version of bulk_export. Yields L9Aggregation objects one by one.
290+
291+
This is more memory efficient for large result sets as it doesn't load
292+
all results into memory at once.
293+
294+
Example:
295+
>>> for aggregation in client.bulk_export_stream([MustQuery(PluginField(Plugin.GitConfigHttpPlugin))]):
296+
... print(aggregation.events[0].ip)
297+
"""
298+
url = f"{self.base_url}/bulk/search"
299+
if queries is None or len(queries) == 0:
300+
serialized_query = EmptyQuery().serialize()
301+
else:
302+
serialized_query = [q.serialize() for q in queries]
303+
serialized_query = " ".join(serialized_query)
304+
serialized_query = f"{serialized_query}"
305+
params = {"q": serialized_query}
306+
r = requests.get(url, params=params, headers=self.headers, stream=True)
307+
if r.status_code != 200:
308+
return
309+
for line in r.iter_lines():
310+
if not line:
311+
continue
312+
try:
313+
json_event = json.loads(line)
314+
yield l9format.L9Aggregation.from_dict(json_event)
315+
except Exception:
316+
pass
317+
318+
def bulk_service_stream(self, queries: list[Query] | None = None):
319+
"""
320+
Streaming version of bulk_service. Yields L9Event objects one by one.
321+
322+
This is more memory efficient for large result sets as it doesn't load
323+
all results into memory at once.
324+
325+
Example:
326+
>>> for event in client.bulk_service_stream([MustQuery(PluginField(Plugin.GitConfigHttpPlugin))]):
327+
... print(event.ip)
328+
"""
329+
url = f"{self.base_url}/bulk/service"
330+
if queries is None or len(queries) == 0:
331+
serialized_query = EmptyQuery().serialize()
332+
else:
333+
serialized_query = [q.serialize() for q in queries]
334+
serialized_query = " ".join(serialized_query)
335+
serialized_query = f"{serialized_query}"
336+
params = {"q": serialized_query}
337+
r = requests.get(url, params=params, headers=self.headers, stream=True)
338+
if r.status_code != 200:
339+
return
340+
for line in r.iter_lines():
341+
if not line:
342+
continue
343+
try:
344+
json_event = json.loads(line)
345+
yield l9format.L9Event.from_dict(json_event)
346+
except Exception:
347+
pass

0 commit comments

Comments
 (0)