ngc 1.5.0
 
Loading...
Searching...
No Matches
ngcbSHM.hpp
Go to the documentation of this file.
1
8#ifndef ngcbSHM_H
9#define ngcbSHM_H
10
11
12#ifndef __cplusplus
13#error This is a C++ include file and cannot be used from plain C
14#endif
15
16#include <stdio.h>
17#include <stdlib.h>
18#include <string.h>
19#include <unistd.h>
20#include <fcntl.h>
21#include <errno.h>
22#include <poll.h>
23#include <sys/socket.h>
24#include <sys/wait.h>
25#include <sys/time.h>
26#include <atomic>
27
29
30#ifdef _PTHREAD_H
31#define NGCB_SHM_THR 1
32#endif
33
37#define ngcbSHM_MAX_CLIENTS 8
38
39/*
40 * Shared memory id special result codes
41 */
42#define ngcbSHM_ID_ERR -1 //< IPC error
43#define ngcbSHM_ID_RETRY -2 //< try again
44#define ngcbSHM_ID_IVLD -3 //< invalid
45
49#define ngcbSHM_VERSION 1
50
56 int semId;
57
60
62 int shmId;
63
65 unsigned int event;
66
68 ngcb_shmevt_t() {semId = -1; reserved1 = 0x0; shmId = -1; event = 0;}
69};
70
76 int semId;
77
79 int retry;
80
82 unsigned int event;
83
86
88 ngcb_shmio_t() {semId = -1; retry = 0; event = 0; reserved = 0x0;}
89};
90
97
99 int frame;
100
103
105 int sx;
106
108 int sy;
109
111 int nx;
112
114 int ny;
115
118
121
123 int bzero;
124
126 int fcnt;
127
129 int fcnt0;
130
133
135 int ndit;
136
138 int flags;
139
142
144 char bytes[308];
145
148 frame = 0;
149 bitPix = 8;
150 sx = 0; sy = 0; nx = 0; ny = 0;
151 endian = 1;
152 fullFrame = 0;
153 bzero = 0;
154 status = 0;
155 fcnt = 0;
156 fcnt0 = 0;
157 ndit = 0;
158 expCnt = 0;
159 flags = 0;
160 memset(bytes, 0, sizeof(bytes));
161 }
162};
163
168public:
171
173 virtual ~ngcbSHMTHR() {}
174
176 virtual void Thread() {};
177
178protected:
179private:
180};
181
182#ifdef NGCB_SHM_THR
186static void *ngcbShmThread(void *arg) {
187 ngcbSHMTHR *evt = static_cast<ngcbSHMTHR *>(arg);
188 sigset_t set;
189
190 sigfillset(&set);
191 pthread_sigmask(SIG_BLOCK, &set, (sigset_t *)NULL);
192
193 if (evt != NULL) evt->Thread();
194
195 return (NULL);
196}
197#endif
198
202class ngcbSHM: public ngcbSHMTHR {
203public:
205 char name[64];
206
209
212 /*
213 * Initialize member data
214 */
215 Init_();
216 initialized_ = 1;
217 }
218
220 explicit ngcbSHM(const char *s) {
221 /*
222 * Initialize member data
223 */
224 Init_();
225
226 /*
227 * Check name
228 */
229 if (s == NULL)
230 {
231 strcpy(ErrMsg(), "invalid name");
232 return;
233 }
234 else if ((strlen(s) < 1) || (strlen(s) >= (sizeof(name) - 1)))
235 {
236 strcpy(ErrMsg(), "invalid name");
237 return;
238 }
239 else if (s[0] == '.')
240 {
241 extPipe_ = 1; // use external pipe (ngcbShmEvt dispatcher)
242 strcpy(name, s + 1);
243 }
244 else if (s[0] == ':')
245 {
246 intPipe_ = 1; // use internal pipe (dispatcher thread)
247 strcpy(name, s + 1);
248 }
249 else
250 {
251 strcpy(name, s);
252 }
253
254 if (strlen(name) == 0)
255 {
256 strcpy(ErrMsg(), "invalid name");
257 return;
258 }
259 else if (extPipe_)
260 {
261 /*
262 * Open external event pipe
263 */
264 regCnt_ = 0; // not used
265 fd_ = OpenPipe_();
266 if (fd_ < 0)
267 {
268 return;
269 }
270 }
271 else
272 {
273 /*
274 * Register to event interface
275 */
276 regCnt_ = 10;
277 if (Register_() != 0)
278 {
279 return;
280 }
281 }
282
283 initialized_ = 1;
284 }
285
287 virtual ~ngcbSHM() {
288 /*
289 * Shutdown event pipes (if applicable)
290 */
291 ClosePipe_();
292
293 if (shm_ != NULL)
294 {
295 /*
296 * Detach from shared memory
297 */
298 shmdt(shm_); shm_ = NULL;
299 }
300
301 if (evt_ != NULL)
302 {
303 /*
304 * Detach from event interface
305 */
306 shmdt(evt_); evt_ = NULL;
307 }
308
309 /*
310 * Check, if we have to cleanup the event resources
311 */
312 if (semEvt_ >= 0)
313 {
314 /*
315 * Detach from events
316 */
317 if (clientId_ >= 0) DetachEvt_();
318
319 if (EvtTryLock_() == 0)
320 {
321 if (NumClient() == 1)
322 {
323 /*
324 * We are the last user - remove the shared memory event
325 */
326 if (evtId_ != -1)
327 {
328 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
329 }
330
331 /*
332 * Remove the semaphore
333 */
334 SemEvtRemove_();
335
336 /*
337 * Remove the IPC-key file
338 */
339 if (strlen(name) > 0)
340 {
341 char fileName[256];
342 snprintf(fileName, sizeof(fileName), "/tmp/ngcshmio_%u_%s",
343 getuid(), name);
344 remove(fileName);
345 }
346 }
347 else
348 {
349 struct sembuf semOp[4];
350
351 /*
352 * Unregister
353 */
354 memset(&semOp[0], 0, sizeof(struct sembuf));
355 semOp[0].sem_num = 1;
356 semOp[0].sem_op = -1;
357 semOp[0].sem_flg = SEM_UNDO;
358 if (semop(semEvt_, semOp, 1) != 0)
359 {
360 if (errno == EIDRM) semEvt_ = -1; // semaphore has been removed
361 }
362 EvtUnlock_();
363 }
364 }
365 }
366 }
367
369 int Initialized() {return (initialized_);}
370
372 char *ErrMsg() {return (erms_);}
373
375 int Fd() {return (fd_);}
376
378 int IntPipe() {return (intPipe_);}
379
381 int NumClient() {
382 if (semEvt_ >= 0)
383 {
384 ngcbSEM_CTL ctl;
385 int val = semctl(semEvt_, 1, GETVAL, ctl);
386 if (val == -1)
387 {
388 sprintf(ErrMsg(), "semaphore error - %s", strerror(errno));
389 return (0);
390 }
391 erms_[0] = '\0';
392 return (val);
393 }
394 else
395 {
396 return (1);
397 }
398 }
399
401 void Remove() {
402 /*
403 * Remove event semaphore set
404 */
405 SemEvtRemove_();
406
407 if ((evtId_ != -1) && (evt_ != NULL))
408 {
409 /*
410 * Delete the shared memory event
411 */
412 shmdt(evt_); evt_ = NULL;
413 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
414 }
415
416 /*
417 * Remove the IPC-key file
418 */
419 if (strlen(name) > 0)
420 {
421 char fileName[256];
422 snprintf(fileName, sizeof(fileName), "/tmp/ngcshmio_%u_%s", getuid(),
423 name);
424 remove(fileName);
425 }
426
427 /*
428 * Interface can no longer be used
429 */
430 initialized_ = 0;
431 }
432
434 int WrInit(int semId) {
435 ngcbSEM_CTL ctl;
436
437 /*
438 * Initialize semaphore
439 */
440 ctl.val = 1; // write lock
441 if (semctl(semId, 0, SETVAL, ctl) == -1)
442 {
443 sprintf(ErrMsg(), "semaphore init failed - %s", strerror(errno));
444 return (-1);
445 }
446 else
447 {
448 ctl.val = 0; // readers unlocked
449 if (semctl(semId, 1, SETVAL, ctl) == -1)
450 {
451 sprintf(ErrMsg(), "semaphore init failed - %s", strerror(errno));
452 return (-1);
453 }
454 }
455
456 return (0);
457 }
458
460 int WrLock(int semId) {
461 if (semop(semId, semOpWrLock_, 3) != 0)
462 {
463 if (errno == EINTR)
464 {
465 while (semop(semId, semOpWrLock_, 3) != 0)
466 {
467 if (errno != EINTR)
468 {
469 sprintf(ErrMsg(), "cannot lock semaphore - %s", strerror(errno));
470 return (-1);
471 }
472 }
473 }
474 else
475 {
476 sprintf(ErrMsg(), "cannot lock semaphore - %s", strerror(errno));
477 return (-1);
478 }
479 }
480
481 return (0);
482 }
483
485 int WrTryLock(int semId) {
486 /*
487 * Try to get write-lock
488 */
489 if (semop(semId, semOpWrTryLock_, 3) != 0)
490 {
491 return (-1);
492 }
493 else
494 {
495 return (0);
496 }
497 }
498
500 int WrUnlock(int semId) {
501 if (semop(semId, semOpWrUnlock_, 1) != 0)
502 {
503 sprintf(ErrMsg(), "semaphore error - %s", strerror(errno));
504 return (-1);
505 }
506 else
507 {
508 return (0);
509 }
510 }
511
513 void RdWait(int semId) {
514 if (semId >= 0)
515 {
516 if (semop(semId, semOpRdWait_, 1) != 0)
517 {
518 if (errno == EINTR)
519 {
520 while (semop(semId, semOpRdWait_, 1) != 0)
521 {
522 if (errno != EINTR)
523 {
524 return;
525 }
526 }
527 }
528 }
529 }
530 }
531
533 int RdCnt(int semId) {
534 ngcbSEM_CTL ctl;
535 return (semctl(semId, 1, GETVAL, ctl));
536 }
537
539 void RdUnlock(int semId) {
540 if (extPipe_)
541 {
542 if ((fd_ >= 0) && !extPipeErr_)
543 {
544 int ack = 1;
545 if (write(fd_, &ack, sizeof(int)) != sizeof(int)) extPipeErr_ = 1;
546 }
547 }
548 else
549 {
550 RdUnlock_(semId);
551 }
552 }
553
556
558 int ShmLock(int shmId, int noWait = 0) {
559 void *m;
560
561 if (shmId < 0)
562 {
563 strcpy(ErrMsg(), "invalid shared memory id");
564 return (-1);
565 }
566
567 /*
568 * Attach to shared memory
569 */
570 m = shmat(shmId, 0, SHM_RDONLY);
571 if (m == (void *)-1)
572 {
573 sprintf(ErrMsg(), "cannot lock shared memory - %s", strerror(errno));
574 return (-1);
575 }
576 else
577 {
578 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(m);
579
580 /*
581 * Get read-lock
582 */
583 if (RdLock_(io->semId, noWait) != 0)
584 {
585 sprintf(ErrMsg(), "cannot lock shared memory - %s", strerror(errno));
586 shmdt(m);
587 return (-1);
588 }
589
590 shmdt(m);
591 return (0);
592 }
593 }
594
596 void ShmNoUndo() {
597 semOpRdLock_[1].sem_flg = 0;
598 semOpRdUnlock_[0].sem_flg = IPC_NOWAIT;
599 }
600
602 int ShmUnlock(int shmId) {
603 void *m;
604
605 if (shmId < 0)
606 {
607 strcpy(ErrMsg(), "invalid shared memory id");
608 return (-1);
609 }
610
611 /*
612 * Attach to shared memory
613 */
614 m = shmat(shmId, 0, SHM_RDONLY);
615 if (m == (void *)-1)
616 {
617 sprintf(ErrMsg(), "cannot lock shared memory - %s", strerror(errno));
618 return (-1);
619 }
620 else
621 {
622 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(m);
623
624 /*
625 * Release read-lock
626 */
627 RdUnlock_(io->semId);
628 shmdt(m);
629 return (0);
630 }
631
632 }
633
635 void *AttachBuffer(int shmId, void *hdr, int noWait = 0) {
636 void *buffer = NULL;
637 ngcb_shmio_t *hdr_io = static_cast<ngcb_shmio_t *>(hdr);
638
639 /*
640 * Reset the retry-flag in any case
641 */
642 hdr_io->retry = 0;
643
644 if (shmId < 0)
645 {
646 strcpy(ErrMsg(), "invalid shared memory id");
647 return (NULL);
648 }
649
650 /*
651 * Attach to shared memory
652 */
653 shm_ = shmat(shmId, 0, SHM_RDONLY);
654 if (shm_ != (void *)-1)
655 {
656 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(shm_);
657
658 /*
659 * Get read-lock
660 */
661 if (RdLock_(io->semId, noWait) != 0)
662 {
663 hdr_io->retry = 1;
664 sprintf(ErrMsg(), "cannot lock buffer - %s", strerror(errno));
665 shmdt(shm_); shm_ = NULL;
666 return (NULL);
667 }
668
669 /*
670 * Get header
671 */
672 memcpy(hdr, shm_, sizeof(ngcb_shmhdr_t));
673
674 /*
675 * Align the event counter. So the Wait() function can ensure
676 * that we do not receive the same event data again.
677 */
678 evtCnt_ = io->event;
679
680 /*
681 * Assign buffer
682 */
683 buffer = (void *)(((char *)shm_) + sizeof(ngcb_shmhdr_t));
684 }
685 else
686 {
687 shm_ = NULL;
688 sprintf(ErrMsg(), "cannot attach to shared memory - %s",
689 strerror(errno));
690 }
691
692 return (buffer);
693 }
694
696 void *TryAttachBuffer(int shmId, void *hdr) {
697 return (AttachBuffer(shmId, hdr, 1));
698 }
699
702 if (shm_ != NULL)
703 {
704 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(shm_);
705
706 /*
707 * Release read-lock
708 */
709 RdUnlock_(io->semId);
710
711 /*
712 * Detach from shared memory
713 */
714 shmdt(shm_); shm_ = NULL;
715 }
716 }
717
719 unsigned int EvtCnt() {return (evtCnt_);}
720
722 unsigned int NextEvt() {
723 unsigned int n = evtCnt_ + 1;
724 if (n == 0) n = 1;
725 return (n);
726 }
727
729 int AttachEvt() {
730 struct sembuf semOp[1];
731 int i;
732
733 if (extPipe_)
734 {
735 if ((fd_ < 0) || extPipeErr_)
736 {
737 strcpy(ErrMsg(), "event interface not initialized");
738 return (-1);
739 }
740
741 return (0);
742 }
743 else if (semEvt_ < 0)
744 {
745 /*
746 * Not yet initialized
747 */
748 strcpy(ErrMsg(), "event interface not initialized");
749 return (-1);
750 }
751 else if (clientId_ >= 0)
752 {
753 /*
754 * We are already attached
755 */
756 return (0);
757 }
758 else if (EvtLock_() != 0)
759 {
760 /*
761 * We cannot proceed
762 */
763 return (-1);
764 }
765
766 cancelled_ = false;
767 memset(&semOp[0], 0, sizeof(struct sembuf));
768 semOp[0].sem_op = -2;
769 semOp[0].sem_flg = IPC_NOWAIT|SEM_UNDO;
770
771 for (i=0;i<ngcbSHM_MAX_CLIENTS;i++)
772 {
773 semOp[0].sem_num = static_cast<unsigned short>(2 + i);
774 if (semop(semEvt_, semOp, 1) == 0)
775 {
776 /*
777 * We are attached
778 */
779 clientId_ = i;
780 EvtUnlock_();
781
782 if (intPipe_)
783 {
784 /*
785 * Start event dispatcher thread
786 */
787 thr_ = 1; // thread interface enabled
788 if (StartThread_() != 0)
789 {
790 thr_ = 0; // thread interface disabled
791 return (-1);
792 }
793 }
794
795 return (0);
796 }
797 else if (errno == EIDRM)
798 {
799 /*
800 * Event semaphore has been removed
801 */
802 semEvt_ = -1;
803 }
804 }
805
806 EvtUnlock_();
807
808 strcpy(ErrMsg(), "too many clients");
809 return (-1);
810 }
811
813 int Wait(int noWait = 0) {
814 int shmId = ngcbSHM_ID_IVLD;
815 if (extPipe_)
816 {
817 if ((fd_ >= 0) && !extPipeErr_)
818 {
819 int ack = 1;
820 int res = WaitPipe_(noWait);
821 if (res < 0)
822 {
823 return (ngcbSHM_ID_ERR);
824 }
825 else if (res == 0)
826 {
827 return (ngcbSHM_ID_RETRY);
828 }
829
830 /*
831 * Read shared memory id
832 */
833 if (read(fd_, &shmId, sizeof(int)) != sizeof(int))
834 {
835 extPipeErr_ = 1;
836 strcpy(ErrMsg(), "error reading event header");
837 return (ngcbSHM_ID_ERR);
838 }
839
840 if (shmId >= 0)
841 {
842 shm_ = shmat(shmId, 0, SHM_RDONLY);
843 if (shm_ != (void *)-1)
844 {
845 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(shm_);
846
847 /*
848 * Get event counter
849 */
850 evtCnt_ = io->event;
851 shmdt(shm_); shm_ = NULL;
852 }
853 else
854 {
855 shm_ = NULL;
856 sprintf(ErrMsg(), "shared memory error - %s", strerror(errno));
857 return (ngcbSHM_ID_ERR);
858 }
859 }
860
861 /*
862 * Send acknowledge
863 */
864 if (write(fd_, &ack, sizeof(int)) != sizeof(int))
865 {
866 extPipeErr_ = 1;
867 strcpy(ErrMsg(), "error sending event acknowledge");
868 return (ngcbSHM_ID_ERR);
869 }
870
871 return (shmId);
872 }
873 else
874 {
875 strcpy(ErrMsg(), "event dispatcher is not running");
876 return (ngcbSHM_ID_ERR);
877 }
878 }
879
880 /*
881 * Wait for next event and lock event
882 */
883 shmId = Wait_(noWait);
884 if ((shmId >= 0) || (shmId == ngcbSHM_ID_IVLD))
885 {
886 /*
887 * Unlock the event
888 */
889 EvtUnlock_();
890 }
891
892 return (shmId);
893 }
894
896 int TryWait() {return (Wait(1));}
897
899 void *GetBuffer(void *hdr, int noWait = 0) {
900 void *buffer = NULL;
901 int shmId = ngcbSHM_ID_IVLD;
902
903 if (extPipe_)
904 {
905 if ((fd_ >= 0) && !extPipeErr_)
906 {
907 if (noWait)
908 {
909 int res = WaitPipe_(1);
910 if (res < 0)
911 {
912 strcpy(ErrMsg(), "error reading event header");
913 return (NULL);
914 }
915 else if (res == 0)
916 {
917 strcpy(ErrMsg(), "busy");
918 return (NULL);
919 }
920 }
921
922 /*
923 * Read shared memory id
924 */
925 if (read(fd_, &shmId, sizeof(int)) != sizeof(int))
926 {
927 extPipeErr_ = 1;
928 strcpy(ErrMsg(), "error reading event header");
929 return (NULL);
930 }
931
932 if (shmId >= 0)
933 {
934 shm_ = shmat(shmId, 0, SHM_RDONLY);
935 if (shm_ != (void *)-1)
936 {
937 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(shm_);
938
939 /*
940 * Get header
941 */
942 memcpy(hdr, shm_, sizeof(ngcb_shmhdr_t));
943 evtCnt_ = io->event;
944
945 /*
946 * Assign buffer
947 */
948 buffer = (void *)(((char *)shm_) + sizeof(ngcb_shmhdr_t));
949 }
950 else
951 {
952 shm_ = NULL;
953 sprintf(ErrMsg(), "cannot attach to shared memory - %s",
954 strerror(errno));
955 return (NULL);
956 }
957 }
958 else
959 {
960 int ack = 1;
961
962 /*
963 * Send acknowledge
964 */
965 if (write(fd_, &ack, sizeof(int)) != sizeof(int))
966 {
967 extPipeErr_ = 1;
968 strcpy(ErrMsg(), "error sending event acknowledge");
969 return (NULL);
970 }
971 }
972
973 return (buffer);
974 }
975 else
976 {
977 strcpy(ErrMsg(), "event dispatcher is not running");
978 return (NULL);
979 }
980 }
981
982 /*
983 * Wait for next event and lock event
984 */
985 shmId = Wait_(noWait);
986 if (shmId >= 0)
987 {
988 /*
989 * Attach buffer and unlock the event
990 */
991 buffer = AttachBuffer(shmId, hdr, 1);
992 EvtUnlock_();
993 }
994 else if (shmId == ngcbSHM_ID_RETRY)
995 {
996 strcpy(ErrMsg(), "busy");
997 }
998 else if (shmId == ngcbSHM_ID_IVLD)
999 {
1000 /*
1001 * Unlock the event
1002 */
1003 EvtUnlock_();
1004 }
1005
1006 return (buffer);
1007 }
1008
1010 void *TryGetBuffer(void *hdr) {return (GetBuffer(hdr, 1));}
1011
1013 int GetHdr(void *hdr, int noWait = 0) {
1014 int shmId = ngcbSHM_ID_IVLD;
1015
1016 if (extPipe_)
1017 {
1018 if ((fd_ >= 0) && !extPipeErr_)
1019 {
1020 if (noWait)
1021 {
1022 int res = WaitPipe_(1);
1023 if (res < 0)
1024 {
1025 strcpy(ErrMsg(), "error reading event header");
1026 return (ngcbSHM_ID_ERR);
1027 }
1028 else if (res == 0)
1029 {
1030 return (ngcbSHM_ID_RETRY);
1031 }
1032 }
1033
1034 /*
1035 * Read shared memory id
1036 */
1037 if (read(fd_, &shmId, sizeof(int)) != sizeof(int))
1038 {
1039 extPipeErr_ = 1;
1040 strcpy(ErrMsg(), "error reading event header");
1041 return (ngcbSHM_ID_ERR);
1042 }
1043
1044 if (shmId >= 0)
1045 {
1046 shm_ = shmat(shmId, 0, SHM_RDONLY);
1047 if (shm_ != (void *)-1)
1048 {
1049 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(shm_);
1050
1051 /*
1052 * Get header
1053 */
1054 memcpy(hdr, shm_, sizeof(ngcb_shmhdr_t));
1055 evtCnt_ = io->event;
1056 shmdt(shm_); shm_ = NULL;
1057 }
1058 else
1059 {
1060 shm_ = NULL;
1061 sprintf(ErrMsg(), "cannot attach to shared memory - %s",
1062 strerror(errno));
1063 return (ngcbSHM_ID_ERR);
1064 }
1065 }
1066 else
1067 {
1068 int ack = 1;
1069
1070 /*
1071 * Send acknowledge
1072 */
1073 if (write(fd_, &ack, sizeof(int)) != sizeof(int))
1074 {
1075 extPipeErr_ = 1;
1076 strcpy(ErrMsg(), "error sending event acknowledge");
1077 return (ngcbSHM_ID_ERR);
1078 }
1079 }
1080
1081 return (shmId);
1082 }
1083 else
1084 {
1085 strcpy(ErrMsg(), "event dispatcher is not running");
1086 return (ngcbSHM_ID_ERR);
1087 }
1088 }
1089
1090 /*
1091 * Wait for next event and lock event
1092 */
1093 shmId = Wait_(noWait);
1094 if (shmId >= 0)
1095 {
1096 /*
1097 * Attach buffer
1098 */
1099 if (AttachBuffer(shmId, hdr, 1) != NULL)
1100 {
1101 /*
1102 * The header is filled - we detach again
1103 */
1104 if (shm_ != NULL)
1105 {
1106 shmdt(shm_); shm_ = NULL;
1107 }
1108 }
1109 else
1110 {
1111 ngcb_shmio_t *io = static_cast<ngcb_shmio_t *>(hdr);
1112 if (io->retry)
1113 {
1114 /*
1115 * We just declare it as invalid
1116 */
1117 shmId = ngcbSHM_ID_IVLD;
1118 }
1119 else
1120 {
1121 /*
1122 * This was a serious error
1123 */
1124 shmId = ngcbSHM_ID_ERR;
1125 }
1126 }
1127
1128 /*
1129 * Unlock the event
1130 */
1131 EvtUnlock_();
1132 }
1133 else if (shmId == ngcbSHM_ID_IVLD)
1134 {
1135 /*
1136 * Unlock the event
1137 */
1138 EvtUnlock_();
1139 }
1140
1141 return (shmId);
1142 }
1143
1145 int TryGetHdr(void *hdr) {return (GetHdr(hdr, 1));}
1146
1148 int Send(int shmId) {
1149 struct sembuf semOp[2];
1150 int i;
1151
1152 /*
1153 * Next event id
1154 */
1155 if (++evtCnt_ == 0) evtCnt_ = 1;
1156
1157 /*
1158 * Write share memory id
1159 */
1160 if (EvtLock_() != 0)
1161 {
1162 return (-1);
1163 }
1164 evt_->shmId = shmId;
1165 evt_->event = evtCnt_;
1166 EvtUnlock_();
1167
1168 /*
1169 * Initialize semaphore operation
1170 */
1171 memset(&semOp[0], 0, sizeof(struct sembuf));
1172 memset(&semOp[1], 0, sizeof(struct sembuf));
1173 semOp[0].sem_op = 0;
1174 semOp[0].sem_flg = IPC_NOWAIT;
1175 semOp[1].sem_op = 1;
1176 semOp[1].sem_flg = 0;
1177
1178 /*
1179 * Send event to all clients
1180 */
1181 for (i=0;i<ngcbSHM_MAX_CLIENTS;i++)
1182 {
1183 semOp[0].sem_num = static_cast<unsigned short>(2 + i);
1184 semOp[1].sem_num = static_cast<unsigned short>(2 + i);
1185 if (semop(semEvt_, &semOp[0], 2) != 0)
1186 {
1187 if (errno == EAGAIN)
1188 {
1189 /*
1190 * Client was not ready
1191 */
1192 continue;
1193 }
1194 else if (errno == EIDRM)
1195 {
1196 /*
1197 * Event semaphore has been removed
1198 */
1199 semEvt_ = -1;
1200 }
1201 sprintf(ErrMsg(), "semaphore error - %s\n", strerror(errno));
1202 return (-1);
1203 }
1204 }
1205
1206 return (0);
1207 }
1208
1210 int Release() {
1211 struct sembuf semOp[2];
1212
1213 if (semEvt_ < 0)
1214 {
1215 strcpy(ErrMsg(), "event interface not initialized");
1216 return (-1);
1217 }
1218 else if (clientId_ < 0)
1219 {
1220 strcpy(ErrMsg(), "not attached to events");
1221 return (-1);
1222 }
1223
1224 if (EvtLock_() == 0)
1225 {
1226 cancelled_ = true;
1227 EvtUnlock_();
1228 }
1229 else
1230 {
1231 cancelled_ = true;
1232 }
1233
1234 /*
1235 * Initialize semaphore operation
1236 */
1237 memset(&semOp[0], 0, sizeof(struct sembuf));
1238 memset(&semOp[1], 0, sizeof(struct sembuf));
1239 semOp[0].sem_op = 0;
1240 semOp[0].sem_flg = IPC_NOWAIT;
1241 semOp[0].sem_num = static_cast<unsigned short>(2 + clientId_);
1242 semOp[1].sem_op = 1;
1243 semOp[1].sem_flg = 0;
1244 semOp[1].sem_num = static_cast<unsigned short>(2 + clientId_);
1245
1246 /*
1247 * Send event
1248 */
1249 if (semop(semEvt_, &semOp[0], 2) != 0)
1250 {
1251 if (errno == EAGAIN)
1252 {
1253 /*
1254 * Event has already been sent but client was not ready
1255 */
1256 return (0);
1257 }
1258 else if (errno == EIDRM)
1259 {
1260 /*
1261 * Event semaphore has been removed
1262 */
1263 semEvt_ = -1;
1264 }
1265 sprintf(ErrMsg(), "semaphore error - %s\n", strerror(errno));
1266 return (-1);
1267 }
1268
1269 return (0);
1270 }
1271
1273 void Invalidate() {
1274 if (EvtLock_() == 0)
1275 {
1276 evt_->shmId = ngcbSHM_ID_ERR;
1277 evt_->event = 0;
1278 EvtUnlock_();
1279 }
1280 }
1281
1283 int RecvEvt() {
1284 int shmId = ngcbSHM_ID_IVLD;
1285 if (intPipe_ && (fd_ >= 0))
1286 {
1287 if (read(fd_, &shmId, sizeof(int)) != sizeof(int))
1288 {
1289 strcpy(ErrMsg(), "error receiving immage event from pipe");
1290 return (ngcbSHM_ID_ERR);
1291 }
1292 }
1293 return (shmId);
1294 }
1295
1297 int SendAck() {
1298 if (intPipe_ && (fd_ >= 0))
1299 {
1300 int ack = 1;
1301 if (write(fd_, &ack, sizeof(int)) != sizeof(int))
1302 {
1303 strcpy(ErrMsg(), "error sending event acknowledge to pipe");
1304 return (-1);
1305 }
1306 }
1307 return (0);
1308 }
1309
1311 void Thread() {
1312 int fd = pipe_;
1313
1314 if ((fd >= 0) && (sharedHeader != NULL))
1315 {
1316 int ack = 1;
1317 int shmId = ngcbSHM_ID_IVLD;
1318
1319 /*
1320 * Thread main-loop
1321 */
1322 while (ack)
1323 {
1324 /*
1325 * Get next event header and acquire read-lock
1326 */
1327 shmId = GetHdr(sharedHeader);
1328 if ((shmId == ngcbSHM_ID_ERR) || (shmId == ngcbSHM_ID_RETRY))
1329 {
1330 /*
1331 * An error occured or the event has been cancelled due to
1332 * an interface shutdown
1333 */
1334 break;
1335 }
1336
1337 if ((shmId >= 0) || (shmId == ngcbSHM_ID_IVLD))
1338 {
1339 if (!sharedHeader->io.retry && (EvtCnt() != 0))
1340 {
1341 /*
1342 * Send shared memory id
1343 */
1344 if (write(fd, &shmId, sizeof(int)) != sizeof(int))
1345 {
1346 break;
1347 }
1348
1349 /*
1350 * Wait for acknowledge
1351 */
1352 if (read(fd, &ack, sizeof(int)) != sizeof(int))
1353 {
1354 break;
1355 }
1356 }
1357
1358 if (shmId >= 0)
1359 {
1360 /*
1361 * Release read-lock
1362 */
1363 RdUnlock(sharedHeader->io.semId);
1364 shmId = ngcbSHM_ID_IVLD;
1365 }
1366 else
1367 {
1368 /*
1369 * Check, if our event pipe is still open
1370 */
1371 if (!pipeOpen_)
1372 {
1373 break;
1374 }
1375 }
1376 }
1377 }
1378
1379 if (shmId >= 0)
1380 {
1381 /*
1382 * Release read-lock
1383 */
1384 RdUnlock(sharedHeader->io.semId);
1385 }
1386 }
1387 }
1388
1389protected:
1390
1391private:
1393 char erms_[256];
1394
1396 int initialized_;
1397
1399 void *shm_;
1400
1402 struct sembuf semOpWrLock_[3];
1403
1405 struct sembuf semOpWrTryLock_[3];
1406
1408 struct sembuf semOpWrUnlock_[1];
1409
1411 struct sembuf semOpRdWait_[1];
1412
1414 struct sembuf semOpRdLock_[2];
1415
1417 struct sembuf semOpRdUnlock_[1];
1418
1420 struct sembuf semOpEvtLock_[2];
1421
1423 struct sembuf semOpEvtTryLock_[2];
1424
1426 struct sembuf semOpEvtUnlock_[1];
1427
1429 int regCnt_;
1430
1432 int evtId_;
1433
1435 ngcb_shmevt_t *evt_;
1436
1438 int semEvt_;
1439
1441 unsigned int evtCnt_;
1442
1444 int clientId_;
1445
1447 int fd_;
1448
1450 int intPipe_;
1451
1453 int pipe_;
1454
1456 int pipeOpen_;
1457
1459 int extPipe_;
1460
1462 int extPipeErr_;
1463
1465 int pid_;
1466
1467#ifdef NGCB_SHM_THR
1469 pthread_t thrId_;
1470#endif
1471
1473 int thr_;
1474
1476 ngcb_shmhdr_t sharedHeader_;
1477
1479 std::atomic_bool cancelled_;
1480
1482 void Init_() {
1483 initialized_ = 0;
1484 name[0] = '\0';
1485 sharedHeader = &sharedHeader_; // link to our private resource
1486 erms_[0] = '\0';
1487 semEvt_ = -1;
1488 evtId_ = -1;
1489 evt_ = NULL;
1490 evtCnt_ = 0;
1491 shm_ = NULL;
1492 clientId_ = -1;
1493 fd_ = -1;
1494 intPipe_ = 0; pipe_ = -1; pipeOpen_ = 0; thr_ = 0;
1495 extPipe_ = 0; extPipeErr_ = 0; pid_ = 0;
1496 cancelled_ = false;
1497 regCnt_ = 0;
1498
1499#ifdef NGCB_SHM_THR
1500 /*
1501 * Initialize thread id
1502 */
1503 thrId_ = pthread_self();
1504#endif
1505
1506 /*
1507 * Semaphore operation:
1508 * - wait until there are no more readers
1509 * - wait until there is no writer
1510 * - increment writer by one
1511 */
1512 memset(&semOpWrLock_[0], 0, sizeof(struct sembuf));
1513 memset(&semOpWrLock_[1], 0, sizeof(struct sembuf));
1514 memset(&semOpWrLock_[2], 0, sizeof(struct sembuf));
1515 semOpWrLock_[0].sem_num = 1;
1516 semOpWrLock_[0].sem_op = 0;
1517 semOpWrLock_[0].sem_flg = 0;
1518 semOpWrLock_[1].sem_num = 0;
1519 semOpWrLock_[1].sem_op = 0;
1520 semOpWrLock_[1].sem_flg = 0;
1521 semOpWrLock_[2].sem_num = 0;
1522 semOpWrLock_[2].sem_op = 1;
1523 semOpWrLock_[2].sem_flg = SEM_UNDO;
1524
1525 /*
1526 * Semaphore operation: same as above with IPC_NOWAIT
1527 */
1528 memset(&semOpWrTryLock_[0], 0, sizeof(struct sembuf));
1529 memset(&semOpWrTryLock_[1], 0, sizeof(struct sembuf));
1530 memset(&semOpWrTryLock_[2], 0, sizeof(struct sembuf));
1531 semOpWrTryLock_[0].sem_num = 1;
1532 semOpWrTryLock_[0].sem_op = 0;
1533 semOpWrTryLock_[0].sem_flg = IPC_NOWAIT;
1534 semOpWrTryLock_[1].sem_num = 0;
1535 semOpWrTryLock_[1].sem_op = 0;
1536 semOpWrTryLock_[1].sem_flg = IPC_NOWAIT;
1537 semOpWrTryLock_[2].sem_num = 0;
1538 semOpWrTryLock_[2].sem_op = 1;
1539 semOpWrTryLock_[2].sem_flg = SEM_UNDO;
1540
1541 /*
1542 * Semaphore operation:
1543 * - check, if we own the write-lock - if not return error
1544 * - release write lock (decrement by one)
1545 */
1546 memset(&semOpWrUnlock_[0], 0, sizeof(struct sembuf));
1547 semOpWrUnlock_[0].sem_num = 0;
1548 semOpWrUnlock_[0].sem_op = -1;
1549 semOpWrUnlock_[0].sem_flg = IPC_NOWAIT|SEM_UNDO;
1550
1551 /*
1552 * Semaphore operation:
1553 * - wait until there are no more readers
1554 */
1555 memset(&semOpRdWait_[0], 0, sizeof(struct sembuf));
1556 semOpRdWait_[0].sem_num = 1;
1557
1558 /*
1559 * Semaphore operation:
1560 * - wait until there is no writer
1561 * - increment readers by one
1562 */
1563 memset(&semOpRdLock_[0], 0, sizeof(struct sembuf));
1564 memset(&semOpRdLock_[1], 0, sizeof(struct sembuf));
1565 semOpRdLock_[0].sem_num = 0;
1566 semOpRdLock_[0].sem_op = 0;
1567 semOpRdLock_[0].sem_flg = 0;
1568 semOpRdLock_[1].sem_num = 1;
1569 semOpRdLock_[1].sem_op = 1;
1570 semOpRdLock_[1].sem_flg = SEM_UNDO;
1571
1572 /*
1573 * Semaphore operation:
1574 * - check, if there is at least one read-lock - if not, return error
1575 * - decrement readers by one
1576 */
1577 memset(&semOpRdUnlock_[0], 0, sizeof(struct sembuf));
1578 semOpRdUnlock_[0].sem_num = 1;
1579 semOpRdUnlock_[0].sem_op = -1;
1580 semOpRdUnlock_[0].sem_flg = IPC_NOWAIT|SEM_UNDO;
1581
1582 /*
1583 * Semaphore operation:
1584 * - lock event access semaphore
1585 */
1586 memset(&semOpEvtLock_[0], 0, sizeof(struct sembuf));
1587 memset(&semOpEvtLock_[1], 0, sizeof(struct sembuf));
1588 semOpEvtLock_[1].sem_num = 0;
1589 semOpEvtLock_[1].sem_op = 1;
1590 semOpEvtLock_[1].sem_flg = SEM_UNDO;
1591
1592 /*
1593 * Semaphore operation:
1594 * - try to lock event access semaphore
1595 */
1596 memset(&semOpEvtTryLock_[0], 0, sizeof(struct sembuf));
1597 memset(&semOpEvtTryLock_[1], 0, sizeof(struct sembuf));
1598 semOpEvtTryLock_[0].sem_num = 0;
1599 semOpEvtTryLock_[0].sem_op = 0;
1600 semOpEvtTryLock_[0].sem_flg = IPC_NOWAIT;
1601 semOpEvtTryLock_[1].sem_num = 0;
1602 semOpEvtTryLock_[1].sem_op = 1;
1603 semOpEvtTryLock_[1].sem_flg = SEM_UNDO;
1604
1605 /*
1606 * Semaphore operation:
1607 * - unlock event access semaphore
1608 */
1609 memset(&semOpEvtUnlock_[0], 0, sizeof(struct sembuf));
1610 semOpEvtUnlock_[0].sem_num = 0;
1611 semOpEvtUnlock_[0].sem_op = -1;
1612 semOpEvtUnlock_[0].sem_flg = IPC_NOWAIT|SEM_UNDO;
1613 }
1614
1616 int Register_() {
1617 int newShm = 0;
1618 void *m;
1619 key_t key = (key_t)-1;
1620
1621 /*
1622 * Ensure that we are not yet bound to any resources
1623 */
1624 if (evt_ != NULL)
1625 {
1626 shmdt(evt_); evt_ = NULL;
1627 }
1628
1629 /*
1630 * Get IPC-key
1631 */
1632 if (strlen(name) > 0)
1633 {
1634 char fileName[256];
1635 int fd;
1636 snprintf(fileName, sizeof(fileName), "/tmp/ngcshmio_%u_%s", getuid(),
1637 name);
1638 fd = open(fileName, O_RDWR|O_CREAT|O_TRUNC, 0644);
1639 if (fd >= 0)
1640 {
1641 close(fd);
1642 key = ftok(fileName, 0x1f);
1643 }
1644 }
1645 else
1646 {
1647 char *evar = getenv("INS_ROOT");
1648 if (evar != NULL) key = ftok(evar, 0x1f);
1649 if (key == -1)
1650 {
1651 evar = getenv("INTROOT");
1652 if (evar != NULL) key = ftok(evar, 0x1f);
1653 }
1654 if (key == -1)
1655 {
1656 evar = getenv("HOME");
1657 if (evar != NULL) key = ftok(evar, 0x1f);
1658 }
1659 }
1660
1661 if (key == (key_t)-1)
1662 {
1663 strcpy(ErrMsg(), "unable to retrieve IPC-key");
1664 return (-1);
1665 }
1666
1667 /*
1668 * Try to access the shared memory event
1669 */
1670 evtId_ = shmget(key, sizeof(ngcb_shmevt_t), 0666);
1671 if (evtId_ == -1)
1672 {
1673 int tmp;
1674
1675 if (errno == EINVAL)
1676 {
1677 /*
1678 * The shared memory size may have changed. We remove the
1679 * exising region and re-create.
1680 */
1681 evtId_ = shmget(key, 0, 0666);
1682 if (evtId_ == -1)
1683 {
1684 sprintf(ErrMsg(), "error creating shared memory - %s",
1685 strerror(errno));
1686 return (-1);
1687 }
1688
1689 /*
1690 * Remove event id and re-create
1691 */
1692 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
1693 }
1694
1695 /*
1696 * Ensure that the event semaphore does not yet exist
1697 */
1698 tmp = semget(key, 0, 0666);
1699 if (tmp >= 0)
1700 {
1701 ngcbSEM_CTL ctl;
1702
1703 /*
1704 * Remove the semaphore
1705 */
1706 ctl.val = 0;
1707 semctl(tmp, 0, IPC_RMID, ctl);
1708 }
1709
1710 /*
1711 * Create a new shared memory region
1712 */
1713 evtId_ = shmget(key, sizeof(ngcb_shmevt_t), 0666|IPC_CREAT|IPC_EXCL);
1714 if (evtId_ == -1)
1715 {
1716 if (errno == EEXIST)
1717 {
1718 /*
1719 * We may have had a race condition here. So we try once
1720 * again to access it.
1721 */
1722 evtId_ = shmget(key, sizeof(ngcb_shmevt_t), 0666);
1723 if (evtId_ == -1)
1724 {
1725 if ((errno == ENOENT) && (regCnt_ > 0))
1726 {
1727 /*
1728 * Retry
1729 */
1730 regCnt_--;
1731 return(Register_());
1732 }
1733 sprintf(ErrMsg(), "error creating shared memory - %s",
1734 strerror(errno));
1735 return (-1);
1736 }
1737 }
1738 else
1739 {
1740 sprintf(ErrMsg(), "error creating shared memory - %s",
1741 strerror(errno));
1742 return (-1);
1743 }
1744 }
1745 else
1746 {
1747 /*
1748 * We actually are the creator
1749 */
1750 newShm = 1;
1751 }
1752 }
1753
1754 /*
1755 * Attach to the shared memory event
1756 */
1757 m = shmat(evtId_, 0, 0);
1758 if (m == (void *)-1)
1759 {
1760 int e = errno;
1761
1762 if (newShm)
1763 {
1764 /*
1765 * Delete the shared memory event again
1766 */
1767 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
1768 }
1769
1770 if ((e == EINVAL) && (regCnt_ > 0))
1771 {
1772 /*
1773 * Retry
1774 */
1775 regCnt_--;
1776 return (Register_());
1777 }
1778 else
1779 {
1780 sprintf(ErrMsg(), "cannot attach to shared memory - %s", strerror(e));
1781 return (-1);
1782 }
1783 }
1784 else
1785 {
1786 evt_ = static_cast<ngcb_shmevt_t *>(m);
1787 }
1788
1789 if (newShm)
1790 {
1791 ngcbSEM_CTL ctl;
1792 ngcb_shmevt_t initStat;
1793 unsigned short semArray[2 + ngcbSHM_MAX_CLIENTS];
1794 int i;
1795
1796 /*
1797 * Initialize the event. This will explicitly mark
1798 * the event semaphore as invalid.
1799 */
1800 memcpy(evt_, &initStat, sizeof(ngcb_shmevt_t));
1801 semArray[0] = 0;
1802 semArray[1] = 0;
1803 for (i=0;i<ngcbSHM_MAX_CLIENTS;i++) semArray[2 + i] = 2;
1804
1805 /*
1806 * Create the event semaphore
1807 */
1808 semEvt_ = semget(key, 2 + ngcbSHM_MAX_CLIENTS, 0666|IPC_CREAT|IPC_EXCL);
1809 if (semEvt_ == -1)
1810 {
1811 sprintf(ErrMsg(), "error creating semaphore - %s", strerror(errno));
1812
1813 /*
1814 * Delete the shared memory event
1815 */
1816 shmdt(evt_); evt_ = NULL;
1817 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
1818 return (-1);
1819 }
1820
1821 /*
1822 * Initialize the event semaphore
1823 */
1824 ctl.array = semArray;
1825 if (semctl(semEvt_, 0, SETALL, ctl) == -1)
1826 {
1827 int e = errno;
1828
1829 sprintf(ErrMsg(), "semaphore error - %s", strerror(e));
1830
1831 if (e == EIDRM)
1832 {
1833 /*
1834 * Semaphore has been removed
1835 */
1836 semEvt_ = -1;
1837 }
1838 else
1839 {
1840 /*
1841 * Delete the event semaphore again
1842 */
1843 SemEvtRemove_();
1844 }
1845
1846 /*
1847 * Delete the shared memory event
1848 */
1849 shmdt(evt_); evt_ = NULL;
1850 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
1851
1852 if ((e == EIDRM) && (regCnt_ > 0))
1853 {
1854 /*
1855 * Retry
1856 */
1857 regCnt_--;
1858 return (Register_());
1859 }
1860 else
1861 {
1862 return (-1);
1863 }
1864 }
1865
1866 /*
1867 * Assign the event semaphore
1868 */
1869 if (EvtLock_() != 0)
1870 {
1871 /*
1872 * Delete the event semaphore again
1873 */
1874 SemEvtRemove_();
1875
1876 /*
1877 * Delete the shared memory event
1878 */
1879 shmdt(evt_); evt_ = NULL;
1880 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
1881 if (regCnt_ > 0)
1882 {
1883 /*
1884 * Retry
1885 */
1886 regCnt_--;
1887 return (Register_());
1888 }
1889 else
1890 {
1891 return (-1);
1892 }
1893 }
1894 evt_->semId = semEvt_;
1895 EvtUnlock_();
1896 }
1897 else
1898 {
1899 /*
1900 * Get event semaphore - this should exist as the associated
1901 * shared memory has already been created
1902 */
1903 semEvt_ = semget(key, 0, 0666);
1904 if (semEvt_ == -1)
1905 {
1906 if ((errno == ENOENT) && (regCnt_ > 0))
1907 {
1908 /*
1909 * We may have had a race condition with the semaphore
1910 * creator. Just retry.
1911 */
1912 shmdt(evt_); evt_ = NULL;
1913 regCnt_--;
1914 return (Register_());
1915 }
1916 else
1917 {
1918 /*
1919 * Any other error causes the registration to fail
1920 */
1921 sprintf(ErrMsg(), "semaphore error - %s", strerror(errno));
1922 return (-1);
1923 }
1924 }
1925
1926 /*
1927 * Verify semaphore - we cannot lock here due to a race condition
1928 * with the creators semctl. This should not matter because only
1929 * a non-faulty matching read indicates that the semaphore is
1930 * both created and initialized.
1931 */
1932 if (semEvt_ != evt_->semId)
1933 {
1934 if (regCnt_ > 0)
1935 {
1936 /*
1937 * Retry
1938 */
1939 semEvt_ = -1;
1940 shmdt(evt_); evt_ = NULL;
1941 regCnt_--;
1942 return (Register_());
1943 }
1944 else
1945 {
1946 strcpy(ErrMsg(), "semaphore mismatch");
1947 return (-1);
1948 }
1949 }
1950
1951 /*
1952 * Align event counter
1953 */
1954 evtCnt_ = evt_->event;
1955 }
1956
1957 /*
1958 * Register to events
1959 */
1960 if (semEvt_ >= 0)
1961 {
1962 struct sembuf semOp[4];
1963
1964 memset(&semOp[0], 0, sizeof(struct sembuf));
1965 memset(&semOp[1], 0, sizeof(struct sembuf));
1966 memset(&semOp[2], 0, sizeof(struct sembuf));
1967 memset(&semOp[3], 0, sizeof(struct sembuf));
1968
1969 semOp[1].sem_num = 0;
1970 semOp[1].sem_op = 1;
1971 semOp[1].sem_flg = SEM_UNDO;
1972
1973 semOp[2].sem_num = 1;
1974 semOp[2].sem_op = 1;
1975 semOp[2].sem_flg = SEM_UNDO;
1976
1977 semOp[3].sem_num = 0;
1978 semOp[3].sem_op = -1;
1979 semOp[3].sem_flg = SEM_UNDO;
1980 if (semop(semEvt_, semOp, 4) != 0)
1981 {
1982 int e = errno;
1983
1984 shmdt(evt_); evt_ = NULL;
1985 if (e == EIDRM) semEvt_ = -1; // semaphore has been removed
1986
1987 if (newShm)
1988 {
1989 /*
1990 * Delete the semaphore again
1991 */
1992 SemEvtRemove_();
1993
1994 /*
1995 * Delete the shared memory
1996 */
1997 shmctl(evtId_, IPC_RMID, 0); evtId_ = -1;
1998 }
1999
2000 if (regCnt_ > 0)
2001 {
2002 /*
2003 * Retry
2004 */
2005 regCnt_--;
2006 return (Register_());
2007 }
2008 else
2009 {
2010 sprintf(ErrMsg(), "semaphore error - %s", strerror(e));
2011 return (-1);
2012 }
2013 }
2014 }
2015
2016 return (0);
2017 }
2018
2020 void DetachEvt_() {
2021 if (!extPipe_ && (semEvt_ >= 0) && (clientId_ >= 0))
2022 {
2023 ngcbSEM_CTL ctl;
2024 ctl.val = 2;
2025 semctl(semEvt_, 2 + clientId_, SETVAL, ctl);
2026 }
2027 }
2028
2030 int EvtLock_() {
2031 if (semEvt_ >= 0)
2032 {
2033 if (semop(semEvt_, semOpEvtLock_, 2) != 0)
2034 {
2035 if (errno == EINTR)
2036 {
2037 while (semop(semEvt_, semOpEvtLock_, 2) != 0)
2038 {
2039 if (errno != EINTR)
2040 {
2041 if (errno == EIDRM) semEvt_ = -1; // removed (and unlocked)
2042 sprintf(ErrMsg(), "cannot lock semaphore - %s", strerror(errno));
2043 return (-1);
2044 }
2045 }
2046 }
2047 else
2048 {
2049 if (errno == EIDRM) semEvt_ = -1; // removed (and unlocked)
2050 sprintf(ErrMsg(), "cannot lock semaphore - %s", strerror(errno));
2051 return (-1);
2052 }
2053 }
2054 }
2055 else
2056 {
2057 strcpy(ErrMsg(), "invalid semaphore");
2058 return (-1);
2059 }
2060
2061 return (0);
2062 }
2063
2065 int EvtTryLock_() {
2066 if (semEvt_ >= 0)
2067 {
2068 if (semop(semEvt_, semOpEvtTryLock_, 2) == 0)
2069 {
2070 return (0);
2071 }
2072 else
2073 {
2074 if (errno == EIDRM) semEvt_ = -1; // removed (and unlocked)
2075 sprintf(ErrMsg(), "cannot lock semaphore - %s", strerror(errno));
2076 return (-1);
2077 }
2078 }
2079
2080 strcpy(ErrMsg(), "invalid semaphore");
2081 return (-1);
2082 }
2083
2085 void EvtUnlock_() {
2086 if (semEvt_ >= 0)
2087 {
2088 if (semop(semEvt_, semOpEvtUnlock_, 1) != 0)
2089 {
2090 if (errno == EIDRM) semEvt_ = -1; // removed (and unlocked)
2091 }
2092 }
2093 }
2094
2096 void SemEvtRemove_() {
2097 if (semEvt_ >= 0)
2098 {
2099 ngcbSEM_CTL ctl;
2100 ctl.val = 0;
2101 semctl(semEvt_, 0, IPC_RMID, ctl); semEvt_ = -1;
2102 }
2103 }
2104
2106 int RdLock_(int semId, int noWait = 0) {
2107 if (semId >= 0)
2108 {
2109 if (noWait) semOpRdLock_[0].sem_flg = IPC_NOWAIT;
2110 else semOpRdLock_[0].sem_flg = 0;
2111
2112 if (semop(semId, semOpRdLock_, 2) != 0)
2113 {
2114 if (errno == EINTR)
2115 {
2116 while (semop(semId, semOpRdLock_, 2) != 0)
2117 {
2118 if (errno != EINTR) return (-1);
2119 }
2120 }
2121 else
2122 {
2123 return (-1);
2124 }
2125 }
2126 }
2127
2128 return (0);
2129 }
2130
2132 int RdUnlock_(int semId) {
2133 if (semId >= 0)
2134 {
2135 if (semop(semId, semOpRdUnlock_, 1) != 0)
2136 {
2137 return (-1);
2138 }
2139 }
2140 return (0);
2141 }
2142
2144 int Wait_(int noWait) {
2145 struct sembuf semOp[3];
2146 int shmId = ngcbSHM_ID_ERR;
2147 int valid = 0;
2148
2149 if (semEvt_ < 0)
2150 {
2151 strcpy(ErrMsg(), "event interface not initialized");
2152 return (ngcbSHM_ID_ERR);
2153 }
2154 else if (clientId_ < 0)
2155 {
2156 strcpy(ErrMsg(), "not attached to events");
2157 return (ngcbSHM_ID_ERR);
2158 }
2159
2160 while (!valid)
2161 {
2162 /*
2163 * Wait for event and then lock event structure
2164 */
2165 memset(&semOp[0], 0, sizeof(struct sembuf));
2166 memset(&semOp[1], 0, sizeof(struct sembuf));
2167 memset(&semOp[2], 0, sizeof(struct sembuf));
2168 semOp[0].sem_num = static_cast<unsigned short>(2 + clientId_);
2169 semOp[0].sem_op = -1;
2170 semOp[2].sem_num = 0;
2171 semOp[2].sem_op = 1;
2172 semOp[2].sem_flg = SEM_UNDO;
2173 if (noWait)
2174 {
2175 semOp[0].sem_flg = IPC_NOWAIT;
2176 semOp[1].sem_flg = IPC_NOWAIT;
2177 }
2178 if (semop(semEvt_, semOp, 3) != 0)
2179 {
2180 if (errno == EAGAIN)
2181 {
2182 if (noWait) return (ngcbSHM_ID_RETRY);
2183 else continue;
2184 }
2185 else if (errno == EINTR)
2186 {
2187 continue;
2188 }
2189 else if (errno == EIDRM)
2190 {
2191 /*
2192 * Event semaphore has been removed
2193 */
2194 semEvt_ = -1;
2195 }
2196
2197 sprintf(ErrMsg(), "semaphore error - %s", strerror(errno));
2198 return (ngcbSHM_ID_ERR);
2199 }
2200
2201 /*
2202 * Check event
2203 */
2204 if (((evt_->event == evtCnt_) || (evt_->event == 0)) || cancelled_)
2205 {
2206 /*
2207 * We still store the event counter
2208 */
2209 evtCnt_ = evt_->event;
2210 if (noWait || cancelled_) shmId = ngcbSHM_ID_RETRY;
2211 cancelled_ = false;
2212 EvtUnlock_();
2213
2214 if (shmId == ngcbSHM_ID_RETRY)
2215 {
2216 /*
2217 * No valid event
2218 */
2219 return (shmId);
2220 }
2221 else
2222 {
2223 /*
2224 * Wait for next event
2225 */
2226 continue;
2227 }
2228 }
2229 else
2230 {
2231 /*
2232 * Store event. We do not unlock here - the public
2233 * function will handle this
2234 */
2235 shmId = evt_->shmId;
2236 evtCnt_ = evt_->event;
2237 if (shmId < 0) shmId = ngcbSHM_ID_IVLD;
2238 valid = 1;
2239 }
2240 }
2241
2242 return (shmId);
2243 }
2244
2246 int OpenPipe_() {
2247 int fd[2];
2248 int ack;
2249 int fdFlags;
2250 char arg[32];
2251
2252 /*
2253 * Socket pair for internal thread event-communication
2254 */
2255 if (socketpair(AF_UNIX, SOCK_STREAM, 0, fd) == -1)
2256 {
2257 sprintf(ErrMsg(), "error opening event pipe - %s", strerror(errno));
2258 return(-1);
2259 }
2260
2261 /*
2262 * We fork here
2263 */
2264 sprintf(arg, "%d", fd[1]);
2265 pid_ = fork();
2266 if (pid_ < 0)
2267 {
2268 /*
2269 * Cannot fork
2270 */
2271 sprintf(ErrMsg(), "error opening event pipe - %s", strerror(errno));
2272 close(fd[0]); close(fd[1]);
2273 return (-1);
2274 }
2275 else if (pid_ == 0)
2276 {
2277 /*
2278 * Close the end we don't need
2279 */
2280 close(fd[0]);
2281
2282 /*
2283 * Execute event dispatcher
2284 */
2285 execlp("ngcbShmEvt", "ngcbShmEvt", name, arg, (char *)NULL);
2286
2287 /*
2288 * If we are still here then the execution has failed
2289 */
2290 close(fd[1]);
2291 _exit(127);
2292 } // child process
2293
2294 /*
2295 * Close the end we don't need
2296 */
2297 close(fd[1]);
2298
2299 /*
2300 * Set close-on-exec flag
2301 */
2302 fdFlags = fcntl(fd[0], F_GETFD);
2303 if (fdFlags >= 0) fcntl(fd[0], F_SETFD, fdFlags | FD_CLOEXEC);
2304
2305 /*
2306 * Read acknowledge from the other side
2307 */
2308 if (read(fd[0], &ack, sizeof(int)) != sizeof(int))
2309 {
2310 strcpy(ErrMsg(), "unable to launch event dispatcher");
2311 close(fd[0]); waitpid(pid_, NULL, 0); pid_ = 0;
2312 return (-1);
2313 }
2314
2315 if (!ack)
2316 {
2317 int size;
2318
2319 /*
2320 * Read back error message
2321 */
2322 if (read(fd[0], &size, sizeof(int)) != sizeof(int))
2323 {
2324 strcpy(ErrMsg(), "event dispatcher connection broken");
2325 close(fd[0]); waitpid(pid_, NULL, 0); pid_ = 0;
2326 return (-1);
2327 }
2328 if (size > 0)
2329 {
2330 size++; // including the terminating NULL-character
2331 if (read(fd[0], ErrMsg(), size) != size)
2332 {
2333 strcpy(ErrMsg(), "event dispatcher connection broken");
2334 }
2335 }
2336 else
2337 {
2338 strcpy(ErrMsg(), "unable to launch event dispatcher");
2339 }
2340 close(fd[0]); waitpid(pid_, NULL, 0); pid_ = 0;
2341 return (-1);
2342 }
2343
2344 return (fd[0]);
2345 }
2346
2348 void ClosePipe_() {
2349 if (fd_ >= 0)
2350 {
2351 pipeOpen_ = 0; close(fd_); fd_ = -1;
2352 }
2353 if (pid_ > 0)
2354 {
2355 kill(pid_, SIGTERM); waitpid(pid_, NULL, 0); pid_ = 0;
2356 }
2357#ifdef NGCB_SHM_THR
2358 if (thr_)
2359 {
2360 Release();
2361 pthread_join(thrId_, NULL);
2362 thr_ = 0;
2363 }
2364#endif
2365 if (pipe_ >= 0)
2366 {
2367 close(pipe_); pipe_ = -1;
2368 }
2369 }
2370
2372 int WaitPipe_(int noWait=0) {
2373 struct pollfd fds;
2374 int res;
2375
2376 if ((fd_ < 0) || extPipeErr_)
2377 {
2378 strcpy(ErrMsg(), "event dispatcher not running");
2379 return (-1);
2380 }
2381
2382 fds.fd = fd_;
2383 fds.events = POLLIN;
2384 fds.revents = 0;
2385 if (noWait)
2386 {
2387 res = poll(&fds, 1, 0);
2388 }
2389 else
2390 {
2391 res = poll(&fds, 1, -1);
2392 }
2393
2394 if (res < 0)
2395 {
2396 /*
2397 * An error occured - file descriptore no longer valid
2398 */
2399 extPipeErr_ = 1;
2400 return (-1);
2401 }
2402 else if (res == 0)
2403 {
2404 /*
2405 * File descriptor valid but not ready
2406 */
2407 return (0);
2408 }
2409 else
2410 {
2411 return (1);
2412 }
2413 }
2414
2416 int StartThread_() {
2417#ifdef NGCB_SHM_THR
2418 int fd[2];
2419 int fdFlags;
2420 pthread_attr_t attr;
2421
2422 /*
2423 * Socket pair for internal thread event-communication
2424 */
2425 if (socketpair(AF_UNIX, SOCK_STREAM, 0, fd) == -1)
2426 {
2427 sprintf(ErrMsg(), "error opening event pipe - %s", strerror(errno));
2428 return(-1);
2429 }
2430
2431 /*
2432 * Set close-on-exec flag
2433 */
2434 fdFlags = fcntl(fd[0], F_GETFD);
2435 if (fdFlags >= 0) fcntl(fd[0], F_SETFD, fdFlags | FD_CLOEXEC);
2436 fdFlags = fcntl(fd[1], F_GETFD);
2437 if (fdFlags >= 0) fcntl(fd[1], F_SETFD, fdFlags | FD_CLOEXEC);
2438
2439 /*
2440 * Set the posix thread attributes
2441 */
2442 pthread_attr_init(&attr);
2443 pthread_attr_setinheritsched(&attr, PTHREAD_INHERIT_SCHED);
2444 pthread_attr_setscope(&attr, PTHREAD_SCOPE_SYSTEM);
2445 pthread_attr_setdetachstate(&attr, PTHREAD_CREATE_JOINABLE);
2446 pthread_attr_setstacksize(&attr, 1024 * 1024);
2447
2448 /*
2449 * Create thread
2450 */
2451 pipe_ = fd[1]; pipeOpen_ = 1;
2452 if (pthread_create(&thrId_, &attr, ngcbShmThread, this) != 0)
2453 {
2454 sprintf(ErrMsg(), "thread creation failed - %s", strerror(errno));
2455 close(fd[0]); close(fd[1]); pipe_ = -1; pipeOpen_ = -1;
2456 return(-1);
2457 }
2458
2459 /*
2460 * We have a valid file descriptor to read events
2461 */
2462 fd_ = fd[0];
2463
2464 return (0);
2465#else
2466 strcpy(ErrMsg(), "event dispatcher thread not installed");
2467 return (-1);
2468#endif
2469 }
2470};
2471
2472#endif
Definition ngcbSHM.hpp:167
ngcbSHMTHR()
Constructor.
Definition ngcbSHM.hpp:170
virtual void Thread()
Thread function.
Definition ngcbSHM.hpp:176
virtual ~ngcbSHMTHR()
Destructor.
Definition ngcbSHM.hpp:173
void * TryAttachBuffer(int shmId, void *hdr)
Try to read-lock and then attach buffer to shared memory.
Definition ngcbSHM.hpp:696
int TryWait()
Try to wait for next event.
Definition ngcbSHM.hpp:896
int NumClient()
Return number of registered clients.
Definition ngcbSHM.hpp:381
int Initialized()
Check, if initialization was OK.
Definition ngcbSHM.hpp:369
int SendAck()
Send event acknowledge to internal pipe.
Definition ngcbSHM.hpp:1297
void ShmNoUndo()
Disable reader-lock undo-function (use with care!)
Definition ngcbSHM.hpp:596
int WrUnlock(int semId)
Writer unlock.
Definition ngcbSHM.hpp:500
void RdUnlock(int semId)
Release read-lock.
Definition ngcbSHM.hpp:539
char * ErrMsg()
Error message.
Definition ngcbSHM.hpp:372
int ShmUnlock(int shmId)
Unlock shared memory (reader)
Definition ngcbSHM.hpp:602
unsigned int NextEvt()
Return next event.
Definition ngcbSHM.hpp:722
void RdWait(int semId)
Wait until there are no more readers.
Definition ngcbSHM.hpp:513
int IntPipe()
Return internal pipe indicator.
Definition ngcbSHM.hpp:378
unsigned int EvtCnt()
Return event counter.
Definition ngcbSHM.hpp:719
int TryGetHdr(void *hdr)
Try to wait for event, try to acquire read-lock, then fill the header.
Definition ngcbSHM.hpp:1145
int ShmLock(int shmId, int noWait=0)
Lock shared memory (reader)
Definition ngcbSHM.hpp:558
ngcb_shmhdr_t * sharedHeader
Shared event header.
Definition ngcbSHM.hpp:208
ngcbSHM()
Constructor.
Definition ngcbSHM.hpp:211
virtual ~ngcbSHM()
Destructor.
Definition ngcbSHM.hpp:287
int Send(int shmId)
Send event.
Definition ngcbSHM.hpp:1148
int Wait(int noWait=0)
Wait for next event.
Definition ngcbSHM.hpp:813
int WrLock(int semId)
Writer lock.
Definition ngcbSHM.hpp:460
int Fd()
Return event file descriptor.
Definition ngcbSHM.hpp:375
void Remove()
Remove shared memory event interface.
Definition ngcbSHM.hpp:401
int RdCnt(int semId)
Return number of readers.
Definition ngcbSHM.hpp:533
ngcbSHM(const char *s)
Constructor with name.
Definition ngcbSHM.hpp:220
void * GetBuffer(void *hdr, int noWait=0)
Wait for next event and return a read-locked buffer.
Definition ngcbSHM.hpp:899
int RecvEvt()
Receive event from internal pipe.
Definition ngcbSHM.hpp:1283
void RdUnlock(ngcb_shmhdr_t &hdr)
Release read-lock.
Definition ngcbSHM.hpp:555
int WrInit(int semId)
Initialize writer lock.
Definition ngcbSHM.hpp:434
void * TryGetBuffer(void *hdr)
Try to wait for next event then try to read-lock the buffer.
Definition ngcbSHM.hpp:1010
int GetHdr(void *hdr, int noWait=0)
Wait for next event, acquire read-lock and fill the header.
Definition ngcbSHM.hpp:1013
int Release()
Release event.
Definition ngcbSHM.hpp:1210
void Thread()
Thread function overloading the thread base class.
Definition ngcbSHM.hpp:1311
void Invalidate()
Invalidate event.
Definition ngcbSHM.hpp:1273
char name[64]
Interface name.
Definition ngcbSHM.hpp:205
void DetachBuffer()
Detach buffer from shared memory and release read-lock.
Definition ngcbSHM.hpp:701
int AttachEvt()
Attach to event handler.
Definition ngcbSHM.hpp:729
int WrTryLock(int semId)
Writer try-lock.
Definition ngcbSHM.hpp:485
void * AttachBuffer(int shmId, void *hdr, int noWait=0)
Attach buffer to shared memory and read-lock.
Definition ngcbSHM.hpp:635
Shared memory built-in database.
#define ngcbSHM_ID_IVLD
Definition ngcbSHM.hpp:44
#define ngcbSHM_MAX_CLIENTS
Definition ngcbSHM.hpp:37
#define ngcbSHM_ID_RETRY
Definition ngcbSHM.hpp:43
#define ngcbSHM_ID_ERR
Definition ngcbSHM.hpp:42
Definition ngcbSHM.hpp:54
unsigned int event
Event.
Definition ngcbSHM.hpp:65
ngcb_shmevt_t()
Constructor.
Definition ngcbSHM.hpp:68
int shmId
Shared memory id.
Definition ngcbSHM.hpp:62
int semId
Semaphore id.
Definition ngcbSHM.hpp:56
int reserved1
Reserved for 64-bit alignment.
Definition ngcbSHM.hpp:59
Definition ngcbSHM.hpp:94
int status
Status.
Definition ngcbSHM.hpp:141
int fcnt0
Initial frame counter.
Definition ngcbSHM.hpp:129
int bitPix
BITPIX.
Definition ngcbSHM.hpp:102
int fullFrame
Full-frame.
Definition ngcbSHM.hpp:120
ngcb_shmio_t io
Event i/o.
Definition ngcbSHM.hpp:96
int sy
Start-y coordinate.
Definition ngcbSHM.hpp:108
int ndit
Number of integrations in pixel data.
Definition ngcbSHM.hpp:135
int ny
Y-dimension.
Definition ngcbSHM.hpp:114
int endian
Endianess (1=little-endian, 0=big-endian)
Definition ngcbSHM.hpp:117
int bzero
BZERO.
Definition ngcbSHM.hpp:123
ngcb_shmhdr_t()
Constructor.
Definition ngcbSHM.hpp:147
int sx
Start-x coordinate.
Definition ngcbSHM.hpp:105
int flags
Flags (BIT0 = recording, ...)
Definition ngcbSHM.hpp:138
int expCnt
Exposure counter.
Definition ngcbSHM.hpp:132
char bytes[308]
General purpose bytes with n * 128 bit data buffer alignment.
Definition ngcbSHM.hpp:144
int nx
X-dimension.
Definition ngcbSHM.hpp:111
int fcnt
Frame counter.
Definition ngcbSHM.hpp:126
int frame
Frame-id.
Definition ngcbSHM.hpp:99
Definition ngcbSHM.hpp:74
ngcb_shmio_t()
Constructor.
Definition ngcbSHM.hpp:88
int reserved
Reserved for 64-bit alignment.
Definition ngcbSHM.hpp:85
unsigned int event
Event id.
Definition ngcbSHM.hpp:82
int retry
Retry.
Definition ngcbSHM.hpp:79
int semId
Semaphore id for buffer r/w-locks.
Definition ngcbSHM.hpp:76
Definition ngcbSHMDB.hpp:46
int val
Definition ngcbSHMDB.hpp:47
unsigned short * array
Definition ngcbSHMDB.hpp:49