otsdaq-mu2e-stm  5.02.01
queue.cc
1 #include <errno.h>
2 #include <fcntl.h>
3 #include <stdarg.h>
4 #include <stdio.h>
5 #include <stdlib.h>
6 #include <string.h>
7 #include <unistd.h>
8 
9 #include <linux/memfd.h>
10 #include <sys/mman.h>
11 #include <sys/syscall.h>
12 #include <sys/types.h>
13 
14 #include "queue.hh"
15 
16 // Convenience wrappers for erroring out
17 static inline void queue_error(const char* fmt, ...)
18 {
19  va_list args;
20  va_start(args, fmt);
21  fprintf(stderr, "queue error: ");
22  vfprintf(stderr, fmt, args);
23  fprintf(stderr, "\n");
24  va_end(args);
25  abort();
26 }
27 static inline void queue_error_errno(const char* fmt, ...)
28 {
29  va_list args;
30  va_start(args, fmt);
31  fprintf(stderr, "queue error: ");
32  vfprintf(stderr, fmt, args);
33  fprintf(stderr, " (errno %d)\n", errno);
34  va_end(args);
35  abort();
36 }
37 
38 // Initialise the queue
39 void queue::init(size_t s)
40 {
41  // We're going to use a trick where we mmap two adjacent pages (in virtual memory) that point to the
42  // same physical memory. This lets us optimize memory access, by virtue of the fact that we don't need
43  // to even worry about wrapping our pointers around until we go through the entire buffer.
44 
45  // Check that the requested size is a multiple of a page. If it isn't, we're in trouble.
46  if(s % getpagesize() != 0)
47  {
48  queue_error("Requested size (%lu) is not a multiple of the page size (%d)",
49  s,
50  getpagesize());
51  }
52 
53  // Create an anonymous file backed by memory
54  if((q.fd = fileno(tmpfile())) == -1)
55  {
56  queue_error_errno("Could not obtain anonymous file");
57  }
58 
59  // Set buffer size
60  if(ftruncate(q.fd, s) != 0)
61  {
62  queue_error_errno("Could not set size of anonymous file");
63  }
64 
65  // Ask mmap for a good address
66  if((q.buffer = (uint8_t*)mmap(
67  NULL, 2 * s, PROT_NONE, MAP_PRIVATE | MAP_ANONYMOUS, -1, 0)) == MAP_FAILED)
68  {
69  queue_error_errno("Could not allocate virtual memory");
70  }
71 
72  // Mmap first region
73  if(mmap(q.buffer, s, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_FIXED, q.fd, 0) ==
74  MAP_FAILED)
75  {
76  queue_error_errno("Could not map buffer into virtual memory");
77  }
78 
79  // Mmap second region, with exact address
80  if(mmap(q.buffer + s, s, PROT_READ | PROT_WRITE, MAP_SHARED | MAP_FIXED, q.fd, 0) ==
81  MAP_FAILED)
82  {
83  queue_error_errno("Could not map buffer into virtual memory");
84  }
85 
86  // Initialize synchronization primitives
87  if(pthread_mutex_init(&q.lock, NULL) != 0)
88  {
89  queue_error_errno("Could not initialize mutex");
90  }
91  if(pthread_cond_init(&q.readable, NULL) != 0)
92  {
93  queue_error_errno("Could not initialize condition variable");
94  }
95  if(pthread_cond_init(&q.writeable, NULL) != 0)
96  {
97  queue_error_errno("Could not initialize condition variable");
98  }
99 
100  // Initialize remaining members
101  q.size = s;
102  q.head = 0;
103  q.tail = 0;
104  // q.head_seq = 0;
105  // q.tail_seq = 0;
106 }
107 
108 // Destroy / unmap the queue
109 void queue::destroy()
110 {
111  if(munmap(q.buffer + q.size, q.size) != 0)
112  {
113  queue_error_errno("Could not unmap buffer");
114  }
115 
116  if(munmap(q.buffer, q.size) != 0)
117  {
118  queue_error_errno("Could not unmap buffer");
119  }
120 
121  if(close(q.fd) != 0)
122  {
123  queue_error_errno("Could not close anonymous file");
124  }
125 
126  if(pthread_mutex_destroy(&q.lock) != 0)
127  {
128  queue_error_errno("Could not destroy mutex");
129  }
130 
131  if(pthread_cond_destroy(&q.readable) != 0)
132  {
133  queue_error_errno("Could not destroy condition variable");
134  }
135 
136  if(pthread_cond_destroy(&q.writeable) != 0)
137  {
138  queue_error_errno("Could not destroy condition variable");
139  }
140 }
141 
142 // Put data in the queue
143 void queue::put(int16_t* data, int16_t data_size)
144 {
145  // Lock the queue mutex
146  pthread_mutex_lock(&q.lock);
147 
148  // Wait for space to become available
149  while(q.size - (q.tail - q.head) < data_size + sizeof(data_size))
150  {
151  pthread_cond_wait(&q.writeable, &q.lock);
152  }
153 
154  // Construct header
155  // message_t m;
156  // m.len = size;
157  // m.seq = q.tail_seq++;
158 
159  // Write message
160  // memcpy(q.buffer + q.tail, &m, sizeof(message_t));
161  memcpy(q.buffer + q.tail, &data_size, sizeof(data_size));
162  // memcpy(&q.buffer[q.tail + sizeof(data_size)], data, data_size );
163  memcpy(q.buffer + q.tail + sizeof(data_size), data, data_size);
164 
165  // Increment write index in bytes
166  q.tail += data_size + sizeof(data_size);
167 
168  // Signal the queue is readable
169  pthread_cond_signal(&q.readable);
170 
171  // Unlock the mutex
172  pthread_mutex_unlock(&q.lock);
173 }
174 
175 // Get data from the queue
176 size_t queue::get(int16_t* data, int16_t max_size)
177 {
178  // Lock the queue mutex
179  pthread_mutex_lock(&q.lock);
180 
181  // Wait for a message that we can successfully consume to reach the front of the queue
182  // message_t m;
183  int16_t data_size;
184  // for(;;) {
185 
186  // Wait for a message to arrive
187  while((q.tail - q.head) == 0)
188  {
189  pthread_cond_wait(&q.readable, &q.lock);
190  }
191 
192  // Read message header
193  // memcpy(&m, &q.buffer[q.head], sizeof(message_t));
194  memcpy(&data_size, q.buffer + q.head, sizeof(data_size));
195 
196  // // Message too long, wait for someone else to consume it
197  // if(m.len > max){
198  // while(q.head_seq == m.seq) {
199  // pthread_cond_wait(&q.writeable, &q.lock);
200  // }
201  // continue;
202  // }
203 
204  // We successfully consumed the header of a suitable message, so proceed
205  // break;
206  // }
207 
208  // Read message body
209  // memcpy(data, &q.buffer[q.head + sizeof(message_t)], m.len);
210  memcpy(data, q.buffer + q.head + sizeof(data_size), data_size);
211 
212  // Consume the message by incrementing the read pointer
213  q.head += data_size + sizeof(data_size);
214  // q.head_seq++;
215 
216  // When read buffer moves into 2nd memory region, we can reset to the 1st region
217  if(q.head >= q.size)
218  {
219  q.head -= q.size;
220  q.tail -= q.size;
221  }
222 
223  // Signal the queue is writable
224  pthread_cond_signal(&q.writeable);
225 
226  // Unlock the mutex
227  pthread_mutex_unlock(&q.lock);
228 
229  return data_size;
230 }
Definition: data.hh:4