use super::*;
use crate::runtime::delivery::inbox::InboxOperation;
#[derive(Debug)]
pub(super) struct RemoteDeliveryState {
owner_control_inbox: SchedulerInbox,
balance_request_node: InboxNode,
}
#[derive(Debug)]
pub(crate) struct PreparedMigrationDelivery {
publication: Option<OwnedCpuRemotePublication>,
core: Option<Arc<ThreadCore>>,
source: CpuId,
target: CpuId,
placement_demand: u64,
}
impl PreparedMigrationDelivery {
pub(crate) fn prepare(
target_remote: &Arc<CpuRemote>,
core: &Arc<ThreadCore>,
source: CpuId,
target: CpuId,
) -> Result<Self, TaskError> {
let publication = target_remote
.begin_owned_publication()
.ok_or(TaskError::CpuOffline(target.as_u32()))?;
if !core.reserve_scheduler_inbox_delivery() {
return Err(TaskError::NotReady);
}
Ok(Self {
publication: Some(publication),
core: Some(Arc::clone(core)),
source,
target,
placement_demand: core.effective_placement_demand(),
})
}
pub(crate) fn refresh_placement_demand(&mut self) {
self.placement_demand = self
.core
.as_ref()
.expect("uncommitted delivery")
.effective_placement_demand();
}
pub(crate) const fn target(&self) -> CpuId {
self.target
}
pub(crate) fn commit(mut self) {
let publication = self
.publication
.take()
.expect("prepared migration must retain its target CPU lease");
let core = self
.core
.take()
.expect("prepared migration must retain its thread delivery lease");
let thread = core.id();
let pointer = Arc::into_raw(core);
let node = unsafe {
Pin::new_unchecked((*pointer).migration_node())
};
let message = InboxMessage::migration_with_payload(
thread,
self.source,
self.target,
thread.generation() as u64,
self.placement_demand,
pointer.expose_provenance(),
);
match publication.publish_owner_control(node, message) {
PublishResult::Published => {}
PublishResult::AlreadyPending => unsafe {
let retained = Arc::from_raw(pointer);
retained.cancel_scheduler_inbox_delivery();
drop(retained);
},
PublishResult::WrongKind => unsafe {
let retained = Arc::from_raw(pointer);
retained.cancel_scheduler_inbox_delivery();
drop(retained);
task_runtime::fatal_invariant(0x4d49_4701, thread.as_u64() as usize);
},
}
}
}
impl Drop for PreparedMigrationDelivery {
fn drop(&mut self) {
if let Some(core) = self.core.take() {
core.cancel_scheduler_inbox_delivery();
}
}
}
impl RemoteDeliveryState {
pub(super) const fn new() -> Self {
Self {
owner_control_inbox: SchedulerInbox::new(InboxKind::OwnerControl),
balance_request_node: InboxNode::new(InboxKind::OwnerControl),
}
}
}
impl CpuRemote {
pub(crate) fn publish_owner_control(
&self,
node: Pin<&'static InboxNode>,
message: InboxMessage,
) -> PublishResult {
let Some(remote_publication) = self.begin_publication() else {
return PublishResult::WrongKind;
};
remote_publication.publish_owner_control(node, message)
}
pub(super) fn publish_owner_control_owned(
&self,
node: Pin<&'static InboxNode>,
message: InboxMessage,
) -> PublishResult {
self.publish_owner_control_owned_observed(node, message).0
}
fn publish_owner_control_owned_observed(
&self,
node: Pin<&'static InboxNode>,
message: InboxMessage,
) -> (
PublishResult,
Option<super::scheduler::SchedulerRequestPublication>,
) {
let _irq = IrqScope::enter();
let _idle_pull_work = self.begin_idle_pull_work();
let migration = message.operation() == InboxOperation::Migration;
if migration {
self.reserve_incoming_migration(message.placement_demand());
}
let (result, head_became_non_empty) = self
.delivery
.owner_control_inbox
.publish_with_head_transition(node, message);
if migration && result != PublishResult::Published {
self.release_incoming_migration_demand(message.placement_demand());
}
let scheduler_publication = if result == PublishResult::Published && head_became_non_empty {
let publication = self.publish_owner_inbox_head_owned();
self.deliver_scheduler_work_owned(publication);
Some(publication)
} else {
None
};
(result, scheduler_publication)
}
pub(crate) fn balance_request_node(&self) -> Pin<&'static InboxNode> {
let node = &self.delivery.balance_request_node as *const InboxNode;
unsafe { Pin::new_unchecked(&*node) }
}
pub(crate) fn owner_control_inbox(&self) -> &SchedulerInbox {
&self.delivery.owner_control_inbox
}
pub(crate) fn has_remote_work(&self) -> bool {
self.delivery.owner_control_inbox.has_pending()
}
}