Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
125 changes: 120 additions & 5 deletions bbot/modules/http.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import time
from http.cookies import SimpleCookie
from urllib.parse import urlparse

Expand Down Expand Up @@ -38,6 +39,13 @@ async def setup(self):
self.store_responses = self.config.get("store_responses", False)
self.client = self.helpers.blasthttp
self.waf_yara_rule = self.helpers.yara.compile_strings(self.helpers.get_waf_strings(), nocase=True)
self._host_cooldowns = {}
self._deferred_events = []
self._429_retry_counts = {}
self._max_429_retries = 3
self._429_default_interval = self.scan.web_config.get("429_sleep_interval", 30)
self._429_max_interval = self.scan.web_config.get("429_max_sleep_interval", 60)
self._wakeup_target = 0
return True

async def filter_event(self, event):
Expand Down Expand Up @@ -86,6 +94,77 @@ def _incoming_dedup_hash(self, event):
urls, url_hash = self.make_url_metadata(event)
return url_hash

@property
def finished(self):
if self._deferred_events:
return False
return super().finished

def is_incoming_duplicate(self, event, add=False):
if getattr(event, "_429_retry", False):
event._429_retry = False
return False, "429 retry"
return super().is_incoming_duplicate(event, add=add)

def _host_is_cooled_down(self, host):
return time.monotonic() < self._host_cooldowns.get(host, 0)

def _set_host_cooldown(self, host, seconds):
resume_time = time.monotonic() + seconds
self._host_cooldowns[host] = max(self._host_cooldowns.get(host, 0), resume_time)

def _parse_retry_after(self, response):
for k, v in response.headers.items():
if k.lower() == "retry-after":
try:
seconds = max(1, int(v))
return min(seconds, self._429_max_interval)
except ValueError:
pass
break
return self._429_default_interval

def _defer_event(self, event):
if event not in self._deferred_events:
self._deferred_events.append(event)
return True
return False

async def _flush_deferred(self):
if not self._deferred_events or self.incoming_event_queue is False:
return
still_deferred = []
flushed = 0
earliest_resume = None
now = time.monotonic()
for event in self._deferred_events:
host = str(event.host)
resume_time = self._host_cooldowns.get(host, 0)
if now >= resume_time:
event._429_retry = True
self.incoming_event_queue.put_nowait(event)
flushed += 1
else:
still_deferred.append(event)
if earliest_resume is None or resume_time < earliest_resume:
earliest_resume = resume_time
self._deferred_events = still_deferred
# prune expired cooldowns
self._host_cooldowns = {h: t for h, t in self._host_cooldowns.items() if t > now}
if flushed:
async with self.event_received:
self.event_received.notify()
if still_deferred and earliest_resume is not None:
if not self._wakeup_target or earliest_resume < self._wakeup_target:
delay = max(0.1, earliest_resume - time.monotonic())
self._wakeup_target = earliest_resume
self.helpers.create_task(self._deferred_wakeup(delay))

async def _deferred_wakeup(self, delay):
await self.helpers.sleep(delay)
self._wakeup_target = 0
await self._flush_deferred()

def _build_headers(self):
"""Build list of (name, value) header tuples from scan config."""
headers = [("User-Agent", self.scan.useragent)]
Expand Down Expand Up @@ -185,6 +264,10 @@ async def handle_batch(self, *events):

for event in events:
urls, url_hash = self.make_url_metadata(event)
host = str(event.host)
if self._host_is_cooled_down(host):
self._defer_event(event)
continue
for url in urls:
stdin[url] = event
if event.type == "OPEN_TCP_PORT":
Expand All @@ -202,6 +285,12 @@ async def handle_batch(self, *events):
paired_probe_urls[schemes["https"]] = key

if not stdin:
if self._deferred_events:
earliest = min(self._host_cooldowns.get(str(e.host), 0) for e in self._deferred_events)
if not self._wakeup_target or earliest < self._wakeup_target:
delay = max(0.1, earliest - time.monotonic())
self._wakeup_target = earliest
self.helpers.create_task(self._deferred_wakeup(delay))
return

headers = self._build_headers()
Expand Down Expand Up @@ -236,14 +325,37 @@ async def resolve_https(key, result):
await self._process_result(result, stdin[result.url])

