use aion_proto::generated;
use tonic::Status;
use crate::worker::registry::WorkerId;
use super::status_from_server_error;
pub(super) fn recorded_capacity(
registry: &crate::worker::ConnectedWorkerRegistry,
worker_id: WorkerId,
) -> Result<std::num::NonZeroU32, Status> {
registry
.worker_by_id(worker_id)
.map_err(|error| status_from_server_error(&error))?
.ok_or_else(|| {
Status::internal("worker left the registry before its registration was acked")
})?
.max_concurrency()
.ok_or_else(|| {
Status::internal(
"gRPC worker was registered with no advertised capacity; the registration \
funnel that refuses this was bypassed",
)
})
}
pub(super) fn register_ack_frame(
worker_id: WorkerId,
namespace: &str,
heartbeat_window: std::time::Duration,
accepted_max_concurrency: std::num::NonZeroU32,
) -> generated::ServerToWorker {
generated::ServerToWorker {
message: Some(generated::server_to_worker::Message::RegisterAck(
generated::RegisterAck {
worker_id: worker_id.value(),
namespace: namespace.to_owned(),
heartbeat_window_ms: u64::try_from(heartbeat_window.as_millis())
.unwrap_or(u64::MAX),
accepted_max_concurrency: accepted_max_concurrency.get(),
},
)),
}
}