use super::*;
pub(super) struct PendingSourceCapture {
connection: String,
request: Request,
pub(super) deadline: Instant,
worker: std::thread::JoinHandle<Result<super::super::application::PreparedSource, Code>>,
}
impl PendingSourceCapture {
pub(super) fn matches_request(&self, connection: &str, request: &Request) -> bool {
self.connection == connection && self.request.request_id == request.request_id
}
pub(super) fn session(&self) -> Option<&str> {
self.request.session_id.as_deref()
}
pub(super) fn join(self) {
let _ = self.worker.join();
}
}
impl Coordinator {
pub(super) fn prepare_source_registration(
&mut self,
connection: &str,
request: &Request,
now: Instant,
) -> Option<Result<Value, Code>> {
if !self.connections[connection].application_profile {
return Some(Err(Code::UnsupportedCapability));
}
if let Err(code) = self.validate_grant(connection, request) {
return Some(Err(code));
}
if self.source_capture.is_some() {
return Some(Err(Code::ConfigurationBusy));
}
if let Err(code) = self
.application
.lock()
.unwrap_or_else(|e| e.into_inner())
.registry
.reserve_source_capture(
request.session_id.as_deref().expect("validated session"),
&request.payload,
)
{
return Some(Err(code));
}
let runtime = Arc::clone(&self.runtime);
let payload = request.payload.clone();
let worker = match std::thread::Builder::new()
.name("magi-resource-capture".into())
.spawn(move || super::super::application::prepare_source(&payload, &runtime))
{
Ok(worker) => worker,
Err(_) => {
self.release_source_reservation();
return Some(Err(Code::InternalError));
}
};
self.source_capture = Some(PendingSourceCapture {
connection: connection.into(),
request: request.clone(),
deadline: now + Duration::from_secs(30),
worker,
});
None
}
pub(super) fn finish_source_capture(&mut self, now: Instant) {
if !self
.source_capture
.as_ref()
.is_some_and(|pending| pending.worker.is_finished())
{
return;
}
let pending = self.source_capture.take().expect("finished capture");
let prepared = pending.worker.join().unwrap_or(Err(Code::InternalError));
let result = (|| {
if now >= pending.deadline {
return Err(Code::RequestTimeout);
}
self.validate_connection(&pending.connection, &pending.request)?;
self.validate_grant(&pending.connection, &pending.request)?;
if self.operations.lookup(
&self.instance,
&self.instance,
pending.request.operation_id.as_deref().unwrap_or_default(),
)["state"]
!= "in_progress"
{
return Err(Code::OperationAlreadyKnown);
}
let attribution = json!({"instance_id":self.instance,"connection_id":pending.connection,
"operation_id":pending.request.operation_id,"grant_generation":pending.request.control.as_ref().expect("validated grant").generation});
self.application
.lock()
.unwrap_or_else(|e| e.into_inner())
.registry
.publish_source(
pending
.request
.session_id
.as_deref()
.expect("validated session"),
prepared?,
&self.runtime,
attribution,
)
})();
self.release_source_reservation();
self.settle_immediate(&pending.request, &result, now);
self.queue(
&pending.connection,
wire::response(
&self.instance,
&pending.connection,
&pending.request,
result,
),
false,
);
}
fn release_source_reservation(&self) {
self.application
.lock()
.unwrap_or_else(|e| e.into_inner())
.registry
.release_source_capture();
}
}