12#include "GaudiKernel/ClassID.h"
13#include "GaudiKernel/FileIncident.h"
14#include "GaudiKernel/GenericAddress.h"
15#include "GaudiKernel/IIncidentSvc.h"
16#include "GaudiKernel/IOpaqueAddress.h"
46 return StatusCode::FAILURE;
64 m_poolSvc->setShareMode(
true);
69 incSvc->addListener(
this,
"StoreCleared", pri);
72 return this->AthenaPoolCnvSvc::initialize();
94 return this->AthenaPoolCnvSvc::finalize();
98 const std::string& openMode) {
99 return AthenaPoolCnvSvc::connectOutput(outputConnectionSpec, openMode);
104 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
107 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
108 return(StatusCode::FAILURE);
113 return(StatusCode::SUCCESS);
117 ATH_MSG_DEBUG(std::format(
"connectOutput SKIPPED for metadata-only server: {}", outputConnectionSpec));
118 return(StatusCode::SUCCESS);
122 return(StatusCode::SUCCESS);
128 std::size_t apend = outputConnectionSpec.find(
'[');
129 if (apend != std::string::npos) {
130 outputConnection += outputConnectionSpec.substr(apend);
132 if (outputConnectionSpec.find(
"[PoolContainerPrefix=" +
m_metadataContainerProp.value() +
"]") != std::string::npos) {
133 return AthenaPoolCnvSvc::connectOutput(outputConnection,
"APPEND");
135 return AthenaPoolCnvSvc::connectOutput(outputConnection);
141 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
145 if (!this->
cleanUp(outputConnection).isSuccess()) {
147 return(StatusCode::FAILURE);
149 return(StatusCode::SUCCESS);
152 ATH_MSG_DEBUG(
"commitOutput SKIPPED for uninitialized server.");
153 return(StatusCode::SUCCESS);
155 std::map<void*, RootType> commitCache;
156 std::string fileName;
159 const char* placementStr =
nullptr;
162 if (
sc.isSuccess() && placementStr !=
nullptr && strlen(placementStr) > 6 && num > 0) {
163 const char * matchedChars = strstr(placementStr,
"[FILE=");
165 ATH_MSG_ERROR(std::format(
"No matching filename in {}", placementStr));
168 fileName = matchedChars;
169 fileName = fileName.substr(6, fileName.find(
']') - 6);
171 ATH_MSG_ERROR(std::format(
"Failed to connectOutput for {}", fileName));
175 bool dataHeaderSeen =
false;
176 std::string dataHeaderID;
178 std::string objName =
"ALL";
179 if (useDetailChronoStat()) {
180 objName = placementStr;
185 std::string_view pStr = placementStr;
186 std::string::size_type cpos = pStr.find (
"[CONT=");
187 if (cpos == std::string::npos) {
188 ATH_MSG_ERROR(std::format(
"No CONT field in placement string: {}", pStr));
189 return StatusCode::FAILURE;
191 std::string tokenStr (pStr.substr(0, cpos));
192 std::string contName (pStr.substr(cpos, std::string::npos));
193 std::string::size_type cl1 = contName.find(
']');
194 if (cl1 == std::string::npos) {
195 ATH_MSG_ERROR(std::format(
"Missing close bracket after CONT field in placement string: {}", pStr));
196 return StatusCode::FAILURE;
198 tokenStr.append(contName, cl1 + 1);
199 contName = contName.substr(6, cl1 - 6);
201 std::string::size_type ppos = pStr.find (
"[PNAME=");
202 if (ppos == std::string::npos) {
203 ATH_MSG_ERROR(std::format(
"No PNAME field in placement string: {}", pStr));
204 return StatusCode::FAILURE;
206 std::string className (pStr.substr(ppos, std::string::npos));
207 std::string::size_type cl2 = className.find(
']');
208 if (cl2 == std::string::npos) {
209 ATH_MSG_ERROR(std::format(
"Missing close bracket after PNAME field in placement string: {}", pStr));
210 return StatusCode::FAILURE;
212 className = className.substr(7, cl2 - 7);
215 const std::string numStr = std::to_string(num);
217 bool foundContainer =
false;
218 std::size_t opPos = contName.find(
'(');
220 foundContainer =
true;
223 if (contName.compare(0, opPos, item) == 0){
224 foundContainer =
true;
230 if (len > 0 && foundContainer && contName[len] ==
'(' ) {
237 memName, {}, memName,
238 "BeginInputMemFile",
"EndInputMemFile");
245 sc = metadataSvc->shmProxy(std::format(
"{}[NUM={}]", pStr, numStr));
246 if (
sc.isRecoverable()) {
248 }
else if (
sc.isFailure()) {
259 if( m_oneDataHeaderForm.value() ) {
260 auto placementWithSwn = [&] {
return std::format(
"{}[SWN={}]", placementStr, num); };
261 if( className ==
"DataHeaderForm_p6" ) {
264 "", placementWithSwn());
265 DHcnv->updateRepRefs(&address,
static_cast<DataObject*
>(obj)).ignore();
271 if (token ==
nullptr) {
275 tokenStr = token->toString();
277 if( className ==
"DataHeader_p6" ) {
280 tokenStr, placementWithSwn());
281 if (!DHcnv->updateRep(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
286 if (className !=
"Token" && className !=
"DataHeaderForm_p6" && !classDesc.
IsFundamental()) {
287 commitCache.insert(std::pair<void*, RootType>(obj, classDesc));
289 placementStr =
nullptr;
293 placement.
fromString(placementStr); placementStr =
nullptr;
295 if (token ==
nullptr) {
299 tokenStr = token->toString();
300 if (className ==
"DataHeader_p6") {
305 if (!DHcnv->updateRep(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
309 dataHeaderSeen =
true;
316 dataHeaderID = std::format(
"{}/{}/{}", token->contID(), numStr, token->dbID().toString());
317 }
else if (dataHeaderSeen) {
318 dataHeaderSeen =
false;
321 if (className ==
"DataHeaderForm_p6") {
324 tokenStr, dataHeaderID);
325 if (!DHcnv->updateRepRefs(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
326 ATH_MSG_ERROR(
"Failed updateRepRefs for obj = " << tokenStr);
331 GenericAddress address(0, 0,
"", dataHeaderID);
332 if (!DHcnv->updateRepRefs(&address,
nullptr).isSuccess()) {
338 if (className !=
"Token" && className !=
"DataHeaderForm_p6" && !classDesc.
IsFundamental()) {
339 commitCache.insert(std::pair<void*, RootType>(obj, classDesc));
346 while (
sc.isRecoverable()) {
349 if (!
sc.isSuccess()) {
355 while (
sc.isRecoverable()) {
358 if (
sc.isFailure()) {
363 if (dataHeaderSeen) {
365 GenericAddress address(0, 0,
"", std::move(dataHeaderID));
366 if (!DHcnv->updateRepRefs(&address,
nullptr).isSuccess()) {
371 placementStr =
nullptr;
372 }
else if (
sc.isSuccess() && placementStr !=
nullptr && strncmp(placementStr,
"stop", 4) == 0) {
373 return(StatusCode::RECOVERABLE);
374 }
else if (
sc.isRecoverable() || num == -1) {
375 return(StatusCode::RECOVERABLE);
377 if (
sc.isFailure() || fileName.empty()) {
382 memName, {}, memName,
383 "BeginInputMemFile",
"EndInputMemFile");
385 if (
sc.isFailure()) {
386 ATH_MSG_INFO(
"All SharedWriter clients stopped - exiting");
390 return(StatusCode::FAILURE);
394 ATH_MSG_DEBUG(std::format(
"commitOutput SKIPPED for metadata-only server: {}", outputConnectionSpec));
395 return(StatusCode::SUCCESS);
397 if (outputConnection.empty()) {
398 outputConnection = std::move(fileName);
400 outputConnection = outputConnectionSpec;
405 std::size_t
merge = outputConnection.find(
"?pmerge=");
406 const std::string baseOutputConnection = outputConnection.substr(0,
merge);
415 StatusCode status = AthenaPoolCnvSvc::commitOutput(outputConnection, doCommit);
416 for (
auto& [ptr, rootType] : commitCache) {
417 rootType.Destruct(ptr);
424 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
427 return(StatusCode::SUCCESS);
433 return(StatusCode::SUCCESS);
442 return AthenaPoolCnvSvc::disconnectOutput(outputConnectionSpec +
m_streamPortString.value());
449 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
453 Token* token =
nullptr;
458 std::string placementStr = placement->
toString();
459 placementStr +=
"[PNAME=";
460 placementStr += classDesc.
Name();
464 while (
sc.isRecoverable()) {
468 if (!
sc.isSuccess()) {
473 const void* buffer =
nullptr;
474 std::size_t nbytes = 0;
476 if (classDesc.
Name() ==
"Token") {
477 nbytes = strlen(
static_cast<const char*
>(obj)) + 1;
481 nbytes = classDesc.
SizeOf();
489 while (
sc.isRecoverable()) {
493 if (own) {
delete []
static_cast<const char*
>(buffer); }
495 if (!
sc.isSuccess()) {
496 ATH_MSG_ERROR(
"Could not share object for: " << placementStr);
501 if (info !=
nullptr) {
503 ATH_MSG_ERROR(
"Could not share dynamic aux store for: " << placementStr);
513 const char* tokenStr =
nullptr;
516 while (
sc.isRecoverable()) {
520 if (!
sc.isSuccess()) {
524 if (!strcmp(tokenStr,
"ABORT")) {
531 tempToken->
fromString(tokenStr); tokenStr =
nullptr;
533 token = tempToken; tempToken =
nullptr;
537 ATH_MSG_DEBUG(
"registerForWrite SKIPPED for uninitialized server, Placement = " << placement->
toString());
540 token = tempToken; tempToken =
nullptr;
546 token = m_poolSvc->registerForWrite(placement, obj, classDesc);
551 token = AthenaPoolCnvSvc::registerForWrite(placement, obj, classDesc);
560 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
565 int num = token->
oid().first;
567 void* buffer =
nullptr;
568 std::size_t nbytes = 0;
570 while (
sc.isRecoverable()) {
573 if (!
sc.isSuccess()) {
581 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
584 std::string className = token->
auxString();
585 className = className.substr(className.find(
"[PNAME="));
586 className = className.substr(7, className.find(
']') - 7);
588 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
593 while (
sc.isRecoverable() && nbytes > 0) {
596 if (
sc.isSuccess() && nbytes > 0) {
598 classId.
fromString(
static_cast<const char*
>(buffer));
600 const std::string contName = std::string(
static_cast<const char*
>(buffer));
603 if (info !=
nullptr) {
604 if (
m_auxDynTool->hasAuxStore(contName, info->clazz().Class() ) && !
m_auxOutput->receiveStore(info->clazz().Class(), obj, num).isSuccess()) {
621 void* buffer =
nullptr;
622 std::size_t nbytes = 0;
623 StatusCode
sc = StatusCode::FAILURE;
628 while (
sc.isRecoverable()) {
633 if (!
sc.isSuccess()) {
638 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
642 while (
sc.isRecoverable() && nbytes > 0) {
645 if (
sc.isSuccess() && nbytes > 0) {
647 classId.
fromString(
static_cast<const char*
>(buffer));
649 const std::string contName = std::string(
static_cast<const char*
>(buffer));
652 if (info !=
nullptr) {
653 if (!
m_auxInput->receiveStore(info->clazz().Class(), obj).isSuccess()) {
664 AthenaPoolCnvSvc::setObjPtr(obj, token);
670 const std::string* par,
671 const unsigned long* ip,
672 IOpaqueAddress*& refpAddress) {
678 addressToken.
setDb(par[0].substr(4));
682 void* buffer =
nullptr;
683 std::size_t nbytes = 0;
685 while (
sc.isRecoverable()) {
689 if (!
sc.isSuccess()) {
691 return(StatusCode::FAILURE);
693 auto token = std::make_unique<Token>();
694 token->fromString(
static_cast<const char*
>(buffer)); buffer =
nullptr;
701 return(StatusCode::SUCCESS);
704 return(StatusCode::RECOVERABLE);
707 if (par[0].compare(0, 3,
"SHM") == 0) {
708 std::unique_ptr<Token> token;
709 token = std::make_unique<Token>();
711 token->setAuxString(
"[PNAME=" + par[2] +
"]");
715 return(StatusCode::SUCCESS);
717 return AthenaPoolCnvSvc::createAddress(svcType, clid, par, ip, refpAddress);
724 const std::string& refAddress,
725 IOpaqueAddress*& refpAddress) {
726 return AthenaPoolCnvSvc::createAddress(svcType, clid, refAddress, refpAddress);
730 auto pos = connection.find(
"?pmerge=");
731 std::string conn = (pos == std::string::npos) ? connection : connection.substr(0, pos);
732 return AthenaPoolCnvSvc::cleanUp(conn);
746 return(StatusCode::FAILURE);
749 m_persSvcPerOutput.setValue(
false);
750 return(StatusCode::SUCCESS);
752 return(StatusCode::RECOVERABLE);
755 return(StatusCode::RECOVERABLE);
764 std::string streamPortSuffix;
767 return(StatusCode::FAILURE);
768 }
else if (!streamPortSuffix.empty()) {
771 ATH_MSG_DEBUG(
"makeClient: Setting conversion service port suffix to " << streamPortSuffix);
776 return(StatusCode::SUCCESS);
779 std::string dummyStr;
785 return(StatusCode::FAILURE);
787 const char* tokenStr =
nullptr;
790 if (
sc.isSuccess() && tokenStr !=
nullptr && strlen(tokenStr) > 0 && num > 0) {
791 ATH_MSG_DEBUG(
"readData: " << tokenStr <<
", for client: " << num);
798 token.
fromString(tokenStr); tokenStr =
nullptr;
800 std::string objName =
"ALL";
801 if (useDetailChronoStat()) {
809 void* buffer =
nullptr;
810 std::size_t nbytes = 0;
813 while (
sc.isRecoverable()) {
816 delete []
static_cast<char*
>(buffer); buffer =
nullptr;
817 if (!
sc.isSuccess()) {
819 return(StatusCode::FAILURE);
822 if (info !=
nullptr) {
825 return(StatusCode::FAILURE);
831 return(StatusCode::FAILURE);
834 std::string returnToken;
836 if( metadataToken ) {
837 returnToken = metadataToken->
toString();
838 metadataToken->
release(); metadataToken =
nullptr;
846 return(StatusCode::FAILURE);
849 return(StatusCode::RECOVERABLE);
851 return(StatusCode::SUCCESS);
856 m_poolSvc->commitCatalog();
857 m_poolSvc->startCatalog();
858 return(StatusCode::SUCCESS);
867 StatusCode
sc = StatusCode::SUCCESS;
868 while (
sc.isSuccess()) {
874 while (
sc.isRecoverable()) {
878 return StatusCode::FAILURE;
889 base_class(name, pSvcLocator) {
#define ATH_CHECK
Evaluate an expression and check for errors.
#define ATH_MSG_DEBUG(x,...)
#define ATH_MSG_ERROR(x,...)
#define ATH_MSG_WARNING(x,...)
#define ATH_MSG_VERBOSE(x,...)
#define ATH_MSG_INFO(x,...)
#define ATH_MSG_FATAL(x,...)
This file contains the class definition for the AthenaPoolSharedIOCnvSvc class.
uint32_t CLID
The Class ID type.
This file contains the class definition for the Placement class (migrated from POOL).
This file contains the class definition for the TokenAddress class.
This file contains the class definition for the Token class (migrated from POOL).
std::map< std::string, int > m_fileCommitCounter
Force SharedWriter to flush data to output file at given intervals, needed by parallel compression.
virtual StatusCode commitOutput(const std::string &outputConnectionSpec, bool doCommit) override
Implementation of IConversionSvc: Commit pending output.
Gaudi::Property< std::string > m_streamPortString
Extension to use ROOT TMemFile for event data, "?pmerge=<host>:<port>" for a TCP socket (default),...
virtual StatusCode finalize() override
Required of all Gaudi Services.
AthenaPoolSharedIOCnvSvc(const std::string &name, ISvcLocator *pSvcLocator)
Standard Service Constructor.
virtual StatusCode cleanUp(const std::string &connection) override
Implement cleanUp to call all registered IAthenaPoolCleanUp cleanUp() function.
virtual void setObjPtr(void *&obj, const Token *token) override
bool m_streamServerActive
StatusCode createAddress(long svcType, const CLID &clid, const std::string *par, const unsigned long *ip, IOpaqueAddress *&refpAddress) override
Create a Generic address using explicit arguments to identify a single object.
std::unique_ptr< RootAuxDynIO::IAuxDynShare > m_auxOutput
StatusCode abortSharedWrClients(int client_n)
Send abort to SharedWriter clients if the server quits on error.
virtual StatusCode makeServer(int num) override
Make this a server.
virtual void handle(const Incident &incident) override
Implementation of IIncidentListener: Handle for EndEvent incidence.
Gaudi::Property< std::string > m_metadataContainerProp
For SharedWriter: To use MetadataSvc to merge data placed in a certain container.
virtual StatusCode disconnectOutput(const std::string &outputConnectionSpec) override
Disconnect to the output connection.
virtual StatusCode readData() override
Read the next data object.
virtual StatusCode commitCatalog() override
Commit Catalog.
virtual StatusCode connectOutput(const std::string &outputConnectionSpec, const std::string &openMode) override
Implementation of IConversionSvc: Connect to the output connection specification with open mode.
std::unique_ptr< RootAuxDynIO::IAuxDynShare > m_auxInput
virtual ~AthenaPoolSharedIOCnvSvc()
Destructor.
Gaudi::Property< std::vector< std::string > > m_metadataContainersAug
Gaudi::Property< bool > m_parallelCompression
Use Athena Object sharing for metadata only, event data is collected and send via ROOT TMemFile.
ToolHandle< IAthenaIPCTool > m_outputStreamingTool
Gaudi::Property< std::map< std::string, int > > m_fileFlushSetting
Gaudi::Property< int > m_streamingTechnology
Use Streaming for selected technologies only.
virtual StatusCode initialize() override
Required of all Gaudi Services.
ToolHandle< IAthenaIPCTool > m_inputStreamingTool
ServiceHandle< IAthenaSerializeSvc > m_serializeSvc
virtual Token * registerForWrite(Placement *placement, const void *obj, const RootType &classDesc) override
Gaudi::Property< int > m_makeStreamingToolClient
Make this instance a Streaming Client during first connect/write automatically.
virtual StatusCode makeClient(int num) override
Make this a client.
std::unique_ptr< RootAuxDynIO::IFactoryTool > m_auxDynTool
This class provides a encapsulation of a GUID/UUID/CLSID/IID data structure (128 bit number).
constexpr void fromString(std::string_view s)
Automatic conversion from string representation.
static const Guid & null() noexcept
NULL-Guid: static class method.
constexpr void toString(std::span< char, StrLen > buf, bool uppercase=true) const noexcept
Automatic conversion to string representation.
This class holds all the necessary information to guide the writing of an object in a physical place.
const std::string & auxString() const
Access auxiliary string.
const std::string & containerName() const
Access container name.
Placement & setFileName(const std::string &fileName)
Set file name.
const std::string toString() const
Retrieve the string representation of the placement.
Placement & setTechnology(int technology)
Set technology type.
const std::string & fileName() const
Access file name.
Placement & fromString(const std::string &from)
Build from the string representation of a placement.
int technology() const
Access technology type.
static TScopeAdapter ByNameNoQuiet(const std::string &name, Bool_t load=kTRUE)
Bool_t IsFundamental() const
std::string Name(unsigned int mod=Reflex::SCOPED) const
void Destruct(void *place) const
This class provides a Generic Transient Address for POOL tokens.
This class provides a token that identifies in a unique way objects on the persistent storage.
const std::string & auxString() const
Access auxiliary string.
Token & setCont(const std::string &cnt)
Set container name.
Token & setDb(const Guid &db)
Set database name.
const std::string & contID() const
Access container identifier.
const Guid & classID() const
Access database identifier.
Token & setClassID(const Guid &cl_id)
Access database identifier.
const std::string toString() const
Retrieve the string representation of the token.
int technology() const
Access technology type.
int release()
Release token: Decrease reference count and eventually delete.
Token & setOid(const OID_t &oid)
Set object identifier.
const OID_t & oid() const
Access object identifier.
Token & fromString(const std::string_view from)
Build from the string representation of a token.
const Guid & dbID() const
Access database identifier.
Token & setAuxString(std::string &&auxString)
Set auxiliary string.
static const RootType forGuid(const Guid &info)
Access classes by Guid.
static Guid guid(const RootType &id)
Determine Guid (normalized string form) from reflection type.
Definition of class DbTypeInfo.
static const DbTypeInfo * create(const std::string &cl_name)
Create type information using name.
static DbType getType(const std::string &name)
Access known storage type object by name.
static const DbType POOL_StorageType
static constexpr CLID ID()