| 11 | |
| 12 | |
| 13 | class SubscriptionService(object): |
| 14 | |
| 15 | def __init__(self, aspace): |
| 16 | self.logger = logging.getLogger(__name__) |
| 17 | self.loop = None |
| 18 | self.aspace = aspace |
| 19 | self.subscriptions = {} |
| 20 | self._sub_id_counter = 77 |
| 21 | self._lock = RLock() |
| 22 | |
| 23 | def set_loop(self, loop): |
| 24 | self.loop = loop |
| 25 | |
| 26 | def create_subscription(self, params, callback): |
| 27 | self.logger.info("create subscription with callback: %s", callback) |
| 28 | result = ua.CreateSubscriptionResult() |
| 29 | result.RevisedPublishingInterval = params.RequestedPublishingInterval |
| 30 | result.RevisedLifetimeCount = params.RequestedLifetimeCount |
| 31 | result.RevisedMaxKeepAliveCount = params.RequestedMaxKeepAliveCount |
| 32 | with self._lock: |
| 33 | self._sub_id_counter += 1 |
| 34 | result.SubscriptionId = self._sub_id_counter |
| 35 | |
| 36 | sub = InternalSubscription(self, result, self.aspace, callback) |
| 37 | sub.start() |
| 38 | self.subscriptions[result.SubscriptionId] = sub |
| 39 | |
| 40 | return result |
| 41 | |
| 42 | def modify_subscription(self, params, callback): |
| 43 | # Requested params are ignored, result = params set during create_subscription. |
| 44 | self.logger.info("modify subscription with callback: %s", callback) |
| 45 | result = ua.ModifySubscriptionResult() |
| 46 | try: |
| 47 | with self._lock: |
| 48 | sub = self.subscriptions[params.SubscriptionId] |
| 49 | result.RevisedPublishingInterval = sub.data.RevisedPublishingInterval |
| 50 | result.RevisedLifetimeCount = sub.data.RevisedLifetimeCount |
| 51 | result.RevisedMaxKeepAliveCount = sub.data.RevisedMaxKeepAliveCount |
| 52 | |
| 53 | return result |
| 54 | except KeyError: |
| 55 | raise utils.ServiceError(ua.StatusCodes.BadSubscriptionIdInvalid) |
| 56 | |
| 57 | def delete_subscriptions(self, ids): |
| 58 | self.logger.info("delete subscriptions: %s", ids) |
| 59 | res = [] |
| 60 | for i in ids: |
| 61 | with self._lock: |
| 62 | if i not in self.subscriptions: |
| 63 | res.append(ua.StatusCode(ua.StatusCodes.BadSubscriptionIdInvalid)) |
| 64 | else: |
| 65 | sub = self.subscriptions.pop(i) |
| 66 | sub.stop() |
| 67 | res.append(ua.StatusCode()) |
| 68 | return res |
| 69 | |
| 70 | def publish(self, acks): |