async for result in iter_batch_results(self.client.request_batch_stream(configs, concurrency=self.threads)):
if result.success and result.response is not None and result.response.status == 429:
url = result.url
host = urlparse(url).hostname
parent_event = stdin[url]
event_key = self._incoming_dedup_hash(parent_event)
retry_count = self._429_retry_counts.get(event_key, 0)
if retry_count >= self._max_429_retries:
self.warning(f"429 from {url} after {self._max_429_retries} retries, giving up")
self._429_retry_counts.pop(event_key, None)
continue
retry_after = self._parse_retry_after(result.response)
self._set_host_cooldown(host, retry_after)
if self._defer_event(parent_event):
self._429_retry_counts[event_key] = retry_count + 1
self.verbose(
f"429 from {host} ({url}), cooling down {retry_after}s (attempt {retry_count + 1}/{self._max_429_retries})"
)
continue

# Non-429 response -- clear any retry tracking for this URL
_retry_parent = stdin.get(result.url)
if _retry_parent is not None:
self._429_retry_counts.pop(self._incoming_dedup_hash(_retry_parent), None)

key = paired_probe_urls.get(result.url)
if key is None:
# Non-paired URL — emit immediately
parent_event = stdin.get(result.url)
if parent_event is None:
# Non-paired URL -- emit immediately
if _retry_parent is None:
self.warning(f"Unable to correlate parent event for: {result.url}")
continue
await self._process_result(result, parent_event)
await self._process_result(result, _retry_parent)
continue

# Paired OPEN_TCP_PORT probe
Expand All @@ -262,6 +374,9 @@ async def resolve_https(key, result):
else:
deferred_https[key] = result

# Stream ended any leftover https had no http result, so emit unconditionally
# Stream ended -- any leftover https had no http result, so emit unconditionally
for key, result in deferred_https.items():
await self._process_result(result, stdin[result.url])

if self._deferred_events:
await self._flush_deferred()
60 changes: 60 additions & 0 deletions bbot/test/test_step_2/module_tests/test_module_http.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
from werkzeug.wrappers import Response

from .base import ModuleTestBase


Expand Down Expand Up @@ -173,3 +175,61 @@ async def setup_after_prep(self, module_test):
def check(self, module_test, events):
# Ensure we received the expected response when the cookie was present
assert [e for e in events if e.type == "URL" and "status-200" in e.tags]


class TestHTTP_429_retry(ModuleTestBase):
"""Test the module's own defer→cooldown→retry→succeed path.

blasthttp_retries=1 means blasthttp makes up to 2 wire attempts per request.
We return 429 for the first 2 requests (exhausting blasthttp's retry),
so the module's 429 handler engages and defers with a cooldown. The 3rd
request (from the module's retry after cooldown) succeeds.
"""

targets = ["http://127.0.0.1:8888"]
modules_overrides = ["http"]
config_overrides = {"web": {"429_sleep_interval": 1, "429_max_sleep_interval": 1}}

async def setup_after_prep(self, module_test):
self.request_count = 0

def handler(request):
self.request_count += 1
if self.request_count <= 2:
return Response("rate limited", status=429, headers={"Retry-After": "1"})
return Response("<html><body>OK</body></html>")

module_test.httpserver.expect_request("/").respond_with_handler(handler)

def check(self, module_test, events):
assert self.request_count >= 3, "Expected at least 3 requests (2 blasthttp attempts + 1 module retry)"
assert any(e.type == "URL" and "status-200" in e.tags for e in events), (
"Expected URL with status-200 after successful retry"
)
assert not any(e.type == "URL" and "status-429" in e.tags for e in events), (
"429 response should not be emitted as a URL event"
)


class TestHTTP_429_max_retries(ModuleTestBase):
targets = ["http://127.0.0.1:8888"]
modules_overrides = ["http"]
config_overrides = {"web": {"429_sleep_interval": 1, "429_max_sleep_interval": 1}}

async def setup_after_prep(self, module_test):
self.request_count = 0

def handler(request):
self.request_count += 1
return Response("rate limited", status=429, headers={"Retry-After": "1"})

module_test.httpserver.expect_request("/").respond_with_handler(handler)

def check(self, module_test, events):
assert self.request_count >= 2, "Expected at least one retry before giving up"
assert not any(e.type == "URL" and "status-429" in e.tags for e in events), (
"Exhausted 429 retries should not emit a URL event"
)
assert not any(e.type == "HTTP_RESPONSE" and e.data.get("status_code") == 429 for e in events), (
"Exhausted 429 retries should not emit an HTTP_RESPONSE event"
)