Commit d92ef7e7 authored by Vitaly Lipatov's avatar Vitaly Lipatov

redis: avoid stream consumer busy loop

parent 7f2a9fb5
...@@ -533,6 +533,7 @@ def connect_redis(): ...@@ -533,6 +533,7 @@ def connect_redis():
redis_client = redis.Redis( redis_client = redis.Redis(
host=redis_server, host=redis_server,
socket_connect_timeout=5, socket_connect_timeout=5,
socket_timeout=10,
socket_keepalive=True, socket_keepalive=True,
health_check_interval=30, health_check_interval=30,
**redis_options, **redis_options,
...@@ -722,17 +723,28 @@ def process_stream_entry(fields): ...@@ -722,17 +723,28 @@ def process_stream_entry(fields):
return success return success
claim_interval_seconds = 60
next_claim_at = 0
read_pending = True
while True: while True:
try: try:
claimed = r.xautoclaim(redis_commands_stream, redis_commands_group, entries = []
redis_commands_consumer, min_idle_time=60000, if read_pending:
start_id='0-0', count=10) entries = r.xreadgroup(redis_commands_group, redis_commands_consumer,
entries = [(redis_commands_stream, claimed[1])] if claimed[1] else []
if not entries:
pending = r.xreadgroup(redis_commands_group, redis_commands_consumer,
{redis_commands_stream: '0'}, count=10) {redis_commands_stream: '0'}, count=10)
entries = pending or r.xreadgroup(redis_commands_group, redis_commands_consumer, read_pending = bool(entries)
{redis_commands_stream: '>'}, count=10, block=5000) if not entries and time.monotonic() >= next_claim_at:
claimed = r.xautoclaim(redis_commands_stream, redis_commands_group,
redis_commands_consumer, min_idle_time=60000,
start_id='0-0', count=10)
entries = [(redis_commands_stream, claimed[1])] if claimed[1] else []
next_claim_at = time.monotonic() + claim_interval_seconds
if not entries:
entries = r.xreadgroup(redis_commands_group, redis_commands_consumer,
{redis_commands_stream: '>'}, count=10, block=5000)
if not entries:
time.sleep(0.1)
for stream, messages in entries: for stream, messages in entries:
for message_id, fields in messages: for message_id, fields in messages:
if process_stream_entry(fields): if process_stream_entry(fields):
...@@ -741,3 +753,5 @@ while True: ...@@ -741,3 +753,5 @@ while True:
log_redis_error("Redis Stream read failed: " + str(error) + "; reconnecting") log_redis_error("Redis Stream read failed: " + str(error) + "; reconnecting")
r = connect_redis() r = connect_redis()
auto_mgr.r = r auto_mgr.r = r
next_claim_at = 0
read_pending = True
...@@ -70,7 +70,10 @@ rg -Fq "'[' + ban_server_ipv6 + ']:81'" gateway/usr/share/eterban/eterban_switch ...@@ -70,7 +70,10 @@ rg -Fq "'[' + ban_server_ipv6 + ']:81'" gateway/usr/share/eterban/eterban_switch
rg -Fq "'80,81,443'" gateway/usr/share/eterban/eterban_switcher.py && \ rg -Fq "'80,81,443'" gateway/usr/share/eterban/eterban_switcher.py && \
rg -Fq "ipset_firehol, 'src', '-p', 'tcp', '-j', 'DNAT'" gateway/usr/share/eterban/eterban_switcher.py && \ rg -Fq "ipset_firehol, 'src', '-p', 'tcp', '-j', 'DNAT'" gateway/usr/share/eterban/eterban_switcher.py && \
rg -Fq "ipset_eterban_1, 'src', '-p', 'tcp', '-j', 'DNAT'" gateway/usr/share/eterban/eterban_switcher.py && \ rg -Fq "ipset_eterban_1, 'src', '-p', 'tcp', '-j', 'DNAT'" gateway/usr/share/eterban/eterban_switcher.py && \
rg -Fq "ipset_eterban_1_ipv6, 'src', '-p', 'tcp', '-j', 'DNAT'" gateway/usr/share/eterban/eterban_switcher.py || { rg -Fq "ipset_eterban_1_ipv6, 'src', '-p', 'tcp', '-j', 'DNAT'" gateway/usr/share/eterban/eterban_switcher.py && \
rg -Fq 'claim_interval_seconds = 60' gateway/usr/share/eterban/eterban_switcher.py && \
rg -Fq 'time.monotonic() >= next_claim_at' gateway/usr/share/eterban/eterban_switcher.py && \
rg -Fq 'socket_timeout=10' gateway/usr/share/eterban/eterban_switcher.py || {
echo 'external IPv4/IPv6 ban redirects must target public port 81' >&2 echo 'external IPv4/IPv6 ban redirects must target public port 81' >&2
exit 1 exit 1
} }
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment