TaskMNG.cpp 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465
  1. #include "TaskMNG.h"
  2. #include "platform.h"
  3. #include "NetMessageDefine.h"
  4. //#include <dlfcn.h>
  5. #include <YD_Public_OCC.h>
  6. #ifdef Q_OS_WIN
  7. #pragma execution_character_set("utf-8")
  8. #endif
  9. TaskMNG* TaskMNG::g_GlobalSingle = NULL;
  10. std::mutex TaskMNG::m_Mutex;
  11. #include "ElectricityVOC.h"
  12. #include <QLibrary>
  13. #include <QCoreApplication>
  14. //单例初始化
  15. typedef CBaseDispatch* (*Create_Sub_Dispatch)();
  16. TaskMNG* TaskMNG::instance()
  17. {
  18. if (g_GlobalSingle == NULL)
  19. {
  20. std::unique_lock<std::mutex> lock(m_Mutex);
  21. if (g_GlobalSingle == NULL)
  22. {
  23. g_GlobalSingle = new (std::nothrow) TaskMNG();
  24. }
  25. }
  26. return g_GlobalSingle;
  27. }
  28. void TaskMNG::destroyInstance()
  29. {
  30. std::unique_lock<std::mutex> lock(m_Mutex);
  31. //NetDll_StopThread();
  32. if (g_GlobalSingle)
  33. {
  34. delete g_GlobalSingle;
  35. g_GlobalSingle = NULL;
  36. }
  37. }
  38. TaskMNG::TaskMNG()
  39. {
  40. nTaskStatus = Status_Null;
  41. }
  42. TaskMNG::~TaskMNG()
  43. {
  44. }
  45. void TaskMNG::InitDB(string strIP, int nType, string strDBname, string strUser, string strPwd)
  46. {
  47. //m_pdb = new YD_DBOperator();
  48. int nPort = 0;
  49. if (nType==1)
  50. {
  51. nPort = 8527;
  52. }
  53. else
  54. {
  55. nPort = 1433;
  56. }
  57. m_strSeverIP = strIP;
  58. m_nDBType = nType;
  59. m_strDBName = strDBname;
  60. m_strUserID = strUser;
  61. m_strPassWord = strPwd;
  62. //m_pdb->SetDBType(nType);
  63. //m_pdb->ConnectDB(strIP,strUser,strPwd,strDBname, nPort);
  64. }
  65. bool TaskMNG::CreateDispatchTask(int courseid, int systype, int trainid, int runlineid, int parm1, int parm2, InitCourInfo* curseinfo)
  66. {
  67. SVRREPLYMSG netret;
  68. memset(&netret, 0, sizeof(netret));
  69. netret.sysType = systype;
  70. string logstr;
  71. bool ret = true;
  72. netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL;
  73. string key = Make2IntKey_2(systype, trainid);
  74. //EnterCriticalSection(&m_pDispatchTasksMappretect);
  75. DISTASKMAP::iterator iter = m_DispatchTasksMap.find(key);
  76. if (iter == m_DispatchTasksMap.end())
  77. {
  78. //std::string Str_Dll =
  79. #ifdef Q_OS_WIN
  80. QString StdPath = QCoreApplication::applicationDirPath() + "/Elc.dll";
  81. #else
  82. QString StdPath = QCoreApplication::applicationDirPath() + "/libElc.so";
  83. #endif
  84. std::string Str_Dll = StdPath.toStdString();
  85. QLibrary myLib;
  86. myLib.setFileName(StdPath);
  87. // 对于Windows系统
  88. // QLibrary myLib("/path/to/your/so/libmylib.so"); // 对于Linux系统
  89. if (!myLib.load())
  90. {
  91. // 获取函数地址
  92. return false;
  93. }
  94. Create_Sub_Dispatch pCreateDispatchTaskMg = (Create_Sub_Dispatch)myLib.resolve("Create_Sub_Dispatch");
  95. CBaseDispatch* pdispatchsys = pCreateDispatchTaskMg();
  96. if (pdispatchsys == NULL)
  97. {
  98. logstr = "调度服务器创建:" + GetSysTypeStr(systype) + "训练号:" + std::to_string(trainid) + "失败\n";
  99. AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_ERROR);
  100. ret = false;
  101. netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL;
  102. nTaskStatus = Status_Error;
  103. }
  104. else
  105. {
  106. ElectricityVOC* pElect = (ElectricityVOC*)pdispatchsys;
  107. //RecTrianPosition(systype, trainid, m_NetServerID);
  108. /*
  109. DATECONFIGMAP::iterator iter = m_DataBaseConfigMap.find(systype);
  110. if (iter == m_DataBaseConfigMap.end())
  111. {
  112. logstr = "调度服务器创建:" + GetSysTypeStr(systype) + "训练号:" + std::to_string(trainid) + "没有找到数据库配置项\n";
  113. AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_NORMAL);
  114. ret = false;
  115. netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL;
  116. }
  117. else*/
  118. m_pdb = new YD_DBOperator();
  119. m_pdb->SetDBType(m_nDBType);
  120. pElect->setDBP(m_pdb,&pLogCtrl);
  121. curseinfo = new InitCourInfo();
  122. if (!pdispatchsys->OpenDB(m_strSeverIP, m_strDBName, m_strUserID, m_strPassWord, courseid, runlineid, parm1, parm2, m_nDBType, curseinfo))
  123. {
  124. pdispatchsys->Exit();
  125. delete pdispatchsys;
  126. pdispatchsys = NULL;
  127. logstr = "调度服务器创建:" + GetSysTypeStr(systype) + "训练号:" + std::to_string(trainid) + "打开数据库任务失败\n";
  128. AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_ERROR);
  129. ret = false;
  130. netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL;
  131. nTaskStatus = Status_Error;
  132. }
  133. else if (!pdispatchsys->OpenWork(trainid, systype, &m_pSendList, &m_pSendGroupList, &m_pSendRealList))
  134. {
  135. pdispatchsys->Exit();
  136. delete pdispatchsys;
  137. pdispatchsys = NULL;
  138. logstr = "调度服务器创建:" + GetSysTypeStr(systype) + "训练号:" + std::to_string(trainid) + "失败\n";
  139. AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_NORMAL);
  140. ret = false;
  141. netret.status = OCC_NETMESSAGE_TRIAN_CREATE_FAIL;
  142. nTaskStatus = Status_Error;
  143. }
  144. else
  145. {
  146. m_DispatchTasksMap.insert(make_pair(key, pdispatchsys));
  147. logstr = "调度服务器创建:" + GetSysTypeStr(systype) + "训练号:" + std::to_string(trainid) + "成功\n";
  148. AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_NORMAL);
  149. netret.status = OCC_NETMESSAGE_TRIAN_CREATE_SUCCESS;
  150. nTaskStatus = Status_Init;
  151. }
  152. }
  153. }
  154. else
  155. {
  156. logstr = "调度服务器创建:" + GetSysTypeStr(systype) + "训练号:" + std::to_string(trainid) + "失败! 该训练已经存在\n";
  157. AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_ERROR);
  158. netret.status = OCC_NETMESSAGE_TRIAN_CREATE_EXIST;
  159. nTaskStatus = Status_Error;
  160. }
  161. //strncpy(netret.ReplyText, logstr.c_str(), sizeof(ALARM_LEN));
  162. // if (curseinfo->coursetype == OCCCOMMAND::eTrainIni) //其他流程单独发送
  163. // {
  164. // NetComand_Send(&netret, OCC_SYS_SERVER_TH, OCCCOMMAND::eSvrReplyMsg, trainid);//创建课程
  165. // }
  166. return ret;
  167. }
  168. bool TaskMNG::DeleteDispatchTask(int courseid, int systype, int trainid)
  169. {
  170. bool returnflag = false;
  171. SVRREPLYMSG netret;
  172. netret.sysType = systype;
  173. string logstr;
  174. bool ret = true;
  175. netret.status = OCC_NETMESSAGE_TRIAN_DELETE_SUCCESS;
  176. string key = Make2IntKey_2(systype, trainid);
  177. /*DISTASKMAP::iterator iter = m_DispatchTasksMap.find(key);
  178. if (iter!=m_DispatchTasksMap.end())*/
  179. bool findflag = false;
  180. //std::unique_lock<std::recursive_mutex> g1(m_pDispatchTasksMappretect);
  181. //EnterCriticalSection(&m_pDispatchTasksMappretect);
  182. for (DISTASKMAP::iterator iter = m_DispatchTasksMap.begin()
  183. ; iter != m_DispatchTasksMap.end()
  184. ; //iter ++
  185. )
  186. {
  187. CBaseDispatch* pdispatchsys = iter->second;
  188. if (pdispatchsys == NULL)
  189. {
  190. m_DispatchTasksMap.erase(iter);
  191. //LeaveCriticalSection(&m_pDispatchTasksMappretect);
  192. nTaskStatus = Status_Null;
  193. return false;
  194. }
  195. else
  196. {
  197. if (pdispatchsys->m_subsystype != systype)
  198. {
  199. iter++;
  200. continue;
  201. }
  202. else
  203. {
  204. courseid = pdispatchsys->m_courseid;
  205. systype = pdispatchsys->m_subsystype;
  206. trainid = pdispatchsys->m_trianid;
  207. }
  208. }
  209. findflag = true;
  210. if (pdispatchsys->StopTask(netret.parm))
  211. {
  212. returnflag = true;
  213. pdispatchsys->Exit();
  214. //DeleteReceiveGroupMsg(trainid, systype);
  215. //Sleep(1000);
  216. delete pdispatchsys;
  217. pdispatchsys = NULL;
  218. netret.status = OCC_NETMESSAGE_TRIAN_DELETE_SUCCESS;
  219. logstr = "调度服务器删除:" + GetSysTypeStr(systype) + "训练号:" + std::to_string(trainid) + "成功\n";
  220. AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_NORMAL);
  221. //DISTASKMAP::iterator tmpiter = iter;
  222. m_DispatchTasksMap.erase(iter);
  223. nTaskStatus = Status_Null;
  224. iter = m_DispatchTasksMap.begin();
  225. }
  226. else
  227. netret.status = OCC_NETMESSAGE_TRIAN_DELETE_FAIL;
  228. }
  229. if (!findflag)
  230. {
  231. netret.status = OCC_NETMESSAGE_TRIAN_DELETE_CANNOETFIND;
  232. }
  233. strncpy(netret.ReplyText, logstr.c_str(), sizeof(ALARM_LEN));
  234. //NetComand_Send(&netret, OCC_SYS_SERVER_TH, OCCCOMMAND::eSvrReplyMsg, trainid);//创建课程
  235. nTaskStatus = Status_Null;
  236. return returnflag;
  237. }
  238. bool TaskMNG::StopDispatchTask(int courseid, int systype, int trainid)
  239. {
  240. SVRREPLYMSG netret;
  241. netret.sysType = systype;
  242. bool ret = true;
  243. string key = Make2IntKey_2(systype, trainid);
  244. //EnterCriticalSection(&m_pDispatchTasksMappretect);
  245. DISTASKMAP::iterator iter = m_DispatchTasksMap.find(key);
  246. if (iter != m_DispatchTasksMap.end())
  247. {
  248. CBaseDispatch* pdispatchsys = iter->second;
  249. if (pdispatchsys == NULL)
  250. {
  251. m_DispatchTasksMap.erase(iter);
  252. //LeaveCriticalSection(&m_pDispatchTasksMappretect);
  253. return false;
  254. }
  255. if (pdispatchsys->StopTask(netret.parm))
  256. {
  257. netret.status = OCC_NETMESSAGE_TRIAN_STOP_SUCCESS;
  258. }
  259. else
  260. {
  261. netret.status = OCC_NETMESSAGE_TRIAN_STOP_FAIL;
  262. }
  263. }
  264. else
  265. {
  266. netret.status = OCC_NETMESSAGE_TRIAN_STOP_CANNOETFIND;
  267. }
  268. nTaskStatus = TaskStatus::Status_Stop;
  269. return false;
  270. }
  271. bool TaskMNG::PauseDispatchTask(int courseid, int systype, int trainid)
  272. {
  273. return false;
  274. }
  275. bool TaskMNG::ResumeDispatchTask(int courseid, int systype, int trainid)
  276. {
  277. return false;
  278. }
  279. bool TaskMNG::RunDispatchTask(int courseid, int systype, int trainid)
  280. {
  281. SVRREPLYMSG netret;
  282. netret.sysType = systype;
  283. bool ret = true;
  284. netret.status = OCC_NETMESSAGE_TRIAN_RUN_SUCCESS;
  285. string key = Make2IntKey_2(systype, trainid);
  286. //std::unique_lock<std::recursive_mutex> g1(m_pDispatchTasksMappretect);
  287. //EnterCriticalSection(&m_pDispatchTasksMappretect);
  288. DISTASKMAP::iterator iter = m_DispatchTasksMap.find(key);
  289. if (iter != m_DispatchTasksMap.end())
  290. {
  291. CBaseDispatch* pdispatchsys = iter->second;
  292. if (pdispatchsys == NULL)
  293. {
  294. m_DispatchTasksMap.erase(iter);
  295. //LeaveCriticalSection(&m_pDispatchTasksMappretect);
  296. nTaskStatus = Status_Null;
  297. return false;
  298. }
  299. if (pdispatchsys->StartTask(netret.parm))
  300. {
  301. //pdispatchsys->m_begintime = time(NULL);
  302. //tm* t = localtime(&pdispatchsys->m_begintime);
  303. //pdispatchsys->m_begintime = CTime::GetCurrentTime().GetTime();
  304. netret.status = OCC_NETMESSAGE_TRIAN_RUN_SUCCESS;
  305. }
  306. else
  307. netret.status = OCC_NETMESSAGE_TRIAN_RUN_FAIL;
  308. }
  309. else
  310. {
  311. netret.status = OCC_NETMESSAGE_TRIAN_RUN_CANNOETFIND;
  312. }
  313. //LeaveCriticalSection(&m_pDispatchTasksMappretect);
  314. nTaskStatus = Status_Run;
  315. //NetComand_Send(&netret, OCC_SYS_SERVER_TH, OCCCOMMAND::eSvrReplyMsg, trainid);//创建课程
  316. return true;
  317. }
  318. //内部互相之间数据传输
  319. bool TaskMNG::NetComand_Receive(OCCCOMMAND& command)
  320. {
  321. if (command.DesSysType == OCC_SYS_SERVER_SVR)
  322. {
  323. return false;
  324. //return NetComandToServer(command);
  325. }
  326. else
  327. {
  328. string key = Make2IntKey_2(command.DesSysType, command.TrainningID);
  329. // if (command.ScrNetID == m_nDebugClientSysID)
  330. // {
  331. // key = Make2IntKey(command.DesSysType, m_nTrainID);
  332. // command.TrainningID = m_nTrainID;
  333. // }
  334. //std::unique_lock<std::recursive_mutex> g1(m_pDispatchTasksMappretect);
  335. //EnterCriticalSection(&m_pDispatchTasksMappretect);
  336. DISTASKMAP::iterator iter = m_DispatchTasksMap.find(key);
  337. if (iter != m_DispatchTasksMap.end())
  338. {
  339. CBaseDispatch* pdispatchsys = iter->second;
  340. if (pdispatchsys == NULL)
  341. {
  342. m_DispatchTasksMap.erase(iter);
  343. //LeaveCriticalSection(&m_pDispatchTasksMappretect);
  344. return false;
  345. }
  346. pdispatchsys->ReceiveMsg(command);
  347. //LeaveCriticalSection(&m_pDispatchTasksMappretect);
  348. return true;
  349. }
  350. //LeaveCriticalSection(&m_pDispatchTasksMappretect);
  351. }
  352. string logstr;
  353. logstr = "未知网络系统消息:DesSysType = " + GetSysTypeStr(command.DesSysType) + "训练号:" + std::to_string(command.TrainningID) + "\n";
  354. //logstr.Format("未知网络系统消息:DesSysType=%d,TrianID=%d", command.DesSysType, command.TrainningID);
  355. //AddLogItem(GetCurTime_ms(), logstr, LOG_LEVEL_NORMAL);
  356. return false;
  357. }
  358. void TaskMNG::NetComand_Send(LPBASENETPACKET pnetmsg, int dessystype, OCCCOMMAND::CmdType cmdtype, int trianid, int nDesID /*= -1*/)
  359. {
  360. OCCCOMMAND NetOrder;
  361. NetOrder.ScrSysType = OCC_SYS_SERVER_SVR;
  362. NetOrder.setPacketValue(pnetmsg, cmdtype);
  363. NetOrder.TrainningID = trianid;
  364. NetOrder.DesSysType = dessystype;
  365. NetOrder.DesNetID = nDesID;
  366. m_pSendList.push(NetOrder);
  367. }
  368. CBaseDispatch* TaskMNG::CreateNewTaskType(int systype)
  369. {
  370. CBaseDispatch* pDispatchSys = NULL;
  371. //创建子模块
  372. auto iters = m_PluginMap.find(systype);
  373. if (iters==m_PluginMap.end())
  374. {
  375. return pDispatchSys;
  376. }
  377. sPlugin* pPlugin = &iters->second;//GetPlugin_ByModuleID(systype);
  378. if (NULL != pPlugin && NULL != pPlugin->bInitOk)
  379. {
  380. pDispatchSys = pPlugin->nCreatePlug();
  381. //pDispatchSys->Module_Init(m_pdb, m_funAddLog, m_pLocalQT);
  382. }
  383. return pDispatchSys;
  384. }
  385. void TaskMNG::AddLogItem(string str1, string str2, int nLevel)
  386. {
  387. m_strLog = str2;
  388. }
  389. std::string TaskMNG::GetCurTime_ms()
  390. {
  391. return "";
  392. }
  393. CBaseDispatch* TaskMNG::pGetEleVOC()
  394. {
  395. auto iters = m_DispatchTasksMap.begin();
  396. if (iters!=m_DispatchTasksMap.end())
  397. {
  398. return iters->second;
  399. }
  400. return nullptr;
  401. }
  402. BOOL LogCtrl::Add_Log(int nLevel, LPCSTR pszFmt, ...)
  403. {
  404. return false;
  405. }
  406. TaskMNG::TaskStatus TaskMNG::getTaskStatus()
  407. {
  408. return nTaskStatus;
  409. }
  410. std::string TaskMNG::getTaskText()
  411. {
  412. return m_strLog;
  413. }