#include <open62541/server_pubsub.h>
#include "ua_pubsub_internal.h"
#ifdef UA_ENABLE_PUBSUB
#ifdef UA_ENABLE_PUBSUB_SKS
#include "ua_pubsub_keystorage.h"
#endif
UA_ReaderGroup *
UA_ReaderGroup_find(UA_PubSubManager *psm, const UA_NodeId id) {
if(!psm)
return NULL;
UA_PubSubConnection *psc;
TAILQ_FOREACH(psc, &psm->connections, listEntry) {
UA_ReaderGroup *rg;
LIST_FOREACH(rg, &psc->readerGroups, listEntry) {
if(UA_NodeId_equal(&id, &rg->head.identifier))
return rg;
}
}
return NULL;
}
UA_StatusCode
UA_ReaderGroupConfig_copy(const UA_ReaderGroupConfig *src,
UA_ReaderGroupConfig *dst) {
memcpy(dst, src, sizeof(UA_ReaderGroupConfig));
UA_StatusCode res = UA_STATUSCODE_GOOD;
res |= UA_String_copy(&src->name, &dst->name);
res |= UA_KeyValueMap_copy(&src->groupProperties, &dst->groupProperties);
res |= UA_String_copy(&src->securityGroupId, &dst->securityGroupId);
res |= UA_ExtensionObject_copy(&src->transportSettings, &dst->transportSettings);
if(res != UA_STATUSCODE_GOOD)
UA_ReaderGroupConfig_clear(dst);
return res;
}
void
UA_ReaderGroupConfig_clear(UA_ReaderGroupConfig *readerGroupConfig) {
UA_String_clear(&readerGroupConfig->name);
UA_KeyValueMap_clear(&readerGroupConfig->groupProperties);
UA_String_clear(&readerGroupConfig->securityGroupId);
UA_ExtensionObject_clear(&readerGroupConfig->transportSettings);
}
#ifdef UA_ENABLE_PUBSUB_SKS
static UA_StatusCode
readerGroupAttachSKSKeystorage(UA_PubSubManager *psm, UA_ReaderGroup *rg) {
if(UA_String_isEmpty(&rg->config.securityGroupId) || !rg->config.securityPolicy)
return UA_STATUSCODE_GOOD;
if(rg->keyStorage)
return UA_STATUSCODE_GOOD;
rg->keyStorage = UA_PubSubKeyStorage_find(psm, rg->config.securityGroupId);
if(rg->keyStorage) {
rg->keyStorage->referenceCount++;
return UA_STATUSCODE_GOOD;
}
rg->keyStorage = (UA_PubSubKeyStorage *)UA_calloc(1, sizeof(UA_PubSubKeyStorage));
if(!rg->keyStorage)
return UA_STATUSCODE_BADOUTOFMEMORY;
UA_StatusCode res =
UA_PubSubKeyStorage_init(psm, rg->keyStorage, &rg->config.securityGroupId,
rg->config.securityPolicy, 0, 0);
if(res != UA_STATUSCODE_GOOD) {
UA_PubSubKeyStorage_delete(psm, rg->keyStorage);
rg->keyStorage = NULL;
return res;
}
rg->keyStorage->referenceCount++;
return UA_STATUSCODE_GOOD;
}
#endif
UA_StatusCode
UA_ReaderGroup_create(UA_PubSubManager *psm, UA_NodeId connectionId,
const UA_ReaderGroupConfig *rgc,
UA_NodeId *readerGroupId) {
if(!psm || !rgc)
return UA_STATUSCODE_BADINVALIDARGUMENT;
UA_PubSubConnection *c = UA_PubSubConnection_find(psm, connectionId);
if(!c)
return UA_STATUSCODE_BADNOTFOUND;
UA_ReaderGroup *newGroup = (UA_ReaderGroup *)UA_calloc(1, sizeof(UA_ReaderGroup));
if(!newGroup)
return UA_STATUSCODE_BADOUTOFMEMORY;
newGroup->head.componentType = UA_PUBSUBCOMPONENT_READERGROUP;
newGroup->linkedConnection = c;
UA_StatusCode retval = UA_ReaderGroupConfig_copy(rgc, &newGroup->config);
if(retval != UA_STATUSCODE_GOOD) {
UA_free(newGroup);
return retval;
}
LIST_INSERT_HEAD(&c->readerGroups, newGroup, listEntry);
c->readerGroupsSize++;
char tmpLogIdStr[128];
mp_snprintf(tmpLogIdStr, 128, "%SReaderGroup %N\t| ",
c->head.logIdString, newGroup->head.identifier);
newGroup->head.logIdString = UA_STRING_ALLOC(tmpLogIdStr);
retval = UA_ReaderGroup_connect(psm, newGroup, true);
if(retval != UA_STATUSCODE_GOOD) {
UA_LOG_ERROR_PUBSUB(psm->logging, newGroup,
"Could not validate the connection parameters");
UA_PubSubComponent_freeWithoutLifecycleCallback(
psm, newGroup, UA_PUBSUBCOMPONENT_READERGROUP);
return retval;
}
#ifdef UA_ENABLE_PUBSUB_SKS
if(rgc->securityMode == UA_MESSAGESECURITYMODE_SIGN ||
rgc->securityMode == UA_MESSAGESECURITYMODE_SIGNANDENCRYPT) {
retval = readerGroupAttachSKSKeystorage(psm, newGroup);
if(retval != UA_STATUSCODE_GOOD) {
UA_LOG_ERROR_PUBSUB(psm->logging, newGroup,
"Attaching the SKS KeyStorage failed");
UA_PubSubComponent_freeWithoutLifecycleCallback(
psm, newGroup, UA_PUBSUBCOMPONENT_READERGROUP);
return retval;
}
}
#endif
#ifdef UA_ENABLE_PUBSUB_INFORMATIONMODEL
retval |= addReaderGroupRepresentation(psm->sc.server, newGroup);
if(retval != UA_STATUSCODE_GOOD) {
UA_PubSubComponent_freeWithoutLifecycleCallback(
psm, newGroup, UA_PUBSUBCOMPONENT_READERGROUP);
return retval;
}
#else
UA_PubSubManager_generateUniqueNodeId(psm, &newGroup->head.identifier);
#endif
UA_Server *server = psm->sc.server;
if(server->config.pubSubConfig.componentLifecycleCallback) {
retval = server->config.pubSubConfig.
componentLifecycleCallback(server, newGroup->head.identifier,
UA_PUBSUBCOMPONENT_READERGROUP, false);
if(retval != UA_STATUSCODE_GOOD) {
UA_PubSubComponent_freeWithoutLifecycleCallback(
psm, newGroup, UA_PUBSUBCOMPONENT_READERGROUP);
return retval;
}
}
UA_LOG_INFO_PUBSUB(psm->logging, newGroup, "ReaderGroup created (State: %s)",
UA_PubSubState_name(newGroup->head.state));
if(rgc->enabled)
UA_PubSubConnection_setPubSubState(psm, c, c->head.state);
if(readerGroupId)
UA_NodeId_copy(&newGroup->head.identifier, readerGroupId);
if(rgc->enabled)
UA_ReaderGroup_setPubSubState(psm, newGroup, UA_PUBSUBSTATE_OPERATIONAL);
return UA_STATUSCODE_GOOD;
}
UA_StatusCode
UA_ReaderGroup_remove(UA_PubSubManager *psm, UA_ReaderGroup *rg) {
UA_LOCK_ASSERT(&psm->sc.server->serviceMutex);
UA_PubSubConnection *connection = rg->linkedConnection;
UA_assert(connection);
UA_Server *server = psm->sc.server;
if(server->config.pubSubConfig.componentLifecycleCallback) {
UA_StatusCode res = server->config.pubSubConfig.
componentLifecycleCallback(server, rg->head.identifier,
UA_PUBSUBCOMPONENT_READERGROUP, true);
if(res != UA_STATUSCODE_GOOD)
return res;
}
rg->deleteFlag = true;
UA_ReaderGroup_setPubSubState(psm, rg, UA_PUBSUBSTATE_DISABLED);
UA_DataSetReader *dsr, *tmp_dsr;
LIST_FOREACH_SAFE(dsr, &rg->readers, listEntry, tmp_dsr) {
UA_DataSetReader_remove(psm, dsr);
}
if(rg->config.securityPolicy && rg->securityPolicyContext) {
UA_PubSubSecurityPolicy *sp = rg->config.securityPolicy;
sp->deleteGroupContext(sp, rg->securityPolicyContext);
rg->securityPolicyContext = NULL;
}
#ifdef UA_ENABLE_PUBSUB_SKS
if(rg->keyStorage) {
UA_PubSubKeyStorage_detachKeyStorage(psm, rg->keyStorage);
rg->keyStorage = NULL;
}
#endif
if(rg->recvChannelsSize == 0) {
LIST_REMOVE(rg, listEntry);
connection->readerGroupsSize--;
rg->linkedConnection = NULL;
#ifdef UA_ENABLE_PUBSUB_INFORMATIONMODEL
deleteNode(psm->sc.server, rg->head.identifier, true);
#endif
UA_LOG_INFO_PUBSUB(psm->logging, rg, "ReaderGroup deleted");
UA_ReaderGroupConfig_clear(&rg->config);
UA_PubSubComponentHead_clear(&rg->head);
UA_free(rg);
}
UA_PubSubConnection_setPubSubState(psm, connection, connection->head.state);
return UA_STATUSCODE_GOOD;
}
UA_StatusCode
UA_ReaderGroup_setPubSubState(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_PubSubState targetState) {
UA_LOCK_ASSERT(&psm->sc.server->serviceMutex);
if(rg->deleteFlag && targetState != UA_PUBSUBSTATE_DISABLED) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"The ReaderGroup is being deleted. Can only be disabled.");
return UA_STATUSCODE_BADINTERNALERROR;
}
UA_Server *server = psm->sc.server;
if(server->config.pubSubConfig.beforeStateChangeCallback) {
server->config.pubSubConfig.
beforeStateChangeCallback(server, rg->head.identifier, &targetState);
}
UA_StatusCode ret = UA_STATUSCODE_GOOD;
UA_PubSubState oldState = rg->head.state;
UA_PubSubConnection *connection = rg->linkedConnection;
UA_Boolean isTransient = rg->head.transientState;
rg->head.transientState = true;
if(rg->config.customStateMachine) {
ret = rg->config.customStateMachine(server, rg->head.identifier, rg->config.context,
&rg->head.state, targetState);
goto finalize_state_machine;
}
switch(targetState) {
case UA_PUBSUBSTATE_DISABLED:
case UA_PUBSUBSTATE_ERROR:
rg->head.state = targetState;
UA_ReaderGroup_disconnect(rg);
rg->hasReceived = false;
break;
case UA_PUBSUBSTATE_PAUSED:
case UA_PUBSUBSTATE_PREOPERATIONAL:
case UA_PUBSUBSTATE_OPERATIONAL:
if(psm->sc.state != UA_LIFECYCLESTATE_STARTED) {
if(oldState != UA_PUBSUBSTATE_PAUSED) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"Cannot enable the ReaderGroup while the "
"server is not running -> Paused State");
}
rg->head.state = UA_PUBSUBSTATE_PAUSED;
UA_ReaderGroup_disconnect(rg);
break;
}
if(connection->head.state != UA_PUBSUBSTATE_OPERATIONAL) {
UA_ReaderGroup_disconnect(rg);
rg->head.state = UA_PUBSUBSTATE_PAUSED;
break;
}
if(UA_ReaderGroup_canConnect(rg))
ret = UA_ReaderGroup_connect(psm, rg, false);
rg->head.state = (rg->hasReceived) ?
UA_PUBSUBSTATE_OPERATIONAL : UA_PUBSUBSTATE_PREOPERATIONAL;
break;
default:
ret = UA_STATUSCODE_BADINTERNALERROR;
break;
}
if(ret != UA_STATUSCODE_GOOD) {
rg->head.state = UA_PUBSUBSTATE_ERROR;
UA_ReaderGroup_disconnect(rg);
rg->hasReceived = false;
}
finalize_state_machine:
rg->head.transientState = isTransient;
if(rg->head.transientState)
return ret;
if(rg->head.state == oldState)
return ret;
UA_LOG_INFO_PUBSUB(psm->logging, rg, "%s -> %s",
UA_PubSubState_name(oldState),
UA_PubSubState_name(rg->head.state));
if(server->config.pubSubConfig.stateChangeCallback)
server->config.pubSubConfig.
stateChangeCallback(server, rg->head.identifier, rg->head.state, ret);
UA_DataSetReader *dsr;
LIST_FOREACH(dsr, &rg->readers, listEntry) {
if(psm->pubSubInitialSetupMode && dsr->config.enabled) {
UA_DataSetReader_setPubSubState(psm, dsr, UA_PUBSUBSTATE_PREOPERATIONAL, UA_STATUSCODE_GOOD);
} else {
UA_DataSetReader_setPubSubState(psm, dsr, dsr->head.state, UA_STATUSCODE_GOOD);
}
}
UA_PubSubManager_setState(psm, psm->sc.state);
return ret;
}
UA_StatusCode
UA_ReaderGroup_setEncryptionKeys(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_UInt32 securityTokenId,
const UA_ByteString signingKey,
const UA_ByteString encryptingKey,
const UA_ByteString keyNonce) {
if(rg->config.encodingMimeType == UA_PUBSUB_ENCODING_JSON) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"JSON encoding is enabled. The message security is "
"only defined for the UADP message mapping.");
return UA_STATUSCODE_BADINTERNALERROR;
}
UA_PubSubSecurityPolicy *sp = rg->config.securityPolicy;
if(!sp) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"No SecurityPolicy configured for the ReaderGroup");
return UA_STATUSCODE_BADINTERNALERROR;
}
if(securityTokenId != rg->securityTokenId) {
rg->securityTokenId = securityTokenId;
rg->nonceSequenceNumber = 1;
}
if(!rg->securityPolicyContext) {
return sp->newGroupContext(sp, &signingKey, &encryptingKey, &keyNonce,
&rg->securityPolicyContext);
}
return sp->setSecurityKeys(sp, rg->securityPolicyContext, &signingKey,
&encryptingKey, &keyNonce);
}
UA_Boolean
UA_ReaderGroup_process(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_NetworkMessage *nm) {
if(rg->head.state != UA_PUBSUBSTATE_OPERATIONAL &&
rg->head.state != UA_PUBSUBSTATE_PREOPERATIONAL)
return false;
rg->hasReceived = true;
UA_ReaderGroup_setPubSubState(psm, rg, rg->head.state);
UA_Boolean processed = false;
UA_DataSetReader *reader, *reader_tmp;
LIST_FOREACH_SAFE(reader, &rg->readers, listEntry, reader_tmp) {
if(reader->head.state != UA_PUBSUBSTATE_OPERATIONAL &&
reader->head.state != UA_PUBSUBSTATE_PREOPERATIONAL)
continue;
UA_StatusCode res = UA_DataSetReader_checkIdentifier(psm, reader, nm);
if(res != UA_STATUSCODE_GOOD)
continue;
if(!rg->hasReceived) {
rg->hasReceived = true;
UA_ReaderGroup_setPubSubState(psm, rg, rg->head.state);
}
processed = true;
UA_LOG_TRACE_PUBSUB(psm->logging, rg, "Processing a NetworkMessage");
if(!nm->payloadHeaderEnabled) {
UA_DataSetReader_process(psm, reader, nm->payload.dataSetMessages);
continue;
}
for(size_t i = 0; i < nm->messageCount; i++) {
if(reader->config.dataSetWriterId == nm->dataSetWriterIds[i])
UA_DataSetReader_process(psm, reader, &nm->payload.dataSetMessages[i]);
}
}
return processed;
}
UA_StatusCode
UA_ReaderGroup_decodeNetworkMessage(UA_PubSubManager *psm,
UA_ReaderGroup *rg,
UA_ByteString buffer,
UA_NetworkMessage *nm) {
PubSubDecodeCtx ctx;
memset(&ctx, 0, sizeof(PubSubDecodeCtx));
ctx.ctx.pos = buffer.data;
ctx.ctx.end = buffer.data + buffer.length;
ctx.ctx.opts.customTypes = psm->sc.server->config.customDataTypes;
UA_StatusCode rv = UA_NetworkMessage_decodeHeaders(&ctx, nm);
if(rv != UA_STATUSCODE_GOOD) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"PubSub receive. decoding headers failed");
UA_NetworkMessage_clear(nm);
return rv;
}
UA_DataSetReader *dsr;
LIST_FOREACH(dsr, &rg->readers, listEntry) {
rv = UA_DataSetReader_checkIdentifier(psm, dsr, nm);
if(rv == UA_STATUSCODE_GOOD)
break;
}
if(!dsr) {
UA_NetworkMessage_clear(nm);
return UA_STATUSCODE_BADNOTFOUND;
}
size_t i = 0;
UA_STACKARRAY(UA_DataSetMessage_EncodingMetaData, emd, rg->readersCount);
memset(emd, 0, sizeof(UA_DataSetMessage_EncodingMetaData) * rg->readersCount);
ctx.eo.metaData = emd;
ctx.eo.metaDataSize = rg->readersCount;
LIST_FOREACH(dsr, &rg->readers, listEntry) {
emd[i].dataSetWriterId = dsr->config.dataSetWriterId;
emd[i].fields = dsr->config.dataSetMetaData.fields;
emd[i].fieldsSize = dsr->config.dataSetMetaData.fieldsSize;
i++;
}
if(!nm->payloadHeaderEnabled) {
rv = UA_NetworkMessage_makeSyntheticPayloadHeader(&ctx.eo, nm);
if(rv != UA_STATUSCODE_GOOD) {
UA_NetworkMessage_clear(nm);
return rv;
}
}
rv = verifyAndDecryptNetworkMessage(psm->logging, buffer, &ctx.ctx, nm, rg);
if(rv != UA_STATUSCODE_GOOD) {
UA_NetworkMessage_clear(nm);
return rv;
}
rv = UA_NetworkMessage_decodePayload(&ctx, nm);
if(rv != UA_STATUSCODE_GOOD) {
UA_NetworkMessage_clear(nm);
return rv;
}
rv = UA_NetworkMessage_decodeFooters(&ctx, nm);
if(rv != UA_STATUSCODE_GOOD) {
UA_NetworkMessage_clear(nm);
return rv;
}
return UA_STATUSCODE_GOOD;
}
#ifdef UA_ENABLE_JSON_ENCODING
UA_StatusCode
UA_ReaderGroup_decodeNetworkMessageJSON(UA_PubSubManager *psm,
UA_ReaderGroup *rg,
UA_ByteString buffer,
UA_NetworkMessage *nm) {
UA_DecodeJsonOptions jo;
memset(&jo, 0, sizeof(jo));
jo.customTypes = psm->sc.server->config.customDataTypes;
UA_NetworkMessage_EncodingOptions eo;
size_t i = 0;
UA_STACKARRAY(UA_DataSetMessage_EncodingMetaData, emd, rg->readersCount);
memset(emd, 0, sizeof(UA_DataSetMessage_EncodingMetaData) * rg->readersCount);
eo.metaData = emd;
eo.metaDataSize = rg->readersCount;
UA_DataSetReader *dsr;
LIST_FOREACH(dsr, &rg->readers, listEntry) {
emd[i].dataSetWriterId = dsr->config.dataSetWriterId;
emd[i].fields = dsr->config.dataSetMetaData.fields;
emd[i].fieldsSize = dsr->config.dataSetMetaData.fieldsSize;
i++;
}
return UA_NetworkMessage_decodeJson(&buffer, nm, &eo, &jo);
}
#endif
static UA_StatusCode
needsDecryption(const UA_Logger *logger,
const UA_NetworkMessage *networkMessage,
const UA_MessageSecurityMode securityMode,
UA_Boolean *doDecrypt) {
UA_StatusCode retval = UA_STATUSCODE_GOOD;
UA_Boolean requiresEncryption = securityMode > UA_MESSAGESECURITYMODE_SIGN;
UA_Boolean isEncrypted = networkMessage->securityHeader.networkMessageEncrypted;
if(isEncrypted && requiresEncryption) {
*doDecrypt = true;
} else if(!isEncrypted && !requiresEncryption) {
*doDecrypt = false;
} else {
if(isEncrypted) {
UA_LOG_ERROR(logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. "
"Message is encrypted but ReaderGroup does not expect encryption");
retval = UA_STATUSCODE_BADSECURITYMODEINSUFFICIENT;
} else {
UA_LOG_ERROR(logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. "
"Message is not encrypted but ReaderGroup requires encryption");
retval = UA_STATUSCODE_BADSECURITYMODEREJECTED;
}
}
return retval;
}
static UA_StatusCode
needsValidation(const UA_Logger *logger,
const UA_NetworkMessage *networkMessage,
const UA_MessageSecurityMode securityMode,
UA_Boolean *doValidate) {
UA_StatusCode retval = UA_STATUSCODE_GOOD;
UA_Boolean isSigned = networkMessage->securityHeader.networkMessageSigned;
UA_Boolean requiresSignature = securityMode > UA_MESSAGESECURITYMODE_NONE;
if(isSigned &&
requiresSignature) {
*doValidate = true;
} else if(!isSigned && !requiresSignature) {
*doValidate = false;
} else {
if(isSigned) {
UA_LOG_ERROR(logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. "
"Message is signed but ReaderGroup does not expect signatures");
retval = UA_STATUSCODE_BADSECURITYMODEINSUFFICIENT;
} else {
UA_LOG_ERROR(logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. "
"Message is not signed but ReaderGroup requires signature");
retval = UA_STATUSCODE_BADSECURITYMODEREJECTED;
}
}
return retval;
}
UA_StatusCode
verifyAndDecryptNetworkMessage(const UA_Logger *logger, UA_ByteString buffer,
Ctx *ctx, UA_NetworkMessage *nm, UA_ReaderGroup *rg) {
UA_MessageSecurityMode securityMode = rg->config.securityMode;
UA_Boolean doValidate = false;
UA_Boolean doDecrypt = false;
UA_StatusCode rv = needsValidation(logger, nm, securityMode, &doValidate);
UA_CHECK_STATUS_WARN(rv, return rv, logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. Validation security mode error");
rv = needsDecryption(logger, nm, securityMode, &doDecrypt);
UA_CHECK_STATUS_WARN(rv, return rv, logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. Decryption security mode error");
if(!doValidate && !doDecrypt)
return UA_STATUSCODE_GOOD;
UA_PubSubSecurityPolicy *sp = rg->config.securityPolicy;
UA_CHECK_MEM_ERROR(sp, return UA_STATUSCODE_BADINVALIDARGUMENT,
logger, UA_LOGCATEGORY_PUBSUB,
"PubSub receive. securityPolicy must be set when security mode"
"is enabled to sign and/or encrypt");
void *cc = rg->securityPolicyContext;
UA_CHECK_MEM_ERROR(cc, return UA_STATUSCODE_BADINVALIDARGUMENT,
logger, UA_LOGCATEGORY_PUBSUB,
"PubSub receive. securityPolicyContext must be initialized "
"when security mode is enabled to sign and/or encrypt");
if(doValidate) {
size_t sigSize = sp->getSignatureSize(sp, cc);
if(buffer.length < sigSize) {
UA_LOG_WARNING(logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. Message too short for signature");
return UA_STATUSCODE_BADSECURITYCHECKSFAILED;
}
UA_ByteString toBeVerified = {buffer.length - sigSize, buffer.data};
UA_ByteString signature = {sigSize, buffer.data + buffer.length - sigSize};
rv = sp->verify(sp, cc, &toBeVerified, &signature);
UA_CHECK_STATUS_WARN(rv, return rv, logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. Signature invalid");
ctx->end -= sigSize;
}
if(doDecrypt) {
const UA_ByteString nonce = {
(size_t)nm->securityHeader.messageNonceSize,
(UA_Byte*)(uintptr_t)nm->securityHeader.messageNonce
};
rv = sp->setMessageNonce(sp, cc, &nonce);
UA_CHECK_STATUS_WARN(rv, return rv, logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. Faulty Nonce set");
UA_ByteString toBeDecrypted = {(uintptr_t)(ctx->end - ctx->pos), ctx->pos};
rv = sp->decrypt(sp, cc, &toBeDecrypted);
UA_CHECK_STATUS_WARN(rv, return rv, logger, UA_LOGCATEGORY_SECURITYPOLICY,
"PubSub receive. Faulty Decryption");
}
return UA_STATUSCODE_GOOD;
}
static UA_StatusCode
UA_ReaderGroup_connectMQTT(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_Boolean validate);
typedef struct {
UA_String profileURI;
UA_String protocol;
UA_Boolean json;
UA_StatusCode (*connectReaderGroup)(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_Boolean validate);
} ReaderGroupProfileMapping;
static ReaderGroupProfileMapping readerGroupProfiles[UA_PUBSUB_PROFILES_SIZE] = {
{UA_STRING_STATIC("http://opcfoundation.org/UA-Profile/Transport/pubsub-udp-uadp"),
UA_STRING_STATIC("udp"), false, NULL},
{UA_STRING_STATIC("http://opcfoundation.org/UA-Profile/Transport/pubsub-mqtt-uadp"),
UA_STRING_STATIC("mqtt"), false, UA_ReaderGroup_connectMQTT},
{UA_STRING_STATIC("http://opcfoundation.org/UA-Profile/Transport/pubsub-mqtt-json"),
UA_STRING_STATIC("mqtt"), true, UA_ReaderGroup_connectMQTT},
{UA_STRING_STATIC("http://opcfoundation.org/UA-Profile/Transport/pubsub-eth-uadp"),
UA_STRING_STATIC("eth"), false, NULL}
};
static void
UA_ReaderGroup_detachConnection(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_ConnectionManager *cm, uintptr_t connectionId) {
for(size_t i = 0; i < UA_PUBSUB_MAXCHANNELS; i++) {
if(rg->recvChannels[i] != connectionId)
continue;
UA_LOG_INFO_PUBSUB(psm->logging, rg, "Detach receive-connection %S %u",
cm->protocol, (unsigned)connectionId);
rg->recvChannels[i] = 0;
rg->recvChannelsSize--;
return;
}
}
static UA_StatusCode
UA_ReaderGroup_attachRecvConnection(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_ConnectionManager *cm, uintptr_t connectionId) {
for(size_t i = 0; i < UA_PUBSUB_MAXCHANNELS; i++) {
if(rg->recvChannels[i] == connectionId)
return UA_STATUSCODE_GOOD;
}
if(rg->recvChannelsSize >= UA_PUBSUB_MAXCHANNELS)
return UA_STATUSCODE_BADINTERNALERROR;
for(size_t i = 0; i < UA_PUBSUB_MAXCHANNELS; i++) {
if(rg->recvChannels[i] != 0)
continue;
UA_LOG_INFO_PUBSUB(psm->logging, rg, "Attach receive-connection %S %u",
cm->protocol, (unsigned)connectionId);
rg->recvChannels[i] = connectionId;
rg->recvChannelsSize++;
break;
}
return UA_STATUSCODE_GOOD;
}
static void
ReaderGroupChannelCallback(UA_ConnectionManager *cm, uintptr_t connectionId,
void *application, void **connectionContext,
UA_ConnectionState state, const UA_KeyValueMap *params,
UA_ByteString msg) {
if(!connectionContext)
return;
UA_ReaderGroup *rg = (UA_ReaderGroup*)*connectionContext;
UA_PubSubManager *psm = (UA_PubSubManager*)application;
UA_Server *server = psm->sc.server;
lockServer(server);
if(state == UA_CONNECTIONSTATE_CLOSING) {
UA_ReaderGroup_detachConnection(psm, rg, cm, connectionId);
if(rg->deleteFlag && rg->recvChannelsSize == 0) {
UA_ReaderGroup_remove(psm, rg);
unlockServer(server);
return;
}
UA_ReaderGroup_setPubSubState(psm, rg, rg->head.state);
UA_PubSubManager_setState(psm, psm->sc.state);
unlockServer(server);
return;
}
UA_StatusCode res = UA_ReaderGroup_attachRecvConnection(psm, rg, cm, connectionId);
if(res != UA_STATUSCODE_GOOD) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"No more space for an additional EventLoop connection");
UA_PubSubConnection *c = rg->linkedConnection;
if(c && c->cm)
c->cm->closeConnection(c->cm, connectionId);
unlockServer(server);
return;
}
UA_ReaderGroup_setPubSubState(psm, rg, rg->head.state);
if(msg.length == 0) {
unlockServer(server);
return;
}
if (rg->head.state != UA_PUBSUBSTATE_OPERATIONAL &&
rg->head.state != UA_PUBSUBSTATE_PREOPERATIONAL) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"Received a message for a disabled ReaderGroup");
unlockServer(server);
return;
}
UA_NetworkMessage nm;
memset(&nm, 0, sizeof(UA_NetworkMessage));
if(rg->config.encodingMimeType == UA_PUBSUB_ENCODING_UADP) {
res = UA_ReaderGroup_decodeNetworkMessage(psm, rg, msg, &nm);
} else {
#ifdef UA_ENABLE_JSON_ENCODING
res = UA_ReaderGroup_decodeNetworkMessageJSON(psm, rg, msg, &nm);
#else
res = UA_STATUSCODE_BADNOTSUPPORTED;
#endif
}
if(res != UA_STATUSCODE_GOOD) {
UA_LOG_WARNING_PUBSUB(psm->logging, rg,
"Verify, decrypt and decode network message failed");
unlockServer(server);
return;
}
UA_ReaderGroup_process(psm, rg, &nm);
UA_NetworkMessage_clear(&nm);
unlockServer(server);
}
static UA_StatusCode
UA_ReaderGroup_connectMQTT(UA_PubSubManager *psm, UA_ReaderGroup *rg,
UA_Boolean validate) {
UA_LOCK_ASSERT(&psm->sc.server->serviceMutex);
UA_PubSubConnection *c = rg->linkedConnection;
UA_NetworkAddressUrlDataType *addressUrl = (UA_NetworkAddressUrlDataType*)
c->config.address.data;
UA_ExtensionObject *ts = &rg->config.transportSettings;
if((ts->encoding != UA_EXTENSIONOBJECT_DECODED &&
ts->encoding != UA_EXTENSIONOBJECT_DECODED_NODELETE) ||
ts->content.decoded.type !=
&UA_TYPES[UA_TYPES_BROKERDATASETREADERTRANSPORTDATATYPE]) {
UA_LOG_ERROR_PUBSUB(psm->logging, rg,
"Wrong TransportSettings type for MQTT");
return UA_STATUSCODE_BADINTERNALERROR;
}
UA_BrokerDataSetReaderTransportDataType *transportSettings =
(UA_BrokerDataSetReaderTransportDataType*)ts->content.decoded.data;
UA_String address;
UA_UInt16 port = 1883;
UA_StatusCode res = UA_parseEndpointUrl(&addressUrl->url, &address, &port, NULL);
if(res != UA_STATUSCODE_GOOD) {
UA_LOG_ERROR_PUBSUB(psm->logging, c, "Could not parse the MQTT network URL");
return res;
}
UA_Boolean listen = true;
UA_KeyValuePair kvp[5];
UA_KeyValueMap kvm = {5, kvp};
kvp[0].key = UA_QUALIFIEDNAME(0, "address");
UA_Variant_setScalar(&kvp[0].value, &address, &UA_TYPES[UA_TYPES_STRING]);
kvp[1].key = UA_QUALIFIEDNAME(0, "subscribe");
UA_Variant_setScalar(&kvp[1].value, &listen, &UA_TYPES[UA_TYPES_BOOLEAN]);
kvp[2].key = UA_QUALIFIEDNAME(0, "port");
UA_Variant_setScalar(&kvp[2].value, &port, &UA_TYPES[UA_TYPES_UINT16]);
kvp[3].key = UA_QUALIFIEDNAME(0, "topic");
UA_Variant_setScalar(&kvp[3].value, &transportSettings->queueName,
&UA_TYPES[UA_TYPES_STRING]);
kvp[4].key = UA_QUALIFIEDNAME(0, "validate");
UA_Variant_setScalar(&kvp[4].value, &validate, &UA_TYPES[UA_TYPES_BOOLEAN]);
res = c->cm->openConnection(c->cm, &kvm, psm, rg, ReaderGroupChannelCallback);
if(res != UA_STATUSCODE_GOOD) {
UA_LOG_ERROR_PUBSUB(psm->logging, rg, "Could not open the MQTT connection");
}
return res;
}
void
UA_ReaderGroup_disconnect(UA_ReaderGroup *rg) {
UA_PubSubConnection *c = rg->linkedConnection;
if(!c)
return;
for(size_t i = 0; i < UA_PUBSUB_MAXCHANNELS; i++) {
if(rg->recvChannels[i] != 0)
c->cm->closeConnection(c->cm, rg->recvChannels[i]);
}
}
UA_Boolean
UA_ReaderGroup_canConnect(UA_ReaderGroup *rg) {
return rg->recvChannelsSize == 0;
}
UA_StatusCode
UA_ReaderGroup_connect(UA_PubSubManager *psm, UA_ReaderGroup *rg, UA_Boolean validate) {
UA_Server *server = psm->sc.server;
UA_LOCK_ASSERT(&server->serviceMutex);
if(rg->config.transportSettings.encoding == UA_EXTENSIONOBJECT_ENCODED_NOBODY)
return UA_STATUSCODE_GOOD;
UA_EventLoop *el = psm->sc.server->config.eventLoop;
if(!el) {
UA_LOG_ERROR_PUBSUB(server->config.logging, rg, "No EventLoop configured");
return UA_STATUSCODE_BADINTERNALERROR;
}
UA_PubSubConnection *c = rg->linkedConnection;
if(!c)
return UA_STATUSCODE_BADINTERNALERROR;
ReaderGroupProfileMapping *profile = NULL;
for(size_t i = 0; i < UA_PUBSUB_PROFILES_SIZE; i++) {
if(!UA_String_equal(&c->config.transportProfileUri,
&readerGroupProfiles[i].profileURI))
continue;
profile = &readerGroupProfiles[i];
break;
}
UA_ConnectionManager *cm = (profile) ? getCM(el, profile->protocol) : NULL;
if(!cm || (c->cm && cm != c->cm)) {
UA_LOG_ERROR_PUBSUB(psm->logging, c,
"The requested profile \"%S\"is not supported",
c->config.transportProfileUri);
return UA_STATUSCODE_BADINTERNALERROR;
}
c->cm = cm;
c->json = profile->json;
return (profile->connectReaderGroup) ?
profile->connectReaderGroup(psm, rg, validate) : UA_STATUSCODE_GOOD;
}
UA_StatusCode
UA_Server_addReaderGroup(UA_Server *server, const UA_NodeId connectionIdentifier,
const UA_ReaderGroupConfig *readerGroupConfig,
UA_NodeId *readerGroupIdentifier) {
if(!server || !readerGroupConfig)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_PubSubManager *psm = getPSM(server);
UA_StatusCode res =
UA_ReaderGroup_create(psm, connectionIdentifier,
readerGroupConfig, readerGroupIdentifier);
unlockServer(server);
return res;
}
UA_StatusCode
UA_Server_removeReaderGroup(UA_Server *server, const UA_NodeId groupIdentifier) {
if(!server)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_StatusCode res = UA_STATUSCODE_GOOD;
UA_PubSubManager *psm = getPSM(server);
UA_ReaderGroup *rg = UA_ReaderGroup_find(psm, groupIdentifier);
if(rg)
UA_ReaderGroup_remove(psm, rg);
else
res = UA_STATUSCODE_BADNOTFOUND;
unlockServer(server);
return res;
}
UA_StatusCode
UA_Server_getReaderGroupConfig(UA_Server *server, const UA_NodeId rgId,
UA_ReaderGroupConfig *config) {
if(!server || !config)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_ReaderGroup *rg = UA_ReaderGroup_find(getPSM(server), rgId);
UA_StatusCode ret = (rg) ?
UA_ReaderGroupConfig_copy(&rg->config, config) : UA_STATUSCODE_BADNOTFOUND;
unlockServer(server);
return ret;
}
UA_StatusCode
UA_Server_getReaderGroupState(UA_Server *server, const UA_NodeId rgId,
UA_PubSubState *state) {
if(!server || !state)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_StatusCode ret = UA_STATUSCODE_BADNOTFOUND;
UA_ReaderGroup *rg = UA_ReaderGroup_find(getPSM(server), rgId);
if(rg) {
*state = rg->head.state;
ret = UA_STATUSCODE_GOOD;
}
unlockServer(server);
return ret;
}
#ifdef UA_ENABLE_PUBSUB_SKS
UA_StatusCode
UA_Server_setReaderGroupActivateKey(UA_Server *server,
const UA_NodeId readerGroupId) {
if(!server)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_PubSubManager *psm = getPSM(server);
UA_ReaderGroup *rg = UA_ReaderGroup_find(psm, readerGroupId);
if(!rg || !rg->keyStorage || !rg->keyStorage->currentItem) {
unlockServer(server);
return UA_STATUSCODE_BADNOTFOUND;
}
UA_StatusCode ret =
UA_PubSubKeyStorage_activateKeyToChannelContext(psm, rg->head.identifier,
rg->config.securityGroupId);
unlockServer(server);
return ret;
}
#endif
UA_StatusCode
UA_Server_enableReaderGroup(UA_Server *server, const UA_NodeId readerGroupId){
if(!server)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_PubSubManager *psm = getPSM(server);
UA_ReaderGroup *rg = UA_ReaderGroup_find(psm, readerGroupId);
UA_StatusCode ret = (rg) ?
UA_ReaderGroup_setPubSubState(psm, rg, UA_PUBSUBSTATE_OPERATIONAL) :
UA_STATUSCODE_BADNOTFOUND;
unlockServer(server);
return ret;
}
UA_StatusCode
UA_Server_disableReaderGroup(UA_Server *server, const UA_NodeId readerGroupId){
if(!server)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_PubSubManager *psm = getPSM(server);
UA_ReaderGroup *rg = UA_ReaderGroup_find(psm, readerGroupId);
UA_StatusCode ret = (rg) ?
UA_ReaderGroup_setPubSubState(psm, rg, UA_PUBSUBSTATE_DISABLED) :
UA_STATUSCODE_BADNOTFOUND;
unlockServer(server);
return ret;
}
UA_StatusCode
UA_Server_setReaderGroupEncryptionKeys(UA_Server *server,
const UA_NodeId readerGroup,
UA_UInt32 securityTokenId,
const UA_ByteString signingKey,
const UA_ByteString encryptingKey,
const UA_ByteString keyNonce) {
if(!server)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_PubSubManager *psm = getPSM(server);
UA_ReaderGroup *rg = UA_ReaderGroup_find(getPSM(server), readerGroup);
UA_StatusCode res = (rg) ?
UA_ReaderGroup_setEncryptionKeys(psm, rg, securityTokenId, signingKey,
encryptingKey, keyNonce) : UA_STATUSCODE_BADNOTFOUND;
unlockServer(server);
return res;
}
UA_StatusCode
UA_Server_updateReaderGroupConfig(UA_Server *server, const UA_NodeId rgId,
const UA_ReaderGroupConfig *config) {
if(!server || !config)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockServer(server);
UA_PubSubManager *psm = getPSM(server);
UA_ReaderGroup *rg = UA_ReaderGroup_find(getPSM(server), rgId);
if(!rg) {
unlockServer(server);
return UA_STATUSCODE_BADNOTFOUND;
}
if(UA_PubSubState_isEnabled(rg->head.state)) {
UA_LOG_ERROR_PUBSUB(psm->logging, rg,
"The ReaderGroup must be disabled to update the config");
unlockServer(server);
return UA_STATUSCODE_BADINTERNALERROR;
}
UA_ReaderGroupConfig oldConfig = rg->config;
UA_StatusCode retval = UA_ReaderGroupConfig_copy(config, &rg->config);
if(retval != UA_STATUSCODE_GOOD) {
unlockServer(server);
return retval;
}
retval = UA_ReaderGroup_connect(psm, rg, true);
if(retval != UA_STATUSCODE_GOOD) {
UA_LOG_ERROR_PUBSUB(psm->logging, rg,
"Could not validate the connection parameters");
goto errout;
}
#ifdef UA_ENABLE_PUBSUB_SKS
if(!UA_String_equal(&rg->config.securityGroupId, &oldConfig.securityGroupId) ||
rg->config.securityMode != oldConfig.securityMode) {
if(rg->keyStorage) {
UA_PubSubKeyStorage_detachKeyStorage(psm, rg->keyStorage);
rg->keyStorage = NULL;
}
if(rg->config.securityMode == UA_MESSAGESECURITYMODE_SIGN ||
rg->config.securityMode == UA_MESSAGESECURITYMODE_SIGNANDENCRYPT) {
retval = readerGroupAttachSKSKeystorage(psm, rg);
if(retval != UA_STATUSCODE_GOOD) {
UA_LOG_ERROR_PUBSUB(psm->logging, rg,
"Attaching the SKS KeyStorage failed");
goto errout;
}
}
}
#endif
UA_ReaderGroupConfig_clear(&oldConfig);
unlockServer(server);
return UA_STATUSCODE_GOOD;
errout:
UA_ReaderGroupConfig_clear(&rg->config);
rg->config = oldConfig;
unlockServer(server);
return retval;
}
#endif