Skip to main content

kaynine_runtime/
actor.rs

1//! Per-session actor: the single writer for one session. Holds the writer
2//! lease, keeps the in-memory current revision, tracks the active run, and
3//! owns the broadcast feed for the session.
4
5use 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
28/// Per-actor shutdown result: (cancelled run ids, interrupted run ids).
29pub(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    /// Mirrors the loop's `cancel_grace`: once cancelled, the run task is
84    /// guaranteed to append its terminal within this bound (+ epsilon), so
85    /// the cancel handler can wait it out instead of detaching.
86    pub(crate) cancel_grace: Duration,
87    /// In-memory steer queue shared with the run task (producer here,
88    /// drained by the run hooks when SteerApplied is persisted).
89    pub(crate) steer_queue: Arc<SteerQueue>,
90    /// Live approval waiters shared with the run task's
91    /// InteractiveApprovalHandler; only populated when the run was started
92    /// with an approval timeout.
93    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    /// Shared with the run task so hook appends and the actor agree on the
101    /// expected revision (single writer, so no conflicts in practice).
102    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    /// Cancels the background lease-renewal task; owned by the actor task.
107    pub(crate) renew_stop: CancellationToken,
108    /// Highest realtime run_seq forwarded for the active run (batch C uses
109    /// this for snapshot.last_run_seq).
110    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, &registry, &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, &registry, &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                    &registry,
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, &registry, &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                            // A lease conflict surfaces on the next append as
190                            // NotLeaseHolder/LeaseExpired; log and keep going.
191                            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, &registry, &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, &registry, &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, &registry, &handle).await;
276    });
277}
278
279/// Stops lease renewal, releases the writer lease, and deregisters the
280/// actor. Idempotent: a second call (loop-exit path after Release/Shutdown
281/// already tore down) is a no-op apart from a harmless extra lease release.
282async 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    /// Clears a finished active run (the run task appends its own terminal
299    /// event, so the actor never awaits completion synchronously).
300    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
308/// Replies with `error` to every command already queued, then deregisters.
309async 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}