Use latest function metadata to check cached function service. (#316)
Watch functions in the router and use that to trigger functionReferenceResolver cache invalidation. We might want to rate-limit syncTriggers in a future change, since we might be triggering it more often than needed.
This commit is contained in:
committed by
Soam Vasani
parent
45c766061a
commit
9d7bd49338
@@ -143,6 +143,7 @@ func (poolMgr *Poolmgr) getServiceForFunctionApi(w http.ResponseWriter, r *http.
|
||||
code, msg := fission.GetHTTPError(err)
|
||||
log.Printf("Error: %v: %v", code, msg)
|
||||
http.Error(w, msg, code)
|
||||
return
|
||||
}
|
||||
|
||||
w.Write([]byte(serviceName))
|
||||
|
||||
@@ -116,3 +116,21 @@ func (frr *functionReferenceResolver) resolveByName(namespace, name string) (*re
|
||||
}
|
||||
return &rr, nil
|
||||
}
|
||||
|
||||
func (frr *functionReferenceResolver) delete(namespace string, fr *fission.FunctionReference) error {
|
||||
nfr := namespacedFunctionReference{
|
||||
namespace: namespace,
|
||||
functionReference: *fr,
|
||||
}
|
||||
return frr.refCache.Delete(nfr)
|
||||
}
|
||||
|
||||
func (frr *functionReferenceResolver) copy() map[namespacedFunctionReference]resolveResult {
|
||||
cache := make(map[namespacedFunctionReference]resolveResult)
|
||||
for k, v := range frr.refCache.Copy() {
|
||||
key := k.(namespacedFunctionReference)
|
||||
val := v.(resolveResult)
|
||||
cache[key] = val
|
||||
}
|
||||
return cache
|
||||
}
|
||||
|
||||
@@ -61,6 +61,7 @@ func (ts *HTTPTriggerSet) subscribeRouter(mr *mutableRouter) {
|
||||
return
|
||||
}
|
||||
go ts.watchTriggers()
|
||||
go ts.watchFunctions()
|
||||
}
|
||||
|
||||
func defaultHomeHandler(w http.ResponseWriter, r *http.Request) {
|
||||
@@ -162,6 +163,48 @@ func (ts *HTTPTriggerSet) watchTriggers() {
|
||||
}
|
||||
}
|
||||
|
||||
func (ts *HTTPTriggerSet) watchFunctions() {
|
||||
rv := ""
|
||||
for {
|
||||
wi, err := ts.fissionClient.Functions(api.NamespaceAll).Watch(api.ListOptions{
|
||||
ResourceVersion: rv,
|
||||
})
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to watch function list: %v", err)
|
||||
}
|
||||
|
||||
for {
|
||||
ev, more := <-wi.ResultChan()
|
||||
if !more {
|
||||
// restart watch from last rv
|
||||
break
|
||||
}
|
||||
if ev.Type == watch.Error {
|
||||
// restart watch from the start
|
||||
rv = ""
|
||||
time.Sleep(time.Second)
|
||||
break
|
||||
}
|
||||
fn := ev.Object.(*tpr.Function)
|
||||
rv = fn.Metadata.ResourceVersion
|
||||
|
||||
// update resolver function reference cache
|
||||
for key, rr := range ts.resolver.copy() {
|
||||
if key.functionReference.Name == fn.Metadata.Name &&
|
||||
rr.functionMetadata.ResourceVersion != fn.Metadata.ResourceVersion {
|
||||
err := ts.resolver.delete(key.namespace, &key.functionReference)
|
||||
if err != nil {
|
||||
log.Printf("Error deleting functionReferenceResolver cache: %v", err)
|
||||
}
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
ts.syncTriggers()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (ts *HTTPTriggerSet) syncTriggers() {
|
||||
log.Printf("Syncing http triggers")
|
||||
|
||||
|
||||
Executable
+62
@@ -0,0 +1,62 @@
|
||||
#!/bin/bash
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
ROOT=$(dirname $0)/../..
|
||||
|
||||
fn=nodejs-hello-$(date +%s)
|
||||
|
||||
# Create a function in nodejs, test it with an HTTP trigger.
|
||||
# Update it and check it's output, the output should be
|
||||
# different from the previous one.
|
||||
|
||||
echo "Pre-test cleanup"
|
||||
fission env delete --name nodejs || true
|
||||
|
||||
echo "Creating nodejs env"
|
||||
fission env create --name nodejs --image fission/node-env
|
||||
trap "fission env delete --name nodejs" EXIT
|
||||
|
||||
echo "Creating function"
|
||||
echo 'module.exports = function(context, callback) { callback(200, "foo!\n"); }' > foo.js
|
||||
fission fn create --name $fn --env nodejs --code foo.js
|
||||
trap "fission fn delete --name $fn" EXIT
|
||||
|
||||
echo "Creating route"
|
||||
fission route create --function $fn --url /$fn --method GET
|
||||
|
||||
echo "Waiting for router to catch up"
|
||||
sleep 10
|
||||
|
||||
echo "Doing an HTTP GET on the function's route"
|
||||
response=$(curl http://$FISSION_ROUTER/$fn)
|
||||
|
||||
echo "Checking for valid response"
|
||||
echo $response | grep -i foo
|
||||
|
||||
# Running a background process to keep access the
|
||||
# function to emulate real online traffic. The router
|
||||
# should be able to update cache under this situation.
|
||||
( watch -n1 curl http://$FISSION_ROUTER/$fn ) > /dev/null 2>&1 &
|
||||
pid=$!
|
||||
|
||||
echo "Updating function"
|
||||
echo 'module.exports = function(context, callback) { callback(200, "bar!\n"); }' > bar.js
|
||||
fission fn update --name $fn --code bar.js
|
||||
trap "fission fn delete --name $fn" EXIT
|
||||
|
||||
echo "Waiting for router to update cache"
|
||||
sleep 10
|
||||
|
||||
echo "Doing an HTTP GET on the function's route"
|
||||
response=$(curl http://$FISSION_ROUTER/$fn)
|
||||
|
||||
echo "Checking for valid response again"
|
||||
echo $response | grep -i bar
|
||||
|
||||
kill -15 $pid
|
||||
|
||||
# crappy cleanup, improve this later
|
||||
kubectl get httptrigger -o name | tail -1 | cut -f2 -d'/' | xargs kubectl delete httptrigger
|
||||
|
||||
echo "All done."
|
||||
Reference in New Issue
Block a user