use super::*;
use super::env::{missing_capability, reject_protocol_command_in_pane, served_socket, session_env};
use super::outbound::send_frame;
use super::projection::{note_verdict, sync_session};
use super::state::{
DispatchInner, DispatchState, SessionSlot, live_sessions, note_beat, rebase_generation,
render_tokens, slot_key_serving_task,
};
use super::transport::{is_revived_connection, note_binding_locked, record_revived_connection};
pub fn dispatch(state: &DispatchState, envelope: &Envelope) -> Result<SessionRef> {
let causality = envelope
.causality
.as_ref()
.context("task envelope missing causality.task")?;
let task_id = causality.task.clone();
let mut inner = state.inner.lock();
if let Some(session) = slot_key_serving_task(&inner, &task_id)
.and_then(|key| inner.sessions.get_mut(&key))
.map(|slot| {
if slot.payload.is_none() {
slot.payload = Some(envelope.clone());
slot.causality = causality.clone();
}
slot.session.clone()
})
{
inner.stall.note_assigned(&task_id, Instant::now());
return Ok(session);
}
if live_sessions(&inner) >= inner.max_sessions as usize {
return Err(anyhow!("max_sessions reached"));
}
let session_id = task_id.clone();
let command = render_tokens(&inner.command, &session_id, &task_id);
reject_protocol_command_in_pane(inner.backend.name(), &command)?;
let env = session_env(
&inner.role,
&session_id,
&task_id,
&inner.relay_required,
inner.relay_count,
&inner.topology,
&served_socket(&inner.workspace),
);
let session = inner.backend.spawn(SpawnSpec {
cwd: inner.workspace.clone(),
task_id: task_id.clone(),
command,
env,
focus: None,
placement: None,
rename: None,
})?;
inner.bridge.track_live(session.clone());
inner.store.open_task(causality, task_cause(causality))?;
rebase_born_session(&inner, &task_id);
feed_created(&inner.bridge, &inner.store, &task_id)?;
feed_dispatched(&inner.bridge, &inner.store, &task_id);
let dropped_at = (!inner.backend.self_driven()).then(Instant::now);
let last_beat = Some(Instant::now());
inner.sessions.insert(
session_id,
SessionSlot {
session: session.clone(),
task_id: Some(task_id.clone()),
ready: false,
payload: Some(envelope.clone()),
msg_id: None,
origin: Some(envelope.from.clone()),
causality: causality.clone(),
dropped_at,
last_beat,
read_only: false,
},
);
inner.stall.note_assigned(&task_id, Instant::now());
Ok(session)
}
fn task_cause(causality: &Causality) -> &'static str {
if causality.parent_task.is_some() {
"relay"
} else {
"root"
}
}
fn rebase_born_session(inner: &DispatchInner, task_id: &str) {
let verdict = rebase_generation(inner, task_id, |stored| {
Observation::initial(stored.isolate_after, stored.terminate_after)
});
match verdict {
Ok(Some(Verdict::Applied(next))) => tracing::info!(
task = %task_id,
generation = next.version.generation,
"the row a re-dispatched session is born onto was rebased onto a new generation"
),
Ok(Some(verdict)) => tracing::warn!(
task = %task_id,
?verdict,
"the row of a re-dispatched session was not rebased"
),
Ok(None) => {}
Err(error) => tracing::warn!(
task = %task_id,
error = %error,
"the row of a re-dispatched session was not rebased"
),
}
}
pub struct ReadyNotice {
pub task_id: String,
pub session_id: String,
pub generation: u64,
pub io: Option<AdapterIo>,
pub capabilities: Vec<Capability>,
}
pub async fn on_ready(state: &DispatchState, notice: ReadyNotice, prose: &str) -> Result<()> {
let ReadyNotice {
task_id,
session_id,
generation,
io,
capabilities,
} = notice;
let (payload, target, session, backend, version) = {
let mut inner = state.inner.lock();
let backend = Arc::clone(&inner.backend);
if let Some(connection) = io.as_ref() {
if is_revived_connection(&inner, connection) {
return Ok(());
}
if !note_binding_locked(&mut inner, &session_id, connection) {
record_revived_connection(
&mut inner,
&session_id,
connection.clone(),
capabilities.clone(),
);
return Ok(());
}
}
let slot = inner
.sessions
.values_mut()
.find(|slot| {
slot.session.task_id == task_id
|| slot
.session
.backend_ref
.get("id")
.and_then(|value| value.as_str())
== Some(session_id.as_str())
})
.ok_or_else(|| anyhow!("unknown session for {task_id}"))?;
if slot.read_only {
return Ok(());
}
let Some(payload) = slot.payload.take() else {
return Ok(());
};
slot.origin = Some(payload.from.clone());
slot.ready = true;
let session = slot.session.clone();
let verdict = feed_ready(&inner.bridge, &inner.store, &task_id)?;
if matches!(verdict, Verdict::Applied(_)) {
note_beat(&mut inner, &task_id, Instant::now());
}
let version = note_verdict(&verdict, &task_id).unwrap_or(Version::new(generation, 0));
(payload, io, session, backend, version)
};
send_frame(
state,
ClientOp::Report(Report::Ready {
task_id: task_id.clone(),
session_id: session_id.clone(),
generation: version.generation,
seq: version.seq,
cluster_ref: None,
}),
)
.await?;
sync_session(state, &task_id).await?;
let text = payload.body.text.clone().unwrap_or_default();
match (backend.self_driven(), target) {
(true, None) => backend.deliver(&session, &task_id, &text),
(true, Some(_)) => Err(anyhow!(
"self-driven session {session_id} unexpectedly has an adapter transport"
)),
(false, Some(target)) if capabilities.contains(&Capability::Inject) => {
let assign = AssignArgs {
envelope: Box::new(payload),
prose: prose.to_string(),
task_id,
generation,
parent: None,
};
target
.notify(AdapterMsg::Host(HostOp::Assign(assign)))
.await
.map_err(|e| anyhow!(e))
}
(false, Some(target)) => target
.notify(AdapterMsg::Host(HostOp::ConfigGet(
onlyne_proto::ConfigGetArgs {
key: format!("stdin:{text}"),
},
)))
.await
.map_err(|e| anyhow!(e)),
(false, None) => Err(anyhow!(
"adapter-backed session {session_id} reported ready without a transport"
)),
}
}
impl DispatchState {
pub async fn hand_session(
&self,
task_id: &str,
io: AdapterIo,
capabilities: Vec<Capability>,
) -> Result<()> {
let prose = self.role_prose();
on_ready(
self,
ReadyNotice {
task_id: task_id.to_string(),
session_id: task_id.to_string(),
generation: self.session_generation(task_id).unwrap_or(1),
io: Some(io),
capabilities,
},
&prose,
)
.await
}
pub async fn hand_staged(&self, session_id: &str) -> Result<bool> {
let self_driven = self.inner.lock().backend.self_driven();
if self_driven {
let prose = self.role_prose();
on_ready(
self,
ReadyNotice {
task_id: session_id.to_string(),
session_id: session_id.to_string(),
generation: self.session_generation(session_id).unwrap_or(1),
io: None,
capabilities: Vec::new(),
},
&prose,
)
.await?;
return Ok(true);
}
let transport = self
.session_transport(session_id)
.or_else(|| self.claim_parked_transport(session_id));
let Some((io, capabilities)) = transport else {
return Ok(false);
};
self.hand_session(session_id, io, capabilities).await?;
Ok(true)
}
pub async fn inject_note(&self, envelope: &Envelope) -> bool {
let Some((task_id, session_id)) = self.ready_session() else {
return false;
};
let Some((io, capabilities)) = self.session_transport(&session_id) else {
return false;
};
if missing_capability(&capabilities, Capability::Inject) {
tracing::debug!(task = %task_id, "plugin takes no message mid-task");
return false;
}
let generation = self.session_generation(&task_id).unwrap_or(1);
let assign = AssignArgs {
envelope: Box::new(envelope.clone()),
prose: self.role_prose(),
task_id,
generation,
parent: None,
};
io.notify(AdapterMsg::Host(HostOp::Assign(assign)))
.await
.is_ok()
}
fn ready_session(&self) -> Option<(String, String)> {
let inner = self.inner.lock();
inner
.sessions
.iter()
.find(|(_, slot)| slot.ready && slot.task_id.is_some() && !slot.read_only)
.map(|(key, slot)| {
(
slot.task_id.clone().unwrap_or_else(|| key.clone()),
key.clone(),
)
})
}
}