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