All complete.

This commit is contained in:
DaLaw2 2024-04-17 19:33:50 +08:00
parent 0b5a680792
commit 8430768c83
23 changed files with 307 additions and 169 deletions

View File

@ -1,9 +1,9 @@
#![allow(non_snake_case)]
use AgentLibrary::management::monitor::Monitor;
use AgentLibrary::management::management::Management;
#[tokio::main]
async fn main() {
Monitor::run().await;
loop {}
Management::run().await;
Management::terminate().await;
}

View File

@ -11,6 +11,7 @@ image = "0.25.1"
common = "0.1.0"
chrono = "0.4.38"
sysinfo = "0.30.10"
async-ctrlc = "1.2.0"
lazy_static = "1.4.0"
serde_json = "1.0.116"
Common = { path = "../Common" }

View File

@ -10,10 +10,13 @@ use tokio::io::AsyncWriteExt;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use tokio::time::{Instant, sleep};
use crate::management::utils::bounding_box::BoundingBox;
use crate::management::utils::model_type::ModelType;
use crate::management::utils::task_result::TaskResult;
use crate::utils::logging::*;
use crate::utils::config::Config;
use crate::connection::packet::Packet;
use crate::management::manager::Manager;
use crate::management::management::Management;
use crate::management::monitor::Monitor;
use crate::utils::clear_unbounded_channel;
use crate::management::utils::task_info::TaskInfo;
@ -33,6 +36,8 @@ use crate::connection::packet::control_acknowledge_packet::ControlAcknowledgePac
use crate::connection::packet::file_transfer_result_packet::FileTransferResultPacket;
use crate::connection::packet::task_info_acknowledge_packet::TaskInfoAcknowledgePacket;
use crate::connection::packet::file_header_acknowledge_packet::FileHeaderAcknowledgePacket;
use crate::connection::packet::result_packet::ResultPacket;
use crate::management::calculate_manager::CalculateManager;
pub struct Agent {
previous_task_uuid: Option<Uuid>,
@ -113,7 +118,7 @@ impl Agent {
let mut timeout_timer = Instant::now();
let timeout_duration = Duration::from_secs(config.control_channel_timeout);
loop {
if Manager::get_state().await == AgentState::Terminate {
if Management::get_state().await == AgentState::Terminate {
return;
}
if timeout_timer.elapsed() > timeout_duration {
@ -143,13 +148,13 @@ impl Agent {
_ = sleep(Duration::from_millis(config.internal_timestamp)) => continue,
}
}
Manager::store_state(AgentState::Terminate).await;
Management::store_state(AgentState::Terminate).await;
}
async fn management(agent: Arc<RwLock<Agent>>) {
loop {
Self::refresh_state(agent.clone()).await;
let state = Manager::get_state().await;
let state = Management::get_state().await;
match state {
AgentState::ProcessTask => Self::process_task(agent.clone()).await,
AgentState::Idle(idle_time) => Self::idle(agent.clone(), Duration::from_secs(idle_time)).await,
@ -167,9 +172,9 @@ impl Agent {
let config = Config::now().await;
let timer = Instant::now();
let timeout_duration = Duration::from_secs(config.control_channel_timeout);
while Manager::get_state().await != AgentState::Terminate {
while Management::get_state().await != AgentState::Terminate {
if timer.elapsed() > timeout_duration {
Manager::store_state(AgentState::Terminate).await;
Management::store_state(AgentState::Terminate).await;
return;
}
let mut agent = agent.write().await;
@ -179,7 +184,7 @@ impl Agent {
Some(packet) => {
clear_unbounded_channel(&mut agent.control_channel_receiver.control_packet).await;
match serde_json::from_slice::<AgentState>(packet.as_data_byte()) {
Ok(state) => Manager::store_state(state).await,
Ok(state) => Management::store_state(state).await,
Err(err) => {
logging_error!("Agent", "Unable to parse packet data", format!("Err: {err}"));
continue;
@ -188,7 +193,7 @@ impl Agent {
},
None => {
logging_notice!("Agent", "Channel has been closed");
Manager::store_state(AgentState::Terminate).await;
Management::store_state(AgentState::Terminate).await;
},
}
},
@ -200,15 +205,27 @@ impl Agent {
}
async fn process_task(agent: Arc<RwLock<Agent>>) {
if let Err(entry) = Self::receive_task(agent.clone()).await {
logging_entry!(entry);
}
if let Err(entry) = Self::inference_task(agent.clone()).await {
let result = match Self::receive_task(agent.clone()).await {
Ok(task_info) => {
match Self::inference_task(task_info).await {
Ok(bounding_box) => Ok(bounding_box),
Err(entry) => {
logging_entry!(entry.clone());
Err(entry.message)
},
}
},
Err(entry) => {
logging_entry!(entry.clone());
Err(entry.message)
},
};
if let Err(entry) = Self::send_result(agent.clone(), result).await {
logging_entry!(entry);
}
}
async fn receive_task(agent: Arc<RwLock<Agent>>) -> Result<(), LogEntry> {
async fn receive_task(agent: Arc<RwLock<Agent>>) -> Result<TaskInfo, LogEntry> {
let task_info = Self::receive_task_info(agent.clone()).await?;
let previous_task_uuid = agent.read().await.previous_task_uuid;
let need_receive_model = if let Some(previous_task_uuid) = previous_task_uuid {
@ -219,10 +236,11 @@ impl Agent {
if need_receive_model {
let model_folder = Path::new(".").join("SavedModel");
Self::receive_file(agent.clone(), &model_folder).await?;
agent.write().await.previous_task_uuid = Some(task_info.uuid);
}
let image_folder = Path::new(".").join("SaveFile");
Self::receive_file(agent.clone(), &image_folder).await?;
Ok(())
Ok(task_info)
}
async fn receive_task_info(agent: Arc<RwLock<Agent>>) -> Result<TaskInfo, LogEntry> {
@ -230,7 +248,7 @@ impl Agent {
let timer = Instant::now();
let timeout_duration = Duration::from_secs(config.data_channel_timeout);
let task_info = loop {
if Manager::get_state().await == AgentState::Terminate {
if Management::get_state().await == AgentState::Terminate {
Err(notice_entry!("Agent", "Terminate. Interrupt current operation"))?;
}
if timer.elapsed() > timeout_duration {
@ -271,7 +289,7 @@ impl Agent {
let timer = Instant::now();
let timeout_duration = Duration::from_secs(config.data_channel_timeout);
let file_header = loop {
if Manager::get_state().await == AgentState::Terminate {
if Management::get_state().await == AgentState::Terminate {
Err(notice_entry!("Agent", "Terminate. Interrupt current operation"))?;
}
if timer.elapsed() > timeout_duration {
@ -307,7 +325,7 @@ impl Agent {
let mut timer = Instant::now();
let timeout_duration = Duration::from_secs(config.data_channel_timeout);
loop {
if Manager::get_state().await == AgentState::Terminate {
if Management::get_state().await == AgentState::Terminate {
Err(notice_entry!("Agent", "Terminate. Interrupt current operation"))?;
}
if timer.elapsed() > timeout_duration {
@ -360,6 +378,8 @@ impl Agent {
for index in 0..file_header.packet_count {
if let Some(block) = file_block.remove(&index) {
sorted_blocks.push(block);
} else {
Err(error_entry!("Agent", "Missing file block"))?
}
}
return Ok(sorted_blocks);
@ -381,15 +401,63 @@ impl Agent {
Ok(())
}
async fn inference_task(agent: Arc<RwLock<Agent>>) -> Result<(), LogEntry> {
Ok(())
async fn inference_task(task_info: TaskInfo) -> Result<Vec<BoundingBox>, LogEntry> {
let model_path = Path::new(".").join("SavedModel").join(task_info.model_filename);
let image_path = Path::new(".").join("SavedFile").join(task_info.image_filename);
return match task_info.model_type {
ModelType::Ultralytics => CalculateManager::ultralytics_inference(model_path, image_path).await,
ModelType::YOLOv4 => CalculateManager::yolov4_inference(model_path, image_path).await,
ModelType::YOLOv7 => CalculateManager::yolov7_inference(model_path, image_path).await,
};
}
async fn send_result(agent: Arc<RwLock<Agent>>, result: Result<Vec<BoundingBox>, String>) -> Result<(), LogEntry> {
let config = Config::now().await;
let result = TaskResult::new(result);
let result_data = serde_json::to_vec(&result)
.map_err(|err| error_entry!("Agent", "Unable to serialize data", format!("Err: {err}")))?;
let timer = Instant::now();
let mut polling_times = 0_u32;
let polling_interval = Duration::from_millis(config.polling_interval);
let timeout_duration = Duration::from_secs(config.control_channel_timeout);
loop {
if Management::get_state().await == AgentState::Terminate {
Err(notice_entry!("Agent", "Terminate. Interrupt current operation"))?;
}
if timer.elapsed() > timeout_duration {
Err(notice_entry!("Agent", "Data Channel timeout"))?;
}
if timer.elapsed() > polling_times * polling_interval {
if let Some(data_channel_sender) = agent.write().await.data_channel_sender.as_mut() {
data_channel_sender.send(ResultPacket::new(result_data.clone())).await;
} else {
Err(warning_entry!("Agent", "Data Channel is not ready"))?;
}
polling_times += 1;
}
if let Some(data_channel_receiver) = agent.write().await.data_channel_receiver.as_mut() {
select! {
packet = data_channel_receiver.result_acknowledge_packet.recv() => {
if let Some(_) = packet {
clear_unbounded_channel(&mut data_channel_receiver.result_acknowledge_packet).await;
return Ok(());
} else {
Err(notice_entry!("Agent", "Channel has been closed"))?;
}
},
_ = sleep(Duration::from_millis(config.internal_timestamp)) => continue,
}
} else {
Err(warning_entry!("Agent", "Data Channel is not ready"))?;
}
}
}
async fn idle(agent: Arc<RwLock<Agent>>, idle_duration: Duration) {
let config = Config::now().await;
let timer = Instant::now();
loop {
if Manager::get_state().await == AgentState::Terminate {
if Management::get_state().await == AgentState::Terminate {
logging_notice!("Agent", "Terminate. Interrupt current operation");
return;
}
@ -397,7 +465,7 @@ impl Agent {
return;
}
let mut agent = agent.write().await;
if let Some(data_channel_receiver) = &mut agent.data_channel_receiver {
if let Some(data_channel_receiver) = agent.data_channel_receiver.as_mut() {
select! {
biased;
_ = data_channel_receiver.alive_packet.recv() => clear_unbounded_channel(&mut data_channel_receiver.alive_packet).await,
@ -407,7 +475,7 @@ impl Agent {
logging_warning!("Agent", "Data Channel is not available.");
return;
}
if let Some(data_channel_sender) = &mut agent.data_channel_sender {
if let Some(data_channel_sender) = agent.data_channel_sender.as_mut() {
data_channel_sender.send(AliveAcknowledgePacket::new()).await;
} else {
logging_warning!("Agent", "Data Channel is not available.");
@ -423,12 +491,12 @@ impl Agent {
let timer = Instant::now();
let timeout_duration = Duration::from_secs(config.control_channel_timeout);
loop {
if Manager::get_state().await == AgentState::Terminate {
if Management::get_state().await == AgentState::Terminate {
logging_notice!("Agent", "Terminate. Interrupt current operation");
return;
}
if timer.elapsed() > timeout_duration {
Manager::store_state(AgentState::Terminate).await;
Management::store_state(AgentState::Terminate).await;
logging_notice!("Agent", "Control channel timout.");
return;
}
@ -447,7 +515,7 @@ impl Agent {
continue;
}
} else {
Manager::store_state(AgentState::Terminate).await;
Management::store_state(AgentState::Terminate).await;
logging_notice!("Agent", "Channel has been closed");
return;
}
@ -456,8 +524,8 @@ impl Agent {
}
}
if let Some(port) = port {
let full_address = format!("{}:{}", config.management_address, port);
match TcpStream::connect(&full_address).await {
let address = format!("{}:{}", config.management_address, port);
match TcpStream::connect(&address).await {
Ok(tcp_stream) => {
let socket_stream = SocketStream::new(tcp_stream);
let (data_channel_sender, data_channel_receiver) = DataChannel::new(socket_stream);

View File

@ -7,10 +7,6 @@ use crate::management::utils::bounding_box::BoundingBox;
pub struct CalculateManager;
impl CalculateManager {
fn new() -> Self {
Self
}
pub async fn ultralytics_inference(model_path: PathBuf, image_path: PathBuf) -> Result<Vec<BoundingBox>, LogEntry> {
#[cfg(target_os = "windows")]
let python = "python";
@ -23,9 +19,7 @@ impl CalculateManager {
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.map_err(|err|
error_entry!("Calculate Manager", "Unable to create process", format!("Err: {err}"))
)?;
.map_err(|err| error_entry!("Calculate Manager", "Unable to create process", format!("Err: {err}")))?;
let output = process.wait_with_output().await
.map_err(|err| error_entry!("Calculate Manager", "An error occurred during process execution", format!("Err: {err}")))?;
if output.status.success() {
@ -39,11 +33,11 @@ impl CalculateManager {
}
}
pub async fn yolov4_inference() -> Result<Vec<BoundingBox>, LogEntry> {
pub async fn yolov4_inference(_model_path: PathBuf, _image_path: PathBuf) -> Result<Vec<BoundingBox>, LogEntry> {
Ok(Vec::new())
}
pub async fn yolov7_inference() -> Result<Vec<BoundingBox>, LogEntry> {
pub async fn yolov7_inference(_model_path: PathBuf, _image_path: PathBuf) -> Result<Vec<BoundingBox>, LogEntry> {
Ok(Vec::new())
}
}

View File

@ -1,18 +1,12 @@
use tokio::fs;
use tokio::fs::File;
use std::path::PathBuf;
use tokio::sync::RwLock;
use tokio::io::AsyncWriteExt;
use lazy_static::lazy_static;
use tokio::process::Command as AsyncCommand;
use crate::utils::logging::*;
use crate::utils::static_files::StaticFiles;
use crate::utils::logging::{Logger, LogLevel};
lazy_static! {
static ref FILE_MANAGER: RwLock<FileManager> = RwLock::new(FileManager {});
}
pub struct FileManager;
impl FileManager {

View File

@ -0,0 +1,102 @@
use std::sync::Arc;
use std::time::Duration;
use async_ctrlc::CtrlC;
use tokio::sync::{RwLock, RwLockReadGuard, RwLockWriteGuard};
use lazy_static::lazy_static;
use tokio::select;
use tokio::time::sleep;
use crate::utils::logging::*;
use crate::connection::socket::management_socket::ManagementSocket;
use crate::management::utils::agent_state::AgentState;
use crate::management::agent::Agent;
use crate::management::file_manager::FileManager;
use crate::utils::config::Config;
lazy_static! {
static ref MANAGER: RwLock<Management> = RwLock::new(Management::new());
}
pub struct Management {
agent: Option<Arc<RwLock<Agent>>>,
state: Option<AgentState>,
terminate: bool,
}
impl Management {
pub fn new() -> Self {
Self {
agent: None,
state: None,
terminate: false,
}
}
pub async fn instance() -> RwLockReadGuard<'static, Self> {
MANAGER.read().await
}
pub async fn instance_mut() -> RwLockWriteGuard<'static, Self> {
MANAGER.write().await
}
pub async fn run() {
FileManager::initialize().await;
tokio::spawn(async move {
Self::hot_reload().await;
});
logging_information!("Management", "Online now");
match CtrlC::new() {
Ok(ctrlc) => ctrlc.await,
Err(err) => logging_emergency!("Management", "Unable to create instance", format!("Err: {err}")),
}
}
pub async fn terminate() {
logging_information!("Management", "Termination in process");
Self::instance_mut().await.terminate = true;
FileManager::cleanup().await;
logging_information!("Management", "Termination complete");
}
pub async fn hot_reload() {
let config = Config::now().await;
while !Self::instance().await.terminate {
let mut management = Self::instance_mut().await;
if management.agent.is_some() {
match management.state {
Some(AgentState::Terminate) => {
management.agent = None;
management.state = Some(AgentState::None);
},
_ => sleep(Duration::from_millis(config.internal_timestamp)).await,
}
} else {
select! {
(socket_stream, _) = ManagementSocket::get_connection() => {
match Agent::new(socket_stream).await {
Ok(agent) => management.agent = Some(Arc::new(RwLock::new(agent))),
Err(entry) => logging_entry!(entry),
}
},
_ = sleep(Duration::from_millis(config.internal_timestamp)) => continue,
}
}
}
}
pub async fn store_state(state: AgentState) {
let mut manager = Self::instance_mut().await;
if let Some(origin_state) = manager.state {
if origin_state != AgentState::Terminate {
manager.state = Some(state);
}
} else {
manager.state = Some(state)
}
}
pub async fn get_state() -> AgentState {
let manager = Self::instance().await;
manager.state.unwrap_or_else(|| AgentState::None)
}
}

View File

@ -1,65 +0,0 @@
use std::sync::Arc;
use tokio::sync::{RwLock, RwLockReadGuard, RwLockWriteGuard};
use lazy_static::lazy_static;
use crate::management::utils::agent_state::AgentState;
use crate::management::agent::Agent;
lazy_static! {
static ref MANAGER: RwLock<Manager> = RwLock::new(Manager::new());
}
pub struct Manager {
agent: Option<Arc<RwLock<Agent>>>,
state: Option<AgentState>,
terminate: bool,
}
impl Manager {
pub fn new() -> Self {
Self {
agent: None,
state: None,
terminate: false,
}
}
pub async fn instance() -> RwLockReadGuard<'static, Self> {
MANAGER.read().await
}
pub async fn instance_mut() -> RwLockWriteGuard<'static, Self> {
MANAGER.write().await
}
pub async fn run() {
}
fn initialize() {
}
pub async fn terminate() {
}
fn cleanup() {
}
pub async fn store_state(state: AgentState) {
let mut manager = Self::instance_mut().await;
if let Some(origin_state) = manager.state {
if origin_state != AgentState::Terminate {
manager.state = Some(state);
}
} else {
manager.state = Some(state)
}
}
pub async fn get_state() -> AgentState {
let manager = Self::instance().await;
manager.state.unwrap_or_else(|| AgentState::None)
}
}

View File

@ -2,5 +2,5 @@ pub mod utils;
pub mod agent;
pub mod calculate_manager;
pub mod file_manager;
pub mod manager;
pub mod management;
pub mod monitor;

42
Cargo.lock generated
View File

@ -15,6 +15,7 @@ name = "AgentLibrary"
version = "0.1.0"
dependencies = [
"Common",
"async-ctrlc",
"chrono",
"common",
"image",
@ -29,7 +30,7 @@ dependencies = [
[[package]]
name = "Common"
version = "1.0.0"
version = "0.1.0"
dependencies = [
"chrono",
"colored",
@ -41,7 +42,7 @@ dependencies = [
[[package]]
name = "Management"
version = "1.0.0"
version = "0.1.0"
dependencies = [
"ManagementLibrary",
"actix-rt",
@ -456,6 +457,15 @@ version = "0.7.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "96d30a06541fbafbc7f82ed10c06164cfbd2c401138f6addd8404629c4b16711"
[[package]]
name = "async-ctrlc"
version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "907279f6e91a51c8ec7cac24711e8308f21da7c10c7700ca2f7e125694ed2df1"
dependencies = [
"ctrlc",
]
[[package]]
name = "atomic_refcell"
version = "0.1.13"
@ -673,6 +683,12 @@ version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "baf1de4339761588bc0619e3cbc0120ee582ebb74b53b4efbf79117bd2da40fd"
[[package]]
name = "cfg_aliases"
version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fd16c4719339c4530435d38e511904438d07cce7950afa3718a84ac36c10e89e"
[[package]]
name = "chrono"
version = "0.4.38"
@ -823,6 +839,16 @@ dependencies = [
"typenum",
]
[[package]]
name = "ctrlc"
version = "3.4.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "672465ae37dc1bc6380a6547a8883d5dd397b0f1faaad4f265726cc7042a5345"
dependencies = [
"nix",
"windows-sys 0.52.0",
]
[[package]]
name = "custom_derive"
version = "0.1.7"
@ -1792,6 +1818,18 @@ version = "1.0.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "650eef8c711430f1a879fdd01d4745a7deea475becfb90269c06775983bbf086"
[[package]]
name = "nix"
version = "0.28.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ab2156c4fce2f8df6c499cc1c763e4394b7482525bf2a9701c9d79d215f519e4"
dependencies = [
"bitflags 2.5.0",
"cfg-if",
"cfg_aliases",
"libc",
]
[[package]]
name = "nom"
version = "7.1.3"

View File

@ -1,14 +1,14 @@
[package]
name = "Common"
version = "1.0.0"
version = "0.1.0"
edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
[dependencies]
chrono = "0.4.37"
chrono = "0.4.38"
colored = "2.1.0"
rust-embed = "8.3.0"
tokio = { version = "1.37.0", features = ["full"] }
serde = { version = "1.0.197", features = ["derive"] }
serde = { version = "1.0.198", features = ["derive"] }
uuid = { version = "1.8.0", features = ["v4", "fast-rng", "macro-diagnostics", "serde"] }
colored = "2.1.0"

View File

@ -6,14 +6,16 @@ use crate::management::utils::model_type::ModelType;
pub struct TaskInfo {
pub uuid: Uuid,
pub model_filename: String,
pub image_filename: String,
pub model_type: ModelType,
}
impl TaskInfo {
pub fn new(uuid: Uuid, model_filename: String, model_type: ModelType) -> Self {
pub fn new(uuid: Uuid, model_filename: String, image_filename: String, model_type: ModelType) -> Self {
Self {
uuid,
model_filename,
image_filename,
model_type,
}
}

View File

@ -7,6 +7,12 @@ pub struct TaskResult {
}
impl TaskResult {
pub fn new(result: Result<Vec<BoundingBox>, String>) -> Self {
Self {
result,
}
}
pub fn into(self) -> Result<Vec<BoundingBox>, String> {
self.result
}

View File

View File

@ -1,6 +1,6 @@
[package]
name = "Management"
version = "1.0.0"
version = "0.1.0"
edition = "2021"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html

View File

@ -1,6 +1,6 @@
#![allow(non_snake_case)]
use ManagementLibrary::management::manager::Management;
use ManagementLibrary::management::management::Management;
#[actix_web::main]
async fn main() {

View File

@ -12,7 +12,6 @@ pub struct AgentSocket {
impl AgentSocket {
pub async fn new() -> Self {
//unimplemented!("Unable to end normally when failed");
let listener = loop {
let config = Config::now().await;
let port = config.agent_listen_port;
@ -30,7 +29,6 @@ impl AgentSocket {
}
pub async fn get_connection(&mut self) -> (SocketStream, SocketAddr) {
//unimplemented!("The loop cannot be ended when there is no connection");
let (stream, address) = loop {
if let Ok(connection) = self.listener.accept().await {
break connection;

View File

@ -258,6 +258,7 @@ impl Agent {
Agent::transfer_task_info(agent.clone(), &image_task).await?;
if need_transfer_model {
Agent::transfer_file(agent.clone(), &image_task.model_filename, &image_task.model_filepath).await?;
agent.write().await.previous_task_uuid = Some(image_task.task_uuid);
}
Agent::transfer_file(agent.clone(), &image_task.image_filename, &image_task.image_filepath).await?;
} else {
@ -270,7 +271,7 @@ impl Agent {
async fn transfer_task_info(agent: Arc<RwLock<Agent>>, image_task: &ImageTask) -> Result<(), LogEntry> {
let uuid = agent.read().await.uuid;
let config = Config::now().await;
let task_info = TaskInfo::new(image_task.task_uuid, image_task.model_filename.clone(), image_task.model_type);
let task_info = TaskInfo::new(image_task.task_uuid, image_task.image_filename.clone(), image_task.model_filename.clone(), image_task.model_type);
let task_info_data = serde_json::to_vec(&task_info)
.map_err(|err| error_entry!("Agent", "Unable to serialize data", format!("Err: {err}")))?;
let timer = Instant::now();

View File

@ -84,6 +84,8 @@ impl FileManager {
}
pub async fn terminate() {
// 需要防止重入函式
// 需要等待其他執行緒結束
logging_information!("File Manager", "Termination in progress");
Self::instance_mut().await.terminate = true;
Self::cleanup().await;
@ -124,9 +126,9 @@ impl FileManager {
_ => {
logging_error!("File Manager", format!("Task {}, unsupported file type", task.uuid));
task.panic("Unsupported file type".to_string()).await;
}
},
}
}
},
None => sleep(Duration::from_millis(config.internal_timestamp)).await,
}
}
@ -161,22 +163,21 @@ impl FileManager {
Ok(_) => {
task.update_unprocessed(1).await;
Self::forward_to_task_manager(task).await;
}
},
Err(err) => {
task.panic("Unable to move file".to_string()).await;
logging_error!("File Manager", "Unable to move file", format!("Source: {}, Destination: {}, Err: {}", source_path.display(), destination_path.display(), err));
}
},
}
}
async fn extract_pre_processing(source_path: &PathBuf, destination_path: &PathBuf, create_folder: &PathBuf) -> Result<(), LogEntry> {
fs::create_dir(&create_folder).await
.map_err(|err| error_entry!("File Manager", "Unable to create folder", format!("Folder:{}, Err: {}", create_folder.display(), err)))?;
.map_err(|err|
error_entry!("File Manager", "Unable to create folder",format!("Folder:{}, Err: {}", create_folder.display(), err)))?;
fs::rename(&source_path, &destination_path).await
.map_err(|err| {
error_entry!("File Manager", "Unable to move file",
format!("Source: {}, Destination: {}, Err: {}", source_path.display(), destination_path.display(), err))
})?;
.map_err(|err|
error_entry!("File Manager", "Unable to move file",format!("Source: {}, Destination: {}, Err: {}", source_path.display(), destination_path.display(), err)))?;
Ok(())
}
@ -190,6 +191,7 @@ impl FileManager {
return;
}
let video_path = destination_path;
let extract_folder = create_folder;
if let Err(entry) = Self::extract_video_info(&video_path).await {
task.panic(entry.message.clone()).await;
logging_entry!(entry);
@ -198,7 +200,7 @@ impl FileManager {
let result = spawn_blocking(move || {
Self::extract_video(video_path)
}).await;
Self::process_extract_result(task, create_folder, result).await;
Self::process_extract_result(task, extract_folder, result).await;
}
async fn extract_video_info(video_path: &PathBuf) -> Result<(), LogEntry> {
@ -218,17 +220,17 @@ impl FileManager {
for field in structure.fields() {
match field.as_str() {
"framerate" => video_info.framerate = structure.get::<gstreamer::Fraction>(field)
.map_or_else(|_| "30/1".to_string(), |f| format!("{number}/{denom}", number = f.numer(), denom = f.denom())),
_ => {}
.map_or_else(|_| "30/1".to_string(), |f| format!("{}/{}", f.numer(), f.denom())),
_ => {},
}
}
}
}
}
let toml_path = video_path.with_extension("toml");
let toml_string = toml::to_string(&video_info)
let video_info = toml::to_string(&video_info)
.map_err(|err| error_entry!("File Manager", "Unable to serialize data", format!("Err: {err}")))?;
fs::write(&toml_path, toml_string).await
fs::write(&toml_path, video_info).await
.map_err(|err| error_entry!("File Manager", "Unable to write file", format!("File:{}, Err: {}", toml_path.display(), err)))?;
Ok(())
}
@ -284,14 +286,14 @@ impl FileManager {
.map_err(|err| error_entry!("File Manager", "Unable to read file", format!("File:{}, Err: {}", zip_path.display(), err)))?;
let mut archive = ZipArchive::new(reader)
.map_err(|err| error_entry!("File Manager", "Unable to create instance", format!("Err: {err}")))?;
let output_folder = zip_path.clone().with_extension("").to_path_buf();
let extract_folder = zip_path.clone().with_extension("").to_path_buf();
for i in 0..archive.len() {
let mut file = archive.by_index(i)
.map_err(|err| error_entry!("File Manager", "An error occurred while reading the file", format!("File: {}, Err: {}", zip_path.display(), err)))?;
if let Some(enclosed_path) = file.enclosed_name() {
if let Some(extension) = enclosed_path.extension() {
if allowed_extensions.contains(&extension.to_str().unwrap_or("")) {
let output_path = output_folder.join(enclosed_path.file_name().unwrap_or_default());
let output_path = extract_folder.join(enclosed_path.file_name().unwrap_or_default());
let mut output_file = File::create(&output_path)
.map_err(|err| error_entry!("File Manager", "Unable to create file", format!("File: {}, Err: {}", output_path.display(), err)))?;
std::io::copy(&mut file, &mut output_file)
@ -310,21 +312,21 @@ impl FileManager {
Ok(count) => {
task.update_unprocessed(count).await;
Self::forward_to_task_manager(task).await;
}
},
Err(entry) => {
task.panic("Unable to read folder".to_string()).await;
logging_entry!(entry);
}
},
}
}
},
Ok(Err(entry)) => {
task.panic(entry.message.clone()).await;
logging_entry!(entry);
}
},
Err(err) => {
task.panic("Panic occurs during execution".to_string()).await;
logging_error!("File Manager", "Panic occurs during execution", format!("Err: {err}"));
}
},
}
}
@ -361,12 +363,11 @@ impl FileManager {
let saved_path = Path::new(".").join("PostProcessing").join(image_task.image_filename.clone());
match Self::draw_bounding_box(image_task, &config, &font).await {
Ok(image) => {
match image.save(&saved_path) {
Ok(_) => Self::forward_to_repository(task).await,
Err(err) => {
task.panic("Unable to write file".to_string()).await;
logging_error!("File Manager", "Unable to write file", format!("File: {}, Err: {}", saved_path.display(), err));
}
if let Err(err) = image.save(&saved_path) {
task.panic("Unable to write file".to_string()).await;
logging_error!("File Manager", "Unable to write file", format!("File: {}, Err: {}", saved_path.display(), err));
} else {
Self::forward_to_repository(task).await;
}
},
Err(entry) => {
@ -426,7 +427,7 @@ impl FileManager {
Self::process_recombination_result(task, result).await;
}
fn recombination_video(video_info_path: PathBuf, frame_folder: PathBuf, target_path: PathBuf) -> Result<(), LogEntry> {
fn recombination_video(video_info_path: PathBuf, frame_folder: PathBuf, saved_path: PathBuf) -> Result<(), LogEntry> {
let toml_str = std::fs::read_to_string(&video_info_path)
.map_err(|err| error_entry!("File Manager", "Unable to read file", format!("File: {}, Err: {}", video_info_path.display(), err)))?;
let video_info: VideoInfo = toml::from_str(&toml_str).unwrap_or_default();
@ -438,13 +439,13 @@ impl FileManager {
"video/x-vp8" => format!("vp8enc target-bitrate={}", bitrate),
_ => format!("x265enc bitrate={}", bitrate),
};
let muxer = match target_path.extension().and_then(OsStr::to_str) {
let muxer = match saved_path.extension().and_then(OsStr::to_str) {
Some("mp4") => "mp4mux",
Some("avi") => "avimux",
Some("mkv") => "matroskamux",
_ => "mp4mux",
};
let pipeline_string = format!("multifilesrc location={:?} index=1 caps=image/png,framerate=(fraction){} ! pngdec ! videoconvert ! {} ! {} ! filesink location={:?}", frame_folder.join("%010d.png"), video_info.framerate, encoder, muxer, target_path);
let pipeline_string = format!("multifilesrc location={:?} index=1 caps=image/png,framerate=(fraction){} ! pngdec ! videoconvert ! {} ! {} ! filesink location={:?}", frame_folder.join("%010d.png"), video_info.framerate, encoder, muxer, saved_path);
let pipeline = gstreamer::parse::launch(&pipeline_string)
.map_err(|err| error_entry!("File Manager", "Unable to create instance", format!("Err: {err}")))?;
let bus = pipeline.bus().ok_or(error_entry!("File Manager", "Unable to create instance"))?;
@ -496,9 +497,7 @@ impl FileManager {
let entry = entry.map_err(|err|
error_entry!("File Manager", "An error occurred while reading the folder", format!("Folder: {}, Err: {}", source_folder.display(), err)))?;
let path = entry.path();
let file_name = path.file_name().ok_or(
error_entry!("File Manager", "Invalid file")
)?.to_string_lossy();
let file_name = path.file_name().ok_or(error_entry!("File Manager", "Invalid file"))?.to_string_lossy();
zip.start_file(file_name.clone(), options)
.map_err(|err| error_entry!("File Manager", "Unable to create file", format!("File: {}, Err: {}", file_name.clone(), err)))?;
let mut file_contents = Vec::new();

View File

@ -2,6 +2,6 @@ pub mod utils;
pub mod agent;
pub mod agent_manager;
pub mod file_manager;
pub mod manager;
pub mod management;
pub mod result_repository;
pub mod task_manager;

View File

@ -1,12 +1,12 @@
use tokio::fs;
use uuid::Uuid;
use std::sync::Arc;
use lazy_static::lazy_static;
use std::path::{Path, PathBuf};
use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use tokio::sync::{RwLock, RwLockReadGuard, RwLockWriteGuard};
use crate::management::agent::Agent;
use crate::utils::logging::*;
use crate::management::agent::Agent;
use crate::management::file_manager::FileManager;
use crate::management::agent_manager::AgentManager;
use crate::management::utils::image_task::ImageTask;
@ -37,10 +37,9 @@ impl TaskManager {
pub async fn add_task(mut task: Task) {
task.change_status(TaskStatus::Processing);
{
let mut task_manager = Self::instance_mut().await;
task_manager.tasks.insert(task.uuid, task.clone());
}
let mut task_manager = Self::instance_mut().await;
task_manager.tasks.insert(task.uuid, task.clone());
drop(task_manager);
tokio::spawn(async move {
Self::distribute_task(task).await;
});
@ -60,6 +59,7 @@ impl TaskManager {
}
pub async fn distribute_task(task: Task) {
// 考慮到刷新效能的不適配問題
let model_filepath = Path::new(".").join("SavedModel").join(task.model_filename.clone());
let vram_usage = Self::estimated_vram_usage(&model_filepath).await;
let filter_agents = AgentManager::filter_agent_by_vram(vram_usage).await;
@ -75,7 +75,7 @@ impl TaskManager {
None => {
logging_error!("Task Manager", "Agent instance does not exist");
0.0
}
},
};
if agent_ram > ram_usage * 0.7 {
if let Some(agent) = AgentManager::get_agent(agent_id).await {
@ -92,7 +92,7 @@ impl TaskManager {
logging_warning!("Task Manager", format!("Task {} cannot be assigned to any agent", task.uuid));
Self::submit_image_task(image_task, false).await;
}
}
},
Some("mp4") | Some("avi") | Some("mkv") | Some("zip") => {
let image_folder = Path::new(".").join("PreProcessing").join(task.media_filename.clone()).with_extension("");
let mut image_folder = match fs::read_dir(&image_folder).await {
@ -140,11 +140,11 @@ impl TaskManager {
}
image_id += 1;
}
}
},
_ => {
Self::task_panic(&task.uuid, "Unsupported file type".to_string()).await;
notice_entry!("File Manager", format!("Task {}, unsupported file type", task.uuid));
}
},
}
}

View File

@ -19,8 +19,8 @@ VisioGrid 是一個使用 Rust 開發的分布式計算平台,專注於圖像
## 安裝和配置
- 複製倉庫:`git clone https://github.com/DaLaw2/VisioGrid`
- 從原碼編譯:
- 编譯管理節點Management需要安裝 GStreamer`cargo build --release --package Management`
- 编譯代理Agent需要安裝 LibTorch`cargo build --release --package Agent`
- 编譯管理節點Management需要安裝 GStreamer`cargo build --release --bin Management`
- 编譯代理Agent需要安裝 LibTorch`cargo build --release --bin Agent`
- 使用 Docker 運行:
- 管理節點Management和代理Agent的 Docker 容器已包含所有必要依賴,無需手動安裝。