ua_server_worker.c 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623
  1. #include "ua_util.h"
  2. #include "ua_server_internal.h"
  3. /**
  4. * There are four types of job execution:
  5. *
  6. * 1. Normal jobs (dispatched to worker threads if multithreading is activated)
  7. *
  8. * 2. Repeated jobs with a repetition interval (dispatched to worker threads)
  9. *
  10. * 3. Mainloop jobs are executed (once) from the mainloop and not in the worker threads. The server
  11. * contains a stack structure where all threads can add mainloop jobs for the next mainloop
  12. * iteration. This is used e.g. to trigger adding and removing repeated jobs without blocking the
  13. * mainloop.
  14. *
  15. * 4. Delayed jobs are executed once in a worker thread. But only when all normal jobs that were
  16. * dispatched earlier have been executed. This is achieved by a counter in the worker threads. We
  17. * compute from the counter if all previous jobs have finished. The delay can be very long, since we
  18. * try to not interfere too much with normal execution. A use case is to eventually free obsolete
  19. * structures that _could_ still be accessed from concurrent threads.
  20. *
  21. * - Remove the entry from the list
  22. * - mark it as "dead" with an atomic operation
  23. * - add a delayed job that frees the memory when all concurrent operations have completed
  24. */
  25. #define MAXTIMEOUT 50000 // max timeout in microsec until the next main loop iteration
  26. #define BATCHSIZE 20 // max number of jobs that are dispatched at once to workers
  27. static void processJobs(UA_Server *server, UA_Job *jobs, size_t jobsSize) {
  28. for(size_t i = 0; i < jobsSize; i++) {
  29. UA_Job *job = &jobs[i];
  30. switch(job->type) {
  31. case UA_JOBTYPE_BINARYMESSAGE:
  32. UA_Server_processBinaryMessage(server, job->job.binaryMessage.connection,
  33. &job->job.binaryMessage.message);
  34. UA_Connection *c = job->job.binaryMessage.connection;
  35. c->releaseBuffer(job->job.binaryMessage.connection, &job->job.binaryMessage.message);
  36. break;
  37. case UA_JOBTYPE_CLOSECONNECTION:
  38. UA_Connection_detachSecureChannel(job->job.closeConnection);
  39. job->job.closeConnection->close(job->job.closeConnection);
  40. break;
  41. case UA_JOBTYPE_METHODCALL:
  42. case UA_JOBTYPE_DELAYEDMETHODCALL:
  43. job->job.methodCall.method(server, job->job.methodCall.data);
  44. break;
  45. default:
  46. UA_LOG_WARNING(server->logger, UA_LOGCATEGORY_SERVER, "Trying to execute a job of unknown type");
  47. break;
  48. }
  49. }
  50. }
  51. /*******************************/
  52. /* Worker Threads and Dispatch */
  53. /*******************************/
  54. #ifdef UA_MULTITHREADING
  55. struct MainLoopJob {
  56. struct cds_lfs_node node;
  57. UA_Job job;
  58. };
  59. /** Entry in the dispatch queue */
  60. struct DispatchJobsList {
  61. struct cds_wfcq_node node; // node for the queue
  62. size_t jobsSize;
  63. UA_Job *jobs;
  64. };
  65. /** Dispatch jobs to workers. Slices the job array up if it contains more than BATCHSIZE items. The jobs
  66. array is freed in the worker threads. */
  67. static void dispatchJobs(UA_Server *server, UA_Job *jobs, size_t jobsSize) {
  68. size_t startIndex = jobsSize; // start at the end
  69. while(jobsSize > 0) {
  70. size_t size = BATCHSIZE;
  71. if(size > jobsSize)
  72. size = jobsSize;
  73. startIndex = startIndex - size;
  74. struct DispatchJobsList *wln = UA_malloc(sizeof(struct DispatchJobsList));
  75. if(startIndex > 0) {
  76. wln->jobs = UA_malloc(size * sizeof(UA_Job));
  77. UA_memcpy(wln->jobs, &jobs[startIndex], size * sizeof(UA_Job));
  78. wln->jobsSize = size;
  79. } else {
  80. /* forward the original array */
  81. wln->jobsSize = size;
  82. wln->jobs = jobs;
  83. }
  84. cds_wfcq_node_init(&wln->node);
  85. cds_wfcq_enqueue(&server->dispatchQueue_head, &server->dispatchQueue_tail, &wln->node);
  86. jobsSize -= size;
  87. }
  88. }
  89. // throwaway struct to bring data into the worker threads
  90. struct workerStartData {
  91. UA_Server *server;
  92. UA_UInt32 **workerCounter;
  93. };
  94. /** Waits until jobs arrive in the dispatch queue and processes them. */
  95. static void * workerLoop(struct workerStartData *startInfo) {
  96. rcu_register_thread();
  97. UA_UInt32 *c = UA_malloc(sizeof(UA_UInt32));
  98. uatomic_set(c, 0);
  99. *startInfo->workerCounter = c;
  100. UA_Server *server = startInfo->server;
  101. UA_free(startInfo);
  102. pthread_mutex_t mutex; // required for the condition variable
  103. pthread_mutex_init(&mutex,0);
  104. pthread_mutex_lock(&mutex);
  105. struct timespec to;
  106. while(*server->running) {
  107. struct DispatchJobsList *wln = (struct DispatchJobsList*)
  108. cds_wfcq_dequeue_blocking(&server->dispatchQueue_head, &server->dispatchQueue_tail);
  109. if(wln) {
  110. processJobs(server, wln->jobs, wln->jobsSize);
  111. UA_free(wln->jobs);
  112. UA_free(wln);
  113. } else {
  114. /* sleep until a work arrives (and wakes up all worker threads) */
  115. clock_gettime(CLOCK_REALTIME, &to);
  116. to.tv_sec += 2;
  117. pthread_cond_timedwait(&server->dispatchQueue_condition, &mutex, &to);
  118. }
  119. uatomic_inc(c); // increase the workerCounter;
  120. }
  121. pthread_mutex_unlock(&mutex);
  122. pthread_mutex_destroy(&mutex);
  123. rcu_barrier(); // wait for all scheduled call_rcu work to complete
  124. rcu_unregister_thread();
  125. /* we need to return _something_ for pthreads */
  126. return UA_NULL;
  127. }
  128. static void emptyDispatchQueue(UA_Server *server) {
  129. while(!cds_wfcq_empty(&server->dispatchQueue_head, &server->dispatchQueue_tail)) {
  130. struct DispatchJobsList *wln = (struct DispatchJobsList*)
  131. cds_wfcq_dequeue_blocking(&server->dispatchQueue_head, &server->dispatchQueue_tail);
  132. processJobs(server, wln->jobs, wln->jobsSize);
  133. UA_free(wln->jobs);
  134. UA_free(wln);
  135. }
  136. }
  137. #endif
  138. /*****************/
  139. /* Repeated Jobs */
  140. /*****************/
  141. struct IdentifiedJob {
  142. UA_Job job;
  143. UA_Guid id;
  144. };
  145. /**
  146. * The RepeatedJobs structure contains an array of jobs that are either executed with the same
  147. * repetition interval. The linked list is sorted, so we can stop traversing when the first element
  148. * has nextTime > now.
  149. */
  150. struct RepeatedJobs {
  151. LIST_ENTRY(RepeatedJobs) pointers; ///> Links to the next list of repeated jobs (with a different) interval
  152. UA_DateTime nextTime; ///> The next time when the jobs are to be executed
  153. UA_UInt32 interval; ///> Interval in 100ns resolution
  154. size_t jobsSize; ///> Number of jobs contained
  155. struct IdentifiedJob jobs[]; ///> The jobs. This is not a pointer, instead the struct is variable sized.
  156. };
  157. /* throwaway struct for the mainloop callback */
  158. struct AddRepeatedJob {
  159. struct IdentifiedJob job;
  160. UA_UInt32 interval;
  161. };
  162. /* internal. call only from the main loop. */
  163. static UA_StatusCode addRepeatedJob(UA_Server *server, struct AddRepeatedJob * UA_RESTRICT arw) {
  164. struct RepeatedJobs *matchingTw = UA_NULL; // add the item here
  165. struct RepeatedJobs *lastTw = UA_NULL; // if there is no repeated job, add a new one this entry
  166. struct RepeatedJobs *tempTw;
  167. /* search for matching entry */
  168. UA_DateTime firstTime = UA_DateTime_now() + arw->interval;
  169. tempTw = LIST_FIRST(&server->repeatedJobs);
  170. while(tempTw) {
  171. if(arw->interval == tempTw->interval) {
  172. matchingTw = tempTw;
  173. break;
  174. }
  175. if(tempTw->nextTime > firstTime)
  176. break;
  177. lastTw = tempTw;
  178. tempTw = LIST_NEXT(lastTw, pointers);
  179. }
  180. if(matchingTw) {
  181. /* append to matching entry */
  182. matchingTw = UA_realloc(matchingTw, sizeof(struct RepeatedJobs) +
  183. (sizeof(struct IdentifiedJob) * (matchingTw->jobsSize + 1)));
  184. if(!matchingTw) {
  185. #ifdef UA_MULTITHREADING
  186. UA_free(arw);
  187. #endif
  188. return UA_STATUSCODE_BADOUTOFMEMORY;
  189. }
  190. /* point the realloced struct */
  191. if(matchingTw->pointers.le_next)
  192. matchingTw->pointers.le_next->pointers.le_prev = &matchingTw->pointers.le_next;
  193. if(matchingTw->pointers.le_prev)
  194. *matchingTw->pointers.le_prev = matchingTw;
  195. } else {
  196. /* create a new entry */
  197. matchingTw = UA_malloc(sizeof(struct RepeatedJobs) + sizeof(struct IdentifiedJob));
  198. if(!matchingTw) {
  199. #ifdef UA_MULTITHREADING
  200. UA_free(arw);
  201. #endif
  202. return UA_STATUSCODE_BADOUTOFMEMORY;
  203. }
  204. matchingTw->jobsSize = 0;
  205. matchingTw->nextTime = firstTime;
  206. matchingTw->interval = arw->interval;
  207. if(lastTw)
  208. LIST_INSERT_AFTER(lastTw, matchingTw, pointers);
  209. else
  210. LIST_INSERT_HEAD(&server->repeatedJobs, matchingTw, pointers);
  211. }
  212. matchingTw->jobs[matchingTw->jobsSize] = arw->job;
  213. matchingTw->jobsSize++;
  214. #ifdef UA_MULTITHREADING
  215. UA_free(arw);
  216. #endif
  217. return UA_STATUSCODE_GOOD;
  218. }
  219. UA_StatusCode UA_Server_addRepeatedJob(UA_Server *server, UA_Job job, UA_UInt32 interval, UA_Guid *jobId) {
  220. /* the interval needs to be at least 5ms */
  221. if(interval < 5)
  222. return UA_STATUSCODE_BADINTERNALERROR;
  223. interval *= 10000; // from ms to 100ns resolution
  224. #ifdef UA_MULTITHREADING
  225. struct AddRepeatedJob *arw = UA_malloc(sizeof(struct AddRepeatedJob));
  226. if(!arw)
  227. return UA_STATUSCODE_BADOUTOFMEMORY;
  228. arw->interval = interval;
  229. arw->job.job = job;
  230. if(jobId) {
  231. arw->job.id = UA_Guid_random(&server->random_seed);
  232. *jobId = arw->job.id;
  233. } else
  234. UA_Guid_init(&arw->job.id);
  235. struct MainLoopJob *mlw = UA_malloc(sizeof(struct MainLoopJob));
  236. if(!mlw) {
  237. UA_free(arw);
  238. return UA_STATUSCODE_BADOUTOFMEMORY;
  239. }
  240. mlw->job = (UA_Job) {
  241. .type = UA_JOBTYPE_METHODCALL,
  242. .job.methodCall = {.data = arw, .method = (void (*)(UA_Server*, void*))addRepeatedJob}};
  243. cds_lfs_push(&server->mainLoopJobs, &mlw->node);
  244. #else
  245. struct AddRepeatedJob arw;
  246. arw.interval = interval;
  247. arw.job.job = job;
  248. if(jobId) {
  249. arw.job.id = UA_Guid_random(&server->random_seed);
  250. *jobId = arw.job.id;
  251. } else
  252. UA_Guid_init(&arw.job.id);
  253. addRepeatedJob(server, &arw);
  254. #endif
  255. return UA_STATUSCODE_GOOD;
  256. }
  257. /* Returns the timeout until the next repeated job in ms */
  258. static UA_UInt16 processRepeatedJobs(UA_Server *server) {
  259. UA_DateTime current = UA_DateTime_now();
  260. struct RepeatedJobs *next = LIST_FIRST(&server->repeatedJobs);
  261. struct RepeatedJobs *tw = UA_NULL;
  262. while(next) {
  263. tw = next;
  264. if(tw->nextTime > current)
  265. break;
  266. next = LIST_NEXT(tw, pointers);
  267. #ifdef UA_MULTITHREADING
  268. // copy the entry and insert at the new location
  269. UA_Job *jobsCopy = UA_malloc(sizeof(UA_Job) * tw->jobsSize);
  270. if(!jobsCopy) {
  271. UA_LOG_ERROR(server->logger, UA_LOGCATEGORY_SERVER, "Not enough memory to dispatch delayed jobs");
  272. break;
  273. }
  274. for(size_t i=0;i<tw->jobsSize;i++)
  275. jobsCopy[i] = tw->jobs[i].job;
  276. dispatchJobs(server, jobsCopy, tw->jobsSize); // frees the job pointer
  277. #else
  278. for(size_t i=0;i<tw->jobsSize;i++)
  279. processJobs(server, &tw->jobs[i].job, 1); // does not free the job ptr
  280. #endif
  281. tw->nextTime += tw->interval;
  282. struct RepeatedJobs *prevTw = tw; // after which tw do we insert?
  283. while(UA_TRUE) {
  284. struct RepeatedJobs *n = LIST_NEXT(prevTw, pointers);
  285. if(!n || n->nextTime > tw->nextTime)
  286. break;
  287. prevTw = n;
  288. }
  289. if(prevTw != tw) {
  290. LIST_REMOVE(tw, pointers);
  291. LIST_INSERT_AFTER(prevTw, tw, pointers);
  292. }
  293. }
  294. // check if the next repeated job is sooner than the usual timeout
  295. struct RepeatedJobs *first = LIST_FIRST(&server->repeatedJobs);
  296. UA_UInt16 timeout = MAXTIMEOUT;
  297. if(first) {
  298. timeout = (first->nextTime - current)/10;
  299. if(timeout > MAXTIMEOUT)
  300. return MAXTIMEOUT;
  301. }
  302. return timeout;
  303. }
  304. /* Call this function only from the main loop! */
  305. static void removeRepeatedJob(UA_Server *server, UA_Guid *jobId) {
  306. struct RepeatedJobs *tw;
  307. LIST_FOREACH(tw, &server->repeatedJobs, pointers) {
  308. for(size_t i = 0; i < tw->jobsSize; i++) {
  309. if(!UA_Guid_equal(jobId, &tw->jobs[i].id))
  310. continue;
  311. if(tw->jobsSize == 1) {
  312. LIST_REMOVE(tw, pointers);
  313. UA_free(tw);
  314. } else {
  315. tw->jobsSize--;
  316. tw->jobs[i] = tw->jobs[tw->jobsSize]; // move the last entry to overwrite
  317. }
  318. goto finish; // ugly break
  319. }
  320. }
  321. finish:
  322. #ifdef UA_MULTITHREADING
  323. UA_free(jobId);
  324. #endif
  325. return;
  326. }
  327. UA_StatusCode UA_Server_removeRepeatedJob(UA_Server *server, UA_Guid jobId) {
  328. #ifdef UA_MULTITHREADING
  329. UA_Guid *idptr = UA_malloc(sizeof(UA_Guid));
  330. if(!idptr)
  331. return UA_STATUSCODE_BADOUTOFMEMORY;
  332. *idptr = jobId;
  333. // dispatch to the mainloopjobs stack
  334. struct MainLoopJob *mlw = UA_malloc(sizeof(struct MainLoopJob));
  335. mlw->job = (UA_Job) {
  336. .type = UA_JOBTYPE_METHODCALL,
  337. .job.methodCall = {.data = idptr, .method = (void (*)(UA_Server*, void*))removeRepeatedJob}};
  338. cds_lfs_push(&server->mainLoopJobs, &mlw->node);
  339. #else
  340. removeRepeatedJob(server, &jobId);
  341. #endif
  342. return UA_STATUSCODE_GOOD;
  343. }
  344. void UA_Server_deleteAllRepeatedJobs(UA_Server *server) {
  345. struct RepeatedJobs *current;
  346. while((current = LIST_FIRST(&server->repeatedJobs))) {
  347. LIST_REMOVE(current, pointers);
  348. UA_free(current);
  349. }
  350. }
  351. /****************/
  352. /* Delayed Jobs */
  353. /****************/
  354. #ifdef UA_MULTITHREADING
  355. #define DELAYEDJOBSSIZE 100 // Collect delayed jobs until we have DELAYEDWORKSIZE items
  356. struct DelayedJobs {
  357. struct DelayedJobs *next;
  358. UA_UInt32 *workerCounters; // initially UA_NULL until the counter are set
  359. UA_UInt32 jobsCount; // the size of the array is DELAYEDJOBSSIZE, the count may be less
  360. UA_Job jobs[DELAYEDJOBSSIZE]; // when it runs full, a new delayedJobs entry is created
  361. };
  362. /* Dispatched as an ordinary job when the DelayedJobs list is full */
  363. static void getCounters(UA_Server *server, struct DelayedJobs *delayed) {
  364. UA_UInt32 *counters = UA_malloc(server->nThreads * sizeof(UA_UInt32));
  365. for(UA_UInt16 i = 0;i<server->nThreads;i++)
  366. counters[i] = *server->workerCounters[i];
  367. delayed->workerCounters = counters;
  368. }
  369. // Call from the main thread only. This is the only function that modifies
  370. // server->delayedWork. processDelayedWorkQueue modifies the "next" (after the
  371. // head).
  372. static void addDelayedJob(UA_Server *server, UA_Job *job) {
  373. struct DelayedJobs *dj = server->delayedJobs;
  374. if(!dj || dj->jobsCount >= DELAYEDJOBSSIZE) {
  375. /* create a new DelayedJobs and add it to the linked list */
  376. dj = UA_malloc(sizeof(struct DelayedJobs));
  377. if(!dj) {
  378. UA_LOG_ERROR(server->logger, UA_LOGCATEGORY_SERVER, "Not enough memory to add a delayed job");
  379. return;
  380. }
  381. dj->jobsCount = 0;
  382. dj->workerCounters = UA_NULL;
  383. dj->next = server->delayedJobs;
  384. server->delayedJobs = dj;
  385. /* dispatch a method that sets the counter for the full list that comes afterwards */
  386. if(dj->next) {
  387. UA_Job *setCounter = UA_malloc(sizeof(UA_Job));
  388. *setCounter = (UA_Job) {.type = UA_JOBTYPE_METHODCALL, .job.methodCall =
  389. {.method = (void (*)(UA_Server*, void*))getCounters, .data = dj->next}};
  390. dispatchJobs(server, setCounter, 1);
  391. }
  392. }
  393. dj->jobs[dj->jobsCount] = *job;
  394. dj->jobsCount++;
  395. }
  396. static void addDelayedJobAsync(UA_Server *server, UA_Job *job) {
  397. addDelayedJob(server, job);
  398. UA_free(job);
  399. }
  400. UA_StatusCode UA_Server_addDelayedJob(UA_Server *server, UA_Job job) {
  401. UA_Job *j = UA_malloc(sizeof(UA_Job));
  402. if(!j)
  403. return UA_STATUSCODE_BADOUTOFMEMORY;
  404. *j = job;
  405. struct MainLoopJob *mlw = UA_malloc(sizeof(struct MainLoopJob));
  406. mlw->job = (UA_Job) {.type = UA_JOBTYPE_METHODCALL, .job.methodCall =
  407. {.data = j, .method = (void (*)(UA_Server*, void*))addDelayedJobAsync}};
  408. cds_lfs_push(&server->mainLoopJobs, &mlw->node);
  409. return UA_STATUSCODE_GOOD;
  410. }
  411. /* Find out which delayed jobs can be executed now */
  412. static void dispatchDelayedJobs(UA_Server *server, void *data /* not used, but needed for the signature*/) {
  413. /* start at the second */
  414. struct DelayedJobs *dw = server->delayedJobs, *beforedw = dw;
  415. if(dw)
  416. dw = dw->next;
  417. /* find the first delayedwork where the counters have been set and have moved */
  418. while(dw) {
  419. if(!dw->workerCounters) {
  420. beforedw = dw;
  421. dw = dw->next;
  422. continue;
  423. }
  424. UA_Boolean allMoved = UA_TRUE;
  425. for(UA_UInt16 i=0;i<server->nThreads;i++) {
  426. if(dw->workerCounters[i] == *server->workerCounters[i]) {
  427. allMoved = UA_FALSE;
  428. break;
  429. }
  430. }
  431. if(allMoved)
  432. break;
  433. beforedw = dw;
  434. dw = dw->next;
  435. }
  436. /* process and free all delayed jobs from here on */
  437. while(dw) {
  438. processJobs(server, dw->jobs, dw->jobsCount);
  439. struct DelayedJobs *next = uatomic_xchg(&beforedw->next, UA_NULL);
  440. UA_free(dw);
  441. UA_free(dw->workerCounters);
  442. dw = next;
  443. }
  444. }
  445. #endif
  446. /********************/
  447. /* Main Server Loop */
  448. /********************/
  449. #ifdef UA_MULTITHREADING
  450. static void processMainLoopJobs(UA_Server *server) {
  451. /* no synchronization required if we only use push and pop_all */
  452. struct cds_lfs_head *head = __cds_lfs_pop_all(&server->mainLoopJobs);
  453. if(!head)
  454. return;
  455. struct MainLoopJob *mlw = (struct MainLoopJob*)&head->node;
  456. struct MainLoopJob *next;
  457. do {
  458. processJobs(server, &mlw->job, 1);
  459. next = (struct MainLoopJob*)mlw->node.next;
  460. UA_free(mlw);
  461. } while((mlw = next));
  462. //UA_free(head);
  463. }
  464. #endif
  465. UA_StatusCode UA_Server_run_startup(UA_Server *server, UA_UInt16 nThreads, UA_Boolean *running) {
  466. #ifdef UA_MULTITHREADING
  467. /* Prepare the worker threads */
  468. server->running = running; // the threads need to access the variable
  469. server->nThreads = nThreads;
  470. pthread_cond_init(&server->dispatchQueue_condition, 0);
  471. server->thr = UA_malloc(nThreads * sizeof(pthread_t));
  472. server->workerCounters = UA_malloc(nThreads * sizeof(UA_UInt32 *));
  473. for(UA_UInt32 i=0;i<nThreads;i++) {
  474. struct workerStartData *startData = UA_malloc(sizeof(struct workerStartData));
  475. startData->server = server;
  476. startData->workerCounter = &server->workerCounters[i];
  477. pthread_create(&server->thr[i], UA_NULL, (void* (*)(void*))workerLoop, startData);
  478. }
  479. /* try to execute the delayed callbacks every 10 sec */
  480. UA_Job processDelayed = {.type = UA_JOBTYPE_METHODCALL,
  481. .job.methodCall = {.method = dispatchDelayedJobs, .data = UA_NULL} };
  482. UA_Server_addRepeatedJob(server, processDelayed, 10000, UA_NULL);
  483. #endif
  484. /* Start the networklayers */
  485. for(size_t i = 0; i < server->networkLayersSize; i++)
  486. server->networkLayers[i].start(server->networkLayers[i].nlHandle, &server->logger);
  487. return UA_STATUSCODE_GOOD;
  488. }
  489. UA_StatusCode UA_Server_run_mainloop(UA_Server *server, UA_Boolean *running) {
  490. #ifdef UA_MULTITHREADING
  491. /* Run Work in the main loop */
  492. processMainLoopJobs(server);
  493. #endif
  494. /* Process repeated work */
  495. UA_UInt16 timeout = processRepeatedJobs(server);
  496. /* Get work from the networklayer */
  497. for(size_t i = 0; i < server->networkLayersSize; i++) {
  498. UA_ServerNetworkLayer *nl = &server->networkLayers[i];
  499. UA_Job *jobs;
  500. UA_Int32 jobsSize;
  501. if(*running) {
  502. if(i == server->networkLayersSize-1)
  503. jobsSize = nl->getJobs(nl->nlHandle, &jobs, timeout);
  504. else
  505. jobsSize = nl->getJobs(nl->nlHandle, &jobs, 0);
  506. } else
  507. jobsSize = server->networkLayers[i].stop(nl->nlHandle, &jobs);
  508. #ifdef UA_MULTITHREADING
  509. /* Filter out delayed work */
  510. for(UA_Int32 k=0;k<jobsSize;k++) {
  511. if(jobs[k].type != UA_JOBTYPE_DELAYEDMETHODCALL)
  512. continue;
  513. addDelayedJob(server, &jobs[k]);
  514. jobs[k].type = UA_JOBTYPE_NOTHING;
  515. }
  516. /* Dispatch work to the worker threads */
  517. dispatchJobs(server, jobs, jobsSize);
  518. /* Trigger sleeping worker threads */
  519. if(jobsSize > 0)
  520. pthread_cond_broadcast(&server->dispatchQueue_condition);
  521. #else
  522. processJobs(server, jobs, jobsSize);
  523. if(jobsSize > 0)
  524. UA_free(jobs);
  525. #endif
  526. }
  527. return UA_STATUSCODE_GOOD;
  528. }
  529. UA_StatusCode UA_Server_run_shutdown(UA_Server *server, UA_UInt16 nThreads){
  530. #ifdef UA_MULTITHREADING
  531. /* Wait for all worker threads to finish */
  532. for(UA_UInt32 i=0;i<nThreads;i++) {
  533. pthread_join(server->thr[i], UA_NULL);
  534. UA_free(server->workerCounters[i]);
  535. }
  536. UA_free(server->workerCounters);
  537. UA_free(server->thr);
  538. /* Manually finish the work still enqueued */
  539. emptyDispatchQueue(server);
  540. /* Process the remaining delayed work */
  541. struct DelayedJobs *dw = server->delayedJobs;
  542. while(dw) {
  543. processJobs(server, dw->jobs, dw->jobsCount);
  544. struct DelayedJobs *next = dw->next;
  545. UA_free(dw->workerCounters);
  546. UA_free(dw);
  547. dw = next;
  548. }
  549. #endif
  550. return UA_STATUSCODE_GOOD;
  551. }
  552. UA_StatusCode UA_Server_run(UA_Server *server, UA_UInt16 nThreads, UA_Boolean *running) {
  553. UA_Server_run_startup(server, nThreads, running);
  554. while(*running) {
  555. UA_Server_run_mainloop(server, running);
  556. }
  557. UA_Server_run_shutdown(server, nThreads);
  558. return UA_STATUSCODE_GOOD;
  559. }