@@ -37,16 +37,31 @@ class STATES:
3737COMPOSITEKEY_NS = '\x00 '
3838EMPTY_KEY_SUBSTITUTE = '\x01 '
3939
40- STATE = STATES .CREATED
41-
42-
4340class Handler :
4441 def __init__ (self , cc_id : str , cc : Chaincode ) -> None :
4542 self .chaincode_id = cc_pb2 .ChaincodeID ()
4643 self .chaincode_id .name = cc_id
4744 self .chaincode = cc
4845 self .msg_queue_handler = None
4946 self .context = None
47+ self .state = STATES .CREATED
48+ self ._pending_tasks = set ()
49+
50+ def _track_task (self , coro ):
51+ """Schedule message handling while surfacing task failures in logs."""
52+ task = asyncio .create_task (coro )
53+ self ._pending_tasks .add (task )
54+
55+ def _done_callback (done_task ):
56+ self ._pending_tasks .discard (done_task )
57+ try :
58+ exc = done_task .exception ()
59+ except asyncio .CancelledError :
60+ return
61+ if exc is not None :
62+ LOGGER .exception ('Unhandled exception while processing peer message' , exc_info = exc )
63+
64+ task .add_done_callback (_done_callback )
5065
5166 async def handle_stub_interaction (self , msg , action = "Invoke" ):
5267 """handle_message calls the Init | Invoke function of the associated chaincode."""
@@ -118,65 +133,51 @@ async def handle_message_ready(self, msg):
118133 await self .handle_stub_interaction (msg , "Invoke" )
119134 return
120135 else :
121- self .context .write (new_error_msg (msg , STATE ))
136+ await self .context .write (new_error_msg (msg , self . state ))
122137
123- def handle_message_established (self , msg ):
138+ async def handle_message_established (self , msg ):
124139 """
125140 handle_message_established handles messages received from the peer when the handler is in the "established" state.
126141 """
127- global STATE
128142 if msg .type != ccshim_pb2 .ChaincodeMessage .READY :
129143 LOGGER .error (f'Chaincode is in "established" state, can only process messages of type "ready", '
130144 f'but received "{ msg .type } "' )
131- # write is an async coroutine on the context
132- try :
133- return asyncio .create_task (self .context .write (new_error_msg (msg , STATE )))
134- except Exception :
135- return
145+ await self .context .write (new_error_msg (msg , self .state ))
136146 else :
137147 LOGGER .info ('Successfully established communication with peer node. State transferred to "ready"' )
138- STATE = STATES .READY
148+ self . state = STATES .READY
139149
140- def handle_message_created (self , msg ):
150+ async def handle_message_created (self , msg ):
141151 """handle_message_created handles messages received from the peer when the handler is in the "created" state."""
142- global STATE
143152 if msg .type != ccshim_pb2 .ChaincodeMessage .REGISTERED :
144153 LOGGER .error (f'Chaincode is in "created" state, can only process messages of type "registered", '
145154 f'but received "{ msg .type } "' )
146- try :
147- return asyncio .create_task (self .context .write (new_error_msg (msg , STATE )))
148- except Exception :
149- return
155+ await self .context .write (new_error_msg (msg , self .state ))
150156 else :
151157 LOGGER .info ('Successfully registered with peer node. State transferred to "established"' )
152- STATE = STATES .ESTABLISHED
158+ self . state = STATES .ESTABLISHED
153159
154160 async def handle_message (self , msg : ccshim_pb2 .ChaincodeMessage ):
155161 """handle_message message handles loop for shim side of chaincode/peer stream."""
156162 LOGGER .warning ('-->> Look out!' )
157- global STATE
158163
159164 # TODO: ?
160165 if msg .type == ccshim_pb2 .ChaincodeMessage .KEEPALIVE :
161166 LOGGER .info ('-| KEEPALIVE' )
162167 return
163168
164- if STATE == STATES .READY :
169+ if self . state == STATES .READY :
165170 await self .handle_message_ready (msg )
166- elif STATE == STATES .ESTABLISHED :
167- self .handle_message_established (msg )
168- elif STATE == STATES .CREATED :
169- self .handle_message_created (msg )
171+ elif self . state == STATES .ESTABLISHED :
172+ await self .handle_message_established (msg )
173+ elif self . state == STATES .CREATED :
174+ await self .handle_message_created (msg )
170175 else :
171- try :
172- asyncio .create_task (self .context .write (new_error_msg (msg , STATE )))
173- except Exception :
174- LOGGER .exception ('Failed to write error message to context' )
176+ await self .context .write (new_error_msg (msg , self .state ))
175177
176178 async def chat_with_peer (self , stream : AsyncIterable [ccshim_pb2 .ChaincodeMessage ], context : grpc .aio .ServicerContext ):
177179 """chat stream for peer-chaincode interactions post connection"""
178- global STATE
179- STATE = STATES .CREATED
180+ self .state = STATES .CREATED
180181
181182 self .context = context
182183 self .msg_queue_handler = MsgQueueHandler (self )
@@ -197,13 +198,18 @@ async def chat_with_peer(self, stream: AsyncIterable[ccshim_pb2.ChaincodeMessage
197198 return ccshim_pb2 .ChaincodeMessage (
198199 type = ccshim_pb2 .ChaincodeMessage .ERROR , payload = err_str .encode (encoding = 'utf-8' ))
199200 else :
200- asyncio .create_task (self .handle_message (receive_message ))
201+ # Keep handling asynchronous so the stream can continue receiving
202+ # response frames needed by in-flight request futures.
203+ self ._track_task (self .handle_message (receive_message ))
201204
202205 LOGGER .info (f'->>>> proposal { receive_message .proposal } ' )
203206 LOGGER .info (f'->>>> payload { receive_message .payload } ' )
204207 LOGGER .info (f'->>>> channel ID { receive_message .channel_id } ' )
205208 LOGGER .info (f'->>>> Tx ID { receive_message .txid } ' )
206209
210+ if self ._pending_tasks :
211+ await asyncio .gather (* self ._pending_tasks , return_exceptions = True )
212+
207213 async def handle_get_state (self , collection , key , channel_id , tx_id ):
208214 msg_pb = ccshim_pb2 .GetState ()
209215 msg_pb .key = key
0 commit comments