Fix router panic when trying to update route (#811)
This commit is contained in:
+56
-36
@@ -38,17 +38,18 @@ type HTTPTriggerSet struct {
|
|||||||
*functionServiceMap
|
*functionServiceMap
|
||||||
*mutableRouter
|
*mutableRouter
|
||||||
|
|
||||||
fissionClient *crd.FissionClient
|
fissionClient *crd.FissionClient
|
||||||
kubeClient *kubernetes.Clientset
|
kubeClient *kubernetes.Clientset
|
||||||
executor *executorClient.Client
|
executor *executorClient.Client
|
||||||
resolver *functionReferenceResolver
|
resolver *functionReferenceResolver
|
||||||
crdClient *rest.RESTClient
|
crdClient *rest.RESTClient
|
||||||
triggers []crd.HTTPTrigger
|
triggers []crd.HTTPTrigger
|
||||||
triggerStore k8sCache.Store
|
triggerStore k8sCache.Store
|
||||||
triggerController k8sCache.Controller
|
triggerController k8sCache.Controller
|
||||||
functions []crd.Function
|
functions []crd.Function
|
||||||
funcStore k8sCache.Store
|
funcStore k8sCache.Store
|
||||||
funcController k8sCache.Controller
|
funcController k8sCache.Controller
|
||||||
|
updateRouterRequestChannel chan struct{}
|
||||||
|
|
||||||
tsRoundTripperParams *tsRoundTripperParams
|
tsRoundTripperParams *tsRoundTripperParams
|
||||||
}
|
}
|
||||||
@@ -56,13 +57,14 @@ type HTTPTriggerSet struct {
|
|||||||
func makeHTTPTriggerSet(fmap *functionServiceMap, fissionClient *crd.FissionClient,
|
func makeHTTPTriggerSet(fmap *functionServiceMap, fissionClient *crd.FissionClient,
|
||||||
kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient *rest.RESTClient, params *tsRoundTripperParams) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) {
|
kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient *rest.RESTClient, params *tsRoundTripperParams) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) {
|
||||||
httpTriggerSet := &HTTPTriggerSet{
|
httpTriggerSet := &HTTPTriggerSet{
|
||||||
functionServiceMap: fmap,
|
functionServiceMap: fmap,
|
||||||
triggers: []crd.HTTPTrigger{},
|
triggers: []crd.HTTPTrigger{},
|
||||||
fissionClient: fissionClient,
|
fissionClient: fissionClient,
|
||||||
kubeClient: kubeClient,
|
kubeClient: kubeClient,
|
||||||
executor: executor,
|
executor: executor,
|
||||||
crdClient: crdClient,
|
crdClient: crdClient,
|
||||||
tsRoundTripperParams: params,
|
updateRouterRequestChannel: make(chan struct{}),
|
||||||
|
tsRoundTripperParams: params,
|
||||||
}
|
}
|
||||||
var tStore, fnStore k8sCache.Store
|
var tStore, fnStore k8sCache.Store
|
||||||
var tController, fnController k8sCache.Controller
|
var tController, fnController k8sCache.Controller
|
||||||
@@ -87,6 +89,7 @@ func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter
|
|||||||
log.Printf("Skipping continuous trigger updates")
|
log.Printf("Skipping continuous trigger updates")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
go ts.updateRouter()
|
||||||
go ts.runWatcher(ctx, ts.funcController)
|
go ts.runWatcher(ctx, ts.funcController)
|
||||||
go ts.runWatcher(ctx, ts.triggerController)
|
go ts.runWatcher(ctx, ts.triggerController)
|
||||||
}
|
}
|
||||||
@@ -192,6 +195,11 @@ func (ts *HTTPTriggerSet) initTriggerController() (k8sCache.Store, k8sCache.Cont
|
|||||||
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
|
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
|
||||||
oldTrigger := oldObj.(*crd.HTTPTrigger)
|
oldTrigger := oldObj.(*crd.HTTPTrigger)
|
||||||
newTrigger := newObj.(*crd.HTTPTrigger)
|
newTrigger := newObj.(*crd.HTTPTrigger)
|
||||||
|
|
||||||
|
if oldTrigger.Metadata.ResourceVersion == newTrigger.Metadata.ResourceVersion {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
go updateIngress(oldTrigger, newTrigger, ts.kubeClient)
|
go updateIngress(oldTrigger, newTrigger, ts.kubeClient)
|
||||||
ts.syncTriggers()
|
ts.syncTriggers()
|
||||||
},
|
},
|
||||||
@@ -211,7 +219,13 @@ func (ts *HTTPTriggerSet) initFunctionController() (k8sCache.Store, k8sCache.Con
|
|||||||
ts.syncTriggers()
|
ts.syncTriggers()
|
||||||
},
|
},
|
||||||
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
|
UpdateFunc: func(oldObj interface{}, newObj interface{}) {
|
||||||
|
oldFn := oldObj.(*crd.Function)
|
||||||
fn := newObj.(*crd.Function)
|
fn := newObj.(*crd.Function)
|
||||||
|
|
||||||
|
if oldFn.Metadata.ResourceVersion == fn.Metadata.ResourceVersion {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
// update resolver function reference cache
|
// update resolver function reference cache
|
||||||
for key, rr := range ts.resolver.copy() {
|
for key, rr := range ts.resolver.copy() {
|
||||||
if key.functionReference.Name == fn.Metadata.Name &&
|
if key.functionReference.Name == fn.Metadata.Name &&
|
||||||
@@ -236,22 +250,28 @@ func (ts *HTTPTriggerSet) runWatcher(ctx context.Context, controller k8sCache.Co
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (ts *HTTPTriggerSet) syncTriggers() {
|
func (ts *HTTPTriggerSet) syncTriggers() {
|
||||||
// get triggers
|
ts.updateRouterRequestChannel <- struct{}{}
|
||||||
latestTriggers := ts.triggerStore.List()
|
}
|
||||||
triggers := make([]crd.HTTPTrigger, len(latestTriggers))
|
|
||||||
for _, t := range latestTriggers {
|
func (ts *HTTPTriggerSet) updateRouter() {
|
||||||
triggers = append(triggers, *t.(*crd.HTTPTrigger))
|
for range ts.updateRouterRequestChannel {
|
||||||
}
|
// get triggers
|
||||||
ts.triggers = triggers
|
latestTriggers := ts.triggerStore.List()
|
||||||
|
triggers := make([]crd.HTTPTrigger, len(latestTriggers))
|
||||||
// get functions
|
for _, t := range latestTriggers {
|
||||||
latestFunctions := ts.funcStore.List()
|
triggers = append(triggers, *t.(*crd.HTTPTrigger))
|
||||||
functions := make([]crd.Function, len(latestFunctions))
|
}
|
||||||
for _, f := range latestFunctions {
|
ts.triggers = triggers
|
||||||
functions = append(functions, *f.(*crd.Function))
|
|
||||||
}
|
// get functions
|
||||||
ts.functions = functions
|
latestFunctions := ts.funcStore.List()
|
||||||
|
functions := make([]crd.Function, len(latestFunctions))
|
||||||
// make a new router and use it
|
for _, f := range latestFunctions {
|
||||||
ts.mutableRouter.updateRouter(ts.getRouter())
|
functions = append(functions, *f.(*crd.Function))
|
||||||
|
}
|
||||||
|
ts.functions = functions
|
||||||
|
|
||||||
|
// make a new router and use it
|
||||||
|
ts.mutableRouter.updateRouter(ts.getRouter())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user