#include "pch.h" #include #include #include "GlobalVarDefine.h" #include "NetMessageDefine.h" #include "CBaseDispatchTaskMg.h" #include #define LOG_PER_LOG_MAX_LEN (1024*1) // 每条日志最大长度 //构造函数 CBaseDispatchTaskMg::CBaseDispatchTaskMg(void) { InitDate(); m_pNetWork = nullptr; m_pDebugModule = nullptr; m_bRollBackEvent = FALSE; m_nCurTrainID = -1; //当前训练号 m_nCurCourseID = -1; //当前课程号 m_nCurRBPointID = 1; //当前回滚记录号 m_NetServerID = -1; //服务器系统号 m_TechServerID = -1; //教员系挺号 m_nViewSysSvrID = -1; //协同服务器ID m_nAllTrainStatus = AllTrain_Status_Stop; m_nEndTrainingTime = -1; //结束训练时间 m_nSysSimulateTime = -1; //系统仿真时间 m_pFrmTool = nullptr; //框架工具接口 //加载调试工具 FunCreate_DebugModule CreateDebugModule = nullptr; #ifdef IN_WINDOWS std::string Str_Dll = GetFolder_Module_2() + "\\OCC.Public.DebugModule.dll"; HINSTANCE hInstance2 = ::LoadLibraryA(Str_Dll.c_str()); CreateDebugModule = (FunCreate_DebugModule)GetProcAddress(hInstance2, "Create_Instance"); #else std::string Str_Dll = GetFolder_Module_2() + "Debug_Module.so"; HINSTANCE hInstance2 = dlopen(Str_Dll.c_str(),RTLD_LAZY); CreateDebugModule = (FunCreate_DebugModule)dlsym(hInstance2,"Create_Instance"); #endif if ( CreateDebugModule ) { m_pDebugModule = (CDebug_Module_I*)CreateDebugModule(); } m_pLogBuf = new char[LOG_PER_LOG_MAX_LEN]; memset(m_cRollBackDes, 0, sizeof(m_cRollBackDes)); //添加命令处理模块列表 Add_AssignCmd_List(OCCCOMMAND::eMaincontrolEnvir, "10,17,12"); Add_AssignCmd_List(OCCCOMMAND::eTSZoneWater_FeedBack, "17"); Add_AssignCmd_List(OCCCOMMAND::eSys_ShijingXietong, "10,17"); Add_AssignCmd_List(OCCCOMMAND::eUpdateTrain, "11,24"); Add_AssignCmd_List(OCCCOMMAND::eTSViewSysCtrl_DHEqmpt, "12"); } //析构函数 CBaseDispatchTaskMg::~CBaseDispatchTaskMg(void) { if ( m_pDebugModule ) { delete m_pDebugModule; m_pDebugModule = nullptr; } delete[]m_pLogBuf; } /******************************************************* 函数名称: StartWork 函数功能: 开始工作 输入参数: ...参数太多不想写了 输出参数: NULL 返 回 值: BOOL - TRUE:成功 FALSE:失败 *******************************************************/ bool CBaseDispatchTaskMg::StartWork( CFrameTool_I* pFrmTool ) { m_pFrmTool = pFrmTool; //1、创建工作线程 m_Thread_DealTCPData = std::thread(std::bind(&CBaseDispatchTaskMg::ThreadFun_DealTCPData, this)); m_Thread_DealUDPData = std::thread(std::bind(&CBaseDispatchTaskMg::ThreadFun_DealUDPData, this)); m_Thread_DealGroupData = std::thread(std::bind(&CBaseDispatchTaskMg::ThreadFun_DealGroupData, this)); m_Thread_RollBackRecord = std::thread(std::bind(&CBaseDispatchTaskMg::ThreadFun_RollBackRecord,this)); //2、加载子调度系统以及数据库配置文件 LoadSub_DispatchSys(); if ( !InitConfigParam() ) { OUT_LOG(LOG_LEVEL_ERROR,"调度服务器初始化数据库配置失败"); return FALSE; } //3、连接Redis数据库服务器 int nResult = Start_Redis_Connect(); if ( -1 == nResult ) { OUT_LOG(LOG_LEVEL_ERROR,"Redis数据库连接失败"); } STARTUP NetOrder; //NetOrder.time = (int)CTime::GetCurrentTime().GetTime(); NetOrder.time = (int)time(nullptr); NetComand_Send(&NetOrder,OCC_SYS_SERVER_ADMIN,OCCCOMMAND::eStartUp); OUT_LOG(LOG_LEVEL_NORMAL, "调度服务器创建消息线程成功"); return TRUE; } /******************************************************* 函数名称: NetComand_Send 函数功能: 网络命令发送 输入参数: pNetPack - 命令内容 nDesSysType - 目标系统类型 eCmdType - 命令类型 nTrainID - 训练号 输出参数: NULL 返回值 : void *******************************************************/ void CBaseDispatchTaskMg::NetComand_Send(LPBASENETPACKET pNetPack, int nDesSysType, OCCCOMMAND::CmdType eCmdType, int nTrainID) { OCCCOMMAND NetCmd; NetCmd.ScrSysType = OCC_SYS_SERVER_SVR; NetCmd.setPacketValue(pNetPack,eCmdType); NetCmd.TrainningID = nTrainID; NetCmd.DesSysType = nDesSysType; m_QI_SendList_TCP.push(NetCmd); } /******************************************************* 函数名称: Start_Redis_Connect 函数功能: 启动Redis数据库连接 输入参数: NULL 输出参数: NULL 返回值 : int - 连接结果:-1:失败 其他:成功 *******************************************************/ int CBaseDispatchTaskMg::Start_Redis_Connect() { #ifdef IN_WINDOWS std::string Str_FileName = GetFolder_Module_2() + "//RedisConfig.ini"; #else std::string Str_FileName = GetFolder_Module_2() + "RedisConfig.ini"; #endif std::string Str_IP = ReadCfg_String(Str_FileName,"REDISSERVER", "ServerIP", "127.0.0.1"); int nPort = ReadCfg_Int (Str_FileName,"REDISSERVER", "ServerPort", 6379); int nResult = m_RedisClient.StartRedis(Str_IP, nPort); return nResult; } /******************************************************* 函数名称: InitDate 函数功能: 初始化数据 输入参数: NULL 输出参数: NULL 返回值 : void *******************************************************/ void CBaseDispatchTaskMg::InitDate() { m_ReceiveMsgThreadOut = true; //线程退出-TCP数据 m_ReceiveRealMsgThreadOut = true; //线程退出-UDP数据 m_ReceiveGroupMsgThreadOut = true; //线程退出-组播数据 m_RollBackRecordThreadOut = true; //线程退出-回滚记录 m_dLastRecordTime = -1.0; m_QI_SendList_TCP.init(); m_QI_SendList_TCP.setSafeCount(100000); m_QI_SendList_TCP.SetQueName("TCP发送队列"); m_QI_GetList_TCP.init(); m_QI_GetList_TCP.setSafeCount(10000); m_QI_GetList_TCP.SetQueName("TCP接收队列"); m_QI_SendList_Group.init(); m_QI_SendList_Group.setSafeCount(10000); m_QI_SendList_Group.SetQueName("组播发送队列"); m_QI_GetList_Group.init(); m_QI_GetList_Group.setSafeCount(10000); m_QI_GetList_Group.SetQueName("组播接收队列"); m_QI_GetList_UDP.init(); m_QI_GetList_UDP.setSafeCount(10000); m_QI_GetList_UDP.SetQueName("UDP接收队列"); m_QI_SendList_UDP.init(); m_QI_SendList_UDP.setSafeCount(10000); m_QI_SendList_UDP.SetQueName("UDP发送队列"); } /******************************************************* 函数名称: NetComandToServer 函数功能: 处理发送给OCC总服务的网络命令 输入参数: NetCmd - 网络命令 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::NetComandToServer( OCCCOMMAND& NetCmd ) { //OUT_LOG(LOG_LEVEL_NORMAL,"<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<<"); //OUT_LOG(LOG_LEVEL_NORMAL, "---------------------------课程管理开始!\n"); try { if ( NetCmd.cmdType == OCCCOMMAND::eTrainIni ) { if ( m_pNetWork ) { m_pNetWork->_ISetTechID(NetCmd.ScrNetID); m_pNetWork->_ISetTrainID(NetCmd.TrainningID); m_TechServerID = NetCmd.ScrNetID; m_nCurTrainID = NetCmd.TrainningID; m_nCurRBPointID = 1; m_nCurCourseID = NetCmd.Net_JY_TrainIni.CourseID; m_pNetWork->_IRecord_Key_System(Record_SysType_TH, m_TechServerID); } TRAINEND te; te.CourseID = NetCmd.Net_JY_TrainIni.CourseID; te.SysID = NetCmd.Net_JY_TrainIni.SysType; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eTrainEnd 消息:CourseID=%d,SysID=%d,TrianID=%d",te.CourseID,te.SysID,NetCmd.TrainningID); StopDispatchTask(te.CourseID,te.SysID, NetCmd.TrainningID); DELTRAIN tb; tb.CourseID = NetCmd.Net_JY_TrainIni.CourseID; tb.SysID = NetCmd.Net_JY_TrainIni.SysType; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eTrainDel 消息:CourseID=%d,SysID=%d,TrianID=%d",tb.CourseID,tb.SysID,NetCmd.TrainningID); DeleteDispatchTask(tb.CourseID,tb.SysID, NetCmd.TrainningID); TRAININI& inc = NetCmd.Net_JY_TrainIni; InitCourInfo tmpcourseinfo; tmpcourseinfo.coursetype = NetCmd.cmdType; tmpcourseinfo.UniqueTrainingID = inc.UniqueTrainingID; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eTrainIni 消息:CourseID=%d,SysType=%d,TrianID=%d,RunLineID=%d,SchalTimeID=%d,begTime=%d",inc.CourseID,inc.SysType,NetCmd.TrainningID,inc.RunLineID,inc.SchalTimeID,inc.begTime); CreateDispatchTask(inc.CourseID,inc.SysType, NetCmd.TrainningID ,inc.RunLineID,inc.SchalTimeID,inc.begTime,&tmpcourseinfo); } else if ( NetCmd.cmdType == OCCCOMMAND::eTrainEnd ) { TRAINEND& te = NetCmd.Net_JY_TrainEnd; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eTrainEnd 消息:CourseID=%d,SysID=%d,TrianID=%d",te.CourseID,te.SysID,NetCmd.TrainningID); StopDispatchTask(te.CourseID,te.SysID, NetCmd.TrainningID); Send_TrainEnd_2Client(NetCmd); } else if ( NetCmd.cmdType == OCCCOMMAND::eTrainBeg ) { TRAINBEG& tb = NetCmd.Net_JY_TrainBeg; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eTrainBeg 消息:CourseID=%d,SysID=%d,TrianID=%d",tb.CourseID,tb.SysID,NetCmd.TrainningID); RunDispatchTask(tb.CourseID,tb.SysID, NetCmd.TrainningID); Send_TrainBeg_2Client(NetCmd); } else if ( NetCmd.cmdType == OCCCOMMAND::eTrainDelete ) { DELTRAIN& tb = NetCmd.Net_JY_TrainDel; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eTrainDel 消息:CourseID=%d,SysID=%d,TrianID=%d",tb.CourseID,tb.SysID,NetCmd.TrainningID); DeleteDispatchTask(tb.CourseID,tb.SysID, NetCmd.TrainningID); } else if ( NetCmd.cmdType == OCCCOMMAND::eSetUserRoll ) { TRAINSETUSERROLL& ur = NetCmd.Net_JY_TrainSetUserRoll; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eSetUserRoll 消息:SysType=%d,TrianID=%d,NetID=%d",ur.SysType,NetCmd.TrainningID,ur.NetID); RecTrianPosition(ur.SysType, NetCmd.TrainningID,ur.NetID); } else if ( NetCmd.cmdType == OCCCOMMAND::eTrainDeleteAll ) { TRAINDELALL& da = NetCmd.Net_JY_TrainiDelAll; OUT_LOG(LOG_LEVEL_NORMAL, "服务器收到删除所有命令!!!!消息:TrianID=%d!!!",da.TrianID); DeleteDispatchTask(0,OCC_SYS_XD_SVR,da.TrianID); DeleteDispatchTask(0,OCC_SYS_LS_SVR,da.TrianID); DeleteDispatchTask(0,OCC_SYS_DD_SVR,da.TrianID); DeleteDispatchTask(0,OCC_SYS_HD_SVR,da.TrianID); SVRREPLYMSG netret; memset(&netret,0,sizeof(netret)); netret.status=OCC_NETMESSAGE_TRIAN_ALLDELETE_SUCCESS; NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg,da.TrianID);//创建课程 OUT_LOG(LOG_LEVEL_NORMAL, "删除成功!!!!!"); } else if( NetCmd.cmdType == OCCCOMMAND::eTrainReStart ) { TRAINDELALL& da = NetCmd.Net_JY_TrainiDelAll; OUT_LOG(LOG_LEVEL_NORMAL, "服务器收到删除所有命令!!!!消息:TrianID=%d!!!", da.TrianID); DeleteDispatchTask(0,OCC_SYS_XD_SVR,da.TrianID); DeleteDispatchTask(0,OCC_SYS_LS_SVR,da.TrianID); DeleteDispatchTask(0,OCC_SYS_DD_SVR,da.TrianID); DeleteDispatchTask(0,OCC_SYS_HD_SVR,da.TrianID); DeleteDispatchTask(0,OCC_SYS_JT_SVR,da.TrianID); SVRREPLYMSG netret; memset(&netret,0,sizeof(netret)); netret.status = OCC_NETMESSAGE_TRIAN_ALLDELETE_SUCCESS; NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg,da.TrianID);//创建课程 OUT_LOG(LOG_LEVEL_NORMAL, "删除成功!!!!!"); #ifdef IN_WINDOWS TCHAR szPath[256]; GetModuleFileName(NULL,szPath,sizeof(szPath)); TCHAR *szcmdline = GetCommandLine(); STARTUPINFO StartUp; GetStartupInfo(&StartUp); PROCESS_INFORMATION info; BOOL bSuccessFlag = CreateProcess(szPath,szcmdline,NULL,NULL,FALSE, NORMAL_PRIORITY_CLASS,NULL,NULL,&StartUp,&info); if( bSuccessFlag ) { ExitProcess(0); ::PostQuitMessage(0); } #endif } else if ( NetCmd.cmdType==OCCCOMMAND::eTrainStart ) { TRAINSTART& ts = NetCmd.Net_JY_TrainiStart; OUT_LOG(LOG_LEVEL_NORMAL,"服务器收到eTrainStart 消息:CourseID=%d,TrianID=%d,RunLineID=%d,SchalTimeID=%d,begTime=%d",ts.CourseID,ts.TrianID,ts.RunLineID,ts.SchalTimeID,ts.begTime); InitCourInfo tmpcourseinfo; tmpcourseinfo.coursetype = OCCCOMMAND::eTrainStart; tmpcourseinfo.startallinfo = NetCmd.Net_JY_TrainiStart; if ( m_pNetWork ) { m_pNetWork->_ISetTechID (NetCmd.ScrNetID ); m_pNetWork->_ISetTrainID(NetCmd.TrainningID); } bool createflag = true; for ( int nIndex = 0; nIndex <100; nIndex++ ) { int systype = ts.UserRoll[nIndex]; nIndex++; int netid = ts.UserRoll[nIndex]; if (systype > 0 && netid >= 0) { RecTrianPosition(systype,ts.TrianID,netid); Send_SetUserRoll_Cmd(netid); } else break; } for ( int count = 0; count < 5; count++ ) //初始化课程数据 { if ( ts.SysType[count] > 0 ) { if (!CreateDispatchTask(ts.CourseID,ts.SysType[count],ts.TrianID ,ts.RunLineID,ts.SchalTimeID,ts.begTime,&tmpcourseinfo)) { createflag = false; } //有轨电车模块特殊处理 //RunDispatchTask(ts.CourseID,ts.SysType[count],ts.TrianID); } else break; } for ( int count = 0; count < 5;count++ ) //运行课程 { if ( ts.SysType[count] > 0 ) { //CreateDispatchTask(ts.CourseID,ts.SysType[count],ts.TrianID,ts.RunLineID,ts.SchalTimeID,ts.begTime); if (!RunDispatchTask(ts.CourseID,ts.SysType[count],ts.TrianID)) { createflag = false; } } else break; } if ( createflag == true ) { SVRREPLYMSG netret; memset(&netret,0,sizeof(netret)); netret.status=OCC_NETMESSAGE_TRIAN_ALLSTART_SUCCESS; NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg,ts.TrianID);//创建课程 OUT_LOG(LOG_LEVEL_NORMAL, "教员一键开始成功!!!"); } else { SVRREPLYMSG netret; memset(&netret,0,sizeof(netret)); netret.status=OCC_NETMESSAGE_TRIAN_ALLSTART_FAIL; NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg,ts.TrianID);//创建课程 OUT_LOG(LOG_LEVEL_NORMAL, "教员一键开始失败!!!"); } OUT_LOG(LOG_LEVEL_NORMAL, "一键开课命令处理完毕!!!!!!"); } else if ( NetCmd.cmdType == OCCCOMMAND::eTrainRollBackCtrl ) { DealWith_RollBack_CtrlCmd(NetCmd); } else { Deal_NetComand(NetCmd); } } catch(...) { OUT_LOG(LOG_LEVEL_ERROR,"系统初始化失败!!!!!"); } //OUT_LOG(LOG_LEVEL_NORMAL, "---------------------------课程管理结束!\n"); //OUT_LOG(LOG_LEVEL_NORMAL, ">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); return true; } /******************************************************* 函数名称: Send_TrainEnd_2Client 函数功能: 转发课程结束命令至客户端 输入参数: NetCmd - 原始命令 输出参数: NULL 返 回 值: void *******************************************************/ bool CBaseDispatchTaskMg::Send_TrainEnd_2Client( OCCCOMMAND NetCmd ) { int nType = NetCmd.Net_JY_TrainEnd.SysID; BOOL bSend = TRUE; switch (nType) { case OCC_SYS_XD_SVR: NetCmd.DesSysType = OCC_SYS_XD_CLIENT; break; case OCC_SYS_HD_SVR: NetCmd.DesSysType = OCC_SYS_HD_CLIENT; break; case OCC_SYS_DD_SVR: NetCmd.DesSysType = OCC_SYS_DD_CLIENT; break; default: bSend = FALSE; break; } if ( bSend ) { m_QI_SendList_TCP.push(NetCmd); } return true; } /******************************************************* 函数名称: Send_TrainBeg_2Client 函数功能: 转发课程开始命令至客户端 输入参数: NetCmd - 原始命令 输出参数: NULL 返 回 值: void *******************************************************/ bool CBaseDispatchTaskMg::Send_TrainBeg_2Client(OCCCOMMAND NetCmd) { int nType = NetCmd.Net_JY_TrainBeg.SysID; BOOL bSend = TRUE; switch (nType) { case OCC_SYS_XD_SVR: NetCmd.DesSysType = OCC_SYS_XD_CLIENT; break; case OCC_SYS_HD_SVR: NetCmd.DesSysType = OCC_SYS_HD_CLIENT; break; case OCC_SYS_DD_SVR: NetCmd.DesSysType = OCC_SYS_DD_CLIENT; break; default: bSend = FALSE; break; } if (bSend) { m_QI_SendList_TCP.push(NetCmd); } return true; } /******************************************************* 函数名称: Add_AssignCmd_List 函数功能: 添加命令分配列表 输入参数: nCmdType - 命令类型 Str_ModuleList - 模块列表 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::Add_AssignCmd_List(int nCmdType, std::string Str_ModuleList) { ASSIGN_CMD AssignCmd; AssignCmd.nCmdType = nCmdType; vector Vec_Temp = str_split_2(Str_ModuleList, ","); for ( size_t nIndex = 0; nIndex < Vec_Temp.size(); nIndex++ ) { int nModuleID = atoi(Vec_Temp[nIndex].c_str()); AssignCmd.VecModuleID.push_back(nModuleID); } m_AssignCmdMap.insert(std::make_pair(nCmdType,AssignCmd)); } /******************************************************* 函数名称: MakeCmd_DealModule 函数功能: 生成命令处理模块列表 输入参数: NetCmd - 网络命令 输出参数: VecModuleID - 命令处理模块列表 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::MakeCmd_DealModule(OCCCOMMAND NetCmd, vector& VecModuleID) { int nCmdType = NetCmd.cmdType; auto itor = m_AssignCmdMap.find(nCmdType); if ( itor != m_AssignCmdMap.end() ) { VecModuleID = itor->second.VecModuleID; } else { VecModuleID.push_back(NetCmd.DesSysType); } } /******************************************************* 函数名称: AssignCmd_ToModule 函数功能: 分发网络命令至处理模块 输入参数: NetCmd - 网络命令 VecModuleID - 命令模块列表 输出参数: NULL 返 回 值: void *******************************************************/ BOOL CBaseDispatchTaskMg::AssignCmd_ToModule(OCCCOMMAND NetCmd, vector VecModuleID) { BOOL bFlag = FALSE; R_Mutex_LOCK(m_MutexDispatchTask); for ( size_t nIndex = 0; nIndex < VecModuleID.size(); nIndex++ ) { string Str_Key = Make2IntKey_2(VecModuleID[nIndex],NetCmd.TrainningID); auto iter = m_DispatchTasksMap.find(Str_Key); if ( iter != m_DispatchTasksMap.end() ) { CBaseDispatch* pDispatchSys = iter->second; if ( pDispatchSys ) { pDispatchSys->ReceiveMsg(NetCmd); bFlag = TRUE; } } } return bFlag; } /******************************************************* 函数名称: NetComand_Receive 函数功能: 网络命令接收-处理+分发 输入参数: NetCmd - 网络命令 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::NetComand_Receive( OCCCOMMAND NetCmd ) { //1、数据只需要在服务端实现数据转发 switch (NetCmd.cmdType) { case OCCCOMMAND::eTelephoneTalkRecord: //通话记录协议 case OCCCOMMAND::eTelephoneSpeechRecognize: //通话语音识别 case OCCCOMMAND::eTelephone_Session: //语音电话_会话 case OCCCOMMAND::eTelephone_Intercom: //语音电话_对讲 { NetCmd.DesNetID = -1; NetCmd.ScrSysType = OCC_SYS_XD_SVR; NetCmd.DesSysType = OCC_SYS_XD_CLIENT; m_QI_SendList_TCP.push(NetCmd); return true; } break; } //2、处理并分发网络数据 if ( NetCmd.DesSysType == OCC_SYS_SERVER_SVR ) { return NetComandToServer(NetCmd); } else { //<1>、存储协同服务器编号 if ( NetCmd.cmdType == OCCCOMMAND::eSys_ShijingXietong && m_pNetWork ) { m_nViewSysSvrID = NetCmd.ScrNetID; m_pNetWork->_IRecord_Key_System(Record_SysType_View, m_nViewSysSvrID); } //<2>、生成命令处理模块ID vector VecModuleID; MakeCmd_DealModule( NetCmd, VecModuleID ); //<3>、将数据分发给各子模块 AssignCmd_ToModule( NetCmd, VecModuleID ); } return false; } /******************************************************* 函数名称: CreateDispatchTask 函数功能: 创建调度任务 输入参数: courseid - 课程号 systype - 系统类型 trainid - 训练号 runlineid - 运营线路编号 parm1 - parm2 - curseinfo - 课程信息 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::CreateDispatchTask( int courseid,int systype,int trainid ,int runlineid,int parm1,int parm2 ,InitCourInfo* curseinfo ) { m_nAllTrainStatus = AllTrain_Status_Init; if ( m_pNetWork && systype == OCC_SYS_XD_SVR ) { int nSimNumberTime = parm2 / 3600 * 10000 + parm2 % 3600 / 60 * 100 + parm2 % 60; m_pNetWork->_IStart_Training(m_nCurTrainID, m_nCurCourseID, nSimNumberTime); m_pNetWork->_IRecord_Training_Sys(Record_SysType_XD); m_pNetWork->_IRecord_Training_Sys(Record_SysType_DD); m_pNetWork->_IRecord_Training_Sys(Record_SysType_HD); } SVRREPLYMSG netret; memset(&netret,0,sizeof(netret)); netret.sysType = systype; bool ret = true; netret.status=OCC_NETMESSAGE_TRIAN_CREATE_FAIL; string key = Make2IntKey_2(systype,trainid); R_Mutex_LOCK(m_MutexDispatchTask); auto iter = m_DispatchTasksMap.find(key); if ( iter == m_DispatchTasksMap.end() ) { CBaseDispatch* pdispatchsys = CreateNewTaskType(systype); if ( !pdispatchsys ) { OUT_LOG(LOG_LEVEL_ERROR,"调度服务器创建: %s 训练号:%d 失败",GetSysTypeStr(systype).c_str(),trainid); ret = false; netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL; } else { RecTrianPosition(systype,trainid,m_NetServerID); auto iter = m_DataBaseConfigMap.find(systype); if ( iter == m_DataBaseConfigMap.end() ) { OUT_LOG(LOG_LEVEL_ERROR,"调度服务器创建: %s 训练号:%d 没有找到数据库配置项",GetSysTypeStr(systype).c_str(), trainid); ret = false; netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL; } else if (!pdispatchsys->OpenDB(iter->second.strSeverIP,iter->second.strDBName ,iter->second.strUserID,iter->second.strPassWord,courseid,runlineid,parm1,parm2,iter->second.nDB_Type,curseinfo)) { pdispatchsys->Exit(); delete pdispatchsys; pdispatchsys=NULL; OUT_LOG(LOG_LEVEL_ERROR,"调度服务器创建: %s 训练号:%d 打开数据库任务失败",GetSysTypeStr(systype).c_str(), trainid); ret=false; netret.status=OCC_NETMESSAGE_TRIAN_CREATE_FAIL; } else if (!pdispatchsys->OpenWork(trainid,systype,&m_QI_SendList_TCP,&m_QI_SendList_Group,&m_QI_SendList_UDP)) { pdispatchsys->Exit(); delete pdispatchsys; pdispatchsys = NULL; OUT_LOG(LOG_LEVEL_ERROR, "调度服务器创建: %s 训练号:%d 失败", GetSysTypeStr(systype).c_str(), trainid); ret=false; netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL; } else { m_DispatchTasksMap.insert(make_pair(key,pdispatchsys)); OUT_LOG(LOG_LEVEL_NORMAL, "调度服务器创建: %s 训练号:%d 成功", GetSysTypeStr(systype).c_str(), trainid); netret.status = OCC_NETMESSAGE_TRIAN_CREATE_SUCCESS; R_Mutex_LOCK(m_MutexTrainConect); NETTRIANNECTMAP::iterator iter = m_NetTrainNectMap.find(trainid); if ( iter == m_NetTrainNectMap.end() ) { NETTRIANDISTRI nettrian; nettrian.nTrainingID = trainid; nettrian.nTrainStartTime = parm2; m_NetTrainNectMap.insert(make_pair(trainid,nettrian)); } else { iter->second.nTrainStartTime = parm2; } AddReceiveGroupMsg(trainid,systype,pdispatchsys); } } } else { OUT_LOG(LOG_LEVEL_ERROR,"调度服务器创建: %s 训练号:%d 失败! 该训练已经存在", GetSysTypeStr(systype).c_str(), trainid); netret.status = OCC_NETMESSAGE_TRIAN_CREATE_EXIST; } Str_To_CharBuf_2("",netret.ReplyText,ALARM_LEN); //if (curseinfo->coursetype == OCCCOMMAND::eTrainIni) //适应老程序特殊处理 //{ // memcpy(&netret,&systype,sizeof(int)); // int tmpstatus = OCC_NETMESSAGE_TRIAN_CREATE_SUCCESS; // memcpy((char*)(&netret)+ sizeof(int),&tmpstatus,sizeof(int)); // /*netret.sysType = systype; // netret.status=OCC_NETMESSAGE_TRIAN_CREATE_SUCCESS;*/ //} if ( curseinfo->coursetype == OCCCOMMAND::eTrainIni && m_pNetWork) //其他流程单独发送 { //NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg,trainid);//创建课程 //修改为立即发送 OCCCOMMAND NetCmd; NetCmd.cmdType = OCCCOMMAND::eSvrReplyMsg; NetCmd.DesSysType = OCC_SYS_SERVER_TH; NetCmd.ScrSysType = OCC_SYS_SERVER_SVR; NetCmd.DesNetID = m_TechServerID; NetCmd.TrainningID = trainid; NetCmd.DateLen = sizeof(SVRREPLYMSG); NetCmd.setPacketValue(&netret,OCCCOMMAND::eSvrReplyMsg); m_pNetWork->_ISend_TCPData_RightNow(NetCmd); } return ret; } /******************************************************* 函数名称: DeleteDispatchTask 函数功能: 删除调度任务 输入参数: nCourseID - 课程号 nSysType - 系统类型 nTrainID - 训练号 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::DeleteDispatchTask( int nCourseID, int nSysType, int nTrainID ) { bool returnflag = false; SVRREPLYMSG netret; netret.sysType = nSysType; netret.status = OCC_NETMESSAGE_TRIAN_DELETE_SUCCESS; BOOL bFindFlag = FALSE; R_Mutex_LOCK(m_MutexDispatchTask); for ( auto iter = m_DispatchTasksMap.begin(); iter != m_DispatchTasksMap.end(); ) { CBaseDispatch* pDispatchSys = iter->second; if ( !pDispatchSys ) { m_DispatchTasksMap.erase(iter); return false; } else { if ( pDispatchSys->m_subsystype != nSysType ) { iter ++; continue; } else { nCourseID= pDispatchSys->m_courseid; nSysType = pDispatchSys->m_subsystype; nTrainID = pDispatchSys->m_trianid; } } //行调训练结束后取结束时间 if ( nSysType == OCC_SYS_XD_SVR ) { m_nEndTrainingTime = pDispatchSys->GetSim_NumberTime(); } bFindFlag = TRUE; if ( pDispatchSys->StopTask(netret.parm) ) { returnflag = true; pDispatchSys->Exit(); DeleteReceiveGroupMsg(nTrainID, nSysType); //Sleep(1000); delete pDispatchSys; pDispatchSys = NULL; netret.status = OCC_NETMESSAGE_TRIAN_DELETE_SUCCESS; OUT_LOG(LOG_LEVEL_ERROR,"调度服务器删除: %s 训练号:%d 成功",GetSysTypeStr(nSysType).c_str(),nTrainID); //DISTASKMAP::iterator tmpiter = iter; m_DispatchTasksMap.erase(iter); if ( nSysType == OCC_SYS_XD_SVR ) { R_Mutex_LOCK(m_MutexTrainConect); auto iter = m_NetTrainNectMap.find(nTrainID); if ( iter != m_NetTrainNectMap.end() ) { m_NetTrainNectMap.erase(iter); } } iter = m_DispatchTasksMap.begin(); } else netret.status = OCC_NETMESSAGE_TRIAN_DELETE_FAIL; } if( !bFindFlag ) { netret.status = OCC_NETMESSAGE_TRIAN_DELETE_CANNOETFIND; } NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg, nTrainID); if ( m_pDebugModule && m_pNetWork && (int)m_DispatchTasksMap.size() == 0 ) { if ( bFindFlag ) { m_pNetWork->_IEnd_Training(m_nCurTrainID, m_nCurCourseID,m_nEndTrainingTime); m_nCurTrainID = -1; m_nCurRBPointID = 1; m_nViewSysSvrID = -1; } m_pDebugModule->Reset_Debug_Tool(); m_nAllTrainStatus = AllTrain_Status_Stop; } return returnflag ; } /******************************************************* 函数名称: StopDispatchTask 函数功能: 停止调度任务 输入参数: nCourseID - 课程号 nSysType - 系统类型 nTrainID - 训练号 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::StopDispatchTask( int nCourseID,int nSysType,int nTrainID ) { SVRREPLYMSG netret; netret.sysType = nSysType; netret.status = OCC_NETMESSAGE_TRIAN_STOP_SUCCESS; string Str_Key = Make2IntKey_2(nSysType, nTrainID); R_Mutex_LOCK(m_MutexDispatchTask); auto iter = m_DispatchTasksMap.find(Str_Key); if ( iter != m_DispatchTasksMap.end() ) { //1、找不到子系统时直接返回 CBaseDispatch* pDispatchSys = iter->second; if ( !pDispatchSys ) { m_DispatchTasksMap.erase(iter); return false; } //2、找到系统后停止子系统的计算 if ( pDispatchSys->StopTask(netret.parm) ) { netret.status = OCC_NETMESSAGE_TRIAN_STOP_SUCCESS; } else { netret.status = OCC_NETMESSAGE_TRIAN_STOP_FAIL; } } else { netret.status=OCC_NETMESSAGE_TRIAN_STOP_CANNOETFIND; } //反馈停止任务结果给教员系统 NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg, nTrainID); return true; } /******************************************************* 函数名称: PauseDispatchTask 函数功能: 暂停调度任务 输入参数: nCourseID - 课程号 nSysType - 系统类型 nTrainID - 训练号 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::PauseDispatchTask( int nCourseID,int nSysType,int nTrainID ) { SVRREPLYMSG netret; netret.sysType = nSysType; netret.status = OCC_NETMESSAGE_TRIAN_PAUSE_SUCCESS; string Str_Key = Make2IntKey_2(nSysType, nTrainID); R_Mutex_LOCK(m_MutexDispatchTask); auto iter = m_DispatchTasksMap.find(Str_Key); if ( iter != m_DispatchTasksMap.end() ) { //1、找不到子系统时直接返回 CBaseDispatch* pDispatchSys = iter->second; if ( !pDispatchSys ) { m_DispatchTasksMap.erase(iter); return false; } //2、找到系统后暂停子系统的计算 if ( pDispatchSys->PauseTask(netret.parm) ) { netret.status = OCC_NETMESSAGE_TRIAN_PAUSE_SUCCESS; } else { netret.status = OCC_NETMESSAGE_TRIAN_PAUSE_FAIL; } } else { netret.status = OCC_NETMESSAGE_TRIAN_PAUSE_CANNOETFIND; } //反馈暂停任务结果给教员系统 NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg, nTrainID); return true; } /******************************************************* 函数名称: ResumeDispatchTask 函数功能: 恢复调度任务 输入参数: nCourseID - 课程号 nSysType - 系统类型 nTrainID - 训练号 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::ResumeDispatchTask( int nCourseID,int nSysType,int nTrainID ) { SVRREPLYMSG netret; netret.sysType = nSysType; netret.status = OCC_NETMESSAGE_TRIAN_RESUME_SUCCESS; string key = Make2IntKey_2(nSysType, nTrainID); R_Mutex_LOCK(m_MutexDispatchTask); auto iter = m_DispatchTasksMap.find(key); if ( iter != m_DispatchTasksMap.end() ) { //1、找不到子系统时直接返回 CBaseDispatch* pDispatchSys = iter->second; if ( !pDispatchSys ) { m_DispatchTasksMap.erase(iter); return false; } //2、找到系统后恢复子系统的计算 if ( pDispatchSys->ResumeTask(netret.parm) ) { netret.status = OCC_NETMESSAGE_TRIAN_RESUME_SUCCESS; } else { netret.status = OCC_NETMESSAGE_TRIAN_RESUME_FAIL; } } else { netret.status = OCC_NETMESSAGE_TRIAN_RESUME_CANNOETFIND; } //反馈恢复任务结果给教员系统 NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg, nTrainID); return true; } /******************************************************* 函数名称: RunDispatchTask 函数功能: 运行调度任务 输入参数: nCourseID - 课程号 nSysType - 系统类型 nTrainID - 训练号 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::RunDispatchTask( int nCourseID,int nSysType,int nTrainID ) { SVRREPLYMSG netret; netret.sysType = nSysType; netret.status = OCC_NETMESSAGE_TRIAN_RUN_SUCCESS; string key = Make2IntKey_2(nSysType, nTrainID); R_Mutex_LOCK(m_MutexDispatchTask); auto iter = m_DispatchTasksMap.find(key); if ( iter != m_DispatchTasksMap.end() ) { CBaseDispatch* pDispatchSys = iter->second; if ( !pDispatchSys ) { m_DispatchTasksMap.erase(iter); return false; } //if (pdispatchsys->m_RunStatus==DISPATCH_RUN_STATUS_INITED) { if ( !InitNetConnectStatus(nTrainID, nSysType) ) { return false; } } if ( pDispatchSys->StartTask(netret.parm) ) { netret.status = OCC_NETMESSAGE_TRIAN_RUN_SUCCESS; } else { netret.status = OCC_NETMESSAGE_TRIAN_RUN_FAIL; } } else { netret.status = OCC_NETMESSAGE_TRIAN_RUN_CANNOETFIND; } //NetComand_Send(&netret,OCC_SYS_SERVER_TH,OCCCOMMAND::eSvrReplyMsg,trainid);//创建课程 //修改为立即发送 OCCCOMMAND NetCmd; NetCmd.cmdType = OCCCOMMAND::eSvrReplyMsg; NetCmd.DesSysType = OCC_SYS_SERVER_TH; NetCmd.ScrSysType = OCC_SYS_SERVER_SVR; NetCmd.DesNetID = m_TechServerID; NetCmd.TrainningID = nTrainID; NetCmd.DateLen = sizeof(SVRREPLYMSG); NetCmd.setPacketValue(&netret,OCCCOMMAND::eSvrReplyMsg); m_pNetWork->_ISend_TCPData_RightNow(NetCmd); return true; } /******************************************************* 函数名称: RecTrianPosition 函数功能: 接收训练系统位置 输入参数: nSysType - 系统类型 trainid - 训练号 sysid - 系统号 输出参数: NULL 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::RecTrianPosition( int nSysType,int nTrainID,int nSysID ) { R_Mutex_LOCK(m_MutexTrainConect); auto iter = m_NetTrainNectMap.find(nTrainID); if ( iter == m_NetTrainNectMap.end() ) { NETTRIANDISTRI nettrian; nettrian.nTrainingID = nTrainID; nettrian.AddSys(nSysType,nSysID); m_NetTrainNectMap.insert(make_pair(nTrainID,nettrian)); } else { iter->second.AddSys(nSysType,nSysID); } return true; } /******************************************************* 函数名称: GetTaskStatusData 函数功能: 获取任务状态数据 输入参数: statusdada - 状态数据 输出参数: statusdada - 状态数据 返回值 : bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatchTaskMg::GetTaskStatusData() { return false; } /******************************************************* 函数名称: GetPlugin_ByModuleID 函数功能: 通过模块号获取模块加载对象指针 输入参数: nModuleID - 模块号 输出参数: NULL 返回值 : sPlugin* - 模块加载对象指针 *******************************************************/ sPlugin* CBaseDispatchTaskMg::GetPlugin_ByModuleID(int nModuleID) { sPlugin* pPlugin = NULL; auto itor = m_PluginMap.find(nModuleID); if ( itor != m_PluginMap.end() ) { pPlugin = &itor->second; } return pPlugin; } /******************************************************* 函数名称: DealWith_RollBack_CtrlCmd 函数功能: 处理回滚控制命令 输入参数: NetCmd - 网络命令 输出参数: NULL 返回值 : void *******************************************************/ void CBaseDispatchTaskMg::DealWith_RollBack_CtrlCmd( OCCCOMMAND NetCmd ) { int nTrainID = NetCmd.Net_JY_RollBackCtrl.nTrainID; int nPointID = NetCmd.Net_JY_RollBackCtrl.nPointID; int nCtrlType = NetCmd.Net_JY_RollBackCtrl.nCtrlType; BOOL bFeedBack = TRUE; BOOL bResult = TRUE; //1、执行回滚控制指令 switch ( nCtrlType ) { case RollCtrl_Type_Prepare: //回滚准备 { bResult = RollBack_Prepare(nTrainID, nPointID); } break; case RollCtrl_Type_Excute: //回滚执行 { bResult = RollBack_Excute(nTrainID, nPointID); } break; case RollCtrl_Type_Finish: //回滚完成 { bResult = RollBack_Finish(nTrainID, nPointID); } break; case RollCtrl_Type_Save: //回滚保存 { if ( !m_bRollBackEvent ) { int nSize = min(sizeof(m_cRollBackDes), sizeof(NetCmd.Net_JY_RollBackCtrl.cPointDes)); memcpy(m_cRollBackDes, NetCmd.Net_JY_RollBackCtrl.cPointDes, nSize); m_bRollBackEvent = TRUE; } bFeedBack = FALSE; } break; case RollCtrl_Type_Event: { RollBack_Prepare(nTrainID, nPointID); RollBack_Excute (nTrainID, nPointID); RollBack_Finish (nTrainID, nPointID); bFeedBack = FALSE; } break; default: break; } //2、反馈回滚控制命令执行状态 if ( bFeedBack ) { int nResult = ( bResult ? 1 : 0 ); FeedBack_RollBack_Cmd(nTrainID, nPointID, nCtrlType, nResult); } } /******************************************************* 函数名称: LoadSub_DispatchSys 函数功能: 加载调度子系统 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::LoadSub_DispatchSys() { #ifdef IN_WINDOWS std::string Str_FileName = GetFolder_Module_2() + "\\SYS_CONFIG\\Module.xls"; #else std::string Str_FileName = GetFolder_Module_2() + "SYS_CONFIG/Module.xls"; #endif // IN_WINDOWS YD_ExcelReader TmpXlsRead; TmpXlsRead.LoadSheet(Str_FileName,"模块列表"); for ( int nIndex = 1; nIndex < TmpXlsRead.GetRowCount(); nIndex++ ) { #ifdef IN_WINDOWS std::string Str_LibName = GetFolder_Module_2() + "\\" + TmpXlsRead.Cell_AsString(nIndex,1); sPlugin Plugin; Plugin.nModuleID = TmpXlsRead.Cell_AsInt(nIndex,0); Plugin.nFileDll = ::LoadLibrary(Str_LibName.c_str()); Plugin.nCreatePlug = (_CreateSubModule)GetProcAddress(Plugin.nFileDll, "Create_Sub_Dispatch"); Plugin.nFreePlug = (_FreeSubModule)GetProcAddress (Plugin.nFileDll, "Free_Sub_Dispatch"); Plugin.Str_ModuleName = TmpXlsRead.Cell_AsString(nIndex,2); Plugin.Check_Is_OK(); m_PluginMap.insert(std::make_pair(Plugin.nModuleID,Plugin)); OUT_LOG(LOG_LEVEL_NORMAL, "%s加载成功:%s", \ TmpXlsRead.Cell_AsString(nIndex,2).c_str(), \ TmpXlsRead.Cell_AsString(nIndex,1).c_str()); #else std::string Str_LibName = GetFolder_Module_2() + TmpXlsRead.Cell_AsString(nIndex, 1); sPlugin Plugin; Plugin.nModuleID = TmpXlsRead.Cell_AsInt(nIndex, 0); Plugin.nFileDll = dlopen(Str_LibName.c_str(), RTLD_LAZY); Plugin.nCreatePlug = (_CreateSubModule)dlsym(Plugin.nFileDll, "Create_Sub_Dispatch"); Plugin.nFreePlug = (_FreeSubModule) dlsym(Plugin.nFileDll, "Free_Sub_Dispatch"); Plugin.Str_ModuleName = coding_conv_ns::gbk_to_utf8_str(TmpXlsRead.Cell_AsString(nIndex,2).c_str()); Plugin.Check_Is_OK(); m_PluginMap.insert(std::make_pair(Plugin.nModuleID, Plugin)); OUT_LOG(LOG_LEVEL_NORMAL, "%s加载成功:%s", \ TmpXlsRead.Cell_AsString(nIndex, 2).c_str(), \ TmpXlsRead.Cell_AsString(nIndex, 1).c_str()); #endif } } /******************************************************* 函数名称: RollBack_Record 函数功能: 回滚记录 输入参数: NULL 输出参数: NULL 返 回 值: BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatchTaskMg::RollBack_Record() { //1、打开回滚数据记录数据库 if ( !m_RollBackTool.Open_SaveDB(m_nCurTrainID,m_nCurRBPointID) ) { return FALSE; } //2、记录回滚数据 BOOL bFlag = TRUE; if ( true ) { R_Mutex_LOCK(m_MutexDispatchTask); for ( auto var : m_DispatchTasksMap ) { //(1)、对象指针为空时直接继续 CBaseDispatch* pDispatchSys = var.second; if ( !pDispatchSys ) continue; //(2)、记录数据并返回记录结果 if (!pDispatchSys->RollBack_Record(&m_RollBackTool)) { bFlag = FALSE; } } } //3、记录成功的前提下将回滚节点信息存入数据库 CBaseDispatch* pSubDispatch = GetSub_Dispatch(OCC_SYS_XD_SVR); if ( bFlag && pSubDispatch ) { int nDate = Get_Number_Date_2(); int nSysTime = GetNumTime_Second_2(); int nSimTime = pSubDispatch->GetSim_NumberTime(); std::string Str_Sql = Str_Format("insert into RollBack_Record ([TrainID],[PointID],[CourseID],[Date],[SimTime],[SysTime],[Description]) values(%d,%d,%d,%d,%d,%d,'%s')", m_nCurTrainID, m_nCurRBPointID, m_nCurCourseID, nDate, nSimTime, nSysTime, m_cRollBackDes); pSubDispatch->Excute_Sql(Str_Sql); } return bFlag; } /******************************************************* 函数名称: RollBack_Excute 函数功能: 回滚执行 输入参数: nTrainID - 训练号 nPointID - 节点号 输出参数: NULL 返回值 : BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatchTaskMg::RollBack_Excute( int nTrainID, int nPointID ) { //1、先检查回滚执行命令参数是否正确 if ( !m_RollBackTool.Check_IsRight_RollBackCmd(nTrainID,nPointID) ) { return FALSE; } //2、执行回滚执行命令 BOOL bFlag = TRUE; R_Mutex_LOCK(m_MutexDispatchTask); for ( auto var : m_DispatchTasksMap ) { //1、对象指针为空时直接继续 CBaseDispatch* pDispatchSys = var.second; if ( !pDispatchSys ) { continue; } //2、执行回滚操作并返回执行结果 if (!pDispatchSys->RollBack_Excute(&m_RollBackTool)) { bFlag = FALSE; } } return bFlag; } /******************************************************* 函数名称: RollBack_Prepare 函数功能: 回滚准备 输入参数: nTrainID - 训练号 nPointID - 节点号 输出参数: NULL 返回值 : BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatchTaskMg::RollBack_Prepare( int nTrainID, int nPointID ) { //1、设置回滚信息并打开数据库 m_RollBackTool.SetRB_TrainID(nTrainID); m_RollBackTool.SetRB_PointID(nPointID); if (!m_RollBackTool.Open_LoadDB(nTrainID, nPointID)) { return FALSE; } //2、准备回滚 BOOL bFlag = TRUE; R_Mutex_LOCK(m_MutexDispatchTask); for ( auto var : m_DispatchTasksMap ) { //1)、对象指针为空时直接继续 CBaseDispatch* pDispatchSys = var.second; if ( !pDispatchSys ) { continue; } //2)、准备回滚并返回准备结果 if ( !pDispatchSys->RollBack_Prepare(&m_RollBackTool) ) { bFlag = FALSE; } } return bFlag; } /******************************************************* 函数名称: RollBack_Finish 函数功能: 回滚-完成 输入参数: nTrainID - 训练号 nPointID - 节点号 输出参数: NULL 返回值 : BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatchTaskMg::RollBack_Finish(int nTrainID, int nPointID) { //1、先检查回滚执行命令参数是否正确 if (!m_RollBackTool.Check_IsRight_RollBackCmd(nTrainID, nPointID)) { return FALSE; } //2、执行回滚完成命令 BOOL bFlag = TRUE; R_Mutex_LOCK(m_MutexDispatchTask); for ( auto var : m_DispatchTasksMap ) { //1、对象指针为空时直接继续 CBaseDispatch* pDispatchSys = var.second; if ( !pDispatchSys ) { continue; } //2、准备回滚并返回准备结果 if (!pDispatchSys->RollBack_Finish(&m_RollBackTool)) { bFlag = FALSE; } } return TRUE; } /******************************************************* 函数名称: RollBack_RecordThread 函数功能: 回滚数据记录线程响应函数 输入参数: NULL 输出参数: NULL 返回值 : BOOL - TRUE:记录成功 FALSE:记录失败 *******************************************************/ BOOL CBaseDispatchTaskMg::RollBack_RecordThread() { //1、记录条件不满足时直接返回FALSE if ( !CheckIs_Record_RollBack() ) { return FALSE; } //2、记录回滚数据 m_dLastRecordTime = GetNowTime(TIMEUNIT_S); BOOL bFlag = RollBack_Record(); //3、反馈回滚记录状态 int nResult = (bFlag ? 1 : 0); FeedBack_RollBack_Cmd(m_nCurTrainID, m_nCurRBPointID, RollCtrl_Type_Save, nResult); m_nCurRBPointID++; return bFlag; } /******************************************************* 函数名称: CheckIs_Record_RollBack 函数功能: 检查是否记录回滚文件 输入参数: NULL 输出参数: NULL 返回值 : BOOL - TRUE:是 FALSE:否 *******************************************************/ BOOL CBaseDispatchTaskMg::CheckIs_Record_RollBack() { BOOL bFlag = FALSE; //1、时间因素,隔一分钟记录一次,或者跨天跨月 double dSysSecondTime = GetNowTime(TIMEUNIT_S); if ( m_dLastRecordTime > 0 && dSysSecondTime - m_dLastRecordTime >= 60.0) { //bFlag = TRUE; } if ( dSysSecondTime < m_dLastRecordTime - 600 ) { //bFlag = TRUE; } //2、初始情况下记录一次 if ( m_dLastRecordTime <= 0 ) { //bFlag = TRUE; } //3、特定事件发生时记录一次 if ( m_bRollBackEvent ) { bFlag = TRUE; m_bRollBackEvent = FALSE; //消除事件 } return bFlag; } /******************************************************* 函数名称: GetSub_Dispatch 函数功能: 获取子调度系统 输入参数: nSysType - 系统类型 输出参数: NULL 返回值 : CBaseDispatch* - 子调度系统指针 *******************************************************/ CBaseDispatch* CBaseDispatchTaskMg::GetSub_Dispatch(int nSysType) { CBaseDispatch* pSubDispatch = NULL; R_Mutex_LOCK(m_MutexDispatchTask); string Str_Key = Make2IntKey_2(nSysType,m_nCurTrainID); auto itor = m_DispatchTasksMap.find(Str_Key); if ( itor != m_DispatchTasksMap.end() ) { pSubDispatch = itor->second; } return pSubDispatch; } /******************************************************* 函数名称: FeedBack_RollBack_Cmd 函数功能: 反馈回滚命令 输入参数: nTrainID - 训练号 nPointID - 节点号 nType - 类型 nResult - 结果 输出参数: NULL 返回值 : void *******************************************************/ void CBaseDispatchTaskMg::FeedBack_RollBack_Cmd(int nTrainID, int nPointID, int nType, int nResult) { TRAIN_ROLLBACK_FEEDBACK RollBackFeedBack; RollBackFeedBack.nTrainID = nTrainID; RollBackFeedBack.nPointID = nPointID; RollBackFeedBack.nCtrlType = nType; RollBackFeedBack.nResult = nResult; NetComand_Send(&RollBackFeedBack, OCC_SYS_SERVER_TH, \ OCCCOMMAND::eTrainRollBackFeedBack,m_nCurTrainID); } /******************************************************* 函数名称: Send_SetUserRoll_Cmd 函数功能: 发送分配角色指令 输入参数: nSysID - 系统号 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::Send_SetUserRoll_Cmd( int nSysID ) { //1、网络指针未空时直接返回 if ( !m_pNetWork ) return; //2、模拟联合教员向客户端发送分配角色命令 TRAINSETUSERROLL TrainSetUserRool; TrainSetUserRool.NetID = m_pNetWork->_IGet_System_ID(); TrainSetUserRool.SysType = OCC_SYS_XD_SVR; OCCCOMMAND NetCmd; NetCmd.cmdType = OCCCOMMAND::eSetUserRoll; NetCmd.DesSysType = OCC_SYS_XD_CLIENT; NetCmd.ScrSysType = OCC_SYS_SERVER_SVR; NetCmd.DesNetID = nSysID; NetCmd.ScrNetID = m_pNetWork->_IGet_System_ID(); NetCmd.TrainningID = m_pNetWork->_IGetTrainID(); NetCmd.DateLen = sizeof(TRAINSETUSERROLL); NetCmd.setPacketValue(&TrainSetUserRool,NetCmd.cmdType); m_pNetWork->_ISend_TCPData_RightNow(NetCmd); } /******************************************************* 函数名称: ThreadFun_DealTCPData 函数功能: 线程函数_处理TCP数据 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::ThreadFun_DealTCPData() { m_ReceiveMsgThreadOut = FALSE; try { OCCCOMMAND NetCmd; while ( !m_ReceiveMsgThreadOut ) { if ( GetRecNetMsg(NetCmd) ) { NetComand_Receive(NetCmd); } else { Thread_Sleep(TIMESPAN_LVN); } } } catch (...) { OUT_LOG(LOG_LEVEL_ERROR,"NetMsgDealThread \ Have a Error"); } } /******************************************************* 函数名称: ThreadFun_DealUDPData 函数功能: 线程函数_处理UDP数据 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::ThreadFun_DealUDPData() { m_ReceiveRealMsgThreadOut = FALSE; try { OCCCOMMAND NetCmd; while (!m_ReceiveRealMsgThreadOut) { if (GetRecNetRealMsg(NetCmd)) { NetComand_Receive(NetCmd); } else { Thread_Sleep(TIMESPAN_LVN); } } } catch (...) { OUT_LOG(LOG_LEVEL_ERROR,"NetRealMsgDealThread \ Have a Error"); } } /******************************************************* 函数名称: ThreadFun_DealGroupData 函数功能: 线程函数_处理组播数据 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::ThreadFun_DealGroupData() { m_ReceiveGroupMsgThreadOut = FALSE; try { OCCCOMMAND NetCmd; while ( !m_ReceiveGroupMsgThreadOut ) { if ( GetRecNetGroupMsg(NetCmd)) { NetComand_Receive(NetCmd); } else { Thread_Sleep(TIMESPAN_LVN); } } } catch (...) { OUT_LOG(LOG_LEVEL_ERROR, "NetGroupMsgDealThread \ Have a Error"); } } /******************************************************* 函数名称: ThreadFun_RollBackRecord 函数功能: 线程函数_处理回滚记录 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatchTaskMg::ThreadFun_RollBackRecord() { m_RollBackRecordThreadOut = FALSE; try { while ( !m_RollBackRecordThreadOut ) { if ( !RollBack_RecordThread() ) { Thread_Sleep(TIMESPAN_LVN); } } } catch (...) { OUT_LOG(LOG_LEVEL_ERROR, "回滚记录错误"); } }