parent
526c8b4d76
commit
025f95b1d6
@ -1,124 +0,0 @@ |
||||
/* |
||||
* Minio Cloud Storage, (C) 2015 Minio, Inc. |
||||
* |
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
* you may not use this file except in compliance with the License. |
||||
* You may obtain a copy of the License at |
||||
* |
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
* |
||||
* Unless required by applicable law or agreed to in writing, software |
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
* See the License for the specific language governing permissions and |
||||
* limitations under the License. |
||||
*/ |
||||
|
||||
package controller |
||||
|
||||
import ( |
||||
"encoding/json" |
||||
"net/http" |
||||
|
||||
jsonrpc "github.com/gorilla/rpc/v2/json" |
||||
"github.com/minio/minio/pkg/auth" |
||||
"github.com/minio/minio/pkg/probe" |
||||
"github.com/minio/minio/pkg/server/rpc" |
||||
) |
||||
|
||||
func closeResp(resp *http.Response) { |
||||
if resp != nil && resp.Body != nil { |
||||
resp.Body.Close() |
||||
} |
||||
} |
||||
|
||||
// GetMemStats get memory status of the server at given url
|
||||
func GetMemStats(url string) ([]byte, *probe.Error) { |
||||
op := RPCOps{ |
||||
Method: "MemStats.Get", |
||||
Request: rpc.Args{Request: ""}, |
||||
} |
||||
req, perr := NewRequest(url, op, http.DefaultTransport) |
||||
if perr != nil { |
||||
return nil, perr.Trace() |
||||
} |
||||
resp, perr := req.Do() |
||||
defer closeResp(resp) |
||||
if perr != nil { |
||||
return nil, perr.Trace() |
||||
} |
||||
var reply rpc.MemStatsReply |
||||
if err := jsonrpc.DecodeClientResponse(resp.Body, &reply); err != nil { |
||||
return nil, probe.NewError(err) |
||||
} |
||||
jsonRespBytes, err := json.MarshalIndent(reply, "", "\t") |
||||
if err != nil { |
||||
return nil, probe.NewError(err) |
||||
} |
||||
return jsonRespBytes, nil |
||||
} |
||||
|
||||
// GetSysInfo get system status of the server at given url
|
||||
func GetSysInfo(url string) ([]byte, *probe.Error) { |
||||
op := RPCOps{ |
||||
Method: "SysInfo.Get", |
||||
Request: rpc.Args{Request: ""}, |
||||
} |
||||
req, perr := NewRequest(url, op, http.DefaultTransport) |
||||
if perr != nil { |
||||
return nil, perr.Trace() |
||||
} |
||||
resp, perr := req.Do() |
||||
defer closeResp(resp) |
||||
if perr != nil { |
||||
return nil, perr.Trace() |
||||
} |
||||
var reply rpc.SysInfoReply |
||||
if err := jsonrpc.DecodeClientResponse(resp.Body, &reply); err != nil { |
||||
return nil, probe.NewError(err) |
||||
} |
||||
jsonRespBytes, err := json.MarshalIndent(reply, "", "\t") |
||||
if err != nil { |
||||
return nil, probe.NewError(err) |
||||
} |
||||
return jsonRespBytes, nil |
||||
} |
||||
|
||||
// GetAuthKeys get access key id and secret access key
|
||||
func GetAuthKeys(url string) ([]byte, *probe.Error) { |
||||
op := RPCOps{ |
||||
Method: "Auth.Get", |
||||
Request: rpc.Args{Request: ""}, |
||||
} |
||||
req, perr := NewRequest(url, op, http.DefaultTransport) |
||||
if perr != nil { |
||||
return nil, perr.Trace() |
||||
} |
||||
resp, perr := req.Do() |
||||
defer closeResp(resp) |
||||
if perr != nil { |
||||
return nil, perr.Trace() |
||||
} |
||||
var reply rpc.AuthReply |
||||
if err := jsonrpc.DecodeClientResponse(resp.Body, &reply); err != nil { |
||||
return nil, probe.NewError(err) |
||||
} |
||||
authConfig := &auth.Config{} |
||||
authConfig.Version = "0.0.1" |
||||
authConfig.Users = make(map[string]*auth.User) |
||||
user := &auth.User{} |
||||
user.Name = "testuser" |
||||
user.AccessKeyID = reply.AccessKeyID |
||||
user.SecretAccessKey = reply.SecretAccessKey |
||||
authConfig.Users[reply.AccessKeyID] = user |
||||
if err := auth.SaveConfig(authConfig); err != nil { |
||||
return nil, err.Trace() |
||||
} |
||||
jsonRespBytes, err := json.MarshalIndent(reply, "", "\t") |
||||
if err != nil { |
||||
return nil, probe.NewError(err) |
||||
} |
||||
return jsonRespBytes, nil |
||||
} |
||||
|
||||
// Add more functions here for other RPC messages
|
@ -0,0 +1,43 @@ |
||||
/* |
||||
* Minio Cloud Storage, (C) 2015 Minio, Inc. |
||||
* |
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
* you may not use this file except in compliance with the License. |
||||
* You may obtain a copy of the License at |
||||
* |
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
* |
||||
* Unless required by applicable law or agreed to in writing, software |
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
* See the License for the specific language governing permissions and |
||||
* limitations under the License. |
||||
*/ |
||||
|
||||
package controller |
||||
|
||||
import ( |
||||
"net/http" |
||||
|
||||
router "github.com/gorilla/mux" |
||||
"github.com/minio/minio/pkg/rpc" |
||||
) |
||||
|
||||
// getRPCHandler rpc handler
|
||||
func getRPCHandler() http.Handler { |
||||
s := rpc.NewServer() |
||||
s.RegisterJSONCodec() |
||||
s.RegisterService(new(rpc.VersionService), "Version") |
||||
s.RegisterService(new(rpc.SysInfoService), "SysInfo") |
||||
s.RegisterService(new(rpc.MemStatsService), "MemStats") |
||||
s.RegisterService(new(rpc.DonutService), "Donut") |
||||
s.RegisterService(new(rpc.AuthService), "Auth") |
||||
// Add new RPC services here
|
||||
return registerRPC(router.NewRouter(), s) |
||||
} |
||||
|
||||
// registerRPC - register rpc handlers
|
||||
func registerRPC(mux *router.Router, s *rpc.Server) http.Handler { |
||||
mux.Handle("/rpc", s) |
||||
return mux |
||||
} |
@ -0,0 +1,67 @@ |
||||
/* |
||||
* Minio Cloud Storage, (C) 2015 Minio, Inc. |
||||
* |
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
* you may not use this file except in compliance with the License. |
||||
* You may obtain a copy of the License at |
||||
* |
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
* |
||||
* Unless required by applicable law or agreed to in writing, software |
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
* See the License for the specific language governing permissions and |
||||
* limitations under the License. |
||||
*/ |
||||
|
||||
package controller |
||||
|
||||
import ( |
||||
"fmt" |
||||
"net" |
||||
"net/http" |
||||
"os" |
||||
"strings" |
||||
|
||||
"github.com/minio/minio/pkg/minhttp" |
||||
"github.com/minio/minio/pkg/probe" |
||||
) |
||||
|
||||
// getRPCServer instance
|
||||
func getRPCServer(rpcHandler http.Handler) (*http.Server, *probe.Error) { |
||||
// Minio server config
|
||||
httpServer := &http.Server{ |
||||
Addr: ":9001", // TODO make this configurable
|
||||
Handler: rpcHandler, |
||||
MaxHeaderBytes: 1 << 20, |
||||
} |
||||
var hosts []string |
||||
addrs, err := net.InterfaceAddrs() |
||||
if err != nil { |
||||
return nil, probe.NewError(err) |
||||
} |
||||
for _, addr := range addrs { |
||||
if addr.Network() == "ip+net" { |
||||
host := strings.Split(addr.String(), "/")[0] |
||||
if ip := net.ParseIP(host); ip.To4() != nil { |
||||
hosts = append(hosts, host) |
||||
} |
||||
} |
||||
} |
||||
for _, host := range hosts { |
||||
fmt.Printf("Starting minio server on: http://%s:9001/rpc, PID: %d\n", host, os.Getpid()) |
||||
} |
||||
return httpServer, nil |
||||
} |
||||
|
||||
func StartController() *probe.Error { |
||||
rpcServer, err := getRPCServer(getRPCHandler()) |
||||
if err != nil { |
||||
return err.Trace() |
||||
} |
||||
// Setting rate limit to 'zero' no ratelimiting implemented
|
||||
if err := minhttp.ListenAndServeLimited(0, rpcServer); err != nil { |
||||
return err.Trace() |
||||
} |
||||
return nil |
||||
} |
Loading…
Reference in new issue