ATLAS Offline Software
Loading...
Searching...
No Matches
JiveXMLServer.cxx
Go to the documentation of this file.
1/*
2 Copyright (C) 2002-2026 CERN for the benefit of the ATLAS collaboration
3*/
4
6
7//tdaq includes
8#include <ers/ers.h>
9
10//JiveXML includes
15
16#include <signal.h>
17#include <ranges>
18
19//Define warning and error
20#define ERS_WARNING( message ) \
21{ \
22 ERS_REPORT_IMPL( ers::warning, ers::Message, message, ); \
23}
24
25#define ERS_ERROR( message ) \
26{ \
27 ERS_REPORT_IMPL( ers::error, ers::Message, message, ); \
28}
29
30namespace JiveXML {
31
36 m_portNumber(port){
37
38 //Make sure ServerThread does not start unexpectedly
39 m_runServerThread = false ;
40
41 //And then start it
42 StartServingThread().ignore();
43
44 //Also register the signal handlers
45 signal( SIGINT , JiveXMLServer::signalHandler );
46 signal( SIGTERM, JiveXMLServer::signalHandler );
47 }
48
53
54 //Just stop the serving thread
55 StopServingThread().ignore();
56 }
57
65
66 ERS_DEBUG(MSG::VERBOSE,"StartServingThread()");
67
68 //The arguments passed on to the server - create new object on the heap that
69 //is persistent through the lifetime of the thread
71
72 //set runServer flag to true, so the thread will start
73 m_runServerThread = true ;
74
75 //create thread itself
76 if ( ( pthread_create(&m_ServerThreadHandle, NULL , &JiveXML::ONCRPCServerThread, (void *) args )) != 0){
77
78 //Create thread failed
79 ERS_WARNING("Thread creation failed");
80 return StatusCode::FAILURE;
81 }
82
83 return StatusCode::SUCCESS;
84
85 }
86
94
95 ERS_DEBUG(MSG::VERBOSE,"StopServingThread()");
96
101
102 //Create a client - this will already cause an update of the file
103 //descriptors on the sockets
104 CLIENT* client = clnt_create("localhost", ONCRPCSERVERPROG,ONCRPCSERVERVERS, "tcp");
105 if (!client){
106 ERS_ERROR("Unable to create shutdown client for local server");
107 return StatusCode::FAILURE;
108 }
109
110 // Next unset the runServerThread flag, which will cause the loop to stop
111 // This needs to happen after client creation, otherwise the next call won't
112 // be serverd anymore
113 m_runServerThread = false ;
114
115 // Now issue the call with a timeout
116 struct timeval timeout; timeout.tv_sec = 1; timeout.tv_usec = 0;
117// xdr_void is defined inconsistently in xdr.h and gets a warning from gcc8.
118#if __GNUC__ >= 8
119# pragma GCC diagnostic push
120# pragma GCC diagnostic ignored "-Wcast-function-type"
121#endif
122#if defined(__clang__) && __clang_major__ >= 19
123# pragma clang diagnostic push
124# pragma clang diagnostic ignored "-Wcast-function-type-mismatch"
125#endif
126 clnt_call(client, NULLPROC, (xdrproc_t)xdr_void, NULL, (xdrproc_t)xdr_void, NULL, timeout);
127#if defined(__clang__) && __clang_major__ >= 19
128# pragma clang diagnostic pop
129#endif
130#if __GNUC__ >= 8
131# pragma GCC diagnostic pop
132#endif
133
134 // A pointer to the return value of the thread
135 void* ret = NULL;
136 // wait till the server thread has finished
137 ERS_INFO("Waiting for server thread to terminate ...");
138 pthread_join(m_ServerThreadHandle, &ret);
139 ERS_INFO(" ... finished server thread");
140
141 //check if there was a return value
142 if (ret){
143 //Get the return value
144 unsigned long NRequests = *(unsigned long*)ret;
145 ERS_DEBUG(MSG::DEBUG,"Server thread stopped after handling " << NRequests << " requests");
146 } else
147 ERS_WARNING("Server thread stopped unexpectedly");
148
149 return StatusCode::SUCCESS;
150 }
151
156 //Store signal and notify thread
157 m_receivedSignal.set_value(signal);
158 }
159
166 auto signal = m_receivedSignal.get_future();
167 //just wait for a signal
168 signal.wait();
169 //Tell why the lock was released
170 ERS_INFO("Reached post-condition after received signal " << signal.get() );
171 }
172
178 //call the signal handler, so we will also reach post condition
179 signalHandler(-1);
180 }
181
186 StatusCode JiveXMLServer::UpdateEventForStream( const EventStreamID& evtStreamID, const std::string & event) {
187
188 ERS_DEBUG(MSG::VERBOSE,"UpdateEventForStream");
189
190 //Check that the event stream id is valid
191 if (!evtStreamID.isValid()){
192 ERS_ERROR("Invalid event stream identifier - cannot add event");
193 return StatusCode::FAILURE;
194 }
195
196 //Make sure we don't have already exceeded the maximum number of streams
197 if (m_eventStreamMap.size() > NSTREAMMAX ){
198 ERS_ERROR("Reached max. allowed number of streams " << NSTREAMMAX << " - cannot add event");
199 return StatusCode::FAILURE;
200 }
201
202 //Make sure the event is not larger than the allowed maximal size
203 if (event.length() > NBYTESMAX ){
204 ERS_ERROR("Event is larger than allowed max. of " << NBYTESMAX << " bytes - cannot add event");
205 return StatusCode::FAILURE;
206 }
207
208 //Make sure we are the only one accessing the data right now, by trying to
209 //obtain a lock. If the lock cannot be obtained after a certain time, an
210 //error is reported
211
212 //Try to obtain the lock within 30 seconds
213 using namespace std::chrono_literals;
214 std::unique_lock lock(m_accessLock, 30s);
215
216 if ( !lock ){
217 ERS_ERROR("Unable to obtain access lock to update event");
218 return StatusCode::FAILURE;
219 }
220
221 //Using std::map::operator[] and std::map::insert() will create a new event
222 //if it did not exist, otherwise just replace the existing entry (making a
223 //copy of the std::string) but would not update the key which holds new
224 //event/run number. Therefore delete existing entry first.
225
226 m_eventStreamMap.erase(evtStreamID);
227 m_eventStreamMap.insert(EventStreamPair(evtStreamID,event));
228
229 ERS_DEBUG(MSG::DEBUG, "Updated stream " << evtStreamID.StreamName()
230 << " with event Nr. " << evtStreamID.EventNumber()
231 << " from run Nr. " << evtStreamID.RunNumber());
232
233 return StatusCode::SUCCESS;
234 }
235
236
245 return 3;
246 }
247
251 std::vector<std::string> JiveXMLServer::GetStreamNames() const {
252
253 //Obtain an exclusive access lock
254 std::scoped_lock lock(m_accessLock);
255
256 return m_eventStreamMap
257 | std::views::keys
258 | std::views::transform(&EventStreamID::StreamName)
259 | std::ranges::to<std::vector>();
260 }
261
265 const EventStreamID JiveXMLServer::GetEventStreamID( const std::string& StreamName) const {
266
267 //Obtain an exclusive access lock
268 std::scoped_lock lock(m_accessLock);
269
270 // Search the entry in the map
271 if (auto MapItr = m_eventStreamMap.find(StreamName); MapItr != m_eventStreamMap.end()) {
272 return MapItr->first;
273 }
274
275 return EventStreamID{""};
276 }
277
281 const std::string JiveXMLServer::GetEvent( const EventStreamID& evtStreamID ) const {
282
283 //Obtain an exclusive access lock
284 std::scoped_lock lock(m_accessLock);
285
286 // Search the entry in the map
287 if (auto MapItr = m_eventStreamMap.find(evtStreamID); MapItr != m_eventStreamMap.end()) {
288 return MapItr->second;
289 }
290
291 return {};
292 }
293
297 void JiveXMLServer::Message( const MSG::Level level, const std::string& msg) const {
298 //Deliver message to the proper stream
299 if (level <= MSG::DEBUG) ERS_REPORT_IMPL( ers::debug, ers::Message, msg, level);
300 if (level == MSG::INFO) ERS_REPORT_IMPL( ers::info, ers::Message, msg, );
301 if (level == MSG::WARNING) ERS_REPORT_IMPL( ers::warning, ers::Message, msg, );
302 if (level == MSG::ERROR) ERS_REPORT_IMPL( ers::error, ers::Message, msg, );
303 if (level >= MSG::FATAL) ERS_REPORT_IMPL( ers::fatal, ers::Message, msg, );
304 }
305
309 MSG::Level JiveXMLServer::LogLevel() const {
310 //set to fixed value for now
311 return MSG::DEBUG;
312 }
313}
virtual void lock()=0
Interface to allow an object to lock itself when made const in SG.
#define ERS_WARNING(message)
#define ERS_ERROR(message)
#define ONCRPCSERVERVERS
#define ONCRPCSERVERPROG
For the client-server communication, each event is uniquely identified by the run number,...
Definition EventStream.h:19
unsigned int RunNumber() const
Definition EventStream.h:42
unsigned long EventNumber() const
Definition EventStream.h:41
const std::string & StreamName() const
Definition EventStream.h:43
virtual StatusCode UpdateEventForStream(const EventStreamID &evtStreamID, const std::string &event) override
Put this event as new current event for stream given by name.
virtual ~JiveXMLServer()
Destructor.
virtual std::vector< std::string > GetStreamNames() const override
get the names of all the streams
virtual void Message(const MSG::Level level, const std::string &msg) const override
This function is exposed to allow using ERS messaging service from other threads.
virtual MSG::Level LogLevel() const override
Get the logging level.
virtual void ServerThreadStopped() override
Callback whenever the server thread is stopped.
static void signalHandler(int signum)
When the signal handler is called, switch the lock to the post condition.
StatusCode StartServingThread()
Start the serving thread.
void Wait()
Wait for the server finish.
EventStreamMap m_eventStreamMap
StatusCode StopServingThread()
Stop the serving thread.
JiveXMLServer(int port=0)
Constructor.
virtual int GetState() const override
get the Status of the application
virtual const std::string GetEvent(const EventStreamID &evtStreamID) const override
get the current event for a particular stream
virtual const EventStreamID GetEventStreamID(const std::string &streamName) const override
get the current EventStreamID for a particular stream
This header is shared inbetween the C-style server thread and the C++ Athena ServerSvc.
const unsigned int NBYTESMAX
std::pair< const EventStreamID, const std::string > EventStreamPair
A map that stores events according to their EventStreamID Due to the way EventStreamID is build,...
Definition EventStream.h:86
const unsigned int NSTREAMMAX
struct ServerThreadArguments_t ServerThreadArguments
Arguments handed over fromt the main (Athena) thread to the server thread.
void * ONCRPCServerThread(void *args)
This is the actual server thread, which takes above arguments.
MsgStream & msg
Definition testRead.cxx:32