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;
456 std::string placementStr = placement->
toString();
457 placementStr +=
"[PNAME=";
458 placementStr += classDesc.
Name();
462 while (
sc.isRecoverable()) {
466 if (!
sc.isSuccess()) {
471 const void* buffer =
nullptr;
472 std::size_t nbytes = 0;
474 if (classDesc.
Name() ==
"Token") {
475 nbytes = strlen(
static_cast<const char*
>(obj)) + 1;
479 nbytes = classDesc.
SizeOf();
487 while (
sc.isRecoverable()) {
491 if (own) {
delete []
static_cast<const char*
>(buffer); }
493 if (!
sc.isSuccess()) {
494 ATH_MSG_ERROR(
"Could not share object for: " << placementStr);
499 if (info !=
nullptr) {
501 ATH_MSG_ERROR(
"Could not share dynamic aux store for: " << placementStr);
511 const char* tokenStr =
nullptr;
514 while (
sc.isRecoverable()) {
518 if (!
sc.isSuccess()) {
522 if (!strcmp(tokenStr,
"ABORT")) {
529 tempToken->
fromString(tokenStr); tokenStr =
nullptr;
531 token = tempToken; tempToken =
nullptr;
535 ATH_MSG_DEBUG(
"registerForWrite SKIPPED for uninitialized server, Placement = " << placement->
toString());
538 token = tempToken; tempToken =
nullptr;
544 token = getPoolSvc()->registerForWrite(placement, obj, classDesc);
549 token = AthenaPoolCnvSvc::registerForWrite(placement, obj, classDesc);
558 ATH_MSG_ERROR(
"Could not make AthenaPoolSharedIOCnvSvc a Share Client");
563 int num = token->
oid().first;
565 void* buffer =
nullptr;
566 std::size_t nbytes = 0;
568 while (
sc.isRecoverable()) {
571 if (!
sc.isSuccess()) {
579 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
582 std::string className = token->
auxString();
583 className = className.substr(className.find(
"[PNAME="));
584 className = className.substr(7, className.find(
']') - 7);
586 obj =
m_serializeSvc->deserialize(buffer, nbytes, cltype); buffer =
nullptr;
591 while (
sc.isRecoverable() && nbytes > 0) {
594 if (
sc.isSuccess() && nbytes > 0) {
596 classId.
fromString(
static_cast<const char*
>(buffer));
598 const std::string contName = std::string(
static_cast<const char*
>(buffer));
601 if (info !=
nullptr) {
602 if (
m_auxDynTool->hasAuxStore(contName, info->clazz().Class() ) && !
m_auxOutput->receiveStore(info->clazz().Class(), obj, num).isSuccess()) {
619 void* buffer =
nullptr;
620 std::size_t nbytes = 0;
621 StatusCode
sc = StatusCode::FAILURE;
626 while (
sc.isRecoverable()) {
631 if (!
sc.isSuccess()) {
639 while (
sc.isRecoverable() && nbytes > 0) {
642 if (
sc.isSuccess() && nbytes > 0) {
644 classId.
fromString(
static_cast<const char*
>(buffer));
646 const std::string contName = std::string(
static_cast<const char*
>(buffer));
649 if (info !=
nullptr) {
650 if (!
m_auxInput->receiveStore(info->clazz().Class(), obj).isSuccess()) {
661 AthenaPoolCnvSvc::setObjPtr(obj, token);
667 const std::string* par,
668 const unsigned long* ip,
669 IOpaqueAddress*& refpAddress) {
675 addressToken.
setDb(par[0].substr(4));
679 void* buffer =
nullptr;
680 std::size_t nbytes = 0;
682 while (
sc.isRecoverable()) {
686 if (!
sc.isSuccess()) {
688 return(StatusCode::FAILURE);
690 auto token = std::make_unique<Token>();
691 token->fromString(
static_cast<const char*
>(buffer)); buffer =
nullptr;
698 return(StatusCode::SUCCESS);
701 return(StatusCode::RECOVERABLE);
704 return AthenaPoolCnvSvc::createAddress(svcType, clid, par, ip, refpAddress);
710 const std::string& refAddress,
711 IOpaqueAddress*& refpAddress) {
712 return AthenaPoolCnvSvc::createAddress(svcType, clid, refAddress, refpAddress);
716 auto pos = connection.find(
"?pmerge=");
717 std::string conn = (pos == std::string::npos) ? connection : connection.substr(0, pos);
718 return AthenaPoolCnvSvc::cleanUp(conn);
732 return(StatusCode::FAILURE);
735 m_persSvcPerOutput.setValue(
false);
736 return(StatusCode::SUCCESS);
738 return(StatusCode::RECOVERABLE);
741 return(StatusCode::RECOVERABLE);
750 std::string streamPortSuffix;
753 return(StatusCode::FAILURE);
756 ATH_MSG_DEBUG(
"makeClient: Setting conversion service port suffix to " << streamPortSuffix);
761 return(StatusCode::SUCCESS);
764 std::string dummyStr;
770 return(StatusCode::FAILURE);
772 const char* tokenStr =
nullptr;
775 if (
sc.isSuccess() && tokenStr !=
nullptr && strlen(tokenStr) > 0 && num > 0) {
776 ATH_MSG_DEBUG(
"readData: " << tokenStr <<
", for client: " << num);
783 token.
fromString(tokenStr); tokenStr =
nullptr;
785 std::string objName =
"ALL";
786 if (useDetailChronoStat()) {
794 void* buffer =
nullptr;
795 std::size_t nbytes = 0;
798 while (
sc.isRecoverable()) {
801 delete []
static_cast<char*
>(buffer); buffer =
nullptr;
802 if (!
sc.isSuccess()) {
804 return(StatusCode::FAILURE);
807 if (info !=
nullptr) {
810 return(StatusCode::FAILURE);
816 return(StatusCode::FAILURE);
819 std::string returnToken;
821 if( metadataToken ) {
822 returnToken = metadataToken->
toString();
823 metadataToken->
release(); metadataToken =
nullptr;
831 return(StatusCode::FAILURE);
834 return(StatusCode::RECOVERABLE);
836 return(StatusCode::SUCCESS);
845 return(StatusCode::SUCCESS);
854 StatusCode
sc = StatusCode::SUCCESS;
855 while (
sc.isSuccess()) {
861 while (
sc.isRecoverable()) {
865 return StatusCode::FAILURE;
876 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).
#define ATLAS_THREAD_SAFE
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.
void commit()
Save catalog to file.
static const DbType POOL_StorageType
static constexpr CLID ID()