@@ -587,6 +587,53 @@ async def websocket_dashboard_handler(request):
587587 Returns:
588588 aiohttp web response.
589589 '''
590+ active_user = await learning_observer .auth .get_active_user (request )
591+ user_context = {
592+ 'user_id' : (active_user or {}).get ('user_id' ),
593+ 'user_email' : (active_user or {}).get ('email' ),
594+ 'user_role' : (active_user or {}).get ('role' ),
595+ }
596+
597+ def _log_protocol_event (event , ** extra ):
598+ '''
599+ Emit structured debug logs describing websocket activity.
600+ '''
601+ payload = {
602+ 'event' : event ,
603+ 'user_id' : user_context .get ('user_id' ),
604+ 'user_email' : user_context .get ('user_email' ),
605+ 'user_role' : user_context .get ('user_role' ),
606+ 'remote' : request .remote ,
607+ 'forwarded_for' : request .headers .get ('X-Forwarded-For' ),
608+ 'request_path' : str (request .rel_url ),
609+ }
610+ payload .update (extra )
611+ try :
612+ debug_log ('communication_protocol_event' , json .dumps (payload , sort_keys = True ))
613+ except TypeError :
614+ # Fall back to logging the raw payload if serialization fails.
615+ debug_log ('communication_protocol_event' , payload )
616+
617+ def _summarize_query (query ):
618+ '''
619+ Provide a compact description of the query for log aggregation.
620+ '''
621+ summary = {}
622+ for key , value in (query or {}).items ():
623+ if not isinstance (value , dict ):
624+ summary [key ] = {'non_dict_value' : repr (value )}
625+ continue
626+ execution_dag = value .get ('execution_dag' )
627+ summary [key ] = {
628+ 'target_exports' : value .get ('target_exports' , []),
629+ 'execution_dag_type' : type (execution_dag ).__name__ ,
630+ 'execution_dag_name' : execution_dag if isinstance (execution_dag , str ) else None ,
631+ 'kwargs' : value .get ('kwargs' , {}),
632+ }
633+ return summary
634+
635+ _log_protocol_event ('connection_opened' )
636+
590637 ws = aiohttp .web .WebSocketResponse (receive_timeout = 0.3 )
591638 await ws .prepare (request )
592639 client_query = None
@@ -661,11 +708,15 @@ async def _drive_generator(generator, dag_kwargs, target=None):
661708 try :
662709 received_params = await ws .receive_json ()
663710 client_query = received_params
711+ _log_protocol_event (
712+ 'query_received' ,
713+ query_summary = _summarize_query (client_query ),
714+ )
664715 # TODO we should validate the client_query structure
665716 except (TypeError , ValueError ):
666717 # these Errors may signal a close
667718 if (await ws .receive ()).type == aiohttp .WSMsgType .CLOSE :
668- debug_log ( "Socket closed!" )
719+ _log_protocol_event ( 'connection_closed' , reason = 'client_close_frame' )
669720 return aiohttp .web .Response ()
670721 except asyncio .exceptions .TimeoutError :
671722 # this is the normal path of the code
@@ -674,7 +725,7 @@ async def _drive_generator(generator, dag_kwargs, target=None):
674725 continue
675726
676727 if ws .closed :
677- debug_log ( "Socket closed." )
728+ _log_protocol_event ( 'connection_closed' , reason = 'websocket_closed_flag' )
678729 return aiohttp .web .Response ()
679730
680731 if client_query != previous_client_query :
0 commit comments