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"
45 return StatusCode::FAILURE;
63 getPoolSvc()->setShareMode(
true);
68 incSvc->addListener(
this,
"StoreCleared", pri);
71 return this->AthenaPoolCnvSvc::initialize();
93 return this->AthenaPoolCnvSvc::finalize();
97 const std::string& openMode) {
98 return AthenaPoolCnvSvc::connectOutput(outputConnectionSpec, openMode);
103 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
106 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
107 return(StatusCode::FAILURE);
112 return(StatusCode::SUCCESS);
116 ATH_MSG_DEBUG(std::format(
"connectOutput SKIPPED for metadata-only server: {}", outputConnectionSpec));
117 return(StatusCode::SUCCESS);
121 return(StatusCode::SUCCESS);
127 std::size_t apend = outputConnectionSpec.find(
'[');
128 if (apend != std::string::npos) {
129 outputConnection += outputConnectionSpec.substr(apend);
131 if (outputConnectionSpec.find(
"[PoolContainerPrefix=" +
m_metadataContainerProp.value() +
"]") != std::string::npos) {
132 return AthenaPoolCnvSvc::connectOutput(outputConnection,
"APPEND");
134 return AthenaPoolCnvSvc::connectOutput(outputConnection);
140 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
144 if (!this->
cleanUp(outputConnection).isSuccess()) {
146 return(StatusCode::FAILURE);
148 return(StatusCode::SUCCESS);
151 ATH_MSG_DEBUG(
"commitOutput SKIPPED for uninitialized server.");
152 return(StatusCode::SUCCESS);
154 std::map<void*, RootType> commitCache;
155 std::string fileName;
158 const char* placementStr =
nullptr;
161 if (
sc.isSuccess() && placementStr !=
nullptr && strlen(placementStr) > 6 && num > 0) {
162 const char * matchedChars = strstr(placementStr,
"[FILE=");
164 ATH_MSG_ERROR(std::format(
"No matching filename in {}", placementStr));
167 fileName = matchedChars;
168 fileName = fileName.substr(6, fileName.find(
']') - 6);
170 ATH_MSG_ERROR(std::format(
"Failed to connectOutput for {}", fileName));
174 bool dataHeaderSeen =
false;
175 std::string dataHeaderID;
177 std::string objName =
"ALL";
178 if (useDetailChronoStat()) {
179 objName = placementStr;
184 std::string_view pStr = placementStr;
185 std::string::size_type cpos = pStr.find (
"[CONT=");
186 if (cpos == std::string::npos) {
187 ATH_MSG_ERROR(std::format(
"No CONT field in placement string: {}", pStr));
188 return StatusCode::FAILURE;
190 std::string tokenStr (pStr.substr(0, cpos));
191 std::string contName (pStr.substr(cpos, std::string::npos));
192 std::string::size_type cl1 = contName.find(
']');
193 if (cl1 == std::string::npos) {
194 ATH_MSG_ERROR(std::format(
"Missing close bracket after CONT field in placement string: {}", pStr));
195 return StatusCode::FAILURE;
197 tokenStr.append(contName, cl1 + 1);
198 contName = contName.substr(6, cl1 - 6);
200 std::string::size_type ppos = pStr.find (
"[PNAME=");
201 if (ppos == std::string::npos) {
202 ATH_MSG_ERROR(std::format(
"No PNAME field in placement string: {}", pStr));
203 return StatusCode::FAILURE;
205 std::string className (pStr.substr(ppos, std::string::npos));
206 std::string::size_type cl2 = className.find(
']');
207 if (cl2 == std::string::npos) {
208 ATH_MSG_ERROR(std::format(
"Missing close bracket after PNAME field in placement string: {}", pStr));
209 return StatusCode::FAILURE;
211 className = className.substr(7, cl2 - 7);
214 const std::string numStr = std::to_string(num);
216 bool foundContainer =
false;
217 std::size_t opPos = contName.find(
'(');
219 foundContainer =
true;
222 if (contName.compare(0, opPos, item) == 0){
223 foundContainer =
true;
229 if (len > 0 && foundContainer && contName[len] ==
'(' ) {
236 memName, {}, memName,
237 "BeginInputMemFile",
"EndInputMemFile");
244 sc = metadataSvc->shmProxy(std::format(
"{}[NUM={}]", pStr, numStr));
245 if (
sc.isRecoverable()) {
247 }
else if (
sc.isFailure()) {
258 if( m_oneDataHeaderForm.value() ) {
259 auto placementWithSwn = [&] {
return std::format(
"{}[SWN={}]", placementStr, num); };
260 if( className ==
"DataHeaderForm_p6" ) {
263 "", placementWithSwn());
264 DHcnv->updateRepRefs(&address,
static_cast<DataObject*
>(obj)).ignore();
270 if (token ==
nullptr) {
274 tokenStr = token->toString();
276 if( className ==
"DataHeader_p6" ) {
279 tokenStr, placementWithSwn());
280 if (!DHcnv->updateRep(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
285 if (className !=
"Token" && className !=
"DataHeaderForm_p6" && !classDesc.
IsFundamental()) {
286 commitCache.insert(std::pair<void*, RootType>(obj, classDesc));
288 placementStr =
nullptr;
292 placement.
fromString(placementStr); placementStr =
nullptr;
294 if (token ==
nullptr) {
298 tokenStr = token->toString();
299 if (className ==
"DataHeader_p6") {
304 if (!DHcnv->updateRep(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
308 dataHeaderSeen =
true;
315 dataHeaderID = std::format(
"{}/{}/{}", token->contID(), numStr, token->dbID().toString());
316 }
else if (dataHeaderSeen) {
317 dataHeaderSeen =
false;
320 if (className ==
"DataHeaderForm_p6") {
323 tokenStr, dataHeaderID);
324 if (!DHcnv->updateRepRefs(&address,
static_cast<DataObject*
>(obj)).isSuccess()) {
325 ATH_MSG_ERROR(
"Failed updateRepRefs for obj = " << tokenStr);
330 GenericAddress address(0, 0,
"", dataHeaderID);
331 if (!DHcnv->updateRepRefs(&address,
nullptr).isSuccess()) {
337 if (className !=
"Token" && className !=
"DataHeaderForm_p6" && !classDesc.
IsFundamental()) {
338 commitCache.insert(std::pair<void*, RootType>(obj, classDesc));
345 while (
sc.isRecoverable()) {
348 if (!
sc.isSuccess()) {
354 while (
sc.isRecoverable()) {
357 if (
sc.isFailure()) {
362 if (dataHeaderSeen) {
364 GenericAddress address(0, 0,
"", std::move(dataHeaderID));
365 if (!DHcnv->updateRepRefs(&address,
nullptr).isSuccess()) {
370 placementStr =
nullptr;
371 }
else if (
sc.isSuccess() && placementStr !=
nullptr && strncmp(placementStr,
"stop", 4) == 0) {
372 return(StatusCode::RECOVERABLE);
373 }
else if (
sc.isRecoverable() || num == -1) {
374 return(StatusCode::RECOVERABLE);
376 if (
sc.isFailure() || fileName.empty()) {
381 memName, {}, memName,
382 "BeginInputMemFile",
"EndInputMemFile");
384 if (
sc.isFailure()) {
385 ATH_MSG_INFO(
"All SharedWriter clients stopped - exiting");
389 return(StatusCode::FAILURE);
393 ATH_MSG_DEBUG(std::format(
"commitOutput SKIPPED for metadata-only server: {}", outputConnectionSpec));
394 return(StatusCode::SUCCESS);
396 if (outputConnection.empty()) {
397 outputConnection = std::move(fileName);
399 outputConnection = outputConnectionSpec;
404 std::size_t
merge = outputConnection.find(
"?pmerge=");
405 const std::string baseOutputConnection = outputConnection.substr(0,
merge);
414 StatusCode status = AthenaPoolCnvSvc::commitOutput(outputConnection, doCommit);
415 for (
auto& [ptr, rootType] : commitCache) {
416 rootType.Destruct(ptr);
423 std::string outputConnection = outputConnectionSpec.substr(0, outputConnectionSpec.find(
'['));
426 return(StatusCode::SUCCESS);
432 return(StatusCode::SUCCESS);
441 return AthenaPoolCnvSvc::disconnectOutput(outputConnectionSpec +
m_streamPortString.value());
448 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
452 Token* token =
nullptr;
457 std::string placementStr = placement->
toString();
458 placementStr +=
"[PNAME=";
459 placementStr += classDesc.
Name();
463 while (
sc.isRecoverable()) {
467 if (!
sc.isSuccess()) {
472 const void* buffer =
nullptr;
473 std::size_t nbytes = 0;
475 if (classDesc.
Name() ==
"Token") {
476 nbytes = strlen(
static_cast<const char*
>(obj)) + 1;
480 nbytes = classDesc.
SizeOf();
488 while (
sc.isRecoverable()) {
492 if (own) {
delete []
static_cast<const char*
>(buffer); }
494 if (!
sc.isSuccess()) {
495 ATH_MSG_ERROR(
"Could not share object for: " << placementStr);
500 if (info !=
nullptr) {
502 ATH_MSG_ERROR(
"Could not share dynamic aux store for: " << placementStr);
512 const char* tokenStr =
nullptr;
515 while (
sc.isRecoverable()) {
519 if (!
sc.isSuccess()) {
523 if (!strcmp(tokenStr,
"ABORT")) {
530 tempToken->
fromString(tokenStr); tokenStr =
nullptr;
532 token = tempToken; tempToken =
nullptr;
536 ATH_MSG_DEBUG(
"registerForWrite SKIPPED for uninitialized server, Placement = " << placement->
toString());
539 token = tempToken; tempToken =
nullptr;
545 token = getPoolSvc()->registerForWrite(placement, obj, classDesc);
550 token = AthenaPoolCnvSvc::registerForWrite(placement, obj, classDesc);
559 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
564 int num = token->
oid().first;
566 void* buffer =
nullptr;
567 std::size_t nbytes = 0;
569 while (
sc.isRecoverable()) {
572 if (!
sc.isSuccess()) {
580 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
583 std::string className = token->
auxString();
584 className = className.substr(className.find(
"[PNAME="));
585 className = className.substr(7, className.find(
']') - 7);
587 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
592 while (
sc.isRecoverable() && nbytes > 0) {
595 if (
sc.isSuccess() && nbytes > 0) {
597 classId.
fromString(
static_cast<const char*
>(buffer));
599 const std::string contName = std::string(
static_cast<const char*
>(buffer));
602 if (info !=
nullptr) {
603 if (
m_auxDynTool->hasAuxStore(contName, info->clazz().Class() ) && !
m_auxOutput->receiveStore(info->clazz().Class(), obj, num).isSuccess()) {
620 void* buffer =
nullptr;
621 std::size_t nbytes = 0;
622 StatusCode
sc = StatusCode::FAILURE;
627 while (
sc.isRecoverable()) {
632 if (!
sc.isSuccess()) {
637 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
641 while (
sc.isRecoverable() && nbytes > 0) {
644 if (
sc.isSuccess() && nbytes > 0) {
646 classId.
fromString(
static_cast<const char*
>(buffer));
648 const std::string contName = std::string(
static_cast<const char*
>(buffer));
651 if (info !=
nullptr) {
652 if (!
m_auxInput->receiveStore(info->clazz().Class(), obj).isSuccess()) {
663 AthenaPoolCnvSvc::setObjPtr(obj, token);
669 const std::string* par,
670 const unsigned long* ip,
671 IOpaqueAddress*& refpAddress) {
677 addressToken.
setDb(par[0].substr(4));
681 void* buffer =
nullptr;
682 std::size_t nbytes = 0;
684 while (
sc.isRecoverable()) {
688 if (!
sc.isSuccess()) {
690 return(StatusCode::FAILURE);
692 auto token = std::make_unique<Token>();
693 token->fromString(
static_cast<const char*
>(buffer)); buffer =
nullptr;
700 return(StatusCode::SUCCESS);
703 return(StatusCode::RECOVERABLE);
706 if (par[0].compare(0, 3,
"SHM") == 0) {
707 std::unique_ptr<Token> token;
708 token = std::make_unique<Token>();
710 token->setAuxString(
"[PNAME=" + par[2] +
"]");
714 return(StatusCode::SUCCESS);
716 return AthenaPoolCnvSvc::createAddress(svcType, clid, par, ip, refpAddress);
723 const std::string& refAddress,
724 IOpaqueAddress*& refpAddress) {
725 return AthenaPoolCnvSvc::createAddress(svcType, clid, refAddress, refpAddress);
729 auto pos = connection.find(
"?pmerge=");
730 std::string conn = (pos == std::string::npos) ? connection : connection.substr(0, pos);
731 return AthenaPoolCnvSvc::cleanUp(conn);
745 return(StatusCode::FAILURE);
748 m_persSvcPerOutput.setValue(
false);
749 return(StatusCode::SUCCESS);
751 return(StatusCode::RECOVERABLE);
754 return(StatusCode::RECOVERABLE);
763 std::string streamPortSuffix;
766 return(StatusCode::FAILURE);
769 ATH_MSG_DEBUG(
"makeClient: Setting conversion service port suffix to " << streamPortSuffix);
774 return(StatusCode::SUCCESS);
777 std::string dummyStr;
783 return(StatusCode::FAILURE);
785 const char* tokenStr =
nullptr;
788 if (
sc.isSuccess() && tokenStr !=
nullptr && strlen(tokenStr) > 0 && num > 0) {
789 ATH_MSG_DEBUG(
"readData: " << tokenStr <<
", for client: " << num);
796 token.
fromString(tokenStr); tokenStr =
nullptr;
798 std::string objName =
"ALL";
799 if (useDetailChronoStat()) {
807 void* buffer =
nullptr;
808 std::size_t nbytes = 0;
811 while (
sc.isRecoverable()) {
814 delete []
static_cast<char*
>(buffer); buffer =
nullptr;
815 if (!
sc.isSuccess()) {
817 return(StatusCode::FAILURE);
820 if (info !=
nullptr) {
823 return(StatusCode::FAILURE);
829 return(StatusCode::FAILURE);
832 std::string returnToken;
834 if( metadataToken ) {
835 returnToken = metadataToken->
toString();
836 metadataToken->
release(); metadataToken =
nullptr;
844 return(StatusCode::FAILURE);
847 return(StatusCode::RECOVERABLE);
849 return(StatusCode::SUCCESS);
854 getPoolSvc()->commitCatalog();
855 getPoolSvc()->startCatalog();
856 return(StatusCode::SUCCESS);
865 StatusCode
sc = StatusCode::SUCCESS;
866 while (
sc.isSuccess()) {
872 while (
sc.isRecoverable()) {
876 return StatusCode::FAILURE;
887 base_class(name, pSvcLocator) {
#define ATH_CHECK
Evaluate an expression and check for errors.
#define ATH_MSG_VERBOSE(x)
#define ATH_MSG_WARNING(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>".
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 TypeH forGuid(const Guid &info)
Access classes by Guid.
static Guid guid(const TypeH &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()