#include "pch.h" #include #include "BaseDispatch.h" #include "GlobalVarDefine.h" #include "SubSys_RollBack_I.h" //构造函数 CBaseDispatch::CBaseDispatch(void) { m_pQI_SendCmdsList = nullptr; m_pQI_SendGroupList = nullptr; m_pQI_SendRealList = nullptr; m_pFrmTool = nullptr; m_pNetWork = nullptr; m_pDebugModule = nullptr; m_pRedisClient = nullptr; m_pDispatchMngI = nullptr; m_pDataBaseOper = nullptr; m_EventF1ThreadTimeDur = -1; m_EventF2ThreadTimeDur = -1; m_EventBeginThreadTimeDur = -1; m_ObjectCPThreadTimeDur = -1; m_dInitTime = -1; m_bF1ThreadFinish = true; m_bF2ThreadFinish = true; m_bBeginThreadFinish = true; m_bObjectCPThreadFinish = true; m_dTaskStartTime = 0.; m_nRunLineID = -1; m_courseid = -1; InitDate(); } //析构函数 CBaseDispatch::~CBaseDispatch(void) { //ReleaseDate(); if ( m_pDataBaseOper ) { delete m_pDataBaseOper; m_pDataBaseOper = nullptr; } } /******************************************************* 函数名称: ReceiveMsg 函数功能: 接收网络消息 输入参数: NetCmd - 网络命令 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::ReceiveMsg( OCCCOMMAND& NetCmd ) { m_QI_GetCommands.push(NetCmd); } /******************************************************* 函数名称: OpenDB 函数功能: 打开数据库 输入参数: strIP - 数据库IP strBDname - 数据库名称 struser - 用户名 strpsw - 密码 courseid - 课程号 runlineid - 运营线路编号 parm1 - 参数1 parm2 - 参数2 curseinfo - 课程初始化信息 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::OpenDB( std::string Str_IP,std::string Str_Name,std::string Str_User ,std::string Str_Psw,int nCourseID,int nRunLineID,int nPara1 ,int nPara2,int nDB_Type,InitCourInfo* pCourseInfo) { m_dTaskStartTime = GetNowTime(TIMEUNIT_S); pCourseInfo->startallinfo.begTime = nPara2; m_nRunLineID = nRunLineID; if ( !m_pDataBaseOper ) { m_pDataBaseOper = new YD_DBOperator(nDB_Type); } if ( !m_pDataBaseOper ) { OUT_LOG(LOG_LEVEL_NORMAL,"【M_%d,T_%d】构建数据库对象失败!!!\n",m_subsystype,m_trianid); return false; } if ( m_pDataBaseOper->ConnectDB(Str_IP,Str_User,Str_Psw,Str_Name) == 0 ) { OUT_LOG(LOG_LEVEL_ERROR, "【M_%d,T_%d】打开数据库失败!!!\n",m_subsystype,m_trianid); return false; } else { OUT_LOG(LOG_LEVEL_NORMAL,"【M_%d,T_%d】打开数据库成功!!!\n",m_subsystype,m_trianid); } ReadCourInfo(nCourseID,nPara1,nPara2,pCourseInfo); OUT_LOG(LOG_LEVEL_NORMAL, "【M_%d,T_%d】应用数据读取开始!!!\n",m_subsystype,m_trianid); if (!InitSubSysDBDate(nCourseID,nPara1,nPara2,pCourseInfo)) { OUT_LOG(LOG_LEVEL_ERROR, "【M_%d,T_%d】读取应用数据失败!!!\n",m_subsystype,m_trianid); return false; } else { OUT_LOG(LOG_LEVEL_NORMAL, "【M_%d,T_%d】读取应用数据成功!!!\n",m_subsystype,m_trianid); } return true; } /******************************************************* 函数名称: ReadCourInfo 函数功能: 读取课程信息 输入参数: nCourseID - 课程编号 nPara1 - 参数1 nPara2 - 参数2 pCourseInfo - 初始化课程信息 输出参数: NULL 返 回 值: BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatch::ReadCourInfo( int nCourseID,int nPara1,int nPara2, InitCourInfo* pCourseInfo ) { BOOL bReadFlag = TRUE; m_curseinfo = *pCourseInfo; //TODO:更多公用课程信息可在此处理 return bReadFlag; } /******************************************************* 函数名称: OpenWork 函数功能: 打开工作任务 输入参数: nTrainID - 训练号 nSysType - 系统类型 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::OpenWork( int nTrainID,int nSysType, CommandQueue* pTCPSendList, CommandQueue* pGroupSendList, CommandQueue* pUDPSendList ) { m_trianid = nTrainID; m_subsystype = nSysType; if (!CreateWorkThread()) { OUT_LOG(LOG_LEVEL_ERROR, "【M_%d,T_%d】调度任务开启工作线程失败",m_subsystype,m_trianid); return false; } else { OUT_LOG(LOG_LEVEL_NORMAL,"【M_%d,T_%d】调度任务开启工作线程成功!!!\n",m_subsystype,m_trianid); } if ( !pTCPSendList ) { return false; } m_pQI_SendCmdsList = pTCPSendList; m_pQI_SendRealList = pUDPSendList; m_pQI_SendGroupList = pGroupSendList; if ( !InitInstance() ) { OUT_LOG(LOG_LEVEL_ERROR,"【M_%d,T_%d】应用数据初始化失败!!!\n",m_subsystype,m_trianid); return false; } else { OUT_LOG(LOG_LEVEL_NORMAL,"【M_%d,T_%d】应用数据初始化成功!!!\n",m_subsystype,m_trianid); } OUT_LOG(LOG_LEVEL_NORMAL, "【M_%d,T_%d】调度任务开启工作成功",m_subsystype,m_trianid); m_dInitTime = GetNowTime(TIMEUNIT_S) - m_dTaskStartTime; return true; } /******************************************************* 函数名称: AddLogItem 函数功能: 添加日志条目 输入参数: strvalue - 日志文本 Pri - 日志等级 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::AddLogItem( std::string strvalue, int Pri ) { if ( m_pFrmTool ) { OUT_LOG(Pri, "【M_%d,T_%d】%s",m_subsystype,m_trianid,strvalue.c_str()); } } /******************************************************* 函数名称: Excute_Sql 函数功能: 指定模块中执行SQL语句 输入参数: Str_Sql - SQL脚本 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::Excute_Sql( std::string Str_Sql ) { if ( m_pDataBaseOper && m_pDataBaseOper->isOpened() ) { m_pDataBaseOper->DirectExecute(Str_Sql); } } /******************************************************* 函数名称: Exit 函数功能: 子系统退出 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::Exit() { m_RunStatus = DISPATCH_RUN_STATUS_STOP; Wait_AllThread_Finish(20); ReleaseDate(); ExitInstance(); Add_Module_Watch(); if ( m_pDataBaseOper ) { m_pDataBaseOper->CloseDB(); } } /******************************************************* 函数名称: Dispatch_ModuleInit 函数功能: 调度模块初始化 输入参数: pDebugModule - 调试工具指针 pNetWork - 网络工具指针 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::Dispatch_ModuleInit( CDispatchMng_I* pDispatchMngI ) { m_pDispatchMngI = pDispatchMngI; if ( m_pDispatchMngI ) { m_pFrmTool = (CFrameTool_I*) m_pDispatchMngI->_IGetPtr_FrameTool(); m_pNetWork = (CNetWorkMng_I*) m_pDispatchMngI->_IGetPtr_NetWorkModule(); m_pDebugModule = (CDebug_Module_I*)m_pDispatchMngI->_IGetPtr_DebugModule(); m_pRedisClient = (YD_RedisClient*) m_pDispatchMngI->_IGetPtr_RedisModule(); } } /******************************************************* 函数名称: CreateWorkThread 函数功能: 创建工作线程 输入参数: NULL 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::CreateWorkThread() { m_Thread_EventF1 = std::thread(std::bind(&CBaseDispatch::ThreadFun_EventF1, this)); m_Thread_EventF2 = std::thread(std::bind(&CBaseDispatch::ThreadFun_EventF2, this)); m_Thread_EventBegin = std::thread(std::bind(&CBaseDispatch::ThreadFun_EventBegin,this)); m_Thread_ObjectCP = std::thread(std::bind(&CBaseDispatch::ThreadFun_ObjectCP, this)); return true; } /******************************************************* 函数名称: ReleaseDate 函数功能: 释放数据 输入参数: NULL 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::ReleaseDate() { m_EventF1ThreadOut = true; m_EventF2ThreadOut = true; m_EventBeginThreadOut = true; m_ObjectCPThreadOut = true; //1、F1线程退出 if ( m_Thread_EventF1.joinable() ) { m_Thread_EventF1.join(); } //2、F2线程退出 if ( m_Thread_EventF2.joinable() ) { m_Thread_EventF2.join(); } //3、实时计算线程退出 if ( m_Thread_EventBegin.joinable() ) { m_Thread_EventBegin.join(); } //4、实时计算线程退出 if ( m_Thread_ObjectCP.joinable() ) { m_Thread_ObjectCP.join(); } OUT_LOG(LOG_LEVEL_ERROR, "【M_%d,T_%d】调度任务清除资源",m_subsystype,m_trianid); return true; } /******************************************************* 函数名称: InitDate 函数功能: 初始化数据 输入参数: NULL 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::InitDate() { m_trianid = 0; m_subsystype = 0; m_EventF1ThreadOut = true; m_EventF2ThreadOut = true; m_EventBeginThreadOut = true; m_ObjectCPThreadOut = true; m_RunStatus = DISPATCH_RUN_STATUS_STOP; m_QI_GetCommands.init(); //初始化接受表 m_QI_GetCommands.setSafeCount(10000); std::string strTmp = Str_Format("系统类型:%d--接收队列",m_subsystype); m_QI_GetCommands.SetQueName(strTmp); _QI_eventF1List.init(); _QI_eventF2List.init(); m_pQI_SendCmdsList = nullptr; return true; } /******************************************************* 函数名称: Add_Module_Watch 函数功能: 添加模块监视数据 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::Add_Module_Watch() { if ( !m_pDebugModule ) return; SUB_MODULE SubModule; SubModule.nID = m_subsystype; SubModule.nStatus = m_RunStatus; SubModule.dTime_Init = m_dInitTime; SubModule.dTime_F1 = m_EventF1ThreadTimeDur; SubModule.dTime_F2 = m_EventF2ThreadTimeDur; SubModule.dTime_CP = m_ObjectCPThreadTimeDur; SubModule.dTime_Event = m_EventBeginThreadTimeDur; SubModule.bFinish_F1 = m_bF1ThreadFinish; SubModule.bFinish_F2 = m_bF2ThreadFinish; SubModule.bFinish_CP = m_bObjectCPThreadFinish; SubModule.bFinish_Event = m_bBeginThreadFinish; SubModule.nUDPSendCnt = (m_pQI_SendRealList ? m_pQI_SendRealList->GetCur_Count() :-1); SubModule.nUDPRecvCnt = (m_pDispatchMngI ? m_pDispatchMngI->_IGetCnt_RecvQue_UDP() :-1); SubModule.nTCPSendCnt = (m_pQI_SendCmdsList ? m_pQI_SendCmdsList->GetCur_Count() :-1); SubModule.nTCPRecvCnt = (m_pDispatchMngI ? m_pDispatchMngI->_IGetCnt_RecvQue_TCP() :-1); SubModule.nSelfRecvCnt = m_QI_GetCommands.GetCur_Count(); m_pDebugModule->Update_SubModule(SubModule); } /******************************************************* 函数名称: RollBack_Prepare 函数功能: 回滚准备 输入参数: pRollBack - 回滚工具指针 输出参数: NULL 返 回 值: BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatch::RollBack_Prepare( CRollBack_ITool* pRollBack ) { //1、设置运行状态为暂停状态 m_RunStatus = DISPATCH_RUN_STATUS_PAUSE; //2、等待全部线程计算完成 if (!Wait_AllThread_Finish()) { m_RunStatus = DISPATCH_RUN_STATUS_RUNNING; return FALSE; } //3、子模块做回滚准备 BOOL bFlag = TRUE; for ( auto var : m_RollBackSysMap ) { if (!var.second->RollBack_Prepare(pRollBack)) { bFlag = FALSE; } } return bFlag; } /******************************************************* 函数名称: RollBack_Excute 函数功能: 回滚执行 输入参数: pRollBack - 回滚工具指针 输出参数: NULL 返 回 值: BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatch::RollBack_Excute (CRollBack_ITool* pRollBack) { BOOL bFlag = TRUE; for ( auto var : m_RollBackSysMap ) { if (!var.second->RollBack_Excute(pRollBack)) { bFlag = FALSE; } } return bFlag; } /******************************************************* 函数名称: RollBack_Finish 函数功能: 回滚完成 输入参数: pRollBack - 回滚工具指针 输出参数: NULL 返 回 值: BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatch::RollBack_Finish(CRollBack_ITool* pRollBack) { BOOL bFlag = TRUE; for ( auto var : m_RollBackSysMap ) { if (!var.second->RollBack_Finish(pRollBack)) { bFlag = FALSE; } } m_RunStatus = DISPATCH_RUN_STATUS_RUNNING; return bFlag; } /******************************************************* 函数名称: RollBack_Record 函数功能: 回滚记录 输入参数: pRollBack - 回滚工具指针 输出参数: NULL 返 回 值: BOOL - TRUE:成功 FALSE:失败 *******************************************************/ BOOL CBaseDispatch::RollBack_Record (CRollBack_ITool* pRollBack) { BOOL bFlag = TRUE; for ( auto var : m_RollBackSysMap ) { if (!var.second->RollBack_Record(pRollBack)) { bFlag = FALSE; } } return bFlag; } /******************************************************* 函数名称: Wait_AllThread_Finish 函数功能: 等待全部线程计算完成 输入参数: nTimeOut_T - 超时时间(秒) 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::Wait_AllThread_Finish( int nTimeOut_T ) { int nStartTime = (int)GetNowTime(TIMEUNIT_S); bool bFlag = ( m_bF1ThreadFinish && m_bF2ThreadFinish && m_bBeginThreadFinish && m_bObjectCPThreadFinish ); while ( !bFlag ) { int nTempTime = (int)GetNowTime(TIMEUNIT_S); if ( nTempTime - nStartTime > nTimeOut_T ) //NOTE:超时直接返回 { return false; } bFlag = ( m_bF1ThreadFinish && m_bF2ThreadFinish && m_bBeginThreadFinish && m_bObjectCPThreadFinish); } return true; } /******************************************************* 函数名称: Send_TrainReset_Cmd 函数功能: 发送训练复位命令 输入参数: NULL 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::Send_TrainReset_Cmd() { int nDesSysType = -1; switch ( m_subsystype ) { case OCC_SYS_LS_SVR: nDesSysType = OCC_SYS_XD_CLIENT; break; case OCC_SYS_DD_SVR: nDesSysType = OCC_SYS_DD_CLIENT; break; case OCC_SYS_HD_SVR: nDesSysType = OCC_SYS_HD_CLIENT; break; default: break; } if ( nDesSysType < 0 || !m_pQI_SendCmdsList ) { return false; } OCCCOMMAND NetCmd; NetCmd.ScrSysType = m_subsystype; NetCmd.DesSysType = nDesSysType; NetCmd.TrainningID = m_trianid; NetCmd.cmdType = OCCCOMMAND::eTrainReset; TRAIN_RESET Train_Reset; Train_Reset.nTrainingID = m_trianid; Train_Reset.nCourseID = m_courseid; Train_Reset.nType = eTrainReset_Type_RollBack; NetCmd.setPacketValue(&Train_Reset,NetCmd.cmdType); m_pQI_SendCmdsList->push(NetCmd); return true; } /******************************************************* 函数名称: Send_DispatchSys_OperR 函数功能: 发送调度系统操作记录 输入参数: NetCmd - 操作命令 nResultCode - 操作结果 输出参数: NULL 返 回 值: bool - true:成功 false:失败 *******************************************************/ bool CBaseDispatch::Send_DispatchSys_OperR( OCCCOMMAND NetCmd,int nResultCode ) { if ( !m_pQI_SendCmdsList || !m_pDispatchMngI ) { return false; } if ( NetCmd.ScrSysType != OCC_SYS_XD_CLIENT && NetCmd.ScrSysType != OCC_SYS_DD_CLIENT && NetCmd.ScrSysType != OCC_SYS_HD_CLIENT && NetCmd.ScrSysType != 28 ) //IBP盘 { return false; } DISPATCH_SYS_OPERRECORD DispatchSysOperR; DispatchSysOperR.nSrcSysID = NetCmd.ScrNetID; DispatchSysOperR.nSimTime = m_pDispatchMngI->_IGetTime_Simulate(); DispatchSysOperR.nSysTime = GetSys_SecondTime_2(); DispatchSysOperR.nDataLen = NetCmd.DateLen; DispatchSysOperR.nCmdType = NetCmd.cmdType; DispatchSysOperR.nCmdResult= nResultCode; int nDataSize = min((unsigned long)NetCmd.DateLen,sizeof(DispatchSysOperR.cCmdPara)); memcpy(DispatchSysOperR.cCmdPara,&NetCmd.Net_MAX_MAXNETINFO,nDataSize); OCCCOMMAND NetCmd2; NetCmd2.DesSysType = OCC_SYS_XD_CLIENT; NetCmd2.ScrSysType = m_subsystype; NetCmd2.TrainningID = m_trianid; NetCmd2.cmdType = OCCCOMMAND::eDispatch_Sys_OperRecord; LPBASENETPACKET TmpPack = (LPBASENETPACKET)&DispatchSysOperR; NetCmd2.setPacketValue(TmpPack,NetCmd2.cmdType); m_pQI_SendCmdsList->push(NetCmd2); return true; } /******************************************************* 函数名称: ThreadFun_EventF1 函数功能: 线程函数_F1 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::ThreadFun_EventF1() { m_EventF1ThreadOut = false; try { while ( !m_EventF1ThreadOut ) { NORMALCMDC cmd; double dTimeS = GetNowTime(TIMEUNIT_MS); if ( _QI_eventF1List.pop(cmd) ) { EventF1Ass(cmd); m_EventF1ThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; } else { m_EventF1ThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; Thread_Sleep(TIMESPAN_LV_NORMAL); } } OUT_LOG(LOG_LEVEL_ERROR,"【M_%d,T_%d】调度任务TreadEventF1线程退出",m_subsystype,m_trianid); } catch (...) { OUT_LOG(LOG_LEVEL_ERROR, "TreadEventF1 Have a Error,Trainid = %d,TrainType= %d", \ m_trianid,m_subsystype); } } /******************************************************* 函数名称: ThreadFun_EventF2 函数功能: 线程函数_F2 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::ThreadFun_EventF2() { m_EventF2ThreadOut = false; try { while ( !m_EventF2ThreadOut ) { NORMALCMDC cmd; double dTimeS = GetNowTime(TIMEUNIT_MS); if ( _QI_eventF2List.pop(cmd) ) { EventF2Ass(cmd); m_EventF2ThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; } else { m_EventF2ThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; Thread_Sleep(TIMESPAN_LVT); } } OUT_LOG(LOG_LEVEL_ERROR, "【M_%d,T_%d】调度任务TreadEventF2线程退出",m_subsystype,m_trianid); } catch (...) { OUT_LOG(LOG_LEVEL_ERROR, "TreadEventF2 Have a Error,Trainid = %d,TrainType= %d", \ m_trianid,m_subsystype); } } /******************************************************* 函数名称: ThreadFun_EventBegin 函数功能: 线程函数_EventBegin 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::ThreadFun_EventBegin() { m_EventBeginThreadOut = false; try { time_t begintime = GetNowTime(TIMEUNIT_S); time_t lasttime = begintime; while (!m_EventBeginThreadOut) { double dTimeS = GetNowTime(TIMEUNIT_MS); time_t spantime = dTimeS - lasttime; if ( EventBeginAss(begintime,spantime) ) { lasttime = GetNowTime(TIMEUNIT_S); m_EventBeginThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; } else { m_EventBeginThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; Thread_Sleep(TIMESPAN_LVT); } } OUT_LOG(LOG_LEVEL_ERROR,"【M_%d,T_%d】调度任务TreadEventBegin线程退出",m_subsystype,m_trianid); } catch (...) { OUT_LOG(LOG_LEVEL_ERROR, "TreadEventBegin Have a Error,Trainid = %d,TrainType= %d", \ m_trianid,m_subsystype); } } /******************************************************* 函数名称: ThreadFun_ObjectCP 函数功能: 线程函数_ObjectCP 输入参数: NULL 输出参数: NULL 返 回 值: void *******************************************************/ void CBaseDispatch::ThreadFun_ObjectCP() { m_ObjectCPThreadOut = false; double dThreadCalTime_CP = GetNowTime(TIMEUNIT_MS); try { //time_t begintime = CTime::GetCurrentTime().GetTime(); time_t begintime = time(nullptr); time_t lasttime = GetNowTime(TIMEUNIT_S); while ( !m_ObjectCPThreadOut ) { time_t spantime = GetNowTime(TIMEUNIT_S) - lasttime; Add_Module_Watch(); double dTimeS = GetNowTime(TIMEUNIT_MS); double dDealT = GetNowTime(TIMEUNIT_MS) - dThreadCalTime_CP; if ( ObjectCPAss(begintime,spantime ) && ( m_subsystype == OCC_SYS_XD_SVR || dDealT >= 1 )) { lasttime = GetNowTime(TIMEUNIT_S); dThreadCalTime_CP = GetNowTime(TIMEUNIT_MS); m_ObjectCPThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; } else { m_ObjectCPThreadTimeDur = GetNowTime(TIMEUNIT_MS) - dTimeS; Thread_Sleep(TIMESPAN_LV_NORMAL); } } OUT_LOG(LOG_LEVEL_ERROR,"【M_%d,T_%d】调度任务TreadObjectCP线程退出",m_subsystype,m_trianid); } catch (...) { OUT_LOG(LOG_LEVEL_ERROR, "TreadObjectCP Have a Error,Trainid = %d,TrainType= %d", \ m_trianid,m_subsystype); } }