@@ -765,9 +765,7 @@ def start(self):
765765 self ._shutdown .clear ()
766766
767767 def run_loop () -> None :
768- loop = asyncio .new_event_loop ()
769- asyncio .set_event_loop (loop )
770- loop .run_until_complete (self ._async_run_loop ())
768+ asyncio .run (self ._async_run_loop ())
771769
772770 self ._logger .info (f"Starting gRPC worker that connects to { self ._host_address } " )
773771 self ._runLoop = Thread (target = run_loop )
@@ -1534,6 +1532,7 @@ def __init__(self,
15341532 self ._received_events : dict [str , list [str | None ]] = {}
15351533 self ._pending_events : dict [str , list [task .CancellableTask [Any ]]] = {}
15361534 self ._new_input : Any | None = None
1535+ self ._new_version : str | None = None
15371536 self ._save_events = False
15381537 self ._encoded_custom_status : str | None = None
15391538 self ._parent_trace_context : pb .TraceContext | None = None
@@ -1642,7 +1641,8 @@ def set_failed(self, ex: Exception | pb.TaskFailureDetails):
16421641 )
16431642 self ._pending_actions [action .id ] = action
16441643
1645- def set_continued_as_new (self , new_input : Any , save_events : bool ):
1644+ def set_continued_as_new (self , new_input : Any , save_events : bool ,
1645+ new_version : str | None = None ):
16461646 if self ._is_complete :
16471647 return
16481648
@@ -1656,6 +1656,7 @@ def set_continued_as_new(self, new_input: Any, save_events: bool):
16561656 # self._pending_actions.clear() # Cancel any pending actions
16571657 self ._completion_status = pb .ORCHESTRATION_STATUS_CONTINUED_AS_NEW
16581658 self ._new_input = new_input
1659+ self ._new_version = new_version
16591660 self ._save_events = save_events
16601661
16611662 def get_actions (self ) -> list [pb .OrchestratorAction ]:
@@ -1680,6 +1681,7 @@ def get_actions(self) -> list[pb.OrchestratorAction]:
16801681 result = self ._data_converter .serialize (self ._new_input ),
16811682 failure_details = None ,
16821683 carryover_events = carryover_events ,
1684+ new_version = self ._new_version ,
16831685 )
16841686 # We must return the existing tasks as well, to capture entity unlocks
16851687 current_actions .append (action )
@@ -2158,11 +2160,12 @@ def send_event(self, instance_id: str, event_name: str, *,
21582160 ),
21592161 )
21602162
2161- def continue_as_new (self , new_input : Any , * , save_events : bool = False ) -> None :
2163+ def continue_as_new (self , new_input : Any , * , save_events : bool = False ,
2164+ new_version : str | None = None ) -> None :
21622165 if self ._is_complete :
21632166 return
21642167
2165- self .set_continued_as_new (new_input , save_events )
2168+ self .set_continued_as_new (new_input , save_events , new_version )
21662169
21672170 def new_uuid (self ) -> str :
21682171 NAMESPACE_UUID : str = "9e952958-5e33-4daf-827f-2fa12937b875"
0 commit comments