X-Git-Url: https://gerrit.o-ran-sc.org/r/gitweb?a=blobdiff_plain;f=pkg%2Fcontrol%2Fcontrol.go;h=4978a8d8921ebbaa76be0874cd51b451b40c07d6;hb=3f7becc8a1e15c5b69dfd95d86b9b7308b3e1c5e;hp=51c84e3ceb811c267c81d4ba144e3ea729fbf2c2;hpb=b31b13f2e2603404d57c2f81260a6e2727aaeb97;p=ric-plt%2Fsubmgr.git diff --git a/pkg/control/control.go b/pkg/control/control.go index 51c84e3..4978a8d 100755 --- a/pkg/control/control.go +++ b/pkg/control/control.go @@ -103,7 +103,7 @@ func NewControl() *Control { //subscriber: subscriber, } c.XappWrapper.Init("") - go xapp.Subscription.Listen(c.SubscriptionHandler, c.QueryHandler) + go xapp.Subscription.Listen(c.SubscriptionHandler, c.QueryHandler, c.SubscriptionDeleteHandler) //go c.subscriber.Listen(c.SubscriptionHandler, c.QueryHandler) return c } @@ -122,7 +122,7 @@ func (c *Control) Run() { //------------------------------------------------------------------- // //------------------------------------------------------------------- -func (c *Control) SubscriptionHandler(stype models.SubscriptionType, params interface{}) (models.SubscriptionResult, error) { +func (c *Control) SubscriptionHandler(stype models.SubscriptionType, params interface{}) (*models.SubscriptionResponse, error) { /* switch p := params.(type) { case *models.ReportParams: @@ -136,7 +136,11 @@ func (c *Control) SubscriptionHandler(stype models.SubscriptionType, params inte case *models.PolicyParams: } */ - return models.SubscriptionResult{}, fmt.Errorf("Subscription rest interface not implemented") + return &models.SubscriptionResponse{}, fmt.Errorf("Subscription rest interface not implemented") +} + +func (c *Control) SubscriptionDeleteHandler(string) error { + return fmt.Errorf("Subscription rest interface not implemented") } func (c *Control) QueryHandler() (models.SubscriptionList, error) { @@ -150,7 +154,7 @@ func (c *Control) QueryHandler() (models.SubscriptionList, error) { func (c *Control) rmrSendToE2T(desc string, subs *Subscription, trans *TransactionSubs) (err error) { params := xapptweaks.NewParams(nil) params.Mtype = trans.GetMtype() - params.SubId = int(subs.GetReqId().Seq) + params.SubId = int(subs.GetReqId().InstanceId) params.Xid = "" params.Meid = subs.GetMeid() params.Src = "" @@ -165,7 +169,7 @@ func (c *Control) rmrSendToXapp(desc string, subs *Subscription, trans *Transact params := xapptweaks.NewParams(nil) params.Mtype = trans.GetMtype() - params.SubId = int(subs.GetReqId().Seq) + params.SubId = int(subs.GetReqId().InstanceId) params.Xid = trans.GetXid() params.Meid = trans.GetMeid() params.Src = "" @@ -218,7 +222,7 @@ func (c *Control) handleXAPPSubscriptionRequest(params *xapptweaks.RMRParams) { return } - trans := c.tracker.NewXappTransaction(xapptweaks.NewRmrEndpoint(params.Src), params.Xid, subReqMsg.RequestId.Seq, params.Meid) + trans := c.tracker.NewXappTransaction(xapptweaks.NewRmrEndpoint(params.Src), params.Xid, subReqMsg.RequestId.InstanceId, params.Meid) if trans == nil { xapp.Logger.Error("XAPP-SubReq: %s", idstring(fmt.Errorf("transaction not created"), params)) return @@ -231,6 +235,7 @@ func (c *Control) handleXAPPSubscriptionRequest(params *xapptweaks.RMRParams) { return } + //TODO handle subscription toward e2term inside AssignToSubscription / hide handleSubscriptionCreate in it? subs, err := c.registry.AssignToSubscription(trans, subReqMsg) if err != nil { xapp.Logger.Error("XAPP-SubReq: %s", idstring(err, trans)) @@ -263,7 +268,7 @@ func (c *Control) handleXAPPSubscriptionRequest(params *xapptweaks.RMRParams) { } } xapp.Logger.Info("XAPP-SubReq: failed %s", idstring(err, trans, subs)) - c.registry.RemoveFromSubscription(subs, trans, 5*time.Second) + //c.registry.RemoveFromSubscription(subs, trans, 5*time.Second) } //------------------------------------------------------------------- @@ -278,7 +283,7 @@ func (c *Control) handleXAPPSubscriptionDeleteRequest(params *xapptweaks.RMRPara return } - trans := c.tracker.NewXappTransaction(xapptweaks.NewRmrEndpoint(params.Src), params.Xid, subDelReqMsg.RequestId.Seq, params.Meid) + trans := c.tracker.NewXappTransaction(xapptweaks.NewRmrEndpoint(params.Src), params.Xid, subDelReqMsg.RequestId.InstanceId, params.Meid) if trans == nil { xapp.Logger.Error("XAPP-SubDelReq: %s", idstring(fmt.Errorf("transaction not created"), params)) return @@ -314,7 +319,8 @@ func (c *Control) handleXAPPSubscriptionDeleteRequest(params *xapptweaks.RMRPara c.rmrSendToXapp("", subs, trans) } - c.registry.RemoveFromSubscription(subs, trans, 5*time.Second) + //TODO handle subscription toward e2term insiged RemoveFromSubscription / hide handleSubscriptionDelete in it? + //c.registry.RemoveFromSubscription(subs, trans, 5*time.Second) } //------------------------------------------------------------------- @@ -369,6 +375,10 @@ func (c *Control) handleSubscriptionCreate(subs *Subscription, parentTrans *Tran xapp.Logger.Debug("SUBS-SubReq: Handling (cached response %s) %s", typeofSubsMessage(subRfMsg), idstring(nil, trans, subs, parentTrans)) } + //Now RemoveFromSubscription in here to avoid race conditions (mostly concerns delete) + if valid == false { + c.registry.RemoveFromSubscription(subs, parentTrans, 5*time.Second) + } parentTrans.SendEvent(subRfMsg, 0) } @@ -393,7 +403,10 @@ func (c *Control) handleSubscriptionDelete(subs *Subscription, parentTrans *Tran } else { subs.mutex.Unlock() } - + //Now RemoveFromSubscription in here to avoid race conditions (mostly concerns delete) + // If parallel deletes ongoing both might pass earlier sendE2TSubscriptionDeleteRequest(...) if + // RemoveFromSubscription locates in caller side (now in handleXAPPSubscriptionDeleteRequest(...)) + c.registry.RemoveFromSubscription(subs, parentTrans, 5*time.Second) parentTrans.SendEvent(nil, 0) } @@ -467,7 +480,7 @@ func (c *Control) handleE2TSubscriptionResponse(params *xapptweaks.RMRParams) { xapp.Logger.Error("MSG-SubResp %s", idstring(err, params)) return } - subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subRespMsg.RequestId.Seq}) + subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subRespMsg.RequestId.InstanceId}) if err != nil { xapp.Logger.Error("MSG-SubResp: %s", idstring(err, params)) return @@ -496,7 +509,7 @@ func (c *Control) handleE2TSubscriptionFailure(params *xapptweaks.RMRParams) { xapp.Logger.Error("MSG-SubFail %s", idstring(err, params)) return } - subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subFailMsg.RequestId.Seq}) + subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subFailMsg.RequestId.InstanceId}) if err != nil { xapp.Logger.Error("MSG-SubFail: %s", idstring(err, params)) return @@ -525,7 +538,7 @@ func (c *Control) handleE2TSubscriptionDeleteResponse(params *xapptweaks.RMRPara xapp.Logger.Error("MSG-SubDelResp: %s", idstring(err, params)) return } - subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subDelRespMsg.RequestId.Seq}) + subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subDelRespMsg.RequestId.InstanceId}) if err != nil { xapp.Logger.Error("MSG-SubDelResp: %s", idstring(err, params)) return @@ -554,7 +567,7 @@ func (c *Control) handleE2TSubscriptionDeleteFailure(params *xapptweaks.RMRParam xapp.Logger.Error("MSG-SubDelFail: %s", idstring(err, params)) return } - subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subDelFailMsg.RequestId.Seq}) + subs, err := c.registry.GetSubscriptionFirstMatch([]uint32{subDelFailMsg.RequestId.InstanceId}) if err != nil { xapp.Logger.Error("MSG-SubDelFail: %s", idstring(err, params)) return