ATLAS Offline Software
Loading...
Searching...
No Matches
AthenaRootSharedWriterSvc Class Reference

This class provides an example for writing event data objects to Pool. More...

#include <AthenaRootSharedWriterSvc.h>

Inheritance diagram for AthenaRootSharedWriterSvc:
Collaboration diagram for AthenaRootSharedWriterSvc:

Public Member Functions

 AthenaRootSharedWriterSvc (const std::string &name, ISvcLocator *pSvcLocator)
 Standard Service Constructor.
virtual ~AthenaRootSharedWriterSvc ()=default
 Destructor.
virtual StatusCode initialize () override
 Gaudi Service Interface method implementations:
virtual StatusCode stop () override
virtual StatusCode finalize () override
virtual StatusCode share (int numClients=0, bool motherClient=false) override

Private Attributes

ServiceHandle< AthenaPoolSharedIOCnvSvcm_cnvSvc {this,"AthenaPoolSharedIOCnvSvc","AthenaPoolSharedIOCnvSvc"}
TServerSocket * m_rootServerSocket
TMonitor * m_rootMonitor
THashTable m_rootMergers
std::unordered_map< TClass *, void * > m_dummyCache
int m_rootClientIndex
int m_rootClientCount
int m_numberOfStreams

Friends

class SvcFactory< AthenaRootSharedWriterSvc >

Detailed Description

This class provides an example for writing event data objects to Pool.

Definition at line 26 of file AthenaRootSharedWriterSvc.h.

Constructor & Destructor Documentation

◆ AthenaRootSharedWriterSvc()

AthenaRootSharedWriterSvc::AthenaRootSharedWriterSvc ( const std::string & name,
ISvcLocator * pSvcLocator )

◆ ~AthenaRootSharedWriterSvc()

virtual AthenaRootSharedWriterSvc::~AthenaRootSharedWriterSvc ( )
virtualdefault

Destructor.

Member Function Documentation

◆ finalize()

StatusCode AthenaRootSharedWriterSvc::finalize ( )
overridevirtual

Definition at line 317 of file AthenaRootSharedWriterSvc.cxx.

317 {
318 ATH_MSG_INFO("in finalize()");
319 delete m_rootMonitor; m_rootMonitor = nullptr;
320 delete m_rootServerSocket; m_rootServerSocket = nullptr;
321 for (auto& [cl, ptr] : m_dummyCache) {
322 if (cl && ptr) {
323 cl->Destructor(ptr, false);
324 }
325 }
326 m_dummyCache.clear();
327 return StatusCode::SUCCESS;
328}
#define ATH_MSG_INFO(x)
std::unordered_map< TClass *, void * > m_dummyCache
cl
print [x.__class__ for x in toList(dqregion.getSubRegions()) ]

◆ initialize()

StatusCode AthenaRootSharedWriterSvc::initialize ( )
overridevirtual

Gaudi Service Interface method implementations:

Definition at line 155 of file AthenaRootSharedWriterSvc.cxx.

155 {
156 ATH_MSG_INFO("in initialize()");
157
158 // Initialize IConversionSvc
159 ATH_CHECK(m_cnvSvc.retrieve());
160 IProperty* propertyServer = dynamic_cast<IProperty*>(m_cnvSvc.get());
161 if (propertyServer == nullptr) {
162 ATH_MSG_ERROR("Unable to cast conversion service to IProperty");
163 return StatusCode::FAILURE;
164 } else {
165 std::string propertyName = "ParallelCompression";
166 bool parallelCompression(false);
167 BooleanProperty parallelCompressionProp(propertyName, parallelCompression);
168 if (propertyServer->getProperty(&parallelCompressionProp).isFailure()) {
169 ATH_MSG_INFO("Conversion service does not have ParallelCompression property");
170 } else if (parallelCompressionProp.value()) {
171 int streamPort = 0;
172 propertyName = "StreamPortString";
173 std::string streamPortString("");
174 StringProperty streamPortStringProp(propertyName, streamPortString);
175 if (propertyServer->getProperty(&streamPortStringProp).isFailure()) {
176 ATH_MSG_INFO("Conversion service does not have StreamPortString property, using default: " << streamPort);
177 } else {
178 streamPort = atoi(streamPortStringProp.value().substr(streamPortStringProp.value().find(':') + 1).c_str());
179 }
180 m_rootServerSocket = new TServerSocket(streamPort, (streamPort == 0 ? false : true), 100, -1, ESocketBindOption::kInaddrLoopback);
181 if (m_rootServerSocket == nullptr || !m_rootServerSocket->IsValid()) {
182 ATH_MSG_FATAL("Could not create ROOT TServerSocket: " << streamPort);
183 return StatusCode::FAILURE;
184 }
185 streamPort = m_rootServerSocket->GetLocalPort();
186 const std::string newStreamPortString{streamPortStringProp.value().substr(0,streamPortStringProp.value().find(':')+1) + std::to_string(streamPort)};
187 if (propertyServer->setProperty(propertyName,newStreamPortString).isFailure()) {
188 ATH_MSG_FATAL("Could not set Conversion Service property " << propertyName << " from " << streamPortString << " to " << newStreamPortString);
189 return StatusCode::FAILURE;
190 }
191 m_rootMonitor = new TMonitor;
193 ATH_MSG_DEBUG("Successfully created ROOT TServerSocket and added it to TMonitor: ready to accept connections, " << streamPort);
194 }
195 }
196 // Count the number of output streams
197 const IAlgManager* algMgr = Gaudi::svcLocator()->as<IAlgManager>();
198 for (const auto& alg : algMgr->getAlgorithms()) {
199 if (alg->type() == "AthenaOutputStream") {
200 ATH_MSG_DEBUG("Counting " << alg->name() << " as an output stream algorithm");
202 }
203 }
204 if (m_numberOfStreams == 0) {
205 ATH_MSG_WARNING("No output stream algorithm found, setting the number of streams to 1");
207 } else {
208 ATH_MSG_INFO("Found a total of " << m_numberOfStreams << " output streams");
209 }
210
211 return StatusCode::SUCCESS;
212}
#define ATH_CHECK
Evaluate an expression and check for errors.
#define ATH_MSG_ERROR(x)
#define ATH_MSG_FATAL(x)
#define ATH_MSG_WARNING(x)
#define ATH_MSG_DEBUG(x)
ServiceHandle< AthenaPoolSharedIOCnvSvc > m_cnvSvc
int atoi(std::string_view str)
Helper functions to unpack numbers decoded in string into integers and doubles The strings are requir...

◆ share()

StatusCode AthenaRootSharedWriterSvc::share ( int numClients = 0,
bool motherClient = false )
overridevirtual

Definition at line 214 of file AthenaRootSharedWriterSvc.cxx.

214 {
215 ATH_MSG_DEBUG("Start commitOutput loop");
216 StatusCode sc = m_cnvSvc->commitOutput("", false);
217
218 // Allow ROOT clients to start up (by setting active clients)
219 // and wait to stop the ROOT server until all clients are done and metadata is written (commitOutput fail).
220 bool anyActiveClients = (m_rootServerSocket != nullptr);
221 while (sc.isSuccess() || sc.isRecoverable() || anyActiveClients) {
222 if (sc.isSuccess()) {
223 ATH_MSG_VERBOSE("Success in commitOutput loop");
224 } else if (m_rootMonitor != nullptr) {
225 TSocket* socket = m_rootMonitor->Select(1);
226 if (socket != nullptr && socket != (TSocket*)-1) {
227 ATH_MSG_DEBUG("ROOT Monitor got: " << socket);
228 if (socket->IsA() == TServerSocket::Class()) {
229 TSocket* client = (static_cast<TServerSocket*>(socket))->Accept();
230 client->Send(m_rootClientIndex, 0);
231 client->Send(1, 1);
234 if (m_rootClientCount < (numClients-1)*m_numberOfStreams + 1) {
235 m_rootMonitor->Add(client);
236 ATH_MSG_INFO("ROOT Monitor add client: " << m_rootClientIndex << ", " << client);
237 } else {
238 ATH_MSG_WARNING("ROOT Monitor do NOT add client: " << m_rootClientIndex << ", " << client);
239 client->Close("force");
241 }
242 } else {
243 TMessage* message = nullptr;
244 Int_t result = socket->Recv(message);
245 if (result < 0) {
246 ATH_MSG_ERROR("ROOT Monitor got an error while receiving the message from the socket: " << result);
247 return StatusCode::FAILURE;
248 }
249 if (message == nullptr) {
250 ATH_MSG_WARNING("ROOT Monitor got no message from socket: " << socket);
251 } else if (message->What() == kMESS_STRING) {
252 char str[64];
253 message->ReadString(str, 64);
254 ATH_MSG_INFO("ROOT Monitor client: " << socket << ", " << str);
255 m_rootMonitor->Remove(socket);
256 ATH_MSG_DEBUG("ROOT Monitor client: " << socket << ", " << socket->GetBytesRecv() << ", " << socket->GetBytesSent());
257 socket->Close();
259 if (m_rootMonitor->GetActive() == 0 || m_rootClientCount == 0) {
260 if (!motherClient) {
261 anyActiveClients = false;
262 ATH_MSG_INFO("ROOT Monitor: No more active clients...");
263 } else {
264 motherClient = false;
265 ATH_MSG_INFO("ROOT Monitor: Mother process is done...");
266 if (!m_cnvSvc->commitCatalog().isSuccess()) {
267 ATH_MSG_FATAL("Failed to commit file catalog.");
268 return StatusCode::FAILURE;
269 }
270 }
271 }
272 } else if (message->What() == kMESS_ANY) {
273 long long length;
274 TString filename;
275 int clientId;
276 message->ReadInt(clientId);
277 message->ReadTString(filename);
278 message->ReadLong64(length);
279 ATH_MSG_DEBUG("ROOT Monitor client: " << socket << ", " << clientId << ": " << filename << ", " << length);
280 std::unique_ptr<TMemFile> transient(new TMemFile(filename, message->Buffer() + message->Length(), length, "UPDATE"));
281 message->SetBufferOffset(message->Length() + length);
282 ParallelFileMerger* info = static_cast<ParallelFileMerger*>(m_rootMergers.FindObject(filename));
283 if (!info) {
284 info = new ParallelFileMerger(filename, &m_dummyCache, transient->GetCompressionSettings());
285 m_rootMergers.Add(info);
286 ATH_MSG_INFO("ROOT Monitor ParallelFileMerger: " << info << ", for: " << filename);
287 }
288 info->MergeTrees(transient.get());
289 }
290 delete message; message = nullptr;
291 }
292 }
293 } else if (m_rootMonitor == nullptr) {
294 usleep(100);
295 }
296 // Once commitOutput failed all legacy clients are finished (writing metadata), do not call again.
297 if (sc.isSuccess() || sc.isRecoverable()) {
298 sc = m_cnvSvc->commitOutput("", false);
299 if (sc.isFailure() && !sc.isRecoverable()) {
300 ATH_MSG_INFO("commitOutput failed, metadata done.");
301 if (anyActiveClients && m_rootClientCount == 0) {
302 ATH_MSG_INFO("ROOT Monitor: No clients, terminating the loop...");
303 anyActiveClients = false;
304 }
305 }
306 }
307 }
308 ATH_MSG_INFO("End commitOutput loop");
309 return StatusCode::SUCCESS;
310}
#define ATH_MSG_VERBOSE(x)
double length(const pvec &v)
static Double_t sc
::StatusCode StatusCode
StatusCode definition for legacy code.

◆ stop()

StatusCode AthenaRootSharedWriterSvc::stop ( )
overridevirtual

Definition at line 312 of file AthenaRootSharedWriterSvc.cxx.

312 {
313 m_rootMergers.Delete();
314 return StatusCode::SUCCESS;
315}

◆ SvcFactory< AthenaRootSharedWriterSvc >

friend class SvcFactory< AthenaRootSharedWriterSvc >
friend

Definition at line 1 of file AthenaRootSharedWriterSvc.h.

Member Data Documentation

◆ m_cnvSvc

ServiceHandle<AthenaPoolSharedIOCnvSvc> AthenaRootSharedWriterSvc::m_cnvSvc {this,"AthenaPoolSharedIOCnvSvc","AthenaPoolSharedIOCnvSvc"}
private

Definition at line 45 of file AthenaRootSharedWriterSvc.h.

45{this,"AthenaPoolSharedIOCnvSvc","AthenaPoolSharedIOCnvSvc"};

◆ m_dummyCache

std::unordered_map<TClass*, void*> AthenaRootSharedWriterSvc::m_dummyCache
private

Definition at line 50 of file AthenaRootSharedWriterSvc.h.

◆ m_numberOfStreams

int AthenaRootSharedWriterSvc::m_numberOfStreams
private

Definition at line 53 of file AthenaRootSharedWriterSvc.h.

◆ m_rootClientCount

int AthenaRootSharedWriterSvc::m_rootClientCount
private

Definition at line 52 of file AthenaRootSharedWriterSvc.h.

◆ m_rootClientIndex

int AthenaRootSharedWriterSvc::m_rootClientIndex
private

Definition at line 51 of file AthenaRootSharedWriterSvc.h.

◆ m_rootMergers

THashTable AthenaRootSharedWriterSvc::m_rootMergers
private

Definition at line 49 of file AthenaRootSharedWriterSvc.h.

◆ m_rootMonitor

TMonitor* AthenaRootSharedWriterSvc::m_rootMonitor
private

Definition at line 48 of file AthenaRootSharedWriterSvc.h.

◆ m_rootServerSocket

TServerSocket* AthenaRootSharedWriterSvc::m_rootServerSocket
private

Definition at line 47 of file AthenaRootSharedWriterSvc.h.


The documentation for this class was generated from the following files: