@@ -760,9 +760,7 @@ def start(self):
760760 self ._shutdown .clear ()
761761
762762 def run_loop () -> None :
763- loop = asyncio .new_event_loop ()
764- asyncio .set_event_loop (loop )
765- loop .run_until_complete (self ._async_run_loop ())
763+ asyncio .run (self ._async_run_loop ())
766764
767765 self ._logger .info (f"Starting gRPC worker that connects to { self ._host_address } " )
768766 self ._runLoop = Thread (target = run_loop )
@@ -1523,6 +1521,7 @@ def __init__(self,
15231521 self ._received_events : dict [str , list [str | None ]] = {}
15241522 self ._pending_events : dict [str , list [task .CancellableTask [Any ]]] = {}
15251523 self ._new_input : Any | None = None
1524+ self ._new_version : str | None = None
15261525 self ._save_events = False
15271526 self ._encoded_custom_status : str | None = None
15281527 self ._parent_trace_context : pb .TraceContext | None = None
@@ -1619,7 +1618,8 @@ def set_failed(self, ex: Exception | pb.TaskFailureDetails):
16191618 )
16201619 self ._pending_actions [action .id ] = action
16211620
1622- def set_continued_as_new (self , new_input : Any , save_events : bool ):
1621+ def set_continued_as_new (self , new_input : Any , save_events : bool ,
1622+ new_version : str | None = None ):
16231623 if self ._is_complete :
16241624 return
16251625
@@ -1633,6 +1633,7 @@ def set_continued_as_new(self, new_input: Any, save_events: bool):
16331633 # self._pending_actions.clear() # Cancel any pending actions
16341634 self ._completion_status = pb .ORCHESTRATION_STATUS_CONTINUED_AS_NEW
16351635 self ._new_input = new_input
1636+ self ._new_version = new_version
16361637 self ._save_events = save_events
16371638
16381639 def get_actions (self ) -> list [pb .OrchestratorAction ]:
@@ -1657,6 +1658,7 @@ def get_actions(self) -> list[pb.OrchestratorAction]:
16571658 result = self ._data_converter .serialize (self ._new_input ),
16581659 failure_details = None ,
16591660 carryover_events = carryover_events ,
1661+ new_version = self ._new_version ,
16601662 )
16611663 # We must return the existing tasks as well, to capture entity unlocks
16621664 current_actions .append (action )
@@ -2135,11 +2137,12 @@ def send_event(self, instance_id: str, event_name: str, *,
21352137 ),
21362138 )
21372139
2138- def continue_as_new (self , new_input : Any , * , save_events : bool = False ) -> None :
2140+ def continue_as_new (self , new_input : Any , * , save_events : bool = False ,
2141+ new_version : str | None = None ) -> None :
21392142 if self ._is_complete :
21402143 return
21412144
2142- self .set_continued_as_new (new_input , save_events )
2145+ self .set_continued_as_new (new_input , save_events , new_version )
21432146
21442147 def new_uuid (self ) -> str :
21452148 NAMESPACE_UUID : str = "9e952958-5e33-4daf-827f-2fa12937b875"
0 commit comments