@@ -384,7 +384,16 @@ def get_messages(self, batch_size: int) -> list[ScheduleMessageItem]:
384384 if len (self .message_pack_cache ) == 0 :
385385 return []
386386 else :
387- return self .message_pack_cache .popleft ()
387+ batch = self .message_pack_cache .popleft ()
388+ if len (batch ) >= batch_size :
389+ logger .debug (
390+ "[REDIS_QUEUE] Dequeued batch. batch_size=%s requested_batch_size=%s cache_packs_remaining=%s stream_count=%s" ,
391+ len (batch ),
392+ batch_size ,
393+ len (self .message_pack_cache ),
394+ len (self .get_stream_keys ()),
395+ )
396+ return batch
388397
389398 def _ensure_consumer_group (self , stream_key ) -> None :
390399 """Ensure the consumer group exists for the stream."""
@@ -449,9 +458,13 @@ def put(
449458 message_id = self ._redis_conn .xadd (
450459 stream_key , message_data , maxlen = self .max_len , approximate = True
451460 )
452-
453- logger .info (
454- f"Added message { message_id } to Redis stream: { message .label } - { message .content [:100 ]} ..."
461+ logger .debug (
462+ "[REDIS_QUEUE] Enqueued message. message_id=%s stream=%s label=%s item_id=%s stream_cache_size=%s" ,
463+ message_id ,
464+ stream_key ,
465+ message .label ,
466+ message .item_id ,
467+ len (self ._stream_keys_cache ),
455468 )
456469
457470 except Exception as e :
@@ -494,7 +507,11 @@ def ack_message(
494507 # Optionally delete the message from the stream to keep it clean
495508 try :
496509 self ._redis_conn .xdel (stream_key , redis_message_id )
497- logger .info (f"Successfully delete acknowledged message { redis_message_id } " )
510+ logger .debug (
511+ "[REDIS_QUEUE] Ack/delete message. redis_message_id=%s stream=%s" ,
512+ redis_message_id ,
513+ stream_key ,
514+ )
498515 except Exception as e :
499516 logger .warning (f"Failed to delete acknowledged message { redis_message_id } : { e } " )
500517
@@ -989,7 +1006,7 @@ def show_task_status(self, stream_key_prefix: str | None = None) -> dict[str, di
9891006 )
9901007 stream_keys = self .get_stream_keys (stream_key_prefix = effective_prefix )
9911008 if not stream_keys :
992- logger .info (f"No Redis streams found for the configured prefix: { effective_prefix } " )
1009+ logger .debug (f"No Redis streams found for the configured prefix: { effective_prefix } " )
9931010 return {}
9941011
9951012 grouped : dict [str , dict [str , int ]] = {}
@@ -1157,7 +1174,7 @@ def connect(self) -> None:
11571174 self ._redis_conn .ping ()
11581175 self ._is_connected = True
11591176 self ._check_xautoclaim_support ()
1160- logger .debug ("Redis connection established successfully" )
1177+ logger .info ("Redis connection established successfully" )
11611178 # Start stream keys refresher when connected
11621179 self ._start_stream_keys_refresh_thread ()
11631180 except Exception as e :
@@ -1174,7 +1191,7 @@ def disconnect(self) -> None:
11741191 self ._stop_stream_keys_refresh_thread ()
11751192 if self ._is_listening :
11761193 self .stop_listening ()
1177- logger .debug ("Disconnected from Redis" )
1194+ logger .info ("Disconnected from Redis" )
11781195
11791196 def __enter__ (self ):
11801197 """Context manager entry."""
@@ -1379,7 +1396,7 @@ def _update_stream_cache_with_log(
13791396 self ._stream_keys_cache = active_stream_keys
13801397 self ._stream_keys_last_refresh = time .time ()
13811398 cache_count = len (self ._stream_keys_cache )
1382- logger .info (
1399+ logger .debug (
13831400 f"Refreshed stream keys cache: { cache_count } active keys, "
13841401 f"{ deleted_count } deleted, { len (candidate_keys )} candidates examined."
13851402 )
0 commit comments