otsdaq-mu2e-stm  5.02.01
queue_safe.hh
1 #ifndef queue_h_
2 #define queue_h_
3 
4 #include <iostream>
5 #include <cstring>
6 #include <thread>
7 #include <atomic>
8 #include <cinttypes>
9 
10 // Data variables header
11 #include "STMDAQ-TestBeam/utils/dataVars.hh"
12 // UDP socket header
13 #include "STMDAQ-TestBeam/utils/UDPsocket.hh"
14 
15 class queue_buffer{
16 
17 public :
18 
19  // Constructor
20  queue_buffer();
21 
22  // // The header for each packet in the buffer
23  // struct packet_header{
24  // int16_t header[MAX_PACKET_LEN] = {};
25  // };
26 
27  // Try to push data to queue
28  int try_push(int chan, int16_t *data, int n);
29 
30  // Push data to queue
31  void push(int chan, int16_t *data, int n);
32 
33  // Try to pull data from queue
34  uint64_t try_pull(int chan, int16_t *&data);
35 
36  // Pull data from queue
37  uint64_t pull(bool *timeout, int chan, int16_t *&data);
38 
39  // Make int32_t from two int16_ts
40  int32_t make_int32_t(int16_t p0, int16_t p1){
41  return (p1 & 0xFFFF) << 16 | p0 & 0xFFFF;
42 
43  }
44 
45  // // Atomic occupancy counters
46  // std::atomic<uint32_t> occupancy_size[CHNUM];
47  // std::atomic<uint32_t> occupancy_data[CHNUM];
48 
49  // // Push tail pointers
50  // struct tail {
51  // // Size buffer pointer (int32_t)
52  // int16_t size_0; // Lower 16 bits
53  // int16_t size_1; // Upper 16 bits
54  // // Data buffer pointer (int32_t)
55  // int16_t data_0; // Lower 16 bits
56  // int16_t data_1; // Upper 16 bits
57  // };
58 
59  // // Pull head pointers
60  // struct head {
61  // // Size buffer pointer (int32_t)
62  // int16_t size_0; // Lower 16 bits
63  // int16_t size_1; // Upper 16 bits
64  // // Data buffer pointer (int32_t)
65  // int16_t data_0; // Lower 16 bits
66  // int16_t data_1; // Upper 16 bits
67  // };
68 
69  // Atomic push tail pointer
70  std::atomic<int64_t> tail[CHNUM];
71  // Atomic pull head pointer
72  std::atomic<int64_t> head[CHNUM];
73  // Atomic buffer occupancy pointer
74  std::atomic<int64_t> occupancy[CHNUM];
75 
76  // // Current tail pointers
77  // tail current_tail[CHNUM];
78  // int32_t current_tail_size[CHNUM];
79  // int32_t current_tail_data[CHNUM];
80  // // Current head pointers
81  // head current_head[CHNUM];
82  // int32_t current_head_size[CHNUM];
83  // int32_t current_head_data[CHNUM];
84 
85  // Current pointers
86  int64_t current_tail[CHNUM];
87  int64_t current_head[CHNUM];
88 
89  // Number written to the buffer
90  int64_t write_num[CHNUM];
91  int64_t read_num[CHNUM];
92  // int32_t write_num_size[CHNUM]; // in push
93  // int32_t write_num_data[CHNUM]; // in push
94  // int32_t read_num_size[CHNUM]; // in push
95  // int32_t read_num_data[CHNUM]; // in push
96 
97  // // Size buffer indeces
98  // static const uint SIZE_INDEX = 0; // Individual datagram size
99  // static const uint ACC_SIZE_INDEX = 1; // Accumulated datagram size
100 
101  // Maximum size of an int32_t
102  static const int32_t INT32_T_MAX = 2147483647;
103 
104  // The max number of datagrams in a buffer
105  // static const uint64_t RING_BUFFER_NUM = 65536;
106  static const int64_t RING_BUFFER_NUM = pow(2,17);
107 
108  // The buffer length
109  static const int64_t RING_BUFFER_LEN = RING_BUFFER_NUM*UDPsocket::MAX_UDP_LEN;
110 
111  // The buffer
112  // uint64_t* size_buffer[CHNUM];
113  // packet_header* header_buffer[CHNUM];
114  int16_t* data_buffer[CHNUM];
115 
116  // Increment function
117  int64_t increment(int n, int x){
118  return (n + x) % RING_BUFFER_LEN;
119  }
120 
121 private :
122 
123 
124 };
125 
126 #endif
127 
Definition: data.hh:4