Refactor controller client package (#1402)

The function implementations of controller client package are
inconsistent. This PR lets functions reuse the functions
that already implemented and able to set additional headers to
request.
This commit is contained in:
Ta-Ching Chen
2019-11-12 18:27:52 +08:00
committed by GitHub
parent d0276f1d52
commit a645a1e197
12 changed files with 120 additions and 110 deletions
+3 -6
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -33,7 +31,7 @@ func (c *Client) CanaryConfigCreate(canaryConf *fv1.CanaryConfig) (*metav1.Objec
return nil, err
}
resp, err := http.Post(c.url("canaryconfigs"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("canaryconfigs", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -57,7 +55,7 @@ func (c *Client) CanaryConfigGet(m *metav1.ObjectMeta) (*fv1.CanaryConfig, error
relativeUrl := fmt.Sprintf("canaryconfigs/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -106,13 +104,12 @@ func (c *Client) CanaryConfigUpdate(canaryConf *fv1.CanaryConfig) (*metav1.Objec
func (c *Client) CanaryConfigDelete(m *metav1.ObjectMeta) error {
relativeUrl := fmt.Sprintf("canaryconfigs/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
return c.delete(relativeUrl)
}
func (c *Client) CanaryConfigList(ns string) ([]fv1.CanaryConfig, error) {
relativeUrl := fmt.Sprintf("canaryconfigs?namespace=%v", ns)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+53 -45
View File
@@ -18,38 +18,56 @@ package client
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"io/ioutil"
"net/http"
"strings"
"time"
"github.com/pkg/errors"
"golang.org/x/net/context/ctxhttp"
ferror "github.com/fission/fission/pkg/error"
"github.com/fission/fission/pkg/info"
)
var (
DefaultRequestHeaders map[string]string
)
type (
Client struct {
Url string
Url string
Headers map[string]string
}
)
func MakeClient(serverUrl string) *Client {
return &Client{Url: strings.TrimSuffix(serverUrl, "/")}
return &Client{
Url: strings.TrimSuffix(serverUrl, "/"),
Headers: DefaultRequestHeaders,
}
}
func (c *Client) create(relativeUrl string, contentType string, payload []byte) (*http.Response, error) {
var reader io.Reader
if len(payload) > 0 {
reader = bytes.NewReader(payload)
}
return c.sendRequest(http.MethodPost, c.v2CrdUrl(relativeUrl), map[string]string{"Content-type": contentType}, reader)
}
func (c *Client) put(relativeUrl string, contentType string, payload []byte) (*http.Response, error) {
var reader io.Reader
if len(payload) > 0 {
reader = bytes.NewReader(payload)
}
return c.sendRequest(http.MethodPut, c.v2CrdUrl(relativeUrl), map[string]string{"Content-type": contentType}, reader)
}
func (c *Client) get(relativeUrl string) (*http.Response, error) {
return c.sendRequest(http.MethodGet, c.v2CrdUrl(relativeUrl), nil, nil)
}
func (c *Client) delete(relativeUrl string) error {
req, err := http.NewRequest("DELETE", c.url(relativeUrl), nil)
if err != nil {
return err
}
resp, err := http.DefaultClient.Do(req)
resp, err := c.sendRequest(http.MethodDelete, c.v2CrdUrl(relativeUrl), nil, nil)
if err != nil {
return err
}
@@ -67,17 +85,33 @@ func (c *Client) delete(relativeUrl string) error {
return nil
}
func (c *Client) put(relativeUrl string, contentType string, body []byte) (*http.Response, error) {
req, err := http.NewRequest("PUT", c.url(relativeUrl), bytes.NewReader(body))
func (c *Client) proxy(method string, relativeUrl string, payload []byte) (*http.Response, error) {
var reader io.Reader
if len(payload) > 0 {
reader = bytes.NewReader(payload)
}
return c.sendRequest(method, c.proxyUrl(relativeUrl), nil, reader)
}
func (c *Client) sendRequest(method string, relativeUrl string, headers map[string]string, reader io.Reader) (*http.Response, error) {
req, err := http.NewRequest(method, relativeUrl, reader)
if err != nil {
return nil, err
}
req.Header.Set("Content-type", contentType)
for _, hs := range []map[string]string{headers, c.Headers} {
for k, v := range hs {
req.Header.Set(k, v)
}
}
return http.DefaultClient.Do(req)
}
func (c *Client) url(relativeUrl string) string {
return c.Url + "/v2/" + relativeUrl
func (c *Client) v2CrdUrl(relativeUrl string) string {
return c.Url + "/v2/" + strings.TrimPrefix(relativeUrl, "/")
}
func (c *Client) proxyUrl(relativeUrl string) string {
return c.Url + "/proxy/" + strings.TrimPrefix(relativeUrl, "/")
}
func (c *Client) handleResponse(resp *http.Response) ([]byte, error) {
@@ -95,29 +129,3 @@ func (c *Client) handleCreateResponse(resp *http.Response) ([]byte, error) {
body, err := ioutil.ReadAll(resp.Body)
return body, err
}
func (c *Client) ServerInfo() (*info.ServerInfo, error) {
url := fmt.Sprintf(c.Url)
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
resp, err := ctxhttp.Get(ctx, &http.Client{}, url)
if err != nil {
return nil, err
}
defer resp.Body.Close()
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, err
}
info := &info.ServerInfo{}
err = json.Unmarshal(body, info)
if err != nil {
return nil, err
}
return info, nil
}
+33 -9
View File
@@ -17,22 +17,26 @@ limitations under the License.
package client
import (
"context"
"encoding/json"
"fmt"
"io/ioutil"
"net/http"
"time"
"github.com/pkg/errors"
"golang.org/x/net/context/ctxhttp"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"github.com/fission/fission/pkg/info"
)
func (c *Client) SecretGet(m *metav1.ObjectMeta) (*apiv1.Secret, error) {
relativeUrl := fmt.Sprintf("secrets/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -56,7 +60,7 @@ func (c *Client) ConfigMapGet(m *metav1.ObjectMeta) (*apiv1.ConfigMap, error) {
relativeUrl := fmt.Sprintf("configmaps/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -77,19 +81,15 @@ func (c *Client) ConfigMapGet(m *metav1.ObjectMeta) (*apiv1.ConfigMap, error) {
}
func (c *Client) GetSvcURL(label string) (string, error) {
url := fmt.Sprintf("%s/proxy/svcname?"+label, c.Url)
resp, err := http.Get(url)
resp, err := c.proxy(http.MethodGet, "svcname?"+label, nil)
if err != nil {
return "", err
}
if resp == nil {
return "", errors.Errorf("failed to find service for given label: %v", label)
}
defer resp.Body.Close()
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return "", err
@@ -99,3 +99,27 @@ func (c *Client) GetSvcURL(label string) (string, error) {
return storageSvc, err
}
func (c *Client) ServerInfo() (*info.ServerInfo, error) {
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
resp, err := ctxhttp.Get(ctx, &http.Client{}, c.Url)
if err != nil {
return nil, err
}
defer resp.Body.Close()
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, err
}
info := &info.ServerInfo{}
err = json.Unmarshal(body, info)
if err != nil {
return nil, err
}
return info, nil
}
+3 -6
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -43,7 +41,7 @@ func (c *Client) EnvironmentCreate(env *fv1.Environment) (*metav1.ObjectMeta, er
return nil, err
}
resp, err := http.Post(c.url("environments"), "application/json", bytes.NewReader(data))
resp, err := c.create("environments", "application/json", data)
if err != nil {
return nil, err
}
@@ -67,7 +65,7 @@ func (c *Client) EnvironmentGet(m *metav1.ObjectMeta) (*fv1.Environment, error)
relativeUrl := fmt.Sprintf("environments/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -117,13 +115,12 @@ func (c *Client) EnvironmentUpdate(env *fv1.Environment) (*metav1.ObjectMeta, er
func (c *Client) EnvironmentDelete(m *metav1.ObjectMeta) error {
relativeUrl := fmt.Sprintf("environments/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
return c.delete(relativeUrl)
}
func (c *Client) EnvironmentList(ns string) ([]fv1.Environment, error) {
relativeUrl := fmt.Sprintf("environments?namespace=%v", ns)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+4 -6
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -38,7 +36,7 @@ func (c *Client) FunctionCreate(f *fv1.Function) (*metav1.ObjectMeta, error) {
return nil, err
}
resp, err := http.Post(c.url("functions"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("functions", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -62,7 +60,7 @@ func (c *Client) FunctionGet(m *metav1.ObjectMeta) (*fv1.Function, error) {
relativeUrl := fmt.Sprintf("functions/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -87,7 +85,7 @@ func (c *Client) FunctionGetRawDeployment(m *metav1.ObjectMeta) ([]byte, error)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
relativeUrl += fmt.Sprintf("&deploymentraw=1")
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -135,7 +133,7 @@ func (c *Client) FunctionDelete(m *metav1.ObjectMeta) error {
func (c *Client) FunctionList(functionNamespace string) ([]fv1.Function, error) {
relativeUrl := fmt.Sprintf("functions?namespace=%v", functionNamespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+3 -5
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -38,7 +36,7 @@ func (c *Client) HTTPTriggerCreate(t *fv1.HTTPTrigger) (*metav1.ObjectMeta, erro
return nil, err
}
resp, err := http.Post(c.url("triggers/http"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("triggers/http", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -62,7 +60,7 @@ func (c *Client) HTTPTriggerGet(m *metav1.ObjectMeta) (*fv1.HTTPTrigger, error)
relativeUrl := fmt.Sprintf("triggers/http/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -121,7 +119,7 @@ func (c *Client) HTTPTriggerDelete(m *metav1.ObjectMeta) error {
func (c *Client) HTTPTriggerList(triggerNamespace string) ([]fv1.HTTPTrigger, error) {
relativeUrl := fmt.Sprintf("triggers/http?namespace=%v", triggerNamespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -39,7 +37,7 @@ func (c *Client) WatchCreate(w *fv1.KubernetesWatchTrigger) (*metav1.ObjectMeta,
return nil, err
}
resp, err := http.Post(c.url("watches"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("watches", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -63,7 +61,7 @@ func (c *Client) WatchGet(m *metav1.ObjectMeta) (*fv1.KubernetesWatchTrigger, er
relativeUrl := fmt.Sprintf("watches/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -84,8 +82,7 @@ func (c *Client) WatchGet(m *metav1.ObjectMeta) (*fv1.KubernetesWatchTrigger, er
}
func (c *Client) WatchUpdate(w *fv1.KubernetesWatchTrigger) (*metav1.ObjectMeta, error) {
return nil, ferror.MakeError(ferror.ErrorNotImplmented,
"watch update not implemented")
return nil, ferror.MakeError(ferror.ErrorNotImplmented, "watch update not implemented")
}
func (c *Client) WatchDelete(m *metav1.ObjectMeta) error {
@@ -96,7 +93,7 @@ func (c *Client) WatchDelete(m *metav1.ObjectMeta) error {
func (c *Client) WatchList(ns string) ([]fv1.KubernetesWatchTrigger, error) {
relativeUrl := fmt.Sprintf("watches?namespace=%v", ns)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+3 -5
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -38,7 +36,7 @@ func (c *Client) MessageQueueTriggerCreate(t *fv1.MessageQueueTrigger) (*metav1.
return nil, err
}
resp, err := http.Post(c.url("triggers/messagequeue"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("triggers/messagequeue", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -62,7 +60,7 @@ func (c *Client) MessageQueueTriggerGet(m *metav1.ObjectMeta) (*fv1.MessageQueue
relativeUrl := fmt.Sprintf("triggers/messagequeue/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -126,7 +124,7 @@ func (c *Client) MessageQueueTriggerList(mqType string, ns string) ([]fv1.Messag
relativeUrl += fmt.Sprintf("?mqtype=%v&namespace=%v", mqType, ns)
}
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+3 -5
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -38,7 +36,7 @@ func (c *Client) PackageCreate(f *fv1.Package) (*metav1.ObjectMeta, error) {
return nil, err
}
resp, err := http.Post(c.url("packages"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("packages", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -62,7 +60,7 @@ func (c *Client) PackageGet(m *metav1.ObjectMeta) (*fv1.Package, error) {
relativeUrl := fmt.Sprintf("packages/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -121,7 +119,7 @@ func (c *Client) PackageDelete(m *metav1.ObjectMeta) error {
func (c *Client) PackageList(pkgNamespace string) ([]fv1.Package, error) {
relativeUrl := fmt.Sprintf("packages?namespace=%v", pkgNamespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+7 -9
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -39,7 +37,7 @@ func (c *Client) RecorderCreate(r *fv1.Recorder) (*metav1.ObjectMeta, error) {
return nil, err
}
resp, err := http.Post(c.url("recorders"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("recorders", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -63,7 +61,7 @@ func (c *Client) RecorderGet(m *metav1.ObjectMeta) (*fv1.Recorder, error) {
relativeUrl := fmt.Sprintf("recorders/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -123,7 +121,7 @@ func (c *Client) RecorderDelete(m *metav1.ObjectMeta) error {
func (c *Client) RecorderList(ns string) ([]fv1.Recorder, error) {
relativeUrl := "recorders"
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -147,7 +145,7 @@ func (c *Client) RecorderList(ns string) ([]fv1.Recorder, error) {
func (c *Client) RecordsByFunction(function string) ([]*redisCache.RecordedEntry, error) {
relativeUrl := fmt.Sprintf("records/function/%v", function)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -170,7 +168,7 @@ func (c *Client) RecordsByFunction(function string) ([]*redisCache.RecordedEntry
func (c *Client) RecordsAll() ([]*redisCache.RecordedEntry, error) {
relativeUrl := "records"
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -193,7 +191,7 @@ func (c *Client) RecordsAll() ([]*redisCache.RecordedEntry, error) {
func (c *Client) RecordsByTrigger(trigger string) ([]*redisCache.RecordedEntry, error) {
relativeUrl := fmt.Sprintf("records/trigger/%v", trigger)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -217,7 +215,7 @@ func (c *Client) RecordsByTime(from string, to string) ([]*redisCache.RecordedEn
relativeUrl := "records/time"
relativeUrl += fmt.Sprintf("?from=%v&to=%v", from, to)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+1 -2
View File
@@ -18,13 +18,12 @@ package client
import (
"encoding/json"
"fmt"
"net/http"
)
func (c *Client) ReplayByReqUID(reqUID string) ([]string, error) {
relativeUrl := fmt.Sprintf("replay/%v", reqUID)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
+3 -5
View File
@@ -17,10 +17,8 @@ limitations under the License.
package client
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -38,7 +36,7 @@ func (c *Client) TimeTriggerCreate(t *fv1.TimeTrigger) (*metav1.ObjectMeta, erro
return nil, err
}
resp, err := http.Post(c.url("triggers/time"), "application/json", bytes.NewReader(reqbody))
resp, err := c.create("triggers/time", "application/json", reqbody)
if err != nil {
return nil, err
}
@@ -62,7 +60,7 @@ func (c *Client) TimeTriggerGet(m *metav1.ObjectMeta) (*fv1.TimeTrigger, error)
relativeUrl := fmt.Sprintf("triggers/time/%v", m.Name)
relativeUrl += fmt.Sprintf("?namespace=%v", m.Namespace)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}
@@ -121,7 +119,7 @@ func (c *Client) TimeTriggerDelete(m *metav1.ObjectMeta) error {
func (c *Client) TimeTriggerList(ns string) ([]fv1.TimeTrigger, error) {
relativeUrl := fmt.Sprintf("triggers/time?namespace=%v", ns)
resp, err := http.Get(c.url(relativeUrl))
resp, err := c.get(relativeUrl)
if err != nil {
return nil, err
}