1use crate::approval::ApprovalShared;
6use crate::run::{self, build_snapshot, SteerQueue};
7use crate::service::{
8 ActiveRunInfo, ApprovalOutcome, ApprovalResolution, CancelOutcome, CancelRequest,
9 ReleaseOutcome, RunAccepted, SessionSnapshot, StartRunRequest, SteerAccepted, SteerRequest,
10 UpdateSessionRequest,
11};
12use kaynine_core::error::KaynineError;
13use kaynine_core::event::{EventEnvelope, RealtimeEvent};
14use kaynine_core::ids::{ModelId, RunId, SessionId};
15use kaynine_core::store::{LeaseOwner, SessionRecord, SessionStore};
16use std::collections::HashMap;
17use std::sync::atomic::{AtomicU64, Ordering};
18use std::sync::{Arc, Mutex};
19use std::time::Duration;
20use tokio::sync::{broadcast, mpsc, oneshot};
21use tokio::task::JoinHandle;
22use tokio_util::sync::CancellationToken;
23
24const LEASE_TTL_SECS: i64 = 60;
25const RENEW_INTERVAL_SECS: u64 = 20;
26const BROADCAST_CAPACITY: usize = 256;
27
28pub(crate) type ActorShutdownOutcome = (Vec<RunId>, Vec<RunId>);
30
31pub(crate) enum ActorCommand {
32 StartRun {
33 request: Box<StartRunRequest>,
34 reply: oneshot::Sender<Result<RunAccepted, KaynineError>>,
35 },
36 Cancel {
37 request: CancelRequest,
38 reply: oneshot::Sender<Result<CancelOutcome, KaynineError>>,
39 },
40 Snapshot {
41 reply: oneshot::Sender<Result<SessionSnapshot, KaynineError>>,
42 },
43 UpdateSession {
44 request: UpdateSessionRequest,
45 reply: oneshot::Sender<Result<SessionRecord, KaynineError>>,
46 },
47 Steer {
48 request: SteerRequest,
49 reply: oneshot::Sender<Result<SteerAccepted, KaynineError>>,
50 },
51 ResolveApproval {
52 request: ApprovalResolution,
53 reply: oneshot::Sender<Result<ApprovalOutcome, KaynineError>>,
54 },
55 Shutdown {
56 grace: Duration,
57 reply: oneshot::Sender<Result<ActorShutdownOutcome, KaynineError>>,
58 },
59 Release {
60 reply: oneshot::Sender<Result<ReleaseOutcome, KaynineError>>,
61 },
62 Subscribe {
63 reply: oneshot::Sender<
64 Result<broadcast::Receiver<EventEnvelope<RealtimeEvent>>, KaynineError>,
65 >,
66 },
67}
68
69#[derive(Clone)]
70pub(crate) struct ActorHandle {
71 pub(crate) actor_id: u64,
72 pub(crate) tx: mpsc::Sender<ActorCommand>,
73}
74
75pub(crate) type ActorRegistry = Arc<Mutex<HashMap<SessionId, ActorHandle>>>;
76
77pub(crate) struct ActiveRun {
78 pub(crate) run_id: RunId,
79 pub(crate) branch_id: kaynine_core::ids::BranchId,
80 pub(crate) model: ModelId,
81 pub(crate) cancel: CancellationToken,
82 pub(crate) join: JoinHandle<()>,
83 pub(crate) cancel_grace: Duration,
87 pub(crate) steer_queue: Arc<SteerQueue>,
90 pub(crate) approval_shared: Arc<ApprovalShared>,
94}
95
96pub(crate) struct ActorState {
97 pub(crate) session_id: SessionId,
98 pub(crate) store: Arc<dyn SessionStore>,
99 pub(crate) owner: LeaseOwner,
100 pub(crate) revision: Arc<tokio::sync::Mutex<u64>>,
103 pub(crate) active: Option<ActiveRun>,
104 pub(crate) event_tx: broadcast::Sender<EventEnvelope<RealtimeEvent>>,
105 pub(crate) shutting: bool,
106 pub(crate) renew_stop: CancellationToken,
108 pub(crate) run_seq: Arc<AtomicU64>,
111}
112
113pub(crate) fn spawn(
114 store: Arc<dyn SessionStore>,
115 registry: ActorRegistry,
116 session_id: SessionId,
117 handle: ActorHandle,
118 mut rx: mpsc::Receiver<ActorCommand>,
119) {
120 tokio::spawn(async move {
121 let owner_id = format!("actor-{}", uuid::Uuid::new_v4());
122 let owner = match store
123 .acquire_lease(&session_id, &owner_id, LEASE_TTL_SECS)
124 .await
125 {
126 Ok(owner) => owner,
127 Err(error) => {
128 fail_pending(rx, error, ®istry, &session_id, &handle).await;
129 return;
130 }
131 };
132
133 match store.recover_session(&session_id, &owner).await {
134 Ok(report) => {
135 tracing::info!(
136 session_id = %session_id,
137 new_revision = report.new_revision,
138 interrupted = report.interrupted_runs.len(),
139 "session actor recovered session"
140 );
141 }
142 Err(error) => {
143 fail_pending(rx, error, ®istry, &session_id, &handle).await;
144 let _ = store.release_lease(&session_id, &owner).await;
145 return;
146 }
147 }
148
149 let revision = match store.get_session(&session_id).await {
150 Ok(Some(session)) => session.current_revision,
151 Ok(None) => {
152 fail_pending(
153 rx,
154 KaynineError::SessionNotFound,
155 ®istry,
156 &session_id,
157 &handle,
158 )
159 .await;
160 let _ = store.release_lease(&session_id, &owner).await;
161 return;
162 }
163 Err(error) => {
164 fail_pending(rx, error, ®istry, &session_id, &handle).await;
165 let _ = store.release_lease(&session_id, &owner).await;
166 return;
167 }
168 };
169
170 let (event_tx, _) = broadcast::channel(BROADCAST_CAPACITY);
171 let renew_stop = CancellationToken::new();
172 {
173 let store = store.clone();
174 let session_id = session_id.clone();
175 let owner = owner.clone();
176 let stop = renew_stop.clone();
177 tokio::spawn(async move {
178 let mut ticker =
179 tokio::time::interval(std::time::Duration::from_secs(RENEW_INTERVAL_SECS));
180 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
181 loop {
182 tokio::select! {
183 _ = stop.cancelled() => return,
184 _ = ticker.tick() => {}
185 }
186 match store.renew_lease(&session_id, &owner, LEASE_TTL_SECS).await {
187 Ok(true) => {}
188 Ok(false) => {
189 tracing::warn!(session_id = %session_id, "writer lease renewal rejected");
192 }
193 Err(error) => {
194 tracing::warn!(session_id = %session_id, ?error, "writer lease renewal failed");
195 }
196 }
197 }
198 });
199 }
200
201 let mut state = ActorState {
202 session_id: session_id.clone(),
203 store: store.clone(),
204 owner: owner.clone(),
205 revision: Arc::new(tokio::sync::Mutex::new(revision)),
206 active: None,
207 event_tx,
208 shutting: false,
209 renew_stop: renew_stop.clone(),
210 run_seq: Arc::new(AtomicU64::new(0)),
211 };
212
213 let mut exiting = false;
214 while let Some(command) = rx.recv().await {
215 match command {
216 ActorCommand::StartRun { request, reply } => {
217 run::handle_start_run(&mut state, request, reply).await;
218 }
219 ActorCommand::Cancel { request, reply } => {
220 run::handle_cancel(&mut state, request, reply).await;
221 }
222 ActorCommand::Snapshot { reply } => {
223 let active = state.take_finished_active().map(|a| ActiveRunInfo {
224 run_id: a.run_id.clone(),
225 branch_id: a.branch_id.clone(),
226 model: a.model.clone(),
227 });
228 let revision = *state.revision.lock().await;
229 let run_seq = (state.run_seq.load(Ordering::SeqCst) > 0)
230 .then(|| state.run_seq.load(Ordering::SeqCst));
231 let snapshot = build_snapshot(
232 state.store.as_ref(),
233 &state.session_id,
234 active,
235 Some(revision),
236 run_seq,
237 )
238 .await;
239 let _ = reply.send(snapshot);
240 }
241 ActorCommand::UpdateSession { request, reply } => {
242 run::handle_update_session(&mut state, request, reply).await;
243 }
244 ActorCommand::Steer { request, reply } => {
245 run::handle_steer(&mut state, request, reply).await;
246 }
247 ActorCommand::ResolveApproval { request, reply } => {
248 run::handle_resolve_approval(&mut state, request, reply).await;
249 }
250 ActorCommand::Shutdown { grace, reply } => {
251 state.shutting = true;
252 let outcome = run::handle_shutdown_actor(&mut state, grace).await;
253 teardown(&mut state, ®istry, &handle).await;
254 let _ = reply.send(Ok(outcome));
255 exiting = true;
256 }
257 ActorCommand::Release { reply } => {
258 if state.active.as_ref().is_some_and(|a| !a.join.is_finished()) {
259 let _ = reply.send(Ok(ReleaseOutcome::RunAlreadyActive));
260 continue;
261 }
262 teardown(&mut state, ®istry, &handle).await;
263 let _ = reply.send(Ok(ReleaseOutcome::Released));
264 exiting = true;
265 }
266 ActorCommand::Subscribe { reply } => {
267 let _ = reply.send(Ok(state.event_tx.subscribe()));
268 }
269 }
270 if exiting {
271 break;
272 }
273 }
274
275 teardown(&mut state, ®istry, &handle).await;
276 });
277}
278
279async fn teardown(state: &mut ActorState, registry: &ActorRegistry, handle: &ActorHandle) {
283 state.renew_stop.cancel();
284 let _ = state
285 .store
286 .release_lease(&state.session_id, &state.owner)
287 .await;
288 let mut actors = registry.lock().expect("actor registry mutex poisoned");
289 if actors
290 .get(&state.session_id)
291 .is_some_and(|current| current.actor_id == handle.actor_id)
292 {
293 actors.remove(&state.session_id);
294 }
295}
296
297impl ActorState {
298 pub(crate) fn take_finished_active(&mut self) -> Option<&ActiveRun> {
301 if self.active.as_ref().is_some_and(|a| a.join.is_finished()) {
302 self.active = None;
303 }
304 self.active.as_ref()
305 }
306}
307
308async fn fail_pending(
310 rx: mpsc::Receiver<ActorCommand>,
311 error: KaynineError,
312 registry: &ActorRegistry,
313 session_id: &SessionId,
314 handle: &ActorHandle,
315) {
316 let mut rx = rx;
317 while let Ok(command) = rx.try_recv() {
318 match command {
319 ActorCommand::StartRun { reply, .. } => {
320 let _ = reply.send(Err(error.clone()));
321 }
322 ActorCommand::Cancel { reply, .. } => {
323 let _ = reply.send(Err(error.clone()));
324 }
325 ActorCommand::Snapshot { reply } => {
326 let _ = reply.send(Err(error.clone()));
327 }
328 ActorCommand::UpdateSession { reply, .. } => {
329 let _ = reply.send(Err(error.clone()));
330 }
331 ActorCommand::Steer { reply, .. } => {
332 let _ = reply.send(Err(error.clone()));
333 }
334 ActorCommand::ResolveApproval { reply, .. } => {
335 let _ = reply.send(Err(error.clone()));
336 }
337 ActorCommand::Shutdown { reply, .. } => {
338 let _ = reply.send(Err(error.clone()));
339 }
340 ActorCommand::Release { reply } => {
341 let _ = reply.send(Ok(ReleaseOutcome::NotFound));
342 }
343 ActorCommand::Subscribe { reply } => {
344 let _ = reply.send(Err(error.clone()));
345 }
346 }
347 }
348 let mut actors = registry.lock().expect("actor registry mutex poisoned");
349 if actors
350 .get(session_id)
351 .is_some_and(|current| current.actor_id == handle.actor_id)
352 {
353 actors.remove(session_id);
354 }
355}