MCPcopy Create free account
hub / github.com/Botloader/botloader / init_worker_handles

Function init_worker_handles

components/scheduler/src/vmworkerpool.rs:369–396  ·  view source on GitHub ↗
(
    pending: PendingWorkerHandle,
    #[cfg(target_family = "windows")] stream: tokio::net::TcpStream,
    #[cfg(target_family = "unix")] stream: tokio::net::UnixStream,
)

Source from the content-addressed store, hash-verified

367}
368
369fn init_worker_handles(
370 pending: PendingWorkerHandle,
371 #[cfg(target_family = "windows")] stream: tokio::net::TcpStream,
372 #[cfg(target_family = "unix")] stream: tokio::net::UnixStream,
373) -> WorkerHandle {
374 let (scheduler_msg_tx, scheduler_msg_rx) = mpsc::unbounded_channel();
375 let (worker_msg_tx, worker_msg_rx) = mpsc::unbounded_channel();
376
377 let (mut reader, mut writer) = stream.into_split();
378
379 tokio::spawn(async move { message_reader(&mut reader, worker_msg_tx).await });
380
381 tokio::spawn(async move { message_writer(&mut writer, scheduler_msg_rx).await });
382
383 WorkerHandle {
384 child: pending.child,
385 tx: scheduler_msg_tx,
386 rx: worker_msg_rx,
387
388 last_active_guild: None,
389 returned_at: Instant::now(),
390 claimed_at: Instant::now(),
391 worker_id: pending.worker_id,
392 priority_index: pending.priority_index,
393
394 session_state: None,
395 }
396}
397
398#[derive(Clone)]
399pub struct WorkerLaunchConfig {

Callers 1

worker_connectedMethod · 0.85

Calls 2

message_readerFunction · 0.85
message_writerFunction · 0.85

Tested by

no test coverage detected