From fd286d8d7a1efcede01036337aad14642ccb9744 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Mon, 19 Oct 2020 16:01:02 -0700 Subject: [PATCH 01/20] update --- go.mod | 5 +-- go.sum | 43 +------------------ pkg/controller/nodes/task/config/config.go | 34 ++++++++++----- pkg/controller/nodes/task/handler.go | 15 +++++-- pkg/controller/nodes/task/handler_test.go | 38 ++++++++++++---- pkg/controller/nodes/task/plugin_config.go | 19 ++++---- .../nodes/task/plugin_config_test.go | 8 ++-- 7 files changed, 83 insertions(+), 79 deletions(-) diff --git a/go.mod b/go.mod index c57043f10..eb2dfe6ce 100644 --- a/go.mod +++ b/go.mod @@ -8,7 +8,6 @@ require ( github.com/Azure/go-autorest/autorest v0.10.0 // indirect github.com/DiSiqueira/GoTree v1.0.1-0.20180907134536-53a8e837f295 github.com/benlaurie/objecthash v0.0.0-20180202135721-d1e3d6079fc1 - github.com/coreos/etcd v3.3.15+incompatible // indirect github.com/coreos/go-oidc v2.2.1+incompatible // indirect github.com/fatih/color v1.9.0 github.com/ghodss/yaml v1.0.0 @@ -23,11 +22,10 @@ require ( github.com/imdario/mergo v0.3.8 // indirect github.com/lyft/datacatalog v0.2.1 github.com/lyft/flyteidl v0.18.9 - github.com/lyft/flyteplugins v0.5.12 + github.com/lyft/flyteplugins v0.5.14-0.20201019203656-dbcbf808fc35 github.com/lyft/flytestdlib v0.3.9 github.com/magiconair/properties v1.8.1 github.com/mattn/go-colorable v0.1.6 // indirect - github.com/mitchellh/go-ps v1.0.0 // indirect github.com/mitchellh/mapstructure v1.1.2 github.com/ncw/swift v1.0.50 // indirect github.com/pkg/errors v0.9.1 @@ -49,7 +47,6 @@ require ( k8s.io/kube-openapi v0.0.0-20200204173128-addea2498afe // indirect k8s.io/utils v0.0.0-20200229041039-0a110f9eb7ab // indirect sigs.k8s.io/controller-runtime v0.5.1 - sigs.k8s.io/testing_frameworks v0.1.2 // indirect sigs.k8s.io/yaml v1.2.0 // indirect ) diff --git a/go.sum b/go.sum index 4206c546b..2a57c1fb3 100644 --- a/go.sum +++ b/go.sum @@ -25,7 +25,6 @@ cloud.google.com/go/storage v1.6.0/go.mod h1:N7U0C8pVQ/+NIKOBQyamJIeKQKkZ+mxpohl dmitri.shuralyov.com/gpu/mtl v0.0.0-20190408044501-666a987793e9/go.mod h1:H6x//7gZCb22OMCxBHrMx7a5I7Hp++hsVxbQ4BYO7hU= github.com/Azure/azure-sdk-for-go v32.5.0+incompatible/go.mod h1:9XXNKU+eRnpl9moKnB4QOLf1HestfXbmab5FXxiDBjc= github.com/Azure/azure-sdk-for-go v38.2.0+incompatible/go.mod h1:9XXNKU+eRnpl9moKnB4QOLf1HestfXbmab5FXxiDBjc= -github.com/Azure/azure-sdk-for-go v39.0.0+incompatible/go.mod h1:9XXNKU+eRnpl9moKnB4QOLf1HestfXbmab5FXxiDBjc= github.com/Azure/azure-sdk-for-go v40.3.0+incompatible h1:NthZg3psrLxvQLN6rVm07pZ9mv2wvGNaBNGQ3fnPvLE= github.com/Azure/azure-sdk-for-go v40.3.0+incompatible/go.mod h1:9XXNKU+eRnpl9moKnB4QOLf1HestfXbmab5FXxiDBjc= github.com/Azure/go-ansiterm v0.0.0-20170929234023-d6e3b3328b78/go.mod h1:LmzpDX56iTiv29bbRTIsUNlaFfuhWRQBWjQdVyAevI8= @@ -91,7 +90,6 @@ github.com/aws/amazon-sagemaker-operator-for-k8s v1.0.1-0.20200410212604-780c48e github.com/aws/amazon-sagemaker-operator-for-k8s v1.0.1-0.20200410212604-780c48ecb21a/go.mod h1:kw+Gl0uvPAMADPoubX+kLx0P7e7zWOr6rc+R7D24pbc= github.com/aws/aws-sdk-go v1.23.4/go.mod h1:KmX6BPdI08NWTb3/sm4ZGu5ShLoqVDhKgpiN924inxo= github.com/aws/aws-sdk-go v1.28.9/go.mod h1:KmX6BPdI08NWTb3/sm4ZGu5ShLoqVDhKgpiN924inxo= -github.com/aws/aws-sdk-go v1.28.11/go.mod h1:KmX6BPdI08NWTb3/sm4ZGu5ShLoqVDhKgpiN924inxo= github.com/aws/aws-sdk-go v1.29.23 h1:wtiGLOzxAP755OfuVTDIy/NbUIYEDxbIbBEDfNhUpeU= github.com/aws/aws-sdk-go v1.29.23/go.mod h1:1KvfttTE3SPKMpo8g2c6jL3ZKfXtFvKscTgahTma5Xg= github.com/aws/aws-sdk-go-v2 v0.20.0 h1:/yefUjgMrda9PNFwWctBU63nL10CJMdBwkAmaQ4w4Hs= @@ -119,10 +117,8 @@ github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGX github.com/cockroachdb/datadriven v0.0.0-20190809214429-80d97fb3cbaa/go.mod h1:zn76sxSg3SzpJ0PPJaLDCu+Bu0Lg3sKTORVIj19EIF8= github.com/coocood/freecache v1.1.0 h1:ENiHOsWdj1BrrlPwblhbn4GdAsMymK3pZORJ+bJGAjA= github.com/coocood/freecache v1.1.0/go.mod h1:ePwxCDzOYvARfHdr1pByNct1at3CoKnsipOHwKlNbzI= -github.com/coreos/bbolt v1.3.1-coreos.6/go.mod h1:iRUV2dpdMOn7Bo10OQBFzIJO9kkE559Wcmn+qkEiiKk= github.com/coreos/bbolt v1.3.2/go.mod h1:iRUV2dpdMOn7Bo10OQBFzIJO9kkE559Wcmn+qkEiiKk= github.com/coreos/etcd v3.3.10+incompatible/go.mod h1:uF7uidLiAD3TWHmW31ZFd/JWoc32PjwdhPthX9715RE= -github.com/coreos/etcd v3.3.15+incompatible/go.mod h1:uF7uidLiAD3TWHmW31ZFd/JWoc32PjwdhPthX9715RE= github.com/coreos/go-etcd v2.0.0+incompatible/go.mod h1:Jez6KQU2B/sWsbdaef3ED8NzMklzPG4d5KIOhIy30Tk= github.com/coreos/go-oidc v2.1.0+incompatible/go.mod h1:CgnwVTmzoESiwO9qyAFEMiHoZ1nMCKZlZ9V6mm3/LKc= github.com/coreos/go-oidc v2.2.1+incompatible h1:mh48q/BqXqgjVHpy2ZY7WnWAbenxRjsz9N1i1YxjHAk= @@ -269,7 +265,6 @@ github.com/golang/mock v1.3.1/go.mod h1:sBzyDLLjw3U8JLTeZvSv8jJB+tU5PVekmnlKIyFU github.com/golang/mock v1.4.0/go.mod h1:UOMv5ysSaYNkG+OFQykRIcU/QvvxJf3p21QfJ2Bt3cw= github.com/golang/mock v1.4.1/go.mod h1:UOMv5ysSaYNkG+OFQykRIcU/QvvxJf3p21QfJ2Bt3cw= github.com/golang/protobuf v0.0.0-20161109072736-4bd1920723d7/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.0.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= @@ -317,7 +312,6 @@ github.com/graymeta/stow v0.2.4/go.mod h1:+0vRL9oMECKjPMP7OeVWl8EIqRCpFwDlth3mrA github.com/graymeta/stow v0.2.5 h1:YFSo4nsAU4Fbi4r/neLIgVYlrMzA1ReDUkdLYTQm/RM= github.com/graymeta/stow v0.2.5/go.mod h1:+0vRL9oMECKjPMP7OeVWl8EIqRCpFwDlth3mrAeV2Kw= github.com/gregjones/httpcache v0.0.0-20170728041850-787624de3eb7/go.mod h1:FecbI9+v66THATjSRHfNgh1IVFe/9kFxbXtjV0ctIMA= -github.com/grpc-ecosystem/go-grpc-middleware v0.0.0-20190222133341-cfaf5686ec79/go.mod h1:FiyG127CGDf3tlThmgyCl78X/SZQqEOJBCDaAfeWzPs= github.com/grpc-ecosystem/go-grpc-middleware v1.0.0/go.mod h1:FiyG127CGDf3tlThmgyCl78X/SZQqEOJBCDaAfeWzPs= github.com/grpc-ecosystem/go-grpc-middleware v1.0.1-0.20190118093823-f849b5445de4/go.mod h1:FiyG127CGDf3tlThmgyCl78X/SZQqEOJBCDaAfeWzPs= github.com/grpc-ecosystem/go-grpc-middleware v1.1.0/go.mod h1:f5nM7jw/oeRSadq3xCzHAvxcr8HZnzsqU6ILg/0NiiE= @@ -325,7 +319,6 @@ github.com/grpc-ecosystem/go-grpc-middleware v1.2.0 h1:0IKlLyQ3Hs9nDaiK5cSHAGmcQ github.com/grpc-ecosystem/go-grpc-middleware v1.2.0/go.mod h1:mJzapYve32yjrKlk9GbyCZHuPgZsrbyIbyKhSzOpg6s= github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0 h1:Ovs26xHkKqVztRpIrF/92BcuyuQ/YW4NSIpoGtfXNho= github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0/go.mod h1:8NvIoxWQoOIhqOTXgfV/d3M/q6VIi02HzZEHgUlZvzk= -github.com/grpc-ecosystem/grpc-gateway v1.3.0/go.mod h1:RSKVYQBd5MCa4OVpNdGskqpgL2+G+NZTnrVHpWWfpdw= github.com/grpc-ecosystem/grpc-gateway v1.9.0/go.mod h1:vNeuVxBJEsws4ogUvrchl83t/GYV9WGTSLVdBhOQFDY= github.com/grpc-ecosystem/grpc-gateway v1.9.5/go.mod h1:vNeuVxBJEsws4ogUvrchl83t/GYV9WGTSLVdBhOQFDY= github.com/grpc-ecosystem/grpc-gateway v1.12.2/go.mod h1:8XEsbTttt/W+VvjtQhLACqCisSPWTxCZ7sBRjU6iH9c= @@ -396,20 +389,12 @@ github.com/lyft/apimachinery v0.0.0-20191031200210-047e3ea32d7f/go.mod h1:llRdnz github.com/lyft/datacatalog v0.2.1 h1:W7LsAjaS297iLCtSH9ZaAAG3YPofwkbbgIaqkfdeM0o= github.com/lyft/datacatalog v0.2.1/go.mod h1:ktrPvzTDUwHO5Lv0hLH38zLHnOJ++rGoAO0iQ/sIPJ4= github.com/lyft/flyteidl v0.17.0/go.mod h1:/zQXxuHO11u/saxTTZc8oYExIGEShXB+xCB1/F1Cu20= -github.com/lyft/flyteidl v0.18.0/go.mod h1:/zQXxuHO11u/saxTTZc8oYExIGEShXB+xCB1/F1Cu20= -github.com/lyft/flyteidl v0.18.7 h1:R8gSt2Tze9BlHbFHZPDPWl630272US+MbSjqoeVkflg= -github.com/lyft/flyteidl v0.18.7/go.mod h1:/zQXxuHO11u/saxTTZc8oYExIGEShXB+xCB1/F1Cu20= github.com/lyft/flyteidl v0.18.9 h1:p9gLp92whTSSOeMGPtZ4tkgsVHNGuBuXXMQ447s0J9E= github.com/lyft/flyteidl v0.18.9/go.mod h1:/zQXxuHO11u/saxTTZc8oYExIGEShXB+xCB1/F1Cu20= -github.com/lyft/flyteplugins v0.5.1 h1:76FpQFohLCy4Eo490sES2empRAi31DJiERfbOSV9pCg= -github.com/lyft/flyteplugins v0.5.1/go.mod h1:8zhqFG9BzbHNQGEXzGYltTJLD+KTmQZkanxXgeFI25c= -github.com/lyft/flyteplugins v0.5.6 h1:4r9aT8XGRcTQk6VCUlJkz9tsSFFPXLVL6C/mVcgCvyQ= -github.com/lyft/flyteplugins v0.5.6/go.mod h1:BRUqCc6HycsFVnhF5Z5MMBHYoqOemdCyfG59lg0B1bY= -github.com/lyft/flyteplugins v0.5.10 h1:mMAtx9PZ1ZdEQ57HiAQsZmF6Rhwo9SWFNnZxYHZY394= -github.com/lyft/flyteplugins v0.5.10/go.mod h1:X17xeh3Sc9ZaAEoZp1Lw6e1lysiToY7aPqYMcstAREI= github.com/lyft/flyteplugins v0.5.12 h1:obz52m9dJ/ununeQJ2OcLZ1z38TE5KMDo8gyH/g8NAg= github.com/lyft/flyteplugins v0.5.12/go.mod h1:UOoiW+rwQdrDDig3bJSxTWoyW8hW5bcqnuRsp+g2zhQ= -github.com/lyft/flytepropeller v0.4.2/go.mod h1:TIiWv/ZP1KOI0mqeUbiMqSn2XuY8O8kn8fQc5tWcaLA= +github.com/lyft/flyteplugins v0.5.14-0.20201019203656-dbcbf808fc35 h1:JWOQEs0sJN/vyW2EhpQfK2VhuGwMcd4MU5JGB9OsWls= +github.com/lyft/flyteplugins v0.5.14-0.20201019203656-dbcbf808fc35/go.mod h1:UOoiW+rwQdrDDig3bJSxTWoyW8hW5bcqnuRsp+g2zhQ= github.com/lyft/flytestdlib v0.3.0/go.mod h1:LJPPJlkFj+wwVWMrQT3K5JZgNhZi2mULsCG4ZYhinhU= github.com/lyft/flytestdlib v0.3.9 h1:NaKp9xkeWWwhVvqTOcR/FqlASy1N2gu/kN7PVe4S7YI= github.com/lyft/flytestdlib v0.3.9/go.mod h1:LJPPJlkFj+wwVWMrQT3K5JZgNhZi2mULsCG4ZYhinhU= @@ -439,7 +424,6 @@ github.com/mattn/go-sqlite3 v1.11.0/go.mod h1:FPy6KqzDD04eiIsT53CuJW3U88zkxoIYsO github.com/matttproud/golang_protobuf_extensions v1.0.1 h1:4hp9jkHxhMHkqkrB3Ix0jegS5sx/RkqARlsWZ6pIwiU= github.com/matttproud/golang_protobuf_extensions v1.0.1/go.mod h1:D8He9yQNgCq6Z5Ld7szi9bcBfOoFv/3dc6xSMkL2PC0= github.com/mitchellh/go-homedir v1.1.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0= -github.com/mitchellh/go-ps v1.0.0/go.mod h1:J4lOc8z8yJs6vUwklHw2XEIiT4z4C40KtWVN3nvg8Pg= github.com/mitchellh/mapstructure v1.1.2 h1:fmNYVwqnSfB9mZU6OS2O6GsXM+wcskZDuKQzvN1EDeE= github.com/mitchellh/mapstructure v1.1.2/go.mod h1:FVVH3fgwuzCH5S8UJGiWEs2h04kUh9fWfEaFds41c1Y= github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= @@ -461,7 +445,6 @@ github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLA github.com/oklog/ulid v1.3.1/go.mod h1:CirwcVhetQ6Lv90oh/F+FBtV6XMibvdAFo93nm5qn4U= github.com/olekukonko/tablewriter v0.0.0-20170122224234-a0225b3f23b5/go.mod h1:vsDQFd/mU46D+Z4whnwzcISnGGzXWMclvtLoiIKAKIo= github.com/onsi/ginkgo v0.0.0-20170829012221-11459a886d9c/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= -github.com/onsi/ginkgo v1.4.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/ginkgo v1.6.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/ginkgo v1.7.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/ginkgo v1.8.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= @@ -469,7 +452,6 @@ github.com/onsi/ginkgo v1.10.1/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+ github.com/onsi/ginkgo v1.11.0 h1:JAKSXpt1YjtLA7YpPiqO9ss6sNXEsPfSGdwN0UHqzrw= github.com/onsi/ginkgo v1.11.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/gomega v0.0.0-20170829124025-dcabb60a477c/go.mod h1:C1qb7wdrVGGVU+Z6iS04AVkA3Q65CEZX59MT0QO5uiA= -github.com/onsi/gomega v1.3.0/go.mod h1:C1qb7wdrVGGVU+Z6iS04AVkA3Q65CEZX59MT0QO5uiA= github.com/onsi/gomega v1.4.2/go.mod h1:ex+gbHU/CVuBBDIJjb2X0qEXbFg53c61hWP/1CpauHY= github.com/onsi/gomega v1.4.3/go.mod h1:ex+gbHU/CVuBBDIJjb2X0qEXbFg53c61hWP/1CpauHY= github.com/onsi/gomega v1.5.0/go.mod h1:ex+gbHU/CVuBBDIJjb2X0qEXbFg53c61hWP/1CpauHY= @@ -499,12 +481,10 @@ github.com/pquerna/cachecontrol v0.0.0-20180517163645-1555304b9b35/go.mod h1:prY github.com/pquerna/ffjson v0.0.0-20190813045741-dac163c6c0a9/go.mod h1:YARuvh7BUWHNhzDq2OM5tzR2RiCcN2D7sapiKyCel/M= github.com/prometheus/client_golang v0.9.0/go.mod h1:7SWBe2y4D6OKWSNQJUaRYU/AaXPKyh/dDVn+NZz0KFw= github.com/prometheus/client_golang v0.9.1/go.mod h1:7SWBe2y4D6OKWSNQJUaRYU/AaXPKyh/dDVn+NZz0KFw= -github.com/prometheus/client_golang v0.9.2/go.mod h1:OsXs2jCmiKlQ1lTBmv21f2mNfw4xf/QclQDMrYNZzcM= github.com/prometheus/client_golang v0.9.3-0.20190127221311-3c4408c8b829/go.mod h1:p2iRAGwDERtqlqzRXnrOVns+ignqQo//hLXqYxZYVNs= github.com/prometheus/client_golang v0.9.3/go.mod h1:/TN21ttK/J9q6uSwhBd54HahCDft0ttaMvbicHlPoso= github.com/prometheus/client_golang v1.0.0/go.mod h1:db9x61etRT2tGnBNRi70OPL5FsnadC4Ky3P0J6CfImo= github.com/prometheus/client_golang v1.3.0/go.mod h1:hJaj2vgQTGQmVCsAACORcieXFeDPbaTKGT+JTgUa3og= -github.com/prometheus/client_golang v1.4.0/go.mod h1:e9GMxYsXl05ICDXkRhurwBS4Q3OK1iX/F2sw+iXX5zU= github.com/prometheus/client_golang v1.5.0 h1:Ctq0iGpCmr3jeP77kbF2UxgvRwzWWz+4Bh9/vJTyg1A= github.com/prometheus/client_golang v1.5.0/go.mod h1:e9GMxYsXl05ICDXkRhurwBS4Q3OK1iX/F2sw+iXX5zU= github.com/prometheus/client_model v0.0.0-20180712105110-5c3871d89910/go.mod h1:MbSGuTsp3dbXC40dX6PRTWyKYBIrTGTE9sqQNg2J8bo= @@ -516,7 +496,6 @@ github.com/prometheus/client_model v0.2.0 h1:uq5h0d+GuxiXLJLNABMgp2qUWDPiLvgCzz2 github.com/prometheus/client_model v0.2.0/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA= github.com/prometheus/common v0.0.0-20180801064454-c7de2306084e/go.mod h1:daVV7qP5qjZbuso7PdcryaAu0sAZbrN9i7WWcTMWvro= github.com/prometheus/common v0.0.0-20181113130724-41aa239b4cce/go.mod h1:daVV7qP5qjZbuso7PdcryaAu0sAZbrN9i7WWcTMWvro= -github.com/prometheus/common v0.0.0-20181126121408-4724e9255275/go.mod h1:daVV7qP5qjZbuso7PdcryaAu0sAZbrN9i7WWcTMWvro= github.com/prometheus/common v0.2.0/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4= github.com/prometheus/common v0.4.0/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4= github.com/prometheus/common v0.4.1/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4= @@ -525,7 +504,6 @@ github.com/prometheus/common v0.9.1 h1:KOMtN28tlbam3/7ZKEYKHhKoJZYYj3gMH4uc62x7X github.com/prometheus/common v0.9.1/go.mod h1:yhUN8i9wzaXS3w1O07YhxHEBxD+W35wd8bs7vj7HSQ4= github.com/prometheus/procfs v0.0.0-20180725123919-05ee40e3a273/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk= github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk= -github.com/prometheus/procfs v0.0.0-20181204211112-1dc9a6cbc91a/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk= github.com/prometheus/procfs v0.0.0-20190117184657-bf6a532e95b1/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk= github.com/prometheus/procfs v0.0.0-20190507164030-5867b95ac084/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA= github.com/prometheus/procfs v0.0.2/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA= @@ -552,7 +530,6 @@ github.com/smartystreets/assertions v0.0.0-20180927180507-b2de0cb4f26d h1:zE9ykE github.com/smartystreets/assertions v0.0.0-20180927180507-b2de0cb4f26d/go.mod h1:OnSkiWE9lh6wB0YB77sQom3nweQdgAjqCqsofrRNTgc= github.com/smartystreets/goconvey v1.6.4 h1:fv0U8FUIMPNf1L9lnHLvLhgicrIVChEkdzIKYqbNC9s= github.com/smartystreets/goconvey v1.6.4/go.mod h1:syvi0/a8iFYH4r/RixwvyeAJjdLS9QV7WQ/tjFTllLA= -github.com/soheilhy/cmux v0.1.3/go.mod h1:IM3LyeVVIOuxMH7sFAkER9+bJ4dT7Ms6E4xg4kGIyLM= github.com/soheilhy/cmux v0.1.4/go.mod h1:IM3LyeVVIOuxMH7sFAkER9+bJ4dT7Ms6E4xg4kGIyLM= github.com/spaolacci/murmur3 v0.0.0-20180118202830-f09979ecbc72 h1:qLC7fQah7D6K1B0ujays3HV9gkFtllcxhzImRR7ArPQ= github.com/spaolacci/murmur3 v0.0.0-20180118202830-f09979ecbc72/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA= @@ -599,7 +576,6 @@ github.com/ugorji/go v1.1.4/go.mod h1:uQMGLiO92mf5W77hV/PUCpI3pbzQx3CRekS0kk+RGr github.com/ugorji/go/codec v0.0.0-20181204163529-d75b2dcb6bc8/go.mod h1:VFNgLljTbGfSG7qAOspJ7OScBnGdDN/yBr0sguwnwf0= github.com/urfave/cli v1.20.0/go.mod h1:70zkFmudgCuE/ngEzBv17Jvp/497gISqfk5gWijbERA= github.com/vektah/gqlparser v1.1.2/go.mod h1:1ycwN7Ij5njmMkPPAOaRFY4rET2Enx7IkVv3vaXspKw= -github.com/xiang90/probing v0.0.0-20160813154853-07dd2e8dfe18/go.mod h1:UETIi67q53MR2AWcXfiuqkDkRtnGDLqkBTpCHuJHxtU= github.com/xiang90/probing v0.0.0-20190116061207-43a291ad63a2/go.mod h1:UETIi67q53MR2AWcXfiuqkDkRtnGDLqkBTpCHuJHxtU= github.com/xordataexchange/crypt v0.0.3-0.20170626215501-b2862e3d0a77/go.mod h1:aYKd//L2LvnjZzWKhF00oedf4jCCReLcmhLdhm1A27Q= go.etcd.io/bbolt v1.3.2/go.mod h1:IbVyRI1SCnLcuJnV2u8VeU0CEYM7e686BmAb1XKL+uU= @@ -614,14 +590,11 @@ go.opencensus.io v0.22.0/go.mod h1:+kGneAE2xo2IficOXnaByMWTGM9T73dGwxeWcUqIpI8= go.opencensus.io v0.22.2/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw= go.opencensus.io v0.22.3 h1:8sGtKOrtQqkN1bp2AtX+misvLIlOmsEsNd+9NIcPEm8= go.opencensus.io v0.22.3/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw= -go.uber.org/atomic v0.0.0-20181018215023-8dc6146f7569/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE= go.uber.org/atomic v1.3.2/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE= go.uber.org/atomic v1.4.0 h1:cxzIVoETapQEqDhQu3QfnvXAV4AlzcvUCxkVUFw3+EU= go.uber.org/atomic v1.4.0/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE= -go.uber.org/multierr v0.0.0-20180122172545-ddea229ff1df/go.mod h1:wR5kodmAFQ0UK8QlbwjlSNy0Z68gJhDJUG5sjR94q/0= go.uber.org/multierr v1.1.0 h1:HoEmRHQPVSqub6w2z2d2EOVs2fjyFRGyofhKuyDq0QI= go.uber.org/multierr v1.1.0/go.mod h1:wR5kodmAFQ0UK8QlbwjlSNy0Z68gJhDJUG5sjR94q/0= -go.uber.org/zap v0.0.0-20180814183419-67bc79d13d15/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q= go.uber.org/zap v1.9.1/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q= go.uber.org/zap v1.10.0 h1:ORx85nbTijNz8ljznvCMR1ZBIPKFn3jQrag10X2AsuM= go.uber.org/zap v1.10.0/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q= @@ -675,13 +648,11 @@ golang.org/x/mod v0.1.1-0.20191105210325-c90efee705ee/go.mod h1:QqPTAvyqsEbceGzB golang.org/x/mod v0.1.1-0.20191107180719-034126e5016b/go.mod h1:QqPTAvyqsEbceGzBzNggFXnrqF1CaUcvgkdR5Ot7KZg= golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/net v0.0.0-20170114055629-f2499483f923/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20180112015858-5ccada7d0a7b/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20180826012351-8a410e7b638d/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20180906233101-161cd47e91fd/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20181005035420-146acd28ed58/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20181114220301-adae6a3d119a/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20181201002055-351d144fa1fc/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20181220203305-927f97764cc3/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20190108225652-1e06a53dbb7e/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/net v0.0.0-20190125091013-d26f9f9a57f3/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= @@ -722,7 +693,6 @@ golang.org/x/sync v0.0.0-20190227155943-e225da77a7e6/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sys v0.0.0-20170830134202-bb24a47a89ea/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20180117170059-2c42eef0765b/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180909124046-d0be0721c37e/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -765,7 +735,6 @@ golang.org/x/sys v0.0.0-20200327173247-9dae0f8f5775/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/text v0.0.0-20160726164857-2910a502d2bf/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= -golang.org/x/text v0.3.1-0.20171227012246-e19ae1496984/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.1-0.20180807135948-17ff2d5776d2/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.2 h1:tW2bmiBqwgJj/UpqtC8EpXEZVYOwU0yG4iWbprSVAcs= golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= @@ -905,7 +874,6 @@ gopkg.in/square/go-jose.v2 v2.4.1/go.mod h1:M9dMgbHiYLoDGQrXy7OpJDJWiKiU//h+vD76 gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7 h1:uRGJdciOHaEIrze2W8Q3AKkepLTh2hOroT7a+7czfdQ= gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7/go.mod h1:dt/ZhP58zS4L8KSrWDmTeBkI65Dw0HsyUHuEVlX15mw= gopkg.in/yaml.v2 v2.0.0-20170812160011-eb3733d160e7/go.mod h1:JAlM8MvJe8wmxCU4Bli9HhUf9+ttbYbLASfIpnQbh74= -gopkg.in/yaml.v2 v2.0.0/go.mod h1:JAlM8MvJe8wmxCU4Bli9HhUf9+ttbYbLASfIpnQbh74= gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v2 v2.2.3/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= @@ -927,16 +895,12 @@ honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWh honnef.co/go/tools v0.0.1-2019.2.3/go.mod h1:a3bituU0lyd329TUQxRnasdCoJDkEUEAqEt0JzvZhAg= honnef.co/go/tools v0.0.1-2020.1.3/go.mod h1:X/FiERA/W4tHapMX5mGpAtMSVEeEUOyHaw9vFzvIQ3k= k8s.io/apiextensions-apiserver v0.0.0-20190409022649-727a075fdec8/go.mod h1:IxkesAMoaCRoLrPJdZNZUQp9NfZnzqaVzLhb2VEQzXE= -k8s.io/apiextensions-apiserver v0.0.0-20190918161926-8f644eb6e783/go.mod h1:xvae1SZB3E17UpV59AWc271W/Ph25N+bjPyR63X6tPY= k8s.io/apiextensions-apiserver v0.17.2 h1:cP579D2hSZNuO/rZj9XFRzwJNYb41DbNANJb6Kolpss= k8s.io/apiextensions-apiserver v0.17.2/go.mod h1:4KdMpjkEjjDI2pPfBA15OscyNldHWdBCfsWMDWAmSTs= -k8s.io/apiserver v0.0.0-20190918160949-bfa5e2e684ad/go.mod h1:XPCXEwhjaFN29a8NldXA901ElnKeKLrLtREO9ZhFyhg= k8s.io/apiserver v0.17.2/go.mod h1:lBmw/TtQdtxvrTk0e2cgtOxHizXI+d0mmGQURIHQZlo= k8s.io/client-go v0.0.0-20191016111102-bec269661e48 h1:C2XVy2z0dV94q9hSSoCuTPp1KOG7IegvbdXuz9VGxoU= k8s.io/client-go v0.0.0-20191016111102-bec269661e48/go.mod h1:hrwktSwYGI4JK+TJA3dMaFyyvHVi/aLarVHpbs8bgCU= -k8s.io/code-generator v0.0.0-20190912054826-cd179ad6a269/go.mod h1:V5BD6M4CyaN5m+VthcclXWsVcT1Hu+glwa1bi3MIsyE= k8s.io/code-generator v0.17.2/go.mod h1:DVmfPQgxQENqDIzVR2ddLXMH34qeszkKSdH/N+s+38s= -k8s.io/component-base v0.0.0-20190918160511-547f6c5d7090/go.mod h1:933PBGtQFJky3TEwYx4aEPZ4IxqhWh3R6DCmzqIn1hA= k8s.io/component-base v0.17.2/go.mod h1:zMPW3g5aH7cHJpKYQ/ZsGMcgbsA/VyhEugF3QT1awLs= k8s.io/gengo v0.0.0-20190128074634-0689ccc1d7d6/go.mod h1:ezvh/TsK7cY6rbqRK0oQQ8IAqLxYwwyPxAX1Pzy0ii0= k8s.io/gengo v0.0.0-20190822140433-26a664648505/go.mod h1:ezvh/TsK7cY6rbqRK0oQQ8IAqLxYwwyPxAX1Pzy0ii0= @@ -965,15 +929,12 @@ rsc.io/binaryregexp v0.2.0/go.mod h1:qTv7/COck+e2FymRvadv62gMdZztPaShugOCi3I+8D8 rsc.io/quote/v3 v3.1.0/go.mod h1:yEA65RcK8LyAZtP9Kv3t0HmxON59tX3rD+tICJqUlj0= rsc.io/sampler v1.3.0/go.mod h1:T1hPZKmBbMNahiBKFy5HrXp6adAjACjK9JXDnKaTXpA= sigs.k8s.io/controller-runtime v0.2.0/go.mod h1:ZHqrRDZi3f6BzONcvlUxkqCKgwasGk5FZrnSv9TVZF4= -sigs.k8s.io/controller-runtime v0.4.0/go.mod h1:ApC79lpY3PHW9xj/w9pj+lYkLgwAAUZwfXkME1Lajns= sigs.k8s.io/controller-runtime v0.5.1 h1:TNidCfVoU/cs2i+9xoTcL/l7yhl0bDhYXU0NCG6wmiE= sigs.k8s.io/controller-runtime v0.5.1/go.mod h1:Uojny7gvg55YLQnEGnPzRE3dC4ik2tRlZJgOUCWXAV4= sigs.k8s.io/structured-merge-diff v0.0.0-20190525122527-15d366b2352e/go.mod h1:wWxsB5ozmmv/SG7nM11ayaAW51xMvak/t1r0CSlcokI= -sigs.k8s.io/structured-merge-diff v0.0.0-20190817042607-6149e4549fca/go.mod h1:IIgPezJWb76P0hotTxzDbWsMYB8APh18qZnxkomBpxA= sigs.k8s.io/structured-merge-diff v1.0.1-0.20191108220359-b1b620dd3f06/go.mod h1:/ULNhyfzRopfcjskuui0cTITekDduZ7ycKN3oUT9R18= sigs.k8s.io/structured-merge-diff/v3 v3.0.0-20200116222232-67a7b8c61874/go.mod h1:PlARxl6Hbt/+BC80dRLi1qAmnMqwqDg62YvvVkZjemw= sigs.k8s.io/testing_frameworks v0.1.1/go.mod h1:VVBKrHmJ6Ekkfz284YKhQePcdycOzNH9qL6ht1zEr/U= -sigs.k8s.io/testing_frameworks v0.1.2/go.mod h1:ToQrwSC3s8Xf/lADdZp3Mktcql9CG0UAmdJG9th5i0w= sigs.k8s.io/yaml v1.1.0/go.mod h1:UJmg0vDUVViEyp3mgSv9WPwZCDxu4rQW1olrI1uml+o= sigs.k8s.io/yaml v1.2.0 h1:kr/MCeFWJWTwyaHoR9c8EjH9OumOmoF9YGiZd7lFm/Q= sigs.k8s.io/yaml v1.2.0/go.mod h1:yfXDCHCao9+ENCvLSE62v9VSji2MKu5jeNfTrofGhJc= diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index b4f45d473..9e4123b4b 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -5,7 +5,6 @@ import ( "time" "github.com/lyft/flytestdlib/config" - "k8s.io/apimachinery/pkg/util/sets" ) //go:generate pflags Config --default-var defaultConfig @@ -14,7 +13,7 @@ const SectionKey = "tasks" var ( defaultConfig = &Config{ - TaskPlugins: TaskPluginConfig{EnabledPlugins: []string{}}, + TaskPlugins: TaskPluginConfig{EnabledPlugins: map[string]EnabledPlugins{}}, MaxPluginPhaseVersions: 100000, BarrierConfig: BarrierConfig{ Enabled: true, @@ -45,8 +44,12 @@ type BarrierConfig struct { CacheTTL config.Duration `json:"cache-ttl" pflag:", Max duration that a barrier would be respected if the process is not restarted. This should account for time required to store the record into persistent storage (across multiple rounds."` } +type EnabledPlugins struct { + DefaultPluginTasks []string `json:"default-plugins-tasks" pflag:",Tasks for which this plugin is the default implementation"` +} + type TaskPluginConfig struct { - EnabledPlugins []string `json:"enabled-plugins" pflag:",Plugins enabled currently"` + EnabledPlugins map[string]EnabledPlugins `json:"enabled-plugins" pflag:",Plugins enabled currently"` } type BackOffConfig struct { @@ -54,14 +57,25 @@ type BackOffConfig struct { MaxDuration config.Duration `json:"max-duration" pflag:",The cap of the backoff duration"` } -func (p TaskPluginConfig) GetEnabledPluginsSet() sets.String { - s := sets.NewString() - for _, e := range p.EnabledPlugins { - cleanedPluginName := strings.Trim(e, " ") - cleanedPluginName = strings.ToLower(cleanedPluginName) - s.Insert(cleanedPluginName) +func cleanString(source string) string { + cleaned := strings.Trim(source, " ") + cleaned = strings.ToLower(cleaned) + return cleaned +} + +func (p TaskPluginConfig) GetEnabledPlugins() map[string]EnabledPlugins { + enabledPlugins := make(map[string]EnabledPlugins) + for pluginName, info := range p.EnabledPlugins { + cleanedDefaultTasks := make([]string, 0, len(info.DefaultPluginTasks)) + for _, taskName := range info.DefaultPluginTasks { + cleanedDefaultTasks = append(cleanedDefaultTasks, cleanString(taskName)) + } + cleanedPluginName := cleanString(pluginName) + enabledPlugins[cleanedPluginName] = EnabledPlugins{ + DefaultPluginTasks: cleanedDefaultTasks, + } } - return s + return enabledPlugins } func GetConfig() *Config { diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index b92dc3039..866f4c61e 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -215,10 +215,19 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { if err != nil { return regErrors.Wrapf(err, "failed to load plugin - %s", p.ID) } + println(fmt.Sprintf("for plugin [%s], registered task types: [%+v] and default task types [%+v]", + p.ID, p.RegisteredTaskTypes, p.DefaultForTaskTypes)) for _, tt := range p.RegisteredTaskTypes { - logger.Infof(ctx, "Plugin [%s] registered for TaskType [%s]", cp.GetID(), tt) - // TODO(katrogan): Make the default task plugin assignment more explicit (https://github.com/lyft/flyte/issues/516) - t.defaultPlugins[tt] = cp + for _, defaultTaskType := range p.DefaultForTaskTypes { + if defaultTaskType == tt { + if existingHandler, alreadyDefaulted := t.defaultPlugins[tt]; alreadyDefaulted { + logger.Warnf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", + tt, existingHandler.GetID(), cp.GetID()) + } + logger.Infof(ctx, "Plugin [%s] registered for TaskType [%s]", cp.GetID(), tt) + t.defaultPlugins[tt] = cp + } + } pluginsForTaskType, ok := t.pluginsForType[tt] if !ok { diff --git a/pkg/controller/nodes/task/handler_test.go b/pkg/controller/nodes/task/handler_test.go index 648893b39..3f17f468d 100644 --- a/pkg/controller/nodes/task/handler_test.go +++ b/pkg/controller/nodes/task/handler_test.go @@ -7,6 +7,8 @@ import ( "testing" "time" + pluginK8sMocks "github.com/lyft/flyteplugins/go/tasks/pluginmachinery/k8s/mocks" + "github.com/lyft/flyteidl/gen/pb-go/flyteidl/admin" "github.com/lyft/flytepropeller/pkg/apis/flyteworkflow/v1alpha1" @@ -28,7 +30,6 @@ import ( "github.com/lyft/flyteplugins/go/tasks/pluginmachinery/io" ioMocks "github.com/lyft/flyteplugins/go/tasks/pluginmachinery/io/mocks" pluginK8s "github.com/lyft/flyteplugins/go/tasks/pluginmachinery/k8s" - pluginK8sMocks "github.com/lyft/flyteplugins/go/tasks/pluginmachinery/k8s/mocks" "github.com/lyft/flytestdlib/promutils" "github.com/lyft/flytestdlib/storage" "github.com/stretchr/testify/assert" @@ -150,41 +151,59 @@ func Test_task_Setup(t *testing.T) { defaultPluginID string } tests := []struct { - name string - registry PluginRegistryIface - fields wantFields - wantErr bool + name string + registry PluginRegistryIface + enabledPluginsConfig map[string]config.EnabledPlugins + fields wantFields + wantErr bool }{ - {"no-plugins", testPluginRegistry{}, wantFields{}, false}, + {"no-plugins", testPluginRegistry{}, map[string]config.EnabledPlugins{}, wantFields{}, false}, {"no-default-only-core", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry}, k8s: []pluginK8s.PluginEntry{}, + }, map[string]config.EnabledPlugins{ + corePluginType: {DefaultPluginTasks: []string{corePluginType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType}, }, false}, {"no-default-only-k8s", testPluginRegistry{ core: []pluginCore.PluginEntry{}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry}, + }, map[string]config.EnabledPlugins{ + k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{k8sPluginType: k8sPluginType}, }, false}, - {"no-default", testPluginRegistry{ - core: []pluginCore.PluginEntry{corePluginEntry}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry}, + {"no-default", testPluginRegistry{}, map[string]config.EnabledPlugins{ + corePluginType: {DefaultPluginTasks: []string{corePluginType}}, + k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, }, wantFields{ - pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, k8sPluginType: k8sPluginType}, + pluginIDs: map[pluginCore.TaskType]string{}, }, false}, {"only-default-core", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry, corePluginEntryDefault}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry}, + }, map[string]config.EnabledPlugins{ + corePluginType: {DefaultPluginTasks: []string{corePluginType}}, + corePluginDefaultType: {DefaultPluginTasks: []string{corePluginDefaultType}}, + k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, corePluginDefaultType: corePluginDefaultType, k8sPluginType: k8sPluginType}, defaultPluginID: corePluginDefaultType, }, false}, {"only-default-k8s", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry}, k8s: []pluginK8s.PluginEntry{k8sPluginEntryDefault}, + }, map[string]config.EnabledPlugins{ + corePluginType: {DefaultPluginTasks: []string{corePluginType}}, + k8sPluginDefaultType: {DefaultPluginTasks: []string{k8sPluginDefaultType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, k8sPluginDefaultType: k8sPluginDefaultType}, defaultPluginID: k8sPluginDefaultType, }, false}, {"default-both", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry, corePluginEntryDefault}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry, k8sPluginEntryDefault}, + }, map[string]config.EnabledPlugins{ + corePluginType: {DefaultPluginTasks: []string{corePluginType}}, + corePluginDefaultType: {DefaultPluginTasks: []string{corePluginDefaultType}}, + k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, + k8sPluginDefaultType: {DefaultPluginTasks: []string{k8sPluginDefaultType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, corePluginDefaultType: corePluginDefaultType, k8sPluginType: k8sPluginType, k8sPluginDefaultType: k8sPluginDefaultType}, defaultPluginID: corePluginDefaultType, @@ -200,6 +219,7 @@ func Test_task_Setup(t *testing.T) { sCtx.On("MetricsScope").Return(promutils.NewTestScope()) tk, err := New(context.TODO(), mocks.NewFakeKubeClient(), &pluginCatalogMocks.Client{}, promutils.NewTestScope()) + tk.cfg.TaskPlugins.EnabledPlugins = tt.enabledPluginsConfig assert.NoError(t, err) tk.pluginRegistry = tt.registry if err := tk.Setup(context.TODO(), sCtx); err != nil { diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index 2f1372010..3bf3b57ca 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -8,7 +8,6 @@ import ( "github.com/lyft/flyteplugins/go/tasks/pluginmachinery/core" "github.com/lyft/flytestdlib/logger" - "k8s.io/apimachinery/pkg/util/sets" "github.com/lyft/flytepropeller/pkg/controller/nodes/task/config" "github.com/lyft/flytepropeller/pkg/controller/nodes/task/k8s" @@ -16,23 +15,25 @@ import ( func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPluginConfig, pr PluginRegistryIface) ([]core.PluginEntry, error) { allPluginsEnabled := false - enabledPlugins := sets.NewString() + enabledPlugins := make(map[string]config.EnabledPlugins) if cfg != nil { - enabledPlugins = cfg.GetEnabledPluginsSet() + enabledPlugins = cfg.GetEnabledPlugins() } - if enabledPlugins.Len() == 0 { + if len(enabledPlugins) == 0 { allPluginsEnabled = true } var finalizedPlugins []core.PluginEntry - logger.Infof(ctx, "Enabled plugins: %v", enabledPlugins.List()) + logger.Infof(ctx, "Enabled plugins: %+v", enabledPlugins) logger.Infof(ctx, "Loading core Plugins, plugin configuration [all plugins enabled: %v]", allPluginsEnabled) for _, cpe := range pr.GetCorePlugins() { id := strings.ToLower(cpe.ID) - if !allPluginsEnabled && !enabledPlugins.Has(id) { + pluginCfg, pluginEnabled := enabledPlugins[id] + if !allPluginsEnabled && !pluginEnabled { logger.Infof(ctx, "Plugin [%s] is DISABLED (not found in enabled plugins list).", id) } else { logger.Infof(ctx, "Plugin [%s] ENABLED", id) + cpe.DefaultForTaskTypes = pluginCfg.DefaultPluginTasks finalizedPlugins = append(finalizedPlugins, cpe) } } @@ -47,7 +48,8 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu for i := range k8sPlugins { kpe := k8sPlugins[i] id := strings.ToLower(kpe.ID) - if !allPluginsEnabled && !enabledPlugins.Has(id) { + pluginConfig, pluginEnabled := enabledPlugins[id] + if !allPluginsEnabled && !pluginEnabled { logger.Infof(ctx, "K8s Plugin [%s] is DISABLED (not found in enabled plugins list).", id) } else { logger.Infof(ctx, "K8s Plugin [%s] is ENABLED.", id) @@ -57,7 +59,8 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu LoadPlugin: func(ctx context.Context, iCtx core.SetupContext) (plugin core.Plugin, e error) { return k8s.NewPluginManagerWithBackOff(ctx, iCtx, kpe, backOffController, monitorIndex) }, - IsDefault: kpe.IsDefault, + IsDefault: kpe.IsDefault, + DefaultForTaskTypes: pluginConfig.DefaultPluginTasks, }) } } diff --git a/pkg/controller/nodes/task/plugin_config_test.go b/pkg/controller/nodes/task/plugin_config_test.go index fd28892ce..3a1417d94 100644 --- a/pkg/controller/nodes/task/plugin_config_test.go +++ b/pkg/controller/nodes/task/plugin_config_test.go @@ -47,13 +47,13 @@ func TestWranglePluginsAndGenerateFinalList(t *testing.T) { args args want want }{ - {"config-no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{coreContainer}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, + {"config-no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{coreContainer: {}}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, {"no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: nil}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, {"no-config-no-plugins", args{}, want{}}, {"no-config-plugins", args{corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{coreContainer, k8sOther}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, - {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{strings.ToUpper(coreContainer), strings.ToUpper(k8sOther)}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, + {"config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{coreContainer: {DefaultPluginTasks: []string{"container"}}, k8sOther: {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{strings.ToUpper(coreContainer): {DefaultPluginTasks: []string{"container"}}, strings.ToUpper(k8sOther): {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { From 70c9375aac65febe3a6a2d4cbe17e267c3b52057 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Wed, 21 Oct 2020 11:06:30 -0700 Subject: [PATCH 02/20] Update pkg/controller/nodes/task/config/config.go Co-authored-by: Haytham AbuelFutuh --- pkg/controller/nodes/task/config/config.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index 9e4123b4b..f423dd095 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -45,7 +45,7 @@ type BarrierConfig struct { } type EnabledPlugins struct { - DefaultPluginTasks []string `json:"default-plugins-tasks" pflag:",Tasks for which this plugin is the default implementation"` + DefaultForTaskTypes []string `json:"default-for-task-types" pflag:",Task types for which this plugin is the default handler."` } type TaskPluginConfig struct { From 7963b3e47eb40fd0f592bdf00cba8420c58b67a4 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Wed, 21 Oct 2020 11:06:39 -0700 Subject: [PATCH 03/20] Update pkg/controller/nodes/task/config/config.go Co-authored-by: Haytham AbuelFutuh --- pkg/controller/nodes/task/config/config.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index f423dd095..deb0e99fa 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -44,7 +44,7 @@ type BarrierConfig struct { CacheTTL config.Duration `json:"cache-ttl" pflag:", Max duration that a barrier would be respected if the process is not restarted. This should account for time required to store the record into persistent storage (across multiple rounds."` } -type EnabledPlugins struct { +type PluginConfig struct { DefaultForTaskTypes []string `json:"default-for-task-types" pflag:",Task types for which this plugin is the default handler."` } From b86f86161a977dd20ece55b563655929645ce67a Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Wed, 21 Oct 2020 15:41:59 -0700 Subject: [PATCH 04/20] sigh --- go.mod | 2 +- go.sum | 2 + pkg/controller/nodes/task/config/config.go | 16 +++---- .../nodes/task/config/config_flags.go | 1 - .../nodes/task/config/config_flags_test.go | 22 ---------- pkg/controller/nodes/task/handler_test.go | 42 +++++++++---------- pkg/controller/nodes/task/plugin_config.go | 6 +-- .../nodes/task/plugin_config_test.go | 8 ++-- 8 files changed, 39 insertions(+), 60 deletions(-) diff --git a/go.mod b/go.mod index eb2dfe6ce..4f8970c93 100644 --- a/go.mod +++ b/go.mod @@ -22,7 +22,7 @@ require ( github.com/imdario/mergo v0.3.8 // indirect github.com/lyft/datacatalog v0.2.1 github.com/lyft/flyteidl v0.18.9 - github.com/lyft/flyteplugins v0.5.14-0.20201019203656-dbcbf808fc35 + github.com/lyft/flyteplugins v0.5.14 github.com/lyft/flytestdlib v0.3.9 github.com/magiconair/properties v1.8.1 github.com/mattn/go-colorable v0.1.6 // indirect diff --git a/go.sum b/go.sum index 2a57c1fb3..935702f77 100644 --- a/go.sum +++ b/go.sum @@ -395,6 +395,8 @@ github.com/lyft/flyteplugins v0.5.12 h1:obz52m9dJ/ununeQJ2OcLZ1z38TE5KMDo8gyH/g8 github.com/lyft/flyteplugins v0.5.12/go.mod h1:UOoiW+rwQdrDDig3bJSxTWoyW8hW5bcqnuRsp+g2zhQ= github.com/lyft/flyteplugins v0.5.14-0.20201019203656-dbcbf808fc35 h1:JWOQEs0sJN/vyW2EhpQfK2VhuGwMcd4MU5JGB9OsWls= github.com/lyft/flyteplugins v0.5.14-0.20201019203656-dbcbf808fc35/go.mod h1:UOoiW+rwQdrDDig3bJSxTWoyW8hW5bcqnuRsp+g2zhQ= +github.com/lyft/flyteplugins v0.5.14 h1:YDJG7d7ZecT6xoe3Jrj99Q4HiCOg6btm9yjKz1LDuOk= +github.com/lyft/flyteplugins v0.5.14/go.mod h1:UOoiW+rwQdrDDig3bJSxTWoyW8hW5bcqnuRsp+g2zhQ= github.com/lyft/flytestdlib v0.3.0/go.mod h1:LJPPJlkFj+wwVWMrQT3K5JZgNhZi2mULsCG4ZYhinhU= github.com/lyft/flytestdlib v0.3.9 h1:NaKp9xkeWWwhVvqTOcR/FqlASy1N2gu/kN7PVe4S7YI= github.com/lyft/flytestdlib v0.3.9/go.mod h1:LJPPJlkFj+wwVWMrQT3K5JZgNhZi2mULsCG4ZYhinhU= diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index deb0e99fa..d9f03da18 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -13,7 +13,7 @@ const SectionKey = "tasks" var ( defaultConfig = &Config{ - TaskPlugins: TaskPluginConfig{EnabledPlugins: map[string]EnabledPlugins{}}, + TaskPlugins: TaskPluginConfig{EnabledPlugins: map[string]PluginConfig{}}, MaxPluginPhaseVersions: 100000, BarrierConfig: BarrierConfig{ Enabled: true, @@ -49,7 +49,7 @@ type PluginConfig struct { } type TaskPluginConfig struct { - EnabledPlugins map[string]EnabledPlugins `json:"enabled-plugins" pflag:",Plugins enabled currently"` + EnabledPlugins map[string]PluginConfig `pflag:"-,"` } type BackOffConfig struct { @@ -63,16 +63,16 @@ func cleanString(source string) string { return cleaned } -func (p TaskPluginConfig) GetEnabledPlugins() map[string]EnabledPlugins { - enabledPlugins := make(map[string]EnabledPlugins) +func (p TaskPluginConfig) GetEnabledPlugins() map[string]PluginConfig { + enabledPlugins := make(map[string]PluginConfig) for pluginName, info := range p.EnabledPlugins { - cleanedDefaultTasks := make([]string, 0, len(info.DefaultPluginTasks)) - for _, taskName := range info.DefaultPluginTasks { + cleanedDefaultTasks := make([]string, 0, len(info.DefaultForTaskTypes)) + for _, taskName := range info.DefaultForTaskTypes { cleanedDefaultTasks = append(cleanedDefaultTasks, cleanString(taskName)) } cleanedPluginName := cleanString(pluginName) - enabledPlugins[cleanedPluginName] = EnabledPlugins{ - DefaultPluginTasks: cleanedDefaultTasks, + enabledPlugins[cleanedPluginName] = PluginConfig{ + DefaultForTaskTypes: cleanedDefaultTasks, } } return enabledPlugins diff --git a/pkg/controller/nodes/task/config/config_flags.go b/pkg/controller/nodes/task/config/config_flags.go index 76707946d..5d572a3e9 100755 --- a/pkg/controller/nodes/task/config/config_flags.go +++ b/pkg/controller/nodes/task/config/config_flags.go @@ -41,7 +41,6 @@ func (Config) mustMarshalJSON(v json.Marshaler) string { // flags is json-name.json-sub-name... etc. func (cfg Config) GetPFlagSet(prefix string) *pflag.FlagSet { cmdFlags := pflag.NewFlagSet("Config", pflag.ExitOnError) - cmdFlags.StringSlice(fmt.Sprintf("%v%v", prefix, "task-plugins.enabled-plugins"), []string{}, "Plugins enabled currently") cmdFlags.Int32(fmt.Sprintf("%v%v", prefix, "max-plugin-phase-versions"), defaultConfig.MaxPluginPhaseVersions, "Maximum number of plugin phase versions allowed for one phase.") cmdFlags.Bool(fmt.Sprintf("%v%v", prefix, "barrier.enabled"), defaultConfig.BarrierConfig.Enabled, "Enable Barrier transitions using inmemory context") cmdFlags.Int(fmt.Sprintf("%v%v", prefix, "barrier.cache-size"), defaultConfig.BarrierConfig.CacheSize, "Max number of barrier to preserve in memory") diff --git a/pkg/controller/nodes/task/config/config_flags_test.go b/pkg/controller/nodes/task/config/config_flags_test.go index b5ebaf283..56a61224f 100755 --- a/pkg/controller/nodes/task/config/config_flags_test.go +++ b/pkg/controller/nodes/task/config/config_flags_test.go @@ -99,28 +99,6 @@ func TestConfig_SetFlags(t *testing.T) { cmdFlags := actual.GetPFlagSet("") assert.True(t, cmdFlags.HasFlags()) - t.Run("Test_task-plugins.enabled-plugins", func(t *testing.T) { - t.Run("DefaultValue", func(t *testing.T) { - // Test that default value is set properly - if vStringSlice, err := cmdFlags.GetStringSlice("task-plugins.enabled-plugins"); err == nil { - assert.Equal(t, []string([]string{}), vStringSlice) - } else { - assert.FailNow(t, err.Error()) - } - }) - - t.Run("Override", func(t *testing.T) { - testValue := join_Config("1,1", ",") - - cmdFlags.Set("task-plugins.enabled-plugins", testValue) - if vStringSlice, err := cmdFlags.GetStringSlice("task-plugins.enabled-plugins"); err == nil { - testDecodeSlice_Config(t, join_Config(vStringSlice, ","), &actual.TaskPlugins.EnabledPlugins) - - } else { - assert.FailNow(t, err.Error()) - } - }) - }) t.Run("Test_max-plugin-phase-versions", func(t *testing.T) { t.Run("DefaultValue", func(t *testing.T) { // Test that default value is set properly diff --git a/pkg/controller/nodes/task/handler_test.go b/pkg/controller/nodes/task/handler_test.go index 3f17f468d..872b5c99c 100644 --- a/pkg/controller/nodes/task/handler_test.go +++ b/pkg/controller/nodes/task/handler_test.go @@ -153,57 +153,57 @@ func Test_task_Setup(t *testing.T) { tests := []struct { name string registry PluginRegistryIface - enabledPluginsConfig map[string]config.EnabledPlugins + enabledPluginsConfig map[string]config.PluginConfig fields wantFields wantErr bool }{ - {"no-plugins", testPluginRegistry{}, map[string]config.EnabledPlugins{}, wantFields{}, false}, + {"no-plugins", testPluginRegistry{}, map[string]config.PluginConfig{}, wantFields{}, false}, {"no-default-only-core", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry}, k8s: []pluginK8s.PluginEntry{}, - }, map[string]config.EnabledPlugins{ - corePluginType: {DefaultPluginTasks: []string{corePluginType}}, + }, map[string]config.PluginConfig{ + corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType}, }, false}, {"no-default-only-k8s", testPluginRegistry{ core: []pluginCore.PluginEntry{}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry}, - }, map[string]config.EnabledPlugins{ - k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, + }, map[string]config.PluginConfig{ + k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{k8sPluginType: k8sPluginType}, }, false}, - {"no-default", testPluginRegistry{}, map[string]config.EnabledPlugins{ - corePluginType: {DefaultPluginTasks: []string{corePluginType}}, - k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, + {"no-default", testPluginRegistry{}, map[string]config.PluginConfig{ + corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, + k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{}, }, false}, {"only-default-core", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry, corePluginEntryDefault}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry}, - }, map[string]config.EnabledPlugins{ - corePluginType: {DefaultPluginTasks: []string{corePluginType}}, - corePluginDefaultType: {DefaultPluginTasks: []string{corePluginDefaultType}}, - k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, + }, map[string]config.PluginConfig{ + corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, + corePluginDefaultType: {DefaultForTaskTypes: []string{corePluginDefaultType}}, + k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, corePluginDefaultType: corePluginDefaultType, k8sPluginType: k8sPluginType}, defaultPluginID: corePluginDefaultType, }, false}, {"only-default-k8s", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry}, k8s: []pluginK8s.PluginEntry{k8sPluginEntryDefault}, - }, map[string]config.EnabledPlugins{ - corePluginType: {DefaultPluginTasks: []string{corePluginType}}, - k8sPluginDefaultType: {DefaultPluginTasks: []string{k8sPluginDefaultType}}, + }, map[string]config.PluginConfig{ + corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, + k8sPluginDefaultType: {DefaultForTaskTypes: []string{k8sPluginDefaultType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, k8sPluginDefaultType: k8sPluginDefaultType}, defaultPluginID: k8sPluginDefaultType, }, false}, {"default-both", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry, corePluginEntryDefault}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry, k8sPluginEntryDefault}, - }, map[string]config.EnabledPlugins{ - corePluginType: {DefaultPluginTasks: []string{corePluginType}}, - corePluginDefaultType: {DefaultPluginTasks: []string{corePluginDefaultType}}, - k8sPluginType: {DefaultPluginTasks: []string{k8sPluginType}}, - k8sPluginDefaultType: {DefaultPluginTasks: []string{k8sPluginDefaultType}}, + }, map[string]config.PluginConfig{ + corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, + corePluginDefaultType: {DefaultForTaskTypes: []string{corePluginDefaultType}}, + k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, + k8sPluginDefaultType: {DefaultForTaskTypes: []string{k8sPluginDefaultType}}, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, corePluginDefaultType: corePluginDefaultType, k8sPluginType: k8sPluginType, k8sPluginDefaultType: k8sPluginDefaultType}, defaultPluginID: corePluginDefaultType, diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index 3bf3b57ca..c53ac681e 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -15,7 +15,7 @@ import ( func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPluginConfig, pr PluginRegistryIface) ([]core.PluginEntry, error) { allPluginsEnabled := false - enabledPlugins := make(map[string]config.EnabledPlugins) + enabledPlugins := make(map[string]config.PluginConfig) if cfg != nil { enabledPlugins = cfg.GetEnabledPlugins() } @@ -33,7 +33,7 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu logger.Infof(ctx, "Plugin [%s] is DISABLED (not found in enabled plugins list).", id) } else { logger.Infof(ctx, "Plugin [%s] ENABLED", id) - cpe.DefaultForTaskTypes = pluginCfg.DefaultPluginTasks + cpe.DefaultForTaskTypes = pluginCfg.DefaultForTaskTypes finalizedPlugins = append(finalizedPlugins, cpe) } } @@ -60,7 +60,7 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu return k8s.NewPluginManagerWithBackOff(ctx, iCtx, kpe, backOffController, monitorIndex) }, IsDefault: kpe.IsDefault, - DefaultForTaskTypes: pluginConfig.DefaultPluginTasks, + DefaultForTaskTypes: pluginConfig.DefaultForTaskTypes, }) } } diff --git a/pkg/controller/nodes/task/plugin_config_test.go b/pkg/controller/nodes/task/plugin_config_test.go index 3a1417d94..76f81b400 100644 --- a/pkg/controller/nodes/task/plugin_config_test.go +++ b/pkg/controller/nodes/task/plugin_config_test.go @@ -47,13 +47,13 @@ func TestWranglePluginsAndGenerateFinalList(t *testing.T) { args args want want }{ - {"config-no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{coreContainer: {}}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, + {"config-no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{coreContainer: {}}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, {"no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: nil}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, {"no-config-no-plugins", args{}, want{}}, {"no-config-plugins", args{corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{coreContainer: {DefaultPluginTasks: []string{"container"}}, k8sOther: {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, - {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.EnabledPlugins{strings.ToUpper(coreContainer): {DefaultPluginTasks: []string{"container"}}, strings.ToUpper(k8sOther): {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, + {"config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{coreContainer: {DefaultForTaskTypes: []string{"container"}}, k8sOther: {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{strings.ToUpper(coreContainer): {DefaultForTaskTypes: []string{"container"}}, strings.ToUpper(k8sOther): {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { From e9167dffee6710ce3aa70796393010f9b350afc4 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Wed, 21 Oct 2020 16:38:15 -0700 Subject: [PATCH 05/20] big sigh --- config.yaml | 12 +++++++++--- pkg/controller/nodes/task/config/config.go | 2 +- pkg/controller/nodes/task/plugin_config.go | 12 ++++++++++-- pkg/controller/workflow/executor_test.go | 3 +++ 4 files changed, 23 insertions(+), 6 deletions(-) diff --git a/config.yaml b/config.yaml index eca2d0923..775bf9065 100644 --- a/config.yaml +++ b/config.yaml @@ -29,9 +29,15 @@ propeller: tasks: task-plugins: enabled-plugins: - - container - - K8S-ARRAY - - qubole-hive-executor + - container: + - "default-for-task-types": + - "container" + - K8S-ARRAY: + - "default-for-task-types": + - "K8S-ARRAY" + - qubole-hive-executor: + - "default-for-task-types": + - "qubole-hive-executor" # Uncomment to enable sagemaker plugin # - sagemaker_training # - sagemaker_hyperparameter_tuning diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index d9f03da18..2887f52f3 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -49,7 +49,7 @@ type PluginConfig struct { } type TaskPluginConfig struct { - EnabledPlugins map[string]PluginConfig `pflag:"-,"` + EnabledPlugins map[string]PluginConfig `json:"enabled-plugins" pflag:"-,"` } type BackOffConfig struct { diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index c53ac681e..cdd12a367 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -31,6 +31,10 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu pluginCfg, pluginEnabled := enabledPlugins[id] if !allPluginsEnabled && !pluginEnabled { logger.Infof(ctx, "Plugin [%s] is DISABLED (not found in enabled plugins list).", id) + } else if allPluginsEnabled { + logger.Infof(ctx, "Plugin [%s] ENABLED", id) + cpe.DefaultForTaskTypes = cpe.RegisteredTaskTypes + finalizedPlugins = append(finalizedPlugins, cpe) } else { logger.Infof(ctx, "Plugin [%s] ENABLED", id) cpe.DefaultForTaskTypes = pluginCfg.DefaultForTaskTypes @@ -53,7 +57,7 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu logger.Infof(ctx, "K8s Plugin [%s] is DISABLED (not found in enabled plugins list).", id) } else { logger.Infof(ctx, "K8s Plugin [%s] is ENABLED.", id) - finalizedPlugins = append(finalizedPlugins, core.PluginEntry{ + plugin := core.PluginEntry{ ID: id, RegisteredTaskTypes: kpe.RegisteredTaskTypes, LoadPlugin: func(ctx context.Context, iCtx core.SetupContext) (plugin core.Plugin, e error) { @@ -61,7 +65,11 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu }, IsDefault: kpe.IsDefault, DefaultForTaskTypes: pluginConfig.DefaultForTaskTypes, - }) + } + if allPluginsEnabled { + plugin.DefaultForTaskTypes = plugin.RegisteredTaskTypes + } + finalizedPlugins = append(finalizedPlugins, plugin) } } return finalizedPlugins, nil diff --git a/pkg/controller/workflow/executor_test.go b/pkg/controller/workflow/executor_test.go index 1c0f4e1ce..61c910c54 100644 --- a/pkg/controller/workflow/executor_test.go +++ b/pkg/controller/workflow/executor_test.go @@ -155,6 +155,7 @@ func createHappyPathTaskExecutor(t assert.TestingT, enableAsserts bool) pluginCo RegisteredTaskTypes: []string{"7"}, LoadPlugin: f, IsDefault: true, + DefaultForTaskTypes: []string{"7"}, } } @@ -184,6 +185,7 @@ func createFailingTaskExecutor(t assert.TestingT) pluginCore.PluginEntry { RegisteredTaskTypes: []string{"7"}, LoadPlugin: f, IsDefault: true, + DefaultForTaskTypes: []string{"7"}, } } @@ -213,6 +215,7 @@ func createTaskExecutorErrorInCheck(t assert.TestingT) pluginCore.PluginEntry { RegisteredTaskTypes: []string{"7"}, LoadPlugin: f, IsDefault: true, + DefaultForTaskTypes: []string{"7"}, } } From 08d6c89d0cfbde789c47c45aba8267cec4eefa93 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Wed, 21 Oct 2020 16:41:15 -0700 Subject: [PATCH 06/20] smaller sigh --- config.yaml | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/config.yaml b/config.yaml index 775bf9065..f310ba03d 100644 --- a/config.yaml +++ b/config.yaml @@ -30,14 +30,14 @@ tasks: task-plugins: enabled-plugins: - container: - - "default-for-task-types": - - "container" + default-for-task-types: + - "container" - K8S-ARRAY: - - "default-for-task-types": - - "K8S-ARRAY" + default-for-task-types: + - "K8S-ARRAY" - qubole-hive-executor: - - "default-for-task-types": - - "qubole-hive-executor" + default-for-task-types: + - "qubole-hive-executor" # Uncomment to enable sagemaker plugin # - sagemaker_training # - sagemaker_hyperparameter_tuning From c5388330a2d2c5117be384ac5899d249dc2f73d1 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Thu, 22 Oct 2020 09:57:44 -0700 Subject: [PATCH 07/20] one more comment --- pkg/controller/nodes/task/handler.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index 866f4c61e..f8f7ffb4e 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -221,7 +221,7 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { for _, defaultTaskType := range p.DefaultForTaskTypes { if defaultTaskType == tt { if existingHandler, alreadyDefaulted := t.defaultPlugins[tt]; alreadyDefaulted { - logger.Warnf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", + logger.Panicf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", tt, existingHandler.GetID(), cp.GetID()) } logger.Infof(ctx, "Plugin [%s] registered for TaskType [%s]", cp.GetID(), tt) From 3b69ab0520f8c1fc58bd3fdacd2c77cc529d99f0 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Thu, 22 Oct 2020 10:35:08 -0700 Subject: [PATCH 08/20] Undo --- config.yaml | 2 +- pkg/controller/nodes/task/config/config.go | 7 +++--- .../nodes/task/config/config_flags.go | 1 + .../nodes/task/config/config_flags_test.go | 22 +++++++++++++++++++ pkg/controller/nodes/task/handler.go | 2 +- pkg/controller/nodes/task/handler_test.go | 12 +++++----- .../nodes/task/plugin_config_test.go | 10 ++++----- pkg/controller/workflow/executor_test.go | 3 --- 8 files changed, 40 insertions(+), 19 deletions(-) diff --git a/config.yaml b/config.yaml index f310ba03d..89763eb85 100644 --- a/config.yaml +++ b/config.yaml @@ -28,7 +28,7 @@ propeller: policy: "ResourceVersionCache" tasks: task-plugins: - enabled-plugins: + plugins-config: - container: default-for-task-types: - "container" diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index 2887f52f3..e5400c702 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -13,7 +13,7 @@ const SectionKey = "tasks" var ( defaultConfig = &Config{ - TaskPlugins: TaskPluginConfig{EnabledPlugins: map[string]PluginConfig{}}, + TaskPlugins: TaskPluginConfig{PluginConfigs: map[string]PluginConfig{}}, MaxPluginPhaseVersions: 100000, BarrierConfig: BarrierConfig{ Enabled: true, @@ -49,7 +49,8 @@ type PluginConfig struct { } type TaskPluginConfig struct { - EnabledPlugins map[string]PluginConfig `json:"enabled-plugins" pflag:"-,"` + EnabledPlugins []string `json:"enabled-plugins" pflag:",Plugins enabled currently"` + PluginConfigs map[string]PluginConfig `json:"plugins-config" pflag:"-,"` } type BackOffConfig struct { @@ -65,7 +66,7 @@ func cleanString(source string) string { func (p TaskPluginConfig) GetEnabledPlugins() map[string]PluginConfig { enabledPlugins := make(map[string]PluginConfig) - for pluginName, info := range p.EnabledPlugins { + for pluginName, info := range p.PluginConfigs { cleanedDefaultTasks := make([]string, 0, len(info.DefaultForTaskTypes)) for _, taskName := range info.DefaultForTaskTypes { cleanedDefaultTasks = append(cleanedDefaultTasks, cleanString(taskName)) diff --git a/pkg/controller/nodes/task/config/config_flags.go b/pkg/controller/nodes/task/config/config_flags.go index 5d572a3e9..76707946d 100755 --- a/pkg/controller/nodes/task/config/config_flags.go +++ b/pkg/controller/nodes/task/config/config_flags.go @@ -41,6 +41,7 @@ func (Config) mustMarshalJSON(v json.Marshaler) string { // flags is json-name.json-sub-name... etc. func (cfg Config) GetPFlagSet(prefix string) *pflag.FlagSet { cmdFlags := pflag.NewFlagSet("Config", pflag.ExitOnError) + cmdFlags.StringSlice(fmt.Sprintf("%v%v", prefix, "task-plugins.enabled-plugins"), []string{}, "Plugins enabled currently") cmdFlags.Int32(fmt.Sprintf("%v%v", prefix, "max-plugin-phase-versions"), defaultConfig.MaxPluginPhaseVersions, "Maximum number of plugin phase versions allowed for one phase.") cmdFlags.Bool(fmt.Sprintf("%v%v", prefix, "barrier.enabled"), defaultConfig.BarrierConfig.Enabled, "Enable Barrier transitions using inmemory context") cmdFlags.Int(fmt.Sprintf("%v%v", prefix, "barrier.cache-size"), defaultConfig.BarrierConfig.CacheSize, "Max number of barrier to preserve in memory") diff --git a/pkg/controller/nodes/task/config/config_flags_test.go b/pkg/controller/nodes/task/config/config_flags_test.go index 56a61224f..b5ebaf283 100755 --- a/pkg/controller/nodes/task/config/config_flags_test.go +++ b/pkg/controller/nodes/task/config/config_flags_test.go @@ -99,6 +99,28 @@ func TestConfig_SetFlags(t *testing.T) { cmdFlags := actual.GetPFlagSet("") assert.True(t, cmdFlags.HasFlags()) + t.Run("Test_task-plugins.enabled-plugins", func(t *testing.T) { + t.Run("DefaultValue", func(t *testing.T) { + // Test that default value is set properly + if vStringSlice, err := cmdFlags.GetStringSlice("task-plugins.enabled-plugins"); err == nil { + assert.Equal(t, []string([]string{}), vStringSlice) + } else { + assert.FailNow(t, err.Error()) + } + }) + + t.Run("Override", func(t *testing.T) { + testValue := join_Config("1,1", ",") + + cmdFlags.Set("task-plugins.enabled-plugins", testValue) + if vStringSlice, err := cmdFlags.GetStringSlice("task-plugins.enabled-plugins"); err == nil { + testDecodeSlice_Config(t, join_Config(vStringSlice, ","), &actual.TaskPlugins.EnabledPlugins) + + } else { + assert.FailNow(t, err.Error()) + } + }) + }) t.Run("Test_max-plugin-phase-versions", func(t *testing.T) { t.Run("DefaultValue", func(t *testing.T) { // Test that default value is set properly diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index f8f7ffb4e..866f4c61e 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -221,7 +221,7 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { for _, defaultTaskType := range p.DefaultForTaskTypes { if defaultTaskType == tt { if existingHandler, alreadyDefaulted := t.defaultPlugins[tt]; alreadyDefaulted { - logger.Panicf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", + logger.Warnf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", tt, existingHandler.GetID(), cp.GetID()) } logger.Infof(ctx, "Plugin [%s] registered for TaskType [%s]", cp.GetID(), tt) diff --git a/pkg/controller/nodes/task/handler_test.go b/pkg/controller/nodes/task/handler_test.go index 872b5c99c..127bf1aef 100644 --- a/pkg/controller/nodes/task/handler_test.go +++ b/pkg/controller/nodes/task/handler_test.go @@ -151,11 +151,11 @@ func Test_task_Setup(t *testing.T) { defaultPluginID string } tests := []struct { - name string - registry PluginRegistryIface - enabledPluginsConfig map[string]config.PluginConfig - fields wantFields - wantErr bool + name string + registry PluginRegistryIface + pluginsConfig map[string]config.PluginConfig + fields wantFields + wantErr bool }{ {"no-plugins", testPluginRegistry{}, map[string]config.PluginConfig{}, wantFields{}, false}, {"no-default-only-core", testPluginRegistry{ @@ -219,7 +219,7 @@ func Test_task_Setup(t *testing.T) { sCtx.On("MetricsScope").Return(promutils.NewTestScope()) tk, err := New(context.TODO(), mocks.NewFakeKubeClient(), &pluginCatalogMocks.Client{}, promutils.NewTestScope()) - tk.cfg.TaskPlugins.EnabledPlugins = tt.enabledPluginsConfig + tk.cfg.TaskPlugins.PluginConfigs = tt.pluginsConfig assert.NoError(t, err) tk.pluginRegistry = tt.registry if err := tk.Setup(context.TODO(), sCtx); err != nil { diff --git a/pkg/controller/nodes/task/plugin_config_test.go b/pkg/controller/nodes/task/plugin_config_test.go index 76f81b400..7705a21e8 100644 --- a/pkg/controller/nodes/task/plugin_config_test.go +++ b/pkg/controller/nodes/task/plugin_config_test.go @@ -47,13 +47,13 @@ func TestWranglePluginsAndGenerateFinalList(t *testing.T) { args args want want }{ - {"config-no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{coreContainer: {}}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, - {"no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: nil}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, + {"config-no-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{coreContainer: {}}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, + {"no-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: nil}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, {"no-config-no-plugins", args{}, want{}}, {"no-config-plugins", args{corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{coreContainer: {DefaultForTaskTypes: []string{"container"}}, k8sOther: {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, - {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: map[string]config.PluginConfig{strings.ToUpper(coreContainer): {DefaultForTaskTypes: []string{"container"}}, strings.ToUpper(k8sOther): {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, + {"config-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{coreContainer: {DefaultForTaskTypes: []string{"container"}}, k8sOther: {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{strings.ToUpper(coreContainer): {DefaultForTaskTypes: []string{"container"}}, strings.ToUpper(k8sOther): {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { diff --git a/pkg/controller/workflow/executor_test.go b/pkg/controller/workflow/executor_test.go index 61c910c54..1c0f4e1ce 100644 --- a/pkg/controller/workflow/executor_test.go +++ b/pkg/controller/workflow/executor_test.go @@ -155,7 +155,6 @@ func createHappyPathTaskExecutor(t assert.TestingT, enableAsserts bool) pluginCo RegisteredTaskTypes: []string{"7"}, LoadPlugin: f, IsDefault: true, - DefaultForTaskTypes: []string{"7"}, } } @@ -185,7 +184,6 @@ func createFailingTaskExecutor(t assert.TestingT) pluginCore.PluginEntry { RegisteredTaskTypes: []string{"7"}, LoadPlugin: f, IsDefault: true, - DefaultForTaskTypes: []string{"7"}, } } @@ -215,7 +213,6 @@ func createTaskExecutorErrorInCheck(t assert.TestingT) pluginCore.PluginEntry { RegisteredTaskTypes: []string{"7"}, LoadPlugin: f, IsDefault: true, - DefaultForTaskTypes: []string{"7"}, } } From c7f38612a090659ab96c7033b4de8dbfa251bebd Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Thu, 22 Oct 2020 10:47:20 -0700 Subject: [PATCH 09/20] deprecate --- pkg/controller/nodes/task/config/config.go | 3 ++- pkg/controller/nodes/task/config/config_flags.go | 2 +- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index e5400c702..d5b2130bc 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -49,7 +49,8 @@ type PluginConfig struct { } type TaskPluginConfig struct { - EnabledPlugins []string `json:"enabled-plugins" pflag:",Plugins enabled currently"` + // DEPRECATED. Use PluginConfigs instead. + EnabledPlugins []string `json:"enabled-plugins" pflag:",deprecated"` PluginConfigs map[string]PluginConfig `json:"plugins-config" pflag:"-,"` } diff --git a/pkg/controller/nodes/task/config/config_flags.go b/pkg/controller/nodes/task/config/config_flags.go index 76707946d..b7e8ad848 100755 --- a/pkg/controller/nodes/task/config/config_flags.go +++ b/pkg/controller/nodes/task/config/config_flags.go @@ -41,7 +41,7 @@ func (Config) mustMarshalJSON(v json.Marshaler) string { // flags is json-name.json-sub-name... etc. func (cfg Config) GetPFlagSet(prefix string) *pflag.FlagSet { cmdFlags := pflag.NewFlagSet("Config", pflag.ExitOnError) - cmdFlags.StringSlice(fmt.Sprintf("%v%v", prefix, "task-plugins.enabled-plugins"), []string{}, "Plugins enabled currently") + cmdFlags.StringSlice(fmt.Sprintf("%v%v", prefix, "task-plugins.enabled-plugins"), []string{}, "deprecated") cmdFlags.Int32(fmt.Sprintf("%v%v", prefix, "max-plugin-phase-versions"), defaultConfig.MaxPluginPhaseVersions, "Maximum number of plugin phase versions allowed for one phase.") cmdFlags.Bool(fmt.Sprintf("%v%v", prefix, "barrier.enabled"), defaultConfig.BarrierConfig.Enabled, "Enable Barrier transitions using inmemory context") cmdFlags.Int(fmt.Sprintf("%v%v", prefix, "barrier.cache-size"), defaultConfig.BarrierConfig.CacheSize, "Max number of barrier to preserve in memory") From 3d1af75ab0f0788ed7e6bd529df072faaba30ef5 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Fri, 23 Oct 2020 11:48:46 -0700 Subject: [PATCH 10/20] offline discussion update --- config.yaml | 17 ++--- pkg/controller/nodes/task/config/config.go | 26 ++++---- pkg/controller/nodes/task/handler.go | 4 +- pkg/controller/nodes/task/handler_test.go | 66 ++++++++++--------- .../nodes/task/plugin_config_test.go | 10 +-- 5 files changed, 62 insertions(+), 61 deletions(-) diff --git a/config.yaml b/config.yaml index 89763eb85..64d2a90fe 100644 --- a/config.yaml +++ b/config.yaml @@ -28,16 +28,13 @@ propeller: policy: "ResourceVersionCache" tasks: task-plugins: - plugins-config: - - container: - default-for-task-types: - - "container" - - K8S-ARRAY: - default-for-task-types: - - "K8S-ARRAY" - - qubole-hive-executor: - default-for-task-types: - - "qubole-hive-executor" + enabled-plugins: + - container + - K8S-ARRAY + - qubole-hive-executor + default-for-task-type: + - container-array: k8s-array + - presto: my-presto # Uncomment to enable sagemaker plugin # - sagemaker_training # - sagemaker_hyperparameter_tuning diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index d5b2130bc..e91208181 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -13,7 +13,7 @@ const SectionKey = "tasks" var ( defaultConfig = &Config{ - TaskPlugins: TaskPluginConfig{PluginConfigs: map[string]PluginConfig{}}, + TaskPlugins: TaskPluginConfig{EnabledPlugins: []string{}, DefaultForTaskTypes: map[string]string{}}, MaxPluginPhaseVersions: 100000, BarrierConfig: BarrierConfig{ Enabled: true, @@ -44,14 +44,10 @@ type BarrierConfig struct { CacheTTL config.Duration `json:"cache-ttl" pflag:", Max duration that a barrier would be respected if the process is not restarted. This should account for time required to store the record into persistent storage (across multiple rounds."` } -type PluginConfig struct { - DefaultForTaskTypes []string `json:"default-for-task-types" pflag:",Task types for which this plugin is the default handler."` -} - type TaskPluginConfig struct { - // DEPRECATED. Use PluginConfigs instead. - EnabledPlugins []string `json:"enabled-plugins" pflag:",deprecated"` - PluginConfigs map[string]PluginConfig `json:"plugins-config" pflag:"-,"` + EnabledPlugins []string `json:"enabled-plugins" pflag:",deprecated"` + // Maps task types to their plugin handler (by ID). + DefaultForTaskTypes map[string]string `json:"default-for-task-types" pflag:"-,"` } type BackOffConfig struct { @@ -59,6 +55,10 @@ type BackOffConfig struct { MaxDuration config.Duration `json:"max-duration" pflag:",The cap of the backoff duration"` } +type PluginConfig struct { + DefaultForTaskTypes []string +} + func cleanString(source string) string { cleaned := strings.Trim(source, " ") cleaned = strings.ToLower(cleaned) @@ -67,10 +67,12 @@ func cleanString(source string) string { func (p TaskPluginConfig) GetEnabledPlugins() map[string]PluginConfig { enabledPlugins := make(map[string]PluginConfig) - for pluginName, info := range p.PluginConfigs { - cleanedDefaultTasks := make([]string, 0, len(info.DefaultForTaskTypes)) - for _, taskName := range info.DefaultForTaskTypes { - cleanedDefaultTasks = append(cleanedDefaultTasks, cleanString(taskName)) + for _, pluginName := range p.EnabledPlugins { + cleanedDefaultTasks := make([]string, 0) + for taskName, taskPluginName := range p.DefaultForTaskTypes { + if taskPluginName == pluginName { + cleanedDefaultTasks = append(cleanedDefaultTasks, cleanString(taskName)) + } } cleanedPluginName := cleanString(pluginName) enabledPlugins[cleanedPluginName] = PluginConfig{ diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index 866f4c61e..246dbca91 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -220,8 +220,8 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { for _, tt := range p.RegisteredTaskTypes { for _, defaultTaskType := range p.DefaultForTaskTypes { if defaultTaskType == tt { - if existingHandler, alreadyDefaulted := t.defaultPlugins[tt]; alreadyDefaulted { - logger.Warnf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", + if existingHandler, alreadyDefaulted := t.defaultPlugins[tt]; alreadyDefaulted && existingHandler.GetID() != cp.GetID() { + logger.Panicf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", tt, existingHandler.GetID(), cp.GetID()) } logger.Infof(ctx, "Plugin [%s] registered for TaskType [%s]", cp.GetID(), tt) diff --git a/pkg/controller/nodes/task/handler_test.go b/pkg/controller/nodes/task/handler_test.go index 127bf1aef..85ab2fb36 100644 --- a/pkg/controller/nodes/task/handler_test.go +++ b/pkg/controller/nodes/task/handler_test.go @@ -151,59 +151,60 @@ func Test_task_Setup(t *testing.T) { defaultPluginID string } tests := []struct { - name string - registry PluginRegistryIface - pluginsConfig map[string]config.PluginConfig - fields wantFields - wantErr bool + name string + registry PluginRegistryIface + enabledPlugins []string + defaultForTaskTypes map[string]string + fields wantFields + wantErr bool }{ - {"no-plugins", testPluginRegistry{}, map[string]config.PluginConfig{}, wantFields{}, false}, + {"no-plugins", testPluginRegistry{}, []string{}, map[string]string{}, wantFields{}, false}, {"no-default-only-core", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry}, k8s: []pluginK8s.PluginEntry{}, - }, map[string]config.PluginConfig{ - corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, - }, wantFields{ - pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType}, - }, false}, + }, []string{corePluginType}, map[string]string{ + corePluginType: corePluginType}, + wantFields{ + pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType}, + }, false}, {"no-default-only-k8s", testPluginRegistry{ core: []pluginCore.PluginEntry{}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry}, - }, map[string]config.PluginConfig{ - k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, - }, wantFields{ - pluginIDs: map[pluginCore.TaskType]string{k8sPluginType: k8sPluginType}, - }, false}, - {"no-default", testPluginRegistry{}, map[string]config.PluginConfig{ - corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, - k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, + }, []string{k8sPluginType}, map[string]string{ + k8sPluginType: k8sPluginType}, + wantFields{ + pluginIDs: map[pluginCore.TaskType]string{k8sPluginType: k8sPluginType}, + }, false}, + {"no-default", testPluginRegistry{}, []string{corePluginType, k8sPluginType}, map[string]string{ + corePluginType: corePluginType, + k8sPluginType: k8sPluginType, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{}, }, false}, {"only-default-core", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry, corePluginEntryDefault}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry}, - }, map[string]config.PluginConfig{ - corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, - corePluginDefaultType: {DefaultForTaskTypes: []string{corePluginDefaultType}}, - k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, + }, []string{corePluginType, corePluginDefaultType, k8sPluginType}, map[string]string{ + corePluginType: corePluginType, + corePluginDefaultType: corePluginDefaultType, + k8sPluginType: k8sPluginType, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, corePluginDefaultType: corePluginDefaultType, k8sPluginType: k8sPluginType}, defaultPluginID: corePluginDefaultType, }, false}, {"only-default-k8s", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry}, k8s: []pluginK8s.PluginEntry{k8sPluginEntryDefault}, - }, map[string]config.PluginConfig{ - corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, - k8sPluginDefaultType: {DefaultForTaskTypes: []string{k8sPluginDefaultType}}, + }, []string{corePluginType, k8sPluginDefaultType}, map[string]string{ + corePluginType: corePluginType, + k8sPluginDefaultType: k8sPluginDefaultType, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, k8sPluginDefaultType: k8sPluginDefaultType}, defaultPluginID: k8sPluginDefaultType, }, false}, {"default-both", testPluginRegistry{ core: []pluginCore.PluginEntry{corePluginEntry, corePluginEntryDefault}, k8s: []pluginK8s.PluginEntry{k8sPluginEntry, k8sPluginEntryDefault}, - }, map[string]config.PluginConfig{ - corePluginType: {DefaultForTaskTypes: []string{corePluginType}}, - corePluginDefaultType: {DefaultForTaskTypes: []string{corePluginDefaultType}}, - k8sPluginType: {DefaultForTaskTypes: []string{k8sPluginType}}, - k8sPluginDefaultType: {DefaultForTaskTypes: []string{k8sPluginDefaultType}}, + }, []string{corePluginType, corePluginDefaultType, k8sPluginType, k8sPluginDefaultType}, map[string]string{ + corePluginType: corePluginType, + corePluginDefaultType: corePluginDefaultType, + k8sPluginType: k8sPluginType, + k8sPluginDefaultType: k8sPluginDefaultType, }, wantFields{ pluginIDs: map[pluginCore.TaskType]string{corePluginType: corePluginType, corePluginDefaultType: corePluginDefaultType, k8sPluginType: k8sPluginType, k8sPluginDefaultType: k8sPluginDefaultType}, defaultPluginID: corePluginDefaultType, @@ -219,7 +220,8 @@ func Test_task_Setup(t *testing.T) { sCtx.On("MetricsScope").Return(promutils.NewTestScope()) tk, err := New(context.TODO(), mocks.NewFakeKubeClient(), &pluginCatalogMocks.Client{}, promutils.NewTestScope()) - tk.cfg.TaskPlugins.PluginConfigs = tt.pluginsConfig + tk.cfg.TaskPlugins.EnabledPlugins = tt.enabledPlugins + tk.cfg.TaskPlugins.DefaultForTaskTypes = tt.defaultForTaskTypes assert.NoError(t, err) tk.pluginRegistry = tt.registry if err := tk.Setup(context.TODO(), sCtx); err != nil { diff --git a/pkg/controller/nodes/task/plugin_config_test.go b/pkg/controller/nodes/task/plugin_config_test.go index 7705a21e8..41bc0757a 100644 --- a/pkg/controller/nodes/task/plugin_config_test.go +++ b/pkg/controller/nodes/task/plugin_config_test.go @@ -47,13 +47,13 @@ func TestWranglePluginsAndGenerateFinalList(t *testing.T) { args args want want }{ - {"config-no-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{coreContainer: {}}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, - {"no-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: nil}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, + {"config-no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{coreContainer}}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, + {"no-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: nil}, backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{}}, {"no-config-no-plugins", args{}, want{}}, {"no-config-plugins", args{corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, - {"config-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{coreContainer: {DefaultForTaskTypes: []string{"container"}}, k8sOther: {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, - {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{PluginConfigs: map[string]config.PluginConfig{strings.ToUpper(coreContainer): {DefaultForTaskTypes: []string{"container"}}, strings.ToUpper(k8sOther): {}}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"empty-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin)}, want{final: sets.NewString(k8sContainer, k8sOther, coreOther, coreContainer)}}, + {"config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{coreContainer, k8sOther}, DefaultForTaskTypes: map[string]string{"container": coreContainer}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, + {"case-differs-config-plugins", args{cfg: &config.TaskPluginConfig{EnabledPlugins: []string{strings.ToUpper(coreContainer), strings.ToUpper(k8sOther)}, DefaultForTaskTypes: map[string]string{"container": coreContainer}}, corePlugins: cpe(coreContainerPlugin, coreOtherPlugin), k8sPlugins: kpe(k8sContainerPlugin, k8sOtherPlugin), backOffCfg: &config.BackOffConfig{BaseSecond: 0, MaxDuration: config2.Duration{Duration: time.Second * 0}}}, want{final: sets.NewString(k8sOther, coreContainer)}}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { From b5b0d6278cddef086b63ffd5523b587bf16e416c Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Fri, 23 Oct 2020 14:15:50 -0700 Subject: [PATCH 11/20] cleanup --- pkg/controller/nodes/task/config/config.go | 14 +++++++------- pkg/controller/nodes/task/handler.go | 2 -- 2 files changed, 7 insertions(+), 9 deletions(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index e91208181..2536fd7e0 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -67,16 +67,16 @@ func cleanString(source string) string { func (p TaskPluginConfig) GetEnabledPlugins() map[string]PluginConfig { enabledPlugins := make(map[string]PluginConfig) + pluginDefaultForTaskType := map[string][]string{} + // Reverse the map. Having the config use task type as a key guarantees only one default plugin can be specified per + // task type but now we need to sort for which tasks a plugin needs to be the default. + for taskName, pluginName := range p.DefaultForTaskTypes { + pluginDefaultForTaskType[pluginName] = append(pluginDefaultForTaskType[pluginName], cleanString(taskName)) + } for _, pluginName := range p.EnabledPlugins { - cleanedDefaultTasks := make([]string, 0) - for taskName, taskPluginName := range p.DefaultForTaskTypes { - if taskPluginName == pluginName { - cleanedDefaultTasks = append(cleanedDefaultTasks, cleanString(taskName)) - } - } cleanedPluginName := cleanString(pluginName) enabledPlugins[cleanedPluginName] = PluginConfig{ - DefaultForTaskTypes: cleanedDefaultTasks, + DefaultForTaskTypes: pluginDefaultForTaskType[pluginName], } } return enabledPlugins diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index 246dbca91..c30382a82 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -215,8 +215,6 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { if err != nil { return regErrors.Wrapf(err, "failed to load plugin - %s", p.ID) } - println(fmt.Sprintf("for plugin [%s], registered task types: [%+v] and default task types [%+v]", - p.ID, p.RegisteredTaskTypes, p.DefaultForTaskTypes)) for _, tt := range p.RegisteredTaskTypes { for _, defaultTaskType := range p.DefaultForTaskTypes { if defaultTaskType == tt { From 2e507c92dd395e457d61e08b08b3c8427a5026bd Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Tue, 27 Oct 2020 11:45:52 -0700 Subject: [PATCH 12/20] more review --- pkg/controller/nodes/task/handler.go | 24 ++++++++++++++++++++++ pkg/controller/nodes/task/plugin_config.go | 7 ------- 2 files changed, 24 insertions(+), 7 deletions(-) diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index c30382a82..92d214ce1 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -204,6 +204,10 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { return err } + // Not every task type will have a default plugin specified in the flytepropeller config. + // That's fine, we resort to using the plugins' static RegisteredTaskTypes as a fallback. + fallbackTaskHandlerMap := make(map[string]map[string]pluginCore.Plugin) + for _, p := range enabledPlugins { // create a new resource registrar proxy for each plugin, and pass it into the plugin's LoadPlugin() via a setup context pluginResourceNamespacePrefix := pluginCore.ResourceNamespace(newResourceManagerBuilder.GetID()).CreateSubNamespace(pluginCore.ResourceNamespace(p.ID)) @@ -233,6 +237,13 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { } pluginsForTaskType[cp.GetID()] = cp t.pluginsForType[tt] = pluginsForTaskType + + fallbackMap, ok := fallbackTaskHandlerMap[tt] + if !ok { + fallbackMap = make(map[string]pluginCore.Plugin) + } + fallbackMap[cp.GetID()] = cp + fallbackTaskHandlerMap[tt] = fallbackMap } if p.IsDefault { if err := t.setDefault(ctx, cp); err != nil { @@ -241,6 +252,19 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { } } + // Read from the fallback task handler map for any remaining tasks without a defaultPlugins registered handler. + for taskType, registeredPlugins := range fallbackTaskHandlerMap { + if _, ok := t.defaultPlugins[taskType]; ok { + break + } + if len(registeredPlugins) != 1 { + logger.Panicf(ctx, "Multiple plugins registered to handle task type: %s. ([%+v])", taskType, registeredPlugins) + } + for _, plugin := range registeredPlugins { + t.defaultPlugins[taskType] = plugin + } + } + rm, err := newResourceManagerBuilder.BuildResourceManager(ctx) if err != nil { logger.Errorf(ctx, "Failed to build a resource manager") diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index cdd12a367..61abb39ec 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -31,10 +31,6 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu pluginCfg, pluginEnabled := enabledPlugins[id] if !allPluginsEnabled && !pluginEnabled { logger.Infof(ctx, "Plugin [%s] is DISABLED (not found in enabled plugins list).", id) - } else if allPluginsEnabled { - logger.Infof(ctx, "Plugin [%s] ENABLED", id) - cpe.DefaultForTaskTypes = cpe.RegisteredTaskTypes - finalizedPlugins = append(finalizedPlugins, cpe) } else { logger.Infof(ctx, "Plugin [%s] ENABLED", id) cpe.DefaultForTaskTypes = pluginCfg.DefaultForTaskTypes @@ -66,9 +62,6 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu IsDefault: kpe.IsDefault, DefaultForTaskTypes: pluginConfig.DefaultForTaskTypes, } - if allPluginsEnabled { - plugin.DefaultForTaskTypes = plugin.RegisteredTaskTypes - } finalizedPlugins = append(finalizedPlugins, plugin) } } From 2dafb0b4b158e486b7d2768f773a871a9d5ad31c Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Fri, 30 Oct 2020 11:49:07 -0700 Subject: [PATCH 13/20] Update pkg/controller/nodes/task/config/config.go Co-authored-by: Haytham AbuelFutuh --- pkg/controller/nodes/task/config/config.go | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index 2536fd7e0..539a0fdc0 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -71,7 +71,12 @@ func (p TaskPluginConfig) GetEnabledPlugins() map[string]PluginConfig { // Reverse the map. Having the config use task type as a key guarantees only one default plugin can be specified per // task type but now we need to sort for which tasks a plugin needs to be the default. for taskName, pluginName := range p.DefaultForTaskTypes { - pluginDefaultForTaskType[pluginName] = append(pluginDefaultForTaskType[pluginName], cleanString(taskName)) + existing, found := pluginDefaultForTaskType[pluginName] + if !found { + existing = make([]string, 0, 1) + } + + pluginDefaultForTaskType[pluginName] = append(existing, cleanString(taskName)) } for _, pluginName := range p.EnabledPlugins { cleanedPluginName := cleanString(pluginName) From e2066d374cc823aa014b5196e032f7228a134d5a Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Fri, 30 Oct 2020 14:24:53 -0700 Subject: [PATCH 14/20] review comments --- pkg/controller/nodes/task/config/config.go | 40 +++++++++++++++++----- pkg/controller/nodes/task/handler.go | 12 ++++--- pkg/controller/nodes/task/plugin_config.go | 6 +++- 3 files changed, 44 insertions(+), 14 deletions(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index 539a0fdc0..377edaa9d 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -1,9 +1,14 @@ package config import ( + "context" + "fmt" "strings" "time" + "github.com/lyft/flytestdlib/logger" + "k8s.io/apimachinery/pkg/util/sets" + "github.com/lyft/flytestdlib/config" ) @@ -65,26 +70,43 @@ func cleanString(source string) string { return cleaned } -func (p TaskPluginConfig) GetEnabledPlugins() map[string]PluginConfig { +func (p TaskPluginConfig) GetEnabledPlugins() (map[string]PluginConfig, error) { enabledPlugins := make(map[string]PluginConfig) pluginDefaultForTaskType := map[string][]string{} - // Reverse the map. Having the config use task type as a key guarantees only one default plugin can be specified per + // Reverse the DefaultForTaskTypes map. Having the config use task type as a key guarantees only one default plugin can be specified per // task type but now we need to sort for which tasks a plugin needs to be the default. for taskName, pluginName := range p.DefaultForTaskTypes { - existing, found := pluginDefaultForTaskType[pluginName] - if !found { - existing = make([]string, 0, 1) - } - - pluginDefaultForTaskType[pluginName] = append(existing, cleanString(taskName)) + existing, found := pluginDefaultForTaskType[pluginName] + if !found { + existing = make([]string, 0, 1) + } + pluginDefaultForTaskType[cleanString(pluginName)] = append(existing, cleanString(taskName)) } + + enabledPluginsNames := sets.NewString() for _, pluginName := range p.EnabledPlugins { cleanedPluginName := cleanString(pluginName) enabledPlugins[cleanedPluginName] = PluginConfig{ DefaultForTaskTypes: pluginDefaultForTaskType[pluginName], } + enabledPluginsNames.Insert(cleanedPluginName) + } + + // All plugins are enabled, nothing further to validate here. + if len(enabledPlugins) == 0 { + return enabledPlugins, nil + } + + // Finally, validate that default plugins for task types only reference enabled plugins + for pluginName, taskTypes := range pluginDefaultForTaskType { + if !enabledPluginsNames.Has(pluginName) { + logger.Errorf(context.TODO(), "Cannot set default plugin [%s] for task types [%+v] when it is not "+ + "configured to be an enabled plugin. Please double check the flytepropeller config.", pluginName, taskTypes) + return nil, fmt.Errorf("cannot set default plugin [%s] for task types [%+v] when it is not "+ + "configured to be an enabled plugin", pluginName, taskTypes) + } } - return enabledPlugins + return enabledPlugins, nil } func GetConfig() *Config { diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index 92d214ce1..db32e969e 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -205,7 +205,7 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { } // Not every task type will have a default plugin specified in the flytepropeller config. - // That's fine, we resort to using the plugins' static RegisteredTaskTypes as a fallback. + // That's fine, we resort to using the plugins' static RegisteredTaskTypes as a fallback further below. fallbackTaskHandlerMap := make(map[string]map[string]pluginCore.Plugin) for _, p := range enabledPlugins { @@ -214,17 +214,20 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { sCtxFinal := newNameSpacedSetupCtx( tSCtx, newResourceManagerBuilder.GetResourceRegistrar(pluginResourceNamespacePrefix)) logger.Infof(ctx, "Loading Plugin [%s] ENABLED", p.ID) - // cp, err := p.LoadPlugin(ctx, tSCtx) cp, err := p.LoadPlugin(ctx, sCtxFinal) if err != nil { return regErrors.Wrapf(err, "failed to load plugin - %s", p.ID) } + // For every default plugin for a task type specified in flytepropeller config we validate that the plugin's + // static definition includes that task type as something it is registered to handle. for _, tt := range p.RegisteredTaskTypes { for _, defaultTaskType := range p.DefaultForTaskTypes { if defaultTaskType == tt { if existingHandler, alreadyDefaulted := t.defaultPlugins[tt]; alreadyDefaulted && existingHandler.GetID() != cp.GetID() { - logger.Panicf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", + logger.Errorf(ctx, "TaskType [%s] has multiple default handlers specified: [%s] and [%s]", tt, existingHandler.GetID(), cp.GetID()) + return regErrors.New(fmt.Sprintf("TaskType [%s] has multiple default handlers specified: [%s] and [%s]", + tt, existingHandler.GetID(), cp.GetID())) } logger.Infof(ctx, "Plugin [%s] registered for TaskType [%s]", cp.GetID(), tt) t.defaultPlugins[tt] = cp @@ -258,7 +261,8 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { break } if len(registeredPlugins) != 1 { - logger.Panicf(ctx, "Multiple plugins registered to handle task type: %s. ([%+v])", taskType, registeredPlugins) + logger.Errorf(ctx, "Multiple plugins registered to handle task type: %s. ([%+v])", taskType, registeredPlugins) + return regErrors.New(fmt.Sprintf("Multiple plugins registered to handle task type: %s. ([%+v])", taskType, registeredPlugins)) } for _, plugin := range registeredPlugins { t.defaultPlugins[taskType] = plugin diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index 61abb39ec..32f0ba2a8 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -16,8 +16,12 @@ import ( func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPluginConfig, pr PluginRegistryIface) ([]core.PluginEntry, error) { allPluginsEnabled := false enabledPlugins := make(map[string]config.PluginConfig) + var err error if cfg != nil { - enabledPlugins = cfg.GetEnabledPlugins() + enabledPlugins, err = cfg.GetEnabledPlugins() + if err != nil { + return nil, err + } } if len(enabledPlugins) == 0 { allPluginsEnabled = true From a8117a001e6ab28b9894a9586b1fe2feeeb5c136 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Fri, 30 Oct 2020 16:18:57 -0700 Subject: [PATCH 15/20] Update pkg/controller/nodes/task/config/config.go Co-authored-by: Haytham AbuelFutuh --- pkg/controller/nodes/task/config/config.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index 377edaa9d..574b91c33 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -70,7 +70,9 @@ func cleanString(source string) string { return cleaned } -func (p TaskPluginConfig) GetEnabledPlugins() (map[string]PluginConfig, error) { +type PluginID = string +type TaskType = string +func (p TaskPluginConfig) GetEnabledPlugins() (map[PluginID]PluginConfig, error) { enabledPlugins := make(map[string]PluginConfig) pluginDefaultForTaskType := map[string][]string{} // Reverse the DefaultForTaskTypes map. Having the config use task type as a key guarantees only one default plugin can be specified per From 6d138e0cb1a29e137d12f4f39e2ea284b10ccbc3 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Fri, 30 Oct 2020 16:44:12 -0700 Subject: [PATCH 16/20] Use plugins config default for task types regardless of whethe all plugins are enabled --- pkg/controller/nodes/task/config/config.go | 46 ++++++++++++---------- pkg/controller/nodes/task/plugin_config.go | 18 ++++----- 2 files changed, 34 insertions(+), 30 deletions(-) diff --git a/pkg/controller/nodes/task/config/config.go b/pkg/controller/nodes/task/config/config.go index 574b91c33..cc7d219be 100644 --- a/pkg/controller/nodes/task/config/config.go +++ b/pkg/controller/nodes/task/config/config.go @@ -60,8 +60,14 @@ type BackOffConfig struct { MaxDuration config.Duration `json:"max-duration" pflag:",The cap of the backoff duration"` } -type PluginConfig struct { - DefaultForTaskTypes []string +type PluginID = string +type TaskType = string + +// Contains the set of enabled plugins for this flytepropeller deployment along with default plugin handlers +// for specific task types. +type PluginsConfigMeta struct { + EnabledPlugins sets.String + AllDefaultForTaskTypes map[PluginID][]TaskType } func cleanString(source string) string { @@ -70,11 +76,14 @@ func cleanString(source string) string { return cleaned } -type PluginID = string -type TaskType = string -func (p TaskPluginConfig) GetEnabledPlugins() (map[PluginID]PluginConfig, error) { - enabledPlugins := make(map[string]PluginConfig) - pluginDefaultForTaskType := map[string][]string{} +func (p TaskPluginConfig) GetEnabledPlugins() (PluginsConfigMeta, error) { + enabledPluginsNames := sets.NewString() + for _, pluginName := range p.EnabledPlugins { + cleanedPluginName := cleanString(pluginName) + enabledPluginsNames.Insert(cleanedPluginName) + } + + pluginDefaultForTaskType := make(map[PluginID][]TaskType) // Reverse the DefaultForTaskTypes map. Having the config use task type as a key guarantees only one default plugin can be specified per // task type but now we need to sort for which tasks a plugin needs to be the default. for taskName, pluginName := range p.DefaultForTaskTypes { @@ -85,18 +94,12 @@ func (p TaskPluginConfig) GetEnabledPlugins() (map[PluginID]PluginConfig, error) pluginDefaultForTaskType[cleanString(pluginName)] = append(existing, cleanString(taskName)) } - enabledPluginsNames := sets.NewString() - for _, pluginName := range p.EnabledPlugins { - cleanedPluginName := cleanString(pluginName) - enabledPlugins[cleanedPluginName] = PluginConfig{ - DefaultForTaskTypes: pluginDefaultForTaskType[pluginName], - } - enabledPluginsNames.Insert(cleanedPluginName) - } - // All plugins are enabled, nothing further to validate here. - if len(enabledPlugins) == 0 { - return enabledPlugins, nil + if enabledPluginsNames.Len() == 0 { + return PluginsConfigMeta{ + EnabledPlugins: enabledPluginsNames, + AllDefaultForTaskTypes: pluginDefaultForTaskType, + }, nil } // Finally, validate that default plugins for task types only reference enabled plugins @@ -104,11 +107,14 @@ func (p TaskPluginConfig) GetEnabledPlugins() (map[PluginID]PluginConfig, error) if !enabledPluginsNames.Has(pluginName) { logger.Errorf(context.TODO(), "Cannot set default plugin [%s] for task types [%+v] when it is not "+ "configured to be an enabled plugin. Please double check the flytepropeller config.", pluginName, taskTypes) - return nil, fmt.Errorf("cannot set default plugin [%s] for task types [%+v] when it is not "+ + return PluginsConfigMeta{}, fmt.Errorf("cannot set default plugin [%s] for task types [%+v] when it is not "+ "configured to be an enabled plugin", pluginName, taskTypes) } } - return enabledPlugins, nil + return PluginsConfigMeta{ + EnabledPlugins: enabledPluginsNames, + AllDefaultForTaskTypes: pluginDefaultForTaskType, + }, nil } func GetConfig() *Config { diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index 32f0ba2a8..677fe47a3 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -15,29 +15,28 @@ import ( func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPluginConfig, pr PluginRegistryIface) ([]core.PluginEntry, error) { allPluginsEnabled := false - enabledPlugins := make(map[string]config.PluginConfig) + pluginsConfigMeta := config.PluginsConfigMeta{} var err error if cfg != nil { - enabledPlugins, err = cfg.GetEnabledPlugins() + pluginsConfigMeta, err = cfg.GetEnabledPlugins() if err != nil { return nil, err } } - if len(enabledPlugins) == 0 { + if pluginsConfigMeta.EnabledPlugins.Len() == 0 { allPluginsEnabled = true } var finalizedPlugins []core.PluginEntry - logger.Infof(ctx, "Enabled plugins: %+v", enabledPlugins) + logger.Infof(ctx, "Enabled plugins: %v", pluginsConfigMeta.EnabledPlugins.List()) logger.Infof(ctx, "Loading core Plugins, plugin configuration [all plugins enabled: %v]", allPluginsEnabled) for _, cpe := range pr.GetCorePlugins() { id := strings.ToLower(cpe.ID) - pluginCfg, pluginEnabled := enabledPlugins[id] - if !allPluginsEnabled && !pluginEnabled { + if !allPluginsEnabled && !pluginsConfigMeta.EnabledPlugins.Has(id) { logger.Infof(ctx, "Plugin [%s] is DISABLED (not found in enabled plugins list).", id) } else { logger.Infof(ctx, "Plugin [%s] ENABLED", id) - cpe.DefaultForTaskTypes = pluginCfg.DefaultForTaskTypes + cpe.DefaultForTaskTypes = pluginsConfigMeta.AllDefaultForTaskTypes[id] finalizedPlugins = append(finalizedPlugins, cpe) } } @@ -52,8 +51,7 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu for i := range k8sPlugins { kpe := k8sPlugins[i] id := strings.ToLower(kpe.ID) - pluginConfig, pluginEnabled := enabledPlugins[id] - if !allPluginsEnabled && !pluginEnabled { + if !allPluginsEnabled && !pluginsConfigMeta.EnabledPlugins.Has(id) { logger.Infof(ctx, "K8s Plugin [%s] is DISABLED (not found in enabled plugins list).", id) } else { logger.Infof(ctx, "K8s Plugin [%s] is ENABLED.", id) @@ -64,7 +62,7 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu return k8s.NewPluginManagerWithBackOff(ctx, iCtx, kpe, backOffController, monitorIndex) }, IsDefault: kpe.IsDefault, - DefaultForTaskTypes: pluginConfig.DefaultForTaskTypes, + DefaultForTaskTypes: pluginsConfigMeta.AllDefaultForTaskTypes[id], } finalizedPlugins = append(finalizedPlugins, plugin) } From c5ecc26dc1a9b974d6f142147cd0ebdb143c00c2 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Mon, 2 Nov 2020 10:45:53 -0800 Subject: [PATCH 17/20] Update pkg/controller/nodes/task/plugin_config.go Co-authored-by: Haytham AbuelFutuh --- pkg/controller/nodes/task/plugin_config.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index 677fe47a3..72d3aeaab 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -15,7 +15,9 @@ import ( func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPluginConfig, pr PluginRegistryIface) ([]core.PluginEntry, error) { allPluginsEnabled := false - pluginsConfigMeta := config.PluginsConfigMeta{} + pluginsConfigMeta := config.PluginsConfigMeta{ + AllDefaultForTaskTypes: map[PluginID][]TaskType{} + } var err error if cfg != nil { pluginsConfigMeta, err = cfg.GetEnabledPlugins() From 99dd7bdac6962394482deda425e32353a4c2346c Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Mon, 2 Nov 2020 10:46:05 -0800 Subject: [PATCH 18/20] Update pkg/controller/nodes/task/handler.go Co-authored-by: Haytham AbuelFutuh --- pkg/controller/nodes/task/handler.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index db32e969e..7303ad08c 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -262,7 +262,7 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { } if len(registeredPlugins) != 1 { logger.Errorf(ctx, "Multiple plugins registered to handle task type: %s. ([%+v])", taskType, registeredPlugins) - return regErrors.New(fmt.Sprintf("Multiple plugins registered to handle task type: %s. ([%+v])", taskType, registeredPlugins)) + return regErrors.New(fmt.Sprintf("Multiple plugins registered to handle task type: %s. ([%+v]). Use default-for-task-type config option to choose the desired plugin.", taskType, registeredPlugins)) } for _, plugin := range registeredPlugins { t.defaultPlugins[taskType] = plugin From b102b8160dafddf0da295ed2ed2331470d106884 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Mon, 2 Nov 2020 10:49:38 -0800 Subject: [PATCH 19/20] review comments --- pkg/controller/nodes/task/handler.go | 5 +++-- pkg/controller/nodes/task/plugin_config.go | 4 +++- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/pkg/controller/nodes/task/handler.go b/pkg/controller/nodes/task/handler.go index db32e969e..6b4c8b777 100644 --- a/pkg/controller/nodes/task/handler.go +++ b/pkg/controller/nodes/task/handler.go @@ -153,6 +153,7 @@ type PluginRegistryIface interface { GetK8sPlugins() []pluginK8s.PluginEntry } +type taskType = string type pluginID = string type Handler struct { @@ -206,7 +207,7 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { // Not every task type will have a default plugin specified in the flytepropeller config. // That's fine, we resort to using the plugins' static RegisteredTaskTypes as a fallback further below. - fallbackTaskHandlerMap := make(map[string]map[string]pluginCore.Plugin) + fallbackTaskHandlerMap := make(map[taskType]map[pluginID]pluginCore.Plugin) for _, p := range enabledPlugins { // create a new resource registrar proxy for each plugin, and pass it into the plugin's LoadPlugin() via a setup context @@ -243,7 +244,7 @@ func (t *Handler) Setup(ctx context.Context, sCtx handler.SetupContext) error { fallbackMap, ok := fallbackTaskHandlerMap[tt] if !ok { - fallbackMap = make(map[string]pluginCore.Plugin) + fallbackMap = make(map[pluginID]pluginCore.Plugin) } fallbackMap[cp.GetID()] = cp fallbackTaskHandlerMap[tt] = fallbackMap diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index 677fe47a3..71631ac1a 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -36,7 +36,9 @@ func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPlu logger.Infof(ctx, "Plugin [%s] is DISABLED (not found in enabled plugins list).", id) } else { logger.Infof(ctx, "Plugin [%s] ENABLED", id) - cpe.DefaultForTaskTypes = pluginsConfigMeta.AllDefaultForTaskTypes[id] + if defaults, ok := pluginsConfigMeta.AllDefaultForTaskTypes[id]; ok { + cpe.DefaultForTaskTypes = defaults + } finalizedPlugins = append(finalizedPlugins, cpe) } } From 171168c6935fedb0fa8d996a83691966627b4788 Mon Sep 17 00:00:00 2001 From: Katrina Rogan Date: Mon, 2 Nov 2020 10:55:19 -0800 Subject: [PATCH 20/20] fix --- pkg/controller/nodes/task/plugin_config.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/controller/nodes/task/plugin_config.go b/pkg/controller/nodes/task/plugin_config.go index c6a41dc1e..d6ab6b4b3 100644 --- a/pkg/controller/nodes/task/plugin_config.go +++ b/pkg/controller/nodes/task/plugin_config.go @@ -16,7 +16,7 @@ import ( func WranglePluginsAndGenerateFinalList(ctx context.Context, cfg *config.TaskPluginConfig, pr PluginRegistryIface) ([]core.PluginEntry, error) { allPluginsEnabled := false pluginsConfigMeta := config.PluginsConfigMeta{ - AllDefaultForTaskTypes: map[PluginID][]TaskType{} + AllDefaultForTaskTypes: map[pluginID][]taskType{}, } var err error if cfg != nil {