|
@@ -43,6 +43,8 @@ UA_Client_Subscriptions_new(UA_Client *client, UA_SubscriptionSettings settings,
|
|
|
}
|
|
|
|
|
|
LIST_INIT(&newSub->monitoredItems);
|
|
|
+ newSub->sequenceNumber = 0;
|
|
|
+ newSub->lastActivity = UA_DateTime_nowMonotonic();
|
|
|
newSub->lifeTime = response.revisedLifetimeCount;
|
|
|
newSub->keepAliveCount = response.revisedMaxKeepAliveCount;
|
|
|
newSub->publishingInterval = response.revisedPublishingInterval;
|
|
@@ -58,8 +60,8 @@ UA_Client_Subscriptions_new(UA_Client *client, UA_SubscriptionSettings settings,
|
|
|
return UA_STATUSCODE_GOOD;
|
|
|
}
|
|
|
|
|
|
-static UA_Client_Subscription *findSubscription(const UA_Client *client, UA_UInt32 subscriptionId)
|
|
|
-{
|
|
|
+static UA_Client_Subscription *
|
|
|
+findSubscription(const UA_Client *client, UA_UInt32 subscriptionId) {
|
|
|
UA_Client_Subscription *sub = NULL;
|
|
|
LIST_FOREACH(sub, &client->subscriptions, listEntry) {
|
|
|
if(sub->subscriptionID == subscriptionId)
|
|
@@ -68,22 +70,15 @@ static UA_Client_Subscription *findSubscription(const UA_Client *client, UA_UInt
|
|
|
return sub;
|
|
|
}
|
|
|
|
|
|
-/* remove the subscription remotely */
|
|
|
UA_StatusCode
|
|
|
UA_Client_Subscriptions_remove(UA_Client *client, UA_UInt32 subscriptionId) {
|
|
|
UA_Client_Subscription *sub = findSubscription(client, subscriptionId);
|
|
|
if(!sub)
|
|
|
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
|
|
|
|
|
|
- UA_StatusCode retval = UA_STATUSCODE_GOOD;
|
|
|
- UA_Client_MonitoredItem *mon, *tmpmon;
|
|
|
- LIST_FOREACH_SAFE(mon, &sub->monitoredItems, listEntry, tmpmon) {
|
|
|
- retval =
|
|
|
- UA_Client_Subscriptions_removeMonitoredItem(client, sub->subscriptionID,
|
|
|
- mon->monitoredItemId);
|
|
|
- if(retval != UA_STATUSCODE_GOOD)
|
|
|
- return retval;
|
|
|
- }
|
|
|
+ /* remove the subscription from the list */
|
|
|
+ /* will be reinserted after if error occurs */
|
|
|
+ LIST_REMOVE(sub, listEntry);
|
|
|
|
|
|
/* remove the subscription remotely */
|
|
|
UA_DeleteSubscriptionsRequest request;
|
|
@@ -91,19 +86,30 @@ UA_Client_Subscriptions_remove(UA_Client *client, UA_UInt32 subscriptionId) {
|
|
|
request.subscriptionIdsSize = 1;
|
|
|
request.subscriptionIds = &sub->subscriptionID;
|
|
|
UA_DeleteSubscriptionsResponse response = UA_Client_Service_deleteSubscriptions(client, request);
|
|
|
- retval = response.responseHeader.serviceResult;
|
|
|
+
|
|
|
+ UA_StatusCode retval = response.responseHeader.serviceResult;
|
|
|
if(retval == UA_STATUSCODE_GOOD && response.resultsSize > 0)
|
|
|
retval = response.results[0];
|
|
|
UA_DeleteSubscriptionsResponse_deleteMembers(&response);
|
|
|
|
|
|
if(retval != UA_STATUSCODE_GOOD && retval != UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID) {
|
|
|
+ /* error occurs re-insert the subscription in the list */
|
|
|
+ LIST_INSERT_HEAD(&client->subscriptions, sub, listEntry);
|
|
|
UA_LOG_INFO(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
"Could not remove subscription %u with error code %s",
|
|
|
sub->subscriptionID, UA_StatusCode_name(retval));
|
|
|
return retval;
|
|
|
}
|
|
|
|
|
|
- UA_Client_Subscriptions_forceDelete(client, sub);
|
|
|
+ /* remove the monitored items locally */
|
|
|
+ UA_Client_MonitoredItem *mon, *tmpmon;
|
|
|
+ LIST_FOREACH_SAFE(mon, &sub->monitoredItems, listEntry, tmpmon) {
|
|
|
+ UA_NodeId_deleteMembers(&mon->monitoredNodeId);
|
|
|
+ LIST_REMOVE(mon, listEntry);
|
|
|
+ UA_free(mon);
|
|
|
+ }
|
|
|
+ UA_free(sub);
|
|
|
+
|
|
|
return UA_STATUSCODE_GOOD;
|
|
|
}
|
|
|
|
|
@@ -178,6 +184,7 @@ addMonitoredItems(UA_Client *client, const UA_UInt32 subscriptionId,
|
|
|
|
|
|
itemResults[i] = result->statusCode;
|
|
|
if(result->statusCode != UA_STATUSCODE_GOOD) {
|
|
|
+ newMonitoredItemIds[i] = 0;
|
|
|
UA_free(newMon);
|
|
|
continue;
|
|
|
}
|
|
@@ -213,7 +220,6 @@ addMonitoredItems(UA_Client *client, const UA_UInt32 subscriptionId,
|
|
|
return retval;
|
|
|
}
|
|
|
|
|
|
-
|
|
|
UA_StatusCode
|
|
|
UA_Client_Subscriptions_addMonitoredItems(UA_Client *client, const UA_UInt32 subscriptionId,
|
|
|
UA_MonitoredItemCreateRequest *items, size_t itemsSize,
|
|
@@ -245,7 +251,6 @@ UA_Client_Subscriptions_addMonitoredItem(UA_Client *client, UA_UInt32 subscripti
|
|
|
return retval | retval_item;
|
|
|
}
|
|
|
|
|
|
-
|
|
|
UA_StatusCode
|
|
|
UA_Client_Subscriptions_addMonitoredEvents(UA_Client *client, const UA_UInt32 subscriptionId,
|
|
|
UA_MonitoredItemCreateRequest *items, size_t itemsSize,
|
|
@@ -359,123 +364,224 @@ UA_Client_Subscriptions_removeMonitoredItem(UA_Client *client, UA_UInt32 subscri
|
|
|
return retval | retval_item;
|
|
|
}
|
|
|
|
|
|
+/* Assume the request is already initialized */
|
|
|
+static UA_StatusCode
|
|
|
+UA_Client_preparePublishRequest(UA_Client *client, UA_PublishRequest *request) {
|
|
|
+ /* Count acks */
|
|
|
+ UA_Client_NotificationsAckNumber *ack;
|
|
|
+ LIST_FOREACH(ack, &client->pendingNotificationsAcks, listEntry)
|
|
|
+ ++request->subscriptionAcknowledgementsSize;
|
|
|
+
|
|
|
+ /* Create the array. Returns a sentinel pointer if the length is zero. */
|
|
|
+ request->subscriptionAcknowledgements = (UA_SubscriptionAcknowledgement*)
|
|
|
+ UA_Array_new(request->subscriptionAcknowledgementsSize,
|
|
|
+ &UA_TYPES[UA_TYPES_SUBSCRIPTIONACKNOWLEDGEMENT]);
|
|
|
+ if(!request->subscriptionAcknowledgements) {
|
|
|
+ request->subscriptionAcknowledgementsSize = 0;
|
|
|
+ return UA_STATUSCODE_BADOUTOFMEMORY;
|
|
|
+ }
|
|
|
+
|
|
|
+ size_t i = 0;
|
|
|
+ UA_Client_NotificationsAckNumber *ack_tmp;
|
|
|
+ LIST_FOREACH_SAFE(ack, &client->pendingNotificationsAcks, listEntry, ack_tmp) {
|
|
|
+ request->subscriptionAcknowledgements[i].sequenceNumber = ack->subAck.sequenceNumber;
|
|
|
+ request->subscriptionAcknowledgements[i].subscriptionId = ack->subAck.subscriptionId;
|
|
|
+ ++i;
|
|
|
+ LIST_REMOVE(ack, listEntry);
|
|
|
+ UA_free(ack);
|
|
|
+ }
|
|
|
+ return UA_STATUSCODE_GOOD;
|
|
|
+}
|
|
|
+
|
|
|
+/* According to OPC Unified Architecture, Part 4 5.13.1.1 i) */
|
|
|
+/* The value 0 is never used for the sequence number */
|
|
|
+static UA_UInt32
|
|
|
+UA_Client_Subscriptions_nextSequenceNumber(UA_UInt32 sequenceNumber) {
|
|
|
+ UA_UInt32 nextSequenceNumber = sequenceNumber + 1;
|
|
|
+ if(nextSequenceNumber == 0)
|
|
|
+ nextSequenceNumber = 1;
|
|
|
+ return nextSequenceNumber;
|
|
|
+}
|
|
|
+
|
|
|
static void
|
|
|
-UA_Client_processPublishResponse(UA_Client *client, UA_PublishRequest *request,
|
|
|
- UA_PublishResponse *response) {
|
|
|
- if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD)
|
|
|
- return;
|
|
|
+processDataChangeNotification(UA_Client *client, UA_Client_Subscription *sub,
|
|
|
+ UA_DataChangeNotification *dataChangeNotification) {
|
|
|
+ for(size_t j = 0; j < dataChangeNotification->monitoredItemsSize; ++j) {
|
|
|
+ UA_MonitoredItemNotification *mitemNot = &dataChangeNotification->monitoredItems[j];
|
|
|
|
|
|
- UA_Client_Subscription *sub = findSubscription(client, response->subscriptionId);
|
|
|
- if(!sub)
|
|
|
- return;
|
|
|
+ /* Find the MonitoredItem */
|
|
|
+ UA_Client_MonitoredItem *mon;
|
|
|
+ LIST_FOREACH(mon, &sub->monitoredItems, listEntry) {
|
|
|
+ if(mon->clientHandle == mitemNot->clientHandle)
|
|
|
+ break;
|
|
|
+ }
|
|
|
|
|
|
- /* Check if the server has acknowledged any of the sent ACKs */
|
|
|
- for(size_t i = 0; i < response->resultsSize && i < request->subscriptionAcknowledgementsSize; ++i) {
|
|
|
- /* remove also acks that are unknown to the server */
|
|
|
- if(response->results[i] != UA_STATUSCODE_GOOD &&
|
|
|
- response->results[i] != UA_STATUSCODE_BADSEQUENCENUMBERUNKNOWN)
|
|
|
+ if(!mon) {
|
|
|
+ UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Could not process a notification with clienthandle %u on subscription %u",
|
|
|
+ mitemNot->clientHandle, sub->subscriptionID);
|
|
|
continue;
|
|
|
+ }
|
|
|
|
|
|
- /* Remove the ack from the list */
|
|
|
- UA_SubscriptionAcknowledgement *orig_ack = &request->subscriptionAcknowledgements[i];
|
|
|
- UA_Client_NotificationsAckNumber *ack;
|
|
|
- LIST_FOREACH(ack, &client->pendingNotificationsAcks, listEntry) {
|
|
|
- if(ack->subAck.subscriptionId == orig_ack->subscriptionId &&
|
|
|
- ack->subAck.sequenceNumber == orig_ack->sequenceNumber) {
|
|
|
- LIST_REMOVE(ack, listEntry);
|
|
|
- UA_free(ack);
|
|
|
- UA_assert(ack != LIST_FIRST(&client->pendingNotificationsAcks));
|
|
|
- break;
|
|
|
- }
|
|
|
+ if(mon->isEventMonitoredItem) {
|
|
|
+ UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "MonitoredItem is configured for Events. But received a "
|
|
|
+ "DataChangeNotification.");
|
|
|
+ continue;
|
|
|
}
|
|
|
+
|
|
|
+ mon->handler.dataChangeHandler(client, mon->monitoredItemId,
|
|
|
+ &mitemNot->value, mon->handlerContext);
|
|
|
}
|
|
|
+}
|
|
|
|
|
|
- /* Process the notification messages */
|
|
|
- UA_NotificationMessage *msg = &response->notificationMessage;
|
|
|
- for(size_t k = 0; k < msg->notificationDataSize; ++k) {
|
|
|
- if(msg->notificationData[k].encoding != UA_EXTENSIONOBJECT_DECODED)
|
|
|
- continue;
|
|
|
+static void
|
|
|
+processEventNotification(UA_Client *client, UA_Client_Subscription *sub,
|
|
|
+ UA_EventNotificationList *eventNotificationList) {
|
|
|
+ for(size_t j = 0; j < eventNotificationList->eventsSize; ++j) {
|
|
|
+ UA_EventFieldList *eventFieldList = &eventNotificationList->events[j];
|
|
|
|
|
|
- /* Handle DataChangeNotification */
|
|
|
- if(msg->notificationData[k].content.decoded.type == &UA_TYPES[UA_TYPES_DATACHANGENOTIFICATION]) {
|
|
|
- UA_DataChangeNotification *dataChangeNotification =
|
|
|
- (UA_DataChangeNotification *)msg->notificationData[k].content.decoded.data;
|
|
|
- for(size_t j = 0; j < dataChangeNotification->monitoredItemsSize; ++j) {
|
|
|
- UA_MonitoredItemNotification *mitemNot = &dataChangeNotification->monitoredItems[j];
|
|
|
-
|
|
|
- /* Find the MonitoredItem */
|
|
|
- UA_Client_MonitoredItem *mon;
|
|
|
- LIST_FOREACH(mon, &sub->monitoredItems, listEntry) {
|
|
|
- if(mon->clientHandle == mitemNot->clientHandle)
|
|
|
- break;
|
|
|
- }
|
|
|
-
|
|
|
- if(!mon) {
|
|
|
- UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
- "Could not process a notification with clienthandle %u on subscription %u",
|
|
|
- mitemNot->clientHandle, sub->subscriptionID);
|
|
|
- continue;
|
|
|
- }
|
|
|
-
|
|
|
- if(mon->isEventMonitoredItem) {
|
|
|
- UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
- "MonitoredItem is configured for Events. But received a "
|
|
|
- "DataChangeNotification.");
|
|
|
- continue;
|
|
|
- }
|
|
|
-
|
|
|
- mon->handler.dataChangeHandler(client, mon->monitoredItemId,
|
|
|
- &mitemNot->value, mon->handlerContext);
|
|
|
- }
|
|
|
+ /* Find the MonitoredItem */
|
|
|
+ UA_Client_MonitoredItem *mon;
|
|
|
+ LIST_FOREACH(mon, &sub->monitoredItems, listEntry) {
|
|
|
+ if(mon->clientHandle == eventFieldList->clientHandle)
|
|
|
+ break;
|
|
|
+ }
|
|
|
+
|
|
|
+ if(!mon) {
|
|
|
+ UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Could not process a notification with clienthandle %u on subscription %u",
|
|
|
+ eventFieldList->clientHandle, sub->subscriptionID);
|
|
|
continue;
|
|
|
}
|
|
|
|
|
|
- /* Handle EventNotification */
|
|
|
- if(msg->notificationData[k].content.decoded.type == &UA_TYPES[UA_TYPES_EVENTNOTIFICATIONLIST]) {
|
|
|
- UA_EventNotificationList *eventNotificationList =
|
|
|
- (UA_EventNotificationList *)msg->notificationData[k].content.decoded.data;
|
|
|
- for(size_t j = 0; j < eventNotificationList->eventsSize; ++j) {
|
|
|
- UA_EventFieldList *eventFieldList = &eventNotificationList->events[j];
|
|
|
-
|
|
|
- /* Find the MonitoredItem */
|
|
|
- UA_Client_MonitoredItem *mon;
|
|
|
- LIST_FOREACH(mon, &sub->monitoredItems, listEntry) {
|
|
|
- if(mon->clientHandle == eventFieldList->clientHandle)
|
|
|
- break;
|
|
|
- }
|
|
|
-
|
|
|
- if(!mon) {
|
|
|
- UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
- "Could not process a notification with clienthandle %u on subscription %u",
|
|
|
- eventFieldList->clientHandle, sub->subscriptionID);
|
|
|
- continue;
|
|
|
- }
|
|
|
-
|
|
|
- if(!mon->isEventMonitoredItem) {
|
|
|
- UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
- "MonitoredItem is configured for DataChanges. But received a "
|
|
|
- "EventNotification.");
|
|
|
- continue;
|
|
|
- }
|
|
|
-
|
|
|
- mon->handler.eventHandler(client, mon->monitoredItemId, eventFieldList->eventFieldsSize,
|
|
|
- eventFieldList->eventFields, mon->handlerContext);
|
|
|
- }
|
|
|
+ if(!mon->isEventMonitoredItem) {
|
|
|
+ UA_LOG_DEBUG(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "MonitoredItem is configured for DataChanges. But received a "
|
|
|
+ "EventNotification.");
|
|
|
+ continue;
|
|
|
}
|
|
|
+
|
|
|
+ mon->handler.eventHandler(client, mon->monitoredItemId, eventFieldList->eventFieldsSize,
|
|
|
+ eventFieldList->eventFields, mon->handlerContext);
|
|
|
}
|
|
|
+}
|
|
|
|
|
|
- /* Add to the list of pending acks */
|
|
|
- UA_Client_NotificationsAckNumber *tmpAck =
|
|
|
- (UA_Client_NotificationsAckNumber*)UA_malloc(sizeof(UA_Client_NotificationsAckNumber));
|
|
|
- if(!tmpAck) {
|
|
|
- UA_LOG_WARNING(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
- "Not enough memory to store the acknowledgement for a publish "
|
|
|
- "message on subscription %u", sub->subscriptionID);
|
|
|
+
|
|
|
+static void
|
|
|
+processNotificationMessage(UA_Client *client, UA_Client_Subscription *sub,
|
|
|
+ UA_ExtensionObject *msg) {
|
|
|
+ if(msg->encoding != UA_EXTENSIONOBJECT_DECODED)
|
|
|
+ return;
|
|
|
+
|
|
|
+ /* Handle DataChangeNotification */
|
|
|
+ if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_DATACHANGENOTIFICATION]) {
|
|
|
+ UA_DataChangeNotification *dataChangeNotification =
|
|
|
+ (UA_DataChangeNotification *)msg->content.decoded.data;
|
|
|
+ processDataChangeNotification(client, sub, dataChangeNotification);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ /* Handle EventNotification */
|
|
|
+ if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_EVENTNOTIFICATIONLIST]) {
|
|
|
+ UA_EventNotificationList *eventNotificationList =
|
|
|
+ (UA_EventNotificationList *)msg->content.decoded.data;
|
|
|
+ processEventNotification(client, sub, eventNotificationList);
|
|
|
return;
|
|
|
}
|
|
|
- tmpAck->subAck.sequenceNumber = msg->sequenceNumber;
|
|
|
- tmpAck->subAck.subscriptionId = sub->subscriptionID;
|
|
|
- LIST_INSERT_HEAD(&client->pendingNotificationsAcks, tmpAck, listEntry);
|
|
|
+}
|
|
|
+
|
|
|
+static void
|
|
|
+processPublishResponse(UA_Client *client, UA_PublishRequest *request,
|
|
|
+ UA_PublishResponse *response) {
|
|
|
+ UA_NotificationMessage *msg = &response->notificationMessage;
|
|
|
+
|
|
|
+ client->currentlyOutStandingPublishRequests--;
|
|
|
+
|
|
|
+ if(response->responseHeader.serviceResult == UA_STATUSCODE_BADTOOMANYPUBLISHREQUESTS){
|
|
|
+ if(client->config.outStandingPublishRequests > 1){
|
|
|
+ client->config.outStandingPublishRequests--;
|
|
|
+ UA_LOG_WARNING(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Too many publishrequest, we reduce outStandingPublishRequests to %d",
|
|
|
+ client->config.outStandingPublishRequests);
|
|
|
+ } else {
|
|
|
+ UA_LOG_ERROR(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Too many publishrequest when outStandingPublishRequests = 1");
|
|
|
+ UA_Client_close(client);
|
|
|
+ }
|
|
|
+ goto cleanup;
|
|
|
+ }
|
|
|
+
|
|
|
+ if(response->responseHeader.serviceResult == UA_STATUSCODE_BADNOSUBSCRIPTION){
|
|
|
+ if(LIST_FIRST(&client->subscriptions)){
|
|
|
+ UA_Client_close(client);
|
|
|
+ UA_LOG_ERROR(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "PublishRequest error : No subscription");
|
|
|
+ }
|
|
|
+ goto cleanup;
|
|
|
+ }
|
|
|
+
|
|
|
+ if(response->responseHeader.serviceResult == UA_STATUSCODE_BADSESSIONIDINVALID){
|
|
|
+ UA_Client_close(client);
|
|
|
+ UA_LOG_ERROR(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Received BadSessionIdInvalid");
|
|
|
+ goto cleanup;
|
|
|
+ }
|
|
|
+
|
|
|
+ UA_Client_Subscription *sub = findSubscription(client, response->subscriptionId);
|
|
|
+ if(!sub)
|
|
|
+ goto cleanup;
|
|
|
+ sub->lastActivity = UA_DateTime_nowMonotonic();
|
|
|
+
|
|
|
+
|
|
|
+ if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD)
|
|
|
+ goto cleanup;
|
|
|
+
|
|
|
+ /* Detect missing message - OPC Unified Architecture, Part 4 5.13.1.1 e) */
|
|
|
+ if((sub->sequenceNumber != msg->sequenceNumber) &&
|
|
|
+ (UA_Client_Subscriptions_nextSequenceNumber(sub->sequenceNumber) != msg->sequenceNumber)) {
|
|
|
+ UA_LOG_ERROR(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Invalid subscritpion sequenceNumber");
|
|
|
+ UA_Client_close(client);
|
|
|
+ goto cleanup;
|
|
|
+ }
|
|
|
+ sub->sequenceNumber = msg->sequenceNumber;
|
|
|
+
|
|
|
+ /* Process the notification messages */
|
|
|
+ for(size_t k = 0; k < msg->notificationDataSize; ++k)
|
|
|
+ processNotificationMessage(client, sub, &msg->notificationData[k]);
|
|
|
+
|
|
|
+ /* Add to the list of pending acks */
|
|
|
+ for(size_t i = 0; i < response->availableSequenceNumbersSize; i++) {
|
|
|
+ if(response->availableSequenceNumbers[i] != msg->sequenceNumber)
|
|
|
+ continue;
|
|
|
+ UA_Client_NotificationsAckNumber *tmpAck = (UA_Client_NotificationsAckNumber*)
|
|
|
+ UA_malloc(sizeof(UA_Client_NotificationsAckNumber));
|
|
|
+ if(!tmpAck) {
|
|
|
+ UA_LOG_WARNING(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Not enough memory to store the acknowledgement for a publish "
|
|
|
+ "message on subscription %u", sub->subscriptionID);
|
|
|
+ break;
|
|
|
+ }
|
|
|
+ tmpAck->subAck.sequenceNumber = msg->sequenceNumber;
|
|
|
+ tmpAck->subAck.subscriptionId = sub->subscriptionID;
|
|
|
+ LIST_INSERT_HEAD(&client->pendingNotificationsAcks, tmpAck, listEntry);
|
|
|
+ break;
|
|
|
+ }
|
|
|
+
|
|
|
+ cleanup:
|
|
|
+ UA_PublishRequest_deleteMembers(request);
|
|
|
+}
|
|
|
+
|
|
|
+static void
|
|
|
+processPublishResponseAsync(UA_Client *client, void *userdata, UA_UInt32 requestId,
|
|
|
+ void *response, const UA_DataType *responseType) {
|
|
|
+ UA_PublishRequest *req = (UA_PublishRequest*)userdata;
|
|
|
+ UA_PublishResponse *res = (UA_PublishResponse*)response;
|
|
|
+ processPublishResponse(client, req, res);
|
|
|
+ UA_PublishRequest_delete(req);
|
|
|
+ /* Fill up the outstanding publish requests */
|
|
|
+ UA_Client_Subscriptions_backgroundPublish(client);
|
|
|
}
|
|
|
|
|
|
UA_StatusCode
|
|
@@ -492,28 +598,20 @@ UA_Client_Subscriptions_manuallySendPublishRequest(UA_Client *client) {
|
|
|
while(moreNotifications) {
|
|
|
UA_PublishRequest request;
|
|
|
UA_PublishRequest_init(&request);
|
|
|
- request.subscriptionAcknowledgementsSize = 0;
|
|
|
-
|
|
|
- UA_Client_NotificationsAckNumber *ack;
|
|
|
- LIST_FOREACH(ack, &client->pendingNotificationsAcks, listEntry)
|
|
|
- ++request.subscriptionAcknowledgementsSize;
|
|
|
- if(request.subscriptionAcknowledgementsSize > 0) {
|
|
|
- request.subscriptionAcknowledgements = (UA_SubscriptionAcknowledgement*)
|
|
|
- UA_malloc(sizeof(UA_SubscriptionAcknowledgement) * request.subscriptionAcknowledgementsSize);
|
|
|
- if(!request.subscriptionAcknowledgements)
|
|
|
- return UA_STATUSCODE_BADOUTOFMEMORY;
|
|
|
- }
|
|
|
+ retval = UA_Client_preparePublishRequest(client, &request);
|
|
|
+ if(retval != UA_STATUSCODE_GOOD)
|
|
|
+ return retval;
|
|
|
|
|
|
- int i = 0;
|
|
|
- LIST_FOREACH(ack, &client->pendingNotificationsAcks, listEntry) {
|
|
|
- request.subscriptionAcknowledgements[i].sequenceNumber = ack->subAck.sequenceNumber;
|
|
|
- request.subscriptionAcknowledgements[i].subscriptionId = ack->subAck.subscriptionId;
|
|
|
- ++i;
|
|
|
- }
|
|
|
+ /* Manually increase the number of sent publish requests. Otherwise we
|
|
|
+ * send out one too many when we process async responses when we wait
|
|
|
+ * for the correct publish response. The
|
|
|
+ * currentlyOutStandingPublishRequests will be reduced during processing
|
|
|
+ * of the response. */
|
|
|
+ client->currentlyOutStandingPublishRequests++;
|
|
|
|
|
|
UA_PublishResponse response = UA_Client_Service_publish(client, request);
|
|
|
- UA_Client_processPublishResponse(client, &request, &response);
|
|
|
-
|
|
|
+ processPublishResponse(client, &request, &response);
|
|
|
+
|
|
|
now = UA_DateTime_nowMonotonic();
|
|
|
if(now > maxDate) {
|
|
|
moreNotifications = UA_FALSE;
|
|
@@ -545,4 +643,52 @@ UA_Client_Subscriptions_clean(UA_Client *client) {
|
|
|
UA_Client_Subscriptions_forceDelete(client, sub); /* force local removal */
|
|
|
}
|
|
|
|
|
|
+UA_StatusCode
|
|
|
+UA_Client_Subscriptions_backgroundPublish(UA_Client *client) {
|
|
|
+ if(client->state < UA_CLIENTSTATE_SESSION)
|
|
|
+ return UA_STATUSCODE_BADSERVERNOTCONNECTED;
|
|
|
+
|
|
|
+ /* The session must have at least one subscription */
|
|
|
+ if(!LIST_FIRST(&client->subscriptions))
|
|
|
+ return UA_STATUSCODE_GOOD;
|
|
|
+
|
|
|
+ while(client->currentlyOutStandingPublishRequests < client->config.outStandingPublishRequests) {
|
|
|
+ UA_PublishRequest *request = UA_PublishRequest_new();
|
|
|
+ if (!request)
|
|
|
+ return UA_STATUSCODE_BADOUTOFMEMORY;
|
|
|
+
|
|
|
+ UA_StatusCode retval = UA_Client_preparePublishRequest(client, request);
|
|
|
+ if(retval != UA_STATUSCODE_GOOD) {
|
|
|
+ UA_PublishRequest_delete(request);
|
|
|
+ return retval;
|
|
|
+ }
|
|
|
+
|
|
|
+ UA_UInt32 requestId;
|
|
|
+ client->currentlyOutStandingPublishRequests++;
|
|
|
+ retval = __UA_Client_AsyncService(client, request, &UA_TYPES[UA_TYPES_PUBLISHREQUEST],
|
|
|
+ processPublishResponseAsync,
|
|
|
+ &UA_TYPES[UA_TYPES_PUBLISHRESPONSE],
|
|
|
+ (void*)request, &requestId);
|
|
|
+ if(retval != UA_STATUSCODE_GOOD) {
|
|
|
+ UA_PublishRequest_delete(request);
|
|
|
+ return retval;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ /* Check subscriptions inactivity */
|
|
|
+ UA_Client_Subscription *sub;
|
|
|
+ LIST_FOREACH(sub, &client->subscriptions, listEntry) {
|
|
|
+ if(((UA_DateTime)(sub->publishingInterval * sub->keepAliveCount + client->config.timeout) *
|
|
|
+ UA_DATETIME_MSEC + sub->lastActivity) < UA_DateTime_nowMonotonic()) {
|
|
|
+ UA_LOG_ERROR(client->config.logger, UA_LOGCATEGORY_CLIENT,
|
|
|
+ "Inactivity for Subscription %d. Closing the connection.",
|
|
|
+ sub->subscriptionID);
|
|
|
+ UA_Client_close(client);
|
|
|
+ return UA_STATUSCODE_BADCONNECTIONCLOSED;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ return UA_STATUSCODE_GOOD;
|
|
|
+}
|
|
|
+
|
|
|
#endif /* UA_ENABLE_SUBSCRIPTIONS */
|