diff --git a/cmd/cloudproxy/main.go b/cmd/cloudproxy/main.go new file mode 100644 index 0000000000..6fec114d90 --- /dev/null +++ b/cmd/cloudproxy/main.go @@ -0,0 +1,95 @@ +// Copyright 2019 Yunion +// +// 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 main + +import ( + "context" + "os" + "os/signal" + "sync" + "syscall" + + "yunion.io/x/log" + + common_app "yunion.io/x/onecloud/pkg/cloudcommon/app" + "yunion.io/x/onecloud/pkg/cloudproxy/agent/worker" + "yunion.io/x/onecloud/pkg/cloudproxy/options" + "yunion.io/x/onecloud/pkg/cloudproxy/service" + "yunion.io/x/onecloud/pkg/util/atexit" +) + +func main() { + defer atexit.Handle() + + nop := true + wg := &sync.WaitGroup{} + ctx := context.Background() + ctx = context.WithValue(ctx, "wg", wg) + ctx, cancelFunc := context.WithCancel(ctx) + + var ( + opts = options.Get() + commonOpts = &opts.CommonOptions + ) + if opts.EnableAPIServer || opts.EnableProxyAgent { + common_app.InitAuth(commonOpts, func() { + log.Infof("Auth complete") + }) + } + + if opts.EnableAPIServer { + nop = false + wg.Add(1) + go func() { + defer wg.Done() + service.StartService() + }() + } + + if opts.EnableProxyAgent { + if opts.EnableAPIServer { + const d = "10m" + log.Infof("set proxy_agent_init_wait to %s", d) + opts.Options.ProxyAgentInitWait = d + } + if err := opts.Options.ValidateThenInit(); err != nil { + log.Fatalf("proxy agent options validation: %v", err) + } + worker := worker.NewWorker(&opts.CommonOptions, &opts.Options) + nop = false + go func() { + worker.Start(ctx) + pid := os.Getpid() + p, err := os.FindProcess(pid) + if err != nil { + log.Fatalf("find process of my pid %d: %v", pid, err) + } + p.Signal(syscall.SIGTERM) + }() + } + if nop { + log.Warningln("nothing to do. Please check configuration") + } else { + go func() { + sigChan := make(chan os.Signal) + signal.Notify(sigChan, syscall.SIGINT) + signal.Notify(sigChan, syscall.SIGTERM) + sig := <-sigChan + log.Infof("signal received: %s", sig) + cancelFunc() + }() + wg.Wait() + } +} diff --git a/pkg/apis/cloudproxy/const.go b/pkg/apis/cloudproxy/const.go new file mode 100644 index 0000000000..e425398877 --- /dev/null +++ b/pkg/apis/cloudproxy/const.go @@ -0,0 +1,44 @@ +// Copyright 2019 Yunion +// +// 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 cloudproxy + +import ( + "yunion.io/x/onecloud/pkg/util/choices" +) + +const ( + PM_SCOPE_VPC = "vpc" + PM_SCOPE_NETWORK = "network" +) + +var PM_SCOPES = choices.NewChoices( + PM_SCOPE_VPC, + PM_SCOPE_NETWORK, +) + +const ( + FORWARD_TYPE_LOCAL = "local" + FORWARD_TYPE_REMOTE = "remote" +) + +var FORWARD_TYPES = choices.NewChoices( + FORWARD_TYPE_LOCAL, + FORWARD_TYPE_REMOTE, +) + +const ( + BindPortMin = 20000 + BindPortMax = 60000 +) diff --git a/pkg/apis/cloudproxy/doc.go b/pkg/apis/cloudproxy/doc.go new file mode 100644 index 0000000000..5d8cb2d9ea --- /dev/null +++ b/pkg/apis/cloudproxy/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 cloudproxy // import "yunion.io/x/onecloud/pkg/apis/cloudproxy" diff --git a/pkg/apis/cloudproxy/forwards.go b/pkg/apis/cloudproxy/forwards.go new file mode 100644 index 0000000000..c21a1417b6 --- /dev/null +++ b/pkg/apis/cloudproxy/forwards.go @@ -0,0 +1,55 @@ +// Copyright 2019 Yunion +// +// 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 cloudproxy + +import ( + "yunion.io/x/onecloud/pkg/apis" +) + +type ForwardCreateInput struct { + apis.VirtualResourceCreateInput + + ProxyEndpointId string + ProxyAgentId string + Type string + BindPortReq int `json:",omitzero"` + RemoteAddr string + RemotePort string + + LastSeenTimeout int `json:",omitzero"` + + Opaque string +} + +type ForwardCreateFromServerInput struct { + ServerId string + + Type string + BindPortReq int `json:",omitzero"` + RemotePort string + + LastSeenTimeout int `json:",omitzero"` +} + +type ForwardHeartbeatInput struct{} + +type ForwardListInput struct { + ProxyAgentId string + ProxyEndpointId string + + Type string + + Opaque string +} diff --git a/pkg/apis/cloudproxy/proxy_agents.go b/pkg/apis/cloudproxy/proxy_agents.go new file mode 100644 index 0000000000..99a57faf4f --- /dev/null +++ b/pkg/apis/cloudproxy/proxy_agents.go @@ -0,0 +1,35 @@ +// Copyright 2019 Yunion +// +// 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 cloudproxy + +import ( + "yunion.io/x/onecloud/pkg/apis" +) + +type ProxyAgentUpdateInput struct { + apis.StandaloneResourceBaseUpdateInput + + BindAddr string `json:"bind_addr"` + AdvertiseAddr string `json:"advertise_addr"` +} + +type ProxyAgentDetails struct { + apis.StandaloneResourceDetails + + Id string `json:"id"` + + BindAddr string `json:"bind_addr"` + AdvertiseAddr string `json:"advertise_addr"` +} diff --git a/pkg/apis/cloudproxy/proxy_endpoints.go b/pkg/apis/cloudproxy/proxy_endpoints.go new file mode 100644 index 0000000000..3371219d8a --- /dev/null +++ b/pkg/apis/cloudproxy/proxy_endpoints.go @@ -0,0 +1,50 @@ +// Copyright 2019 Yunion +// +// 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 cloudproxy + +import ( + "yunion.io/x/onecloud/pkg/apis" +) + +type ProxyEndpointCreateInput struct { + apis.VirtualResourceCreateInput + + User string + Host string + Port int `json:",omitzero"` + PrivateKey string + + IntranetIpAddr string +} + +type ProxyEndpointCreateFromServerInput struct { + ServerId string +} + +type ProxyEndpointUpdateInput struct { + apis.VirtualResourceBaseUpdateInput + + User string `json:"user"` + Host string `json:"host"` + Port *int `json:"port"` + PrivateKey string `json:"private_key"` +} + +type ProxyEndpointListInput struct { + apis.VirtualResourceBaseUpdateInput + + VpcId string + NetworkId string +} diff --git a/pkg/cloudproxy/agent/models/doc.go b/pkg/cloudproxy/agent/models/doc.go new file mode 100644 index 0000000000..22023092c7 --- /dev/null +++ b/pkg/cloudproxy/agent/models/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 models // import "yunion.io/x/onecloud/pkg/proxyagent/models" diff --git a/pkg/cloudproxy/agent/models/models.go b/pkg/cloudproxy/agent/models/models.go new file mode 100644 index 0000000000..4298118968 --- /dev/null +++ b/pkg/cloudproxy/agent/models/models.go @@ -0,0 +1,43 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + proxy_models "yunion.io/x/onecloud/pkg/cloudproxy/models" +) + +type ProxyEndpoint struct { + proxy_models.SProxyEndpoint + + Forwards Forwards `json:"-"` +} + +func (el *ProxyEndpoint) Copy() *ProxyEndpoint { + return &ProxyEndpoint{ + SProxyEndpoint: el.SProxyEndpoint, + } +} + +type Forward struct { + proxy_models.SForward + + ProxyEndpoint *ProxyEndpoint +} + +func (el *Forward) Copy() *Forward { + return &Forward{ + SForward: el.SForward, + } +} diff --git a/pkg/cloudproxy/agent/models/modelset.go b/pkg/cloudproxy/agent/models/modelset.go new file mode 100644 index 0000000000..2918ec212a --- /dev/null +++ b/pkg/cloudproxy/agent/models/modelset.go @@ -0,0 +1,90 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/apihelper" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + mcclient_modulebase "yunion.io/x/onecloud/pkg/mcclient/modulebase" + mcclient_modules "yunion.io/x/onecloud/pkg/mcclient/modules/cloudproxy" +) + +type ( + ProxyEndpoints map[string]*ProxyEndpoint + Forwards map[string]*Forward +) + +func (set ProxyEndpoints) ModelManager() mcclient_modulebase.IBaseManager { + return &mcclient_modules.ProxyEndpoints +} + +func (set ProxyEndpoints) NewModel() db.IModel { + return &ProxyEndpoint{} +} + +func (set ProxyEndpoints) AddModel(i db.IModel) { + m := i.(*ProxyEndpoint) + set[m.Id] = m +} + +func (set ProxyEndpoints) Copy() apihelper.IModelSet { + setCopy := ProxyEndpoints{} + for id, el := range set { + setCopy[id] = el.Copy() + } + return setCopy +} + +func (ms ProxyEndpoints) joinForwards(subEntries Forwards) bool { + correct := true + for _, subEntry := range subEntries { + epId := subEntry.ProxyEndpointId + m, ok := ms[epId] + if !ok { + log.Warningf("proxy_endpoint_id %s of forward %s(%s) is not present", epId, subEntry.Name, subEntry.Id) + correct = false + continue + } + subEntry.ProxyEndpoint = m + if m.Forwards == nil { + m.Forwards = Forwards{} + } + m.Forwards[subEntry.Id] = subEntry + } + return correct +} + +func (set Forwards) ModelManager() mcclient_modulebase.IBaseManager { + return &mcclient_modules.Forwards +} + +func (set Forwards) NewModel() db.IModel { + return &Forward{} +} + +func (set Forwards) AddModel(i db.IModel) { + m := i.(*Forward) + set[m.Id] = m +} + +func (set Forwards) Copy() apihelper.IModelSet { + setCopy := Forwards{} + for id, el := range set { + setCopy[id] = el.Copy() + } + return setCopy +} diff --git a/pkg/cloudproxy/agent/models/modelsets.go b/pkg/cloudproxy/agent/models/modelsets.go new file mode 100644 index 0000000000..379622b984 --- /dev/null +++ b/pkg/cloudproxy/agent/models/modelsets.go @@ -0,0 +1,106 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "time" + + "yunion.io/x/onecloud/pkg/apihelper" +) + +type ModelSetsMaxUpdatedAt struct { + ProxyEndpoints time.Time + Forwards time.Time +} + +func NewModelSetsMaxUpdatedAt() *ModelSetsMaxUpdatedAt { + return &ModelSetsMaxUpdatedAt{ + ProxyEndpoints: apihelper.PseudoZeroTime, + Forwards: apihelper.PseudoZeroTime, + } +} + +type ModelSets struct { + ProxyEndpoints ProxyEndpoints + Forwards Forwards +} + +func NewModelSets() *ModelSets { + return &ModelSets{ + ProxyEndpoints: ProxyEndpoints{}, + Forwards: Forwards{}, + } +} + +func (mss *ModelSets) ModelSetList() []apihelper.IModelSet { + // it's ordered this way to favour creation, not deletion + return []apihelper.IModelSet{ + mss.ProxyEndpoints, + mss.Forwards, + } +} + +func (mss *ModelSets) NewEmpty() apihelper.IModelSets { + return NewModelSets() +} + +func (mss *ModelSets) copy_() *ModelSets { + mssCopy := &ModelSets{ + ProxyEndpoints: mss.ProxyEndpoints.Copy().(ProxyEndpoints), + Forwards: mss.Forwards.Copy().(Forwards), + } + return mssCopy +} + +func (mss *ModelSets) Copy() apihelper.IModelSets { + return mss.copy_() +} + +func (mss *ModelSets) CopyJoined() apihelper.IModelSets { + mssCopy := mss.copy_() + mssCopy.join() + return mssCopy +} + +func (mss *ModelSets) ApplyUpdates(mssNews apihelper.IModelSets) apihelper.ModelSetsUpdateResult { + r := apihelper.ModelSetsUpdateResult{ + Changed: false, + Correct: true, + } + mssList := mss.ModelSetList() + mssNewsList := mssNews.ModelSetList() + for i, mss := range mssList { + mssNews := mssNewsList[i] + msR := apihelper.ModelSetApplyUpdates(mss, mssNews) + if !r.Changed && msR.Changed { + r.Changed = true + } + } + if r.Changed { + r.Correct = mss.join() + } + return r +} + +func (mss *ModelSets) join() bool { + var p []bool + p = append(p, mss.ProxyEndpoints.joinForwards(mss.Forwards)) + for _, b := range p { + if !b { + return false + } + } + return true +} diff --git a/pkg/cloudproxy/agent/options/doc.go b/pkg/cloudproxy/agent/options/doc.go new file mode 100644 index 0000000000..70db70ed62 --- /dev/null +++ b/pkg/cloudproxy/agent/options/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 options // import "yunion.io/x/onecloud/pkg/proxyagent/options" diff --git a/pkg/cloudproxy/agent/options/options.go b/pkg/cloudproxy/agent/options/options.go new file mode 100644 index 0000000000..32c269691d --- /dev/null +++ b/pkg/cloudproxy/agent/options/options.go @@ -0,0 +1,54 @@ +// Copyright 2019 Yunion +// +// 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 options + +import ( + "fmt" + "time" +) + +type Options struct { + ProxyAgentId string + ProxyAgentInitWait string `help:"duration to try and wait for init" default:"15s"` + proxyAgentInitWaitDuration time.Duration + + APISyncInterval int `default:"10"` + APIListBatchSize int `default:"1024"` +} + +func (opts *Options) GetProxyAgentInitWaitDuration() time.Duration { + return opts.proxyAgentInitWaitDuration +} + +func (opts *Options) ValidateThenInit() error { + if opts.ProxyAgentId == "" { + return fmt.Errorf("empty proxy_agent_id") + } + + if d, err := time.ParseDuration(opts.ProxyAgentInitWait); err != nil { + return fmt.Errorf("parse proxy_agent_init_wait: %v", err) + } else { + opts.proxyAgentInitWaitDuration = d + } + + if opts.APIListBatchSize <= 20 { + opts.APIListBatchSize = 20 + } + if opts.APISyncInterval <= 10 { + opts.APISyncInterval = 10 + } + + return nil +} diff --git a/pkg/cloudproxy/agent/ssh/client.go b/pkg/cloudproxy/agent/ssh/client.go new file mode 100644 index 0000000000..1a03dfed8f --- /dev/null +++ b/pkg/cloudproxy/agent/ssh/client.go @@ -0,0 +1,345 @@ +// Copyright 2019 Yunion +// +// 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 ssh + +import ( + "context" + "fmt" + "net" + "sync" + "time" + + "golang.org/x/crypto/ssh" + + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/sets" +) + +type addrMap map[string]interface{} +type portMap map[int]addrMap + +func (pm portMap) contains(port int, addr string) bool { + am, ok := pm[port] + if !ok { + return false + } + return am.contains(addr) +} + +func (pm portMap) get(port int, addr string) interface{} { + am, ok := pm[port] + if ok { + return am.get(addr) + } + return nil +} + +func (pm portMap) set(port int, addr string, v interface{}) { + am, ok := pm[port] + if !ok { + am = addrMap{} + pm[port] = am + } + am.set(addr, v) +} + +func (pm portMap) delete(port int, addr string) { + if am, ok := pm[port]; ok { + am.delete(addr) + } +} + +func (am addrMap) contains(addr string) bool { + const ( + ip4wild = "0.0.0.0" + ip6wild = "::" + ) + _, ok := am[addr] + if ok { + return true + } + if _, ok := am[ip4wild]; ok { + return true + } + if _, ok := am[ip6wild]; ok { + return true + } + return false +} + +func (am addrMap) get(addr string) interface{} { + return am[addr] +} + +func (am addrMap) set(addr string, v interface{}) { + am[addr] = v +} + +func (am addrMap) delete(addr string) { + delete(am, addr) +} + +type Client struct { + cc *ClientConfig + c *ssh.Client + + stopc chan sets.Empty + stopcEx *sync.Mutex + stopcc bool + + wakec chan sets.Empty + lfc chan LocalForwardReq + rfc chan RemoteForwardReq + + lfclosec chan LocalForwardReq + rfclosec chan RemoteForwardReq + + localForwards portMap + remoteForwards portMap +} + +func NewClient(cc *ClientConfig) *Client { + c := &Client{ + cc: cc, + + stopc: make(chan sets.Empty), + stopcEx: &sync.Mutex{}, + + wakec: make(chan sets.Empty), + lfc: make(chan LocalForwardReq), + rfc: make(chan RemoteForwardReq), + + lfclosec: make(chan LocalForwardReq), + rfclosec: make(chan RemoteForwardReq), + + localForwards: portMap{}, + remoteForwards: portMap{}, + } + return c +} + +func (c *Client) Stop(ctx context.Context) { + c.stopcEx.Lock() + defer c.stopcEx.Unlock() + if !c.stopcc { + close(c.stopc) + c.stopcc = true + } +} + +func (c *Client) Start(ctx context.Context) { + pingT := time.NewTimer(17 * time.Second) + pingFailCount := 0 + const pingMaxFail = 3 + + const ( + stateInit = iota + stateOK + ) + state := stateInit + stateC := make(chan int) + stateInitRetryInterval := 7 * time.Second + stateInitRetryT := time.NewTicker(stateInitRetryInterval) + for { + switch state { + case stateOK: + // check forwards and start + case stateInit: + if sshc, err := c.connect(ctx); err != nil { + log.Errorf("ssh connect: %v", err) + } else { + c.c = sshc + state = stateOK + go func() { + defer c.c.Conn.Close() + + err := c.c.Conn.Wait() + if err != nil { + log.Errorf("ssh client conn: %v", err) + } + select { + case stateC <- stateInit: + case <-ctx.Done(): + } + }() + } + } + + select { + case req := <-c.lfc: + if c.c != nil { + c.localForward(ctx, req) + } + case req := <-c.rfc: + if c.c != nil { + c.remoteForward(ctx, req) + } + case req := <-c.lfclosec: + c.localForwardClose(ctx, req) + case req := <-c.rfclosec: + c.remoteForwardClose(ctx, req) + case <-c.wakec: + break + case <-pingT.C: + //TODO ping check + //ping fail + if pingFailCount > pingMaxFail { + state = stateInit + } + case newState := <-stateC: + state = newState + case <-stateInitRetryT.C: + case <-c.stopc: + if c.c != nil { + c.c.Conn.Close() + } + return + case <-ctx.Done(): + return + } + } +} + +func (c *Client) connect(ctx context.Context) (*ssh.Client, error) { + sshc, err := c.cc.NewClient(ctx) + return sshc, err +} + +func (c *Client) LocalForward(ctx context.Context, req LocalForwardReq) { + select { + case c.lfc <- req: + case <-ctx.Done(): + } +} + +func (c *Client) localForward(ctx context.Context, req LocalForwardReq) { + if err := c.localForward_(ctx, req); err != nil { + log.Errorf("local forward: %v", err) + } +} + +func (c *Client) localForward_(ctx context.Context, req LocalForwardReq) error { + // check LocalAddr/LocalPort existence + if c.localForwards.contains(req.LocalPort, req.LocalAddr) { + return errors.Errorf("local addr occupied: %s:%d", req.LocalAddr, req.LocalPort) + } + + addr := net.JoinHostPort(req.LocalAddr, fmt.Sprintf("%d", req.LocalPort)) + listener, err := net.Listen("tcp", addr) + if err != nil { + return errors.Wrapf(err, "tcp listen %s", addr) + } + fwd := &forwarder{ + listener: listener, + + dial: c.c.Dial, + dialAddr: req.RemoteAddr, + dialPort: req.RemotePort, + + done: c.localForwardDone, + doneAddr: req.LocalAddr, + donePort: req.LocalPort, + + tick: req.Tick, + tickCb: req.TickCb, + } + + c.localForwards.set(req.LocalPort, req.LocalAddr, fwd) + go fwd.Start(ctx) + return nil +} + +func (c *Client) localForwardDone(laddr string, lport int) { + c.localForwards.delete(lport, laddr) +} + +func (c *Client) RemoteForward(ctx context.Context, req RemoteForwardReq) { + select { + case c.rfc <- req: + case <-ctx.Done(): + } +} + +func (c *Client) remoteForward(ctx context.Context, req RemoteForwardReq) { + if err := c.remoteForward_(ctx, req); err != nil { + log.Errorf("remote forward: %v", err) + } +} + +func (c *Client) remoteForward_(ctx context.Context, req RemoteForwardReq) error { + // check RemoteAddr/RemotePort existence + if c.remoteForwards.contains(req.RemotePort, req.RemoteAddr) { + return errors.Errorf("remote addr occupied: %s:%d", req.RemoteAddr, req.RemotePort) + } + + addr := net.JoinHostPort(req.RemoteAddr, fmt.Sprintf("%d", req.RemotePort)) + listener, err := c.c.Listen("tcp", addr) + if err != nil { + return errors.Wrapf(err, "ssh listen %s", addr) + } + + fwd := &forwarder{ + listener: listener, + + dial: net.Dial, + dialAddr: req.LocalAddr, + dialPort: req.LocalPort, + + done: c.remoteForwardDone, + doneAddr: req.RemoteAddr, + donePort: req.RemotePort, + + tick: req.Tick, + tickCb: req.TickCb, + } + c.remoteForwards.set(req.RemotePort, req.RemoteAddr, fwd) + go fwd.Start(ctx) + return nil +} + +func (c *Client) remoteForwardDone(raddr string, rport int) { + c.remoteForwards.delete(rport, raddr) +} + +func (c *Client) LocalForwardClose(ctx context.Context, req LocalForwardReq) { + select { + case c.lfclosec <- req: + case <-ctx.Done(): + } +} + +func (c *Client) localForwardClose(ctx context.Context, req LocalForwardReq) { + v := c.localForwards.get(req.LocalPort, req.LocalAddr) + if v != nil { + fwd := v.(*forwarder) + fwd.Stop(ctx) + } +} + +func (c *Client) RemoteForwardClose(ctx context.Context, req RemoteForwardReq) { + select { + case c.rfclosec <- req: + case <-ctx.Done(): + } +} + +func (c *Client) remoteForwardClose(ctx context.Context, req RemoteForwardReq) { + v := c.remoteForwards.get(req.RemotePort, req.RemoteAddr) + if v != nil { + fwd := v.(*forwarder) + fwd.Stop(ctx) + } +} diff --git a/pkg/cloudproxy/agent/ssh/client_config.go b/pkg/cloudproxy/agent/ssh/client_config.go new file mode 100644 index 0000000000..c5f8c9b691 --- /dev/null +++ b/pkg/cloudproxy/agent/ssh/client_config.go @@ -0,0 +1,62 @@ +// Copyright 2019 Yunion +// +// 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 ssh + +import ( + "context" + "fmt" + "net" + + "golang.org/x/crypto/ssh" + + "yunion.io/x/pkg/errors" +) + +type ClientConfig struct { + User string + Host string + Port int + Key string +} + +func (cc *ClientConfig) NewClient(ctx context.Context) (*ssh.Client, error) { + signer, err := ssh.ParsePrivateKey([]byte(cc.Key)) + if err != nil { + return nil, errors.Wrap(err, "parse ssh key") + } + sshcc := &ssh.ClientConfig{ + User: cc.User, + Auth: []ssh.AuthMethod{ + ssh.PublicKeys(signer), + }, + HostKeyCallback: ssh.InsecureIgnoreHostKey(), + } + + addr := net.JoinHostPort(cc.Host, fmt.Sprintf("%d", cc.Port)) + d := &net.Dialer{} + netconn, err := d.DialContext(ctx, "tcp", addr) + if err != nil { + return nil, errors.Wrap(err, "net dial") + } + + sshconn, chans, reqs, err := ssh.NewClientConn(netconn, addr, sshcc) + if err != nil { + netconn.Close() + return nil, errors.Wrap(err, "ssh new client conn") + } + + sshc := ssh.NewClient(sshconn, chans, reqs) + return sshc, nil +} diff --git a/pkg/cloudproxy/agent/ssh/clientset.go b/pkg/cloudproxy/agent/ssh/clientset.go new file mode 100644 index 0000000000..05a75de780 --- /dev/null +++ b/pkg/cloudproxy/agent/ssh/clientset.go @@ -0,0 +1,184 @@ +// Copyright 2019 Yunion +// +// 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 ssh + +import ( + "context" +) + +type epClientSet struct { + cc ClientConfig + clients []*Client + + mark bool +} + +func (epcs *epClientSet) clearMark() { + epcs.mark = false +} + +func (epcs *epClientSet) setMark() { + epcs.mark = true +} + +func (epcs *epClientSet) getMark() bool { + return epcs.mark +} + +func (epcs *epClientSet) stop(ctx context.Context) { + for _, client := range epcs.clients { + client.Stop(ctx) + } +} + +type epClients map[string]*epClientSet // key: epKey + +type ClientSet struct { + epClients epClients +} + +func NewClientSet() *ClientSet { + cs := &ClientSet{ + epClients: epClients{}, + } + return cs +} + +func (cs *ClientSet) ClearAllMark() { + for _, epcs := range cs.epClients { + epcs.clearMark() + } +} + +func (cs *ClientSet) ResetIfChanged(ctx context.Context, epKey string, cc ClientConfig) bool { + epcs, ok := cs.epClients[epKey] + if ok { + if epcs.cc != cc { + epcs.stop(ctx) + delete(cs.epClients, epKey) + return true + } + epcs.setMark() + } + return false +} + +func (cs *ClientSet) AddIfNotExist(ctx context.Context, epKey string, cc ClientConfig) bool { + epcs, ok := cs.epClients[epKey] + if !ok { + epcs := &epClientSet{ + cc: cc, + } + epcs.setMark() + cs.epClients[epKey] = epcs + return true + } + epcs.setMark() + return false +} + +func (cs *ClientSet) ResetUnmarked(ctx context.Context) { + for epKey, epcs := range cs.epClients { + if !epcs.getMark() { + epcs.stop(ctx) + delete(cs.epClients, epKey) + } + } +} + +func (cs *ClientSet) ForwardKeySet() ForwardKeySet { + fks := ForwardKeySet{} + for epKey, epcs := range cs.epClients { + for _, client := range epcs.clients { + fks.addByPortMap(epKey, ForwardKeyTypeL, client.localForwards) + fks.addByPortMap(epKey, ForwardKeyTypeR, client.remoteForwards) + } + } + return fks +} + +func (cs *ClientSet) LocalForward(ctx context.Context, epKey string, req LocalForwardReq) { + client, created := cs.getOrCreateClient(epKey, ForwardKeyTypeL) + if created { + go client.Start(ctx) + } + client.LocalForward(ctx, req) +} + +func (cs *ClientSet) RemoteForward(ctx context.Context, epKey string, req RemoteForwardReq) { + client, created := cs.getOrCreateClient(epKey, ForwardKeyTypeR) + if created { + go client.Start(ctx) + } + client.RemoteForward(ctx, req) +} + +func (cs *ClientSet) CloseForward(ctx context.Context, fk ForwardKey) { + client := cs.getClient(fk.EpKey, fk.Type) + if client == nil { + return + } + switch fk.Type { + case ForwardKeyTypeL: + client.LocalForwardClose(ctx, LocalForwardReq{ + LocalAddr: fk.KeyAddr, + LocalPort: fk.KeyPort, + }) + case ForwardKeyTypeR: + client.RemoteForwardClose(ctx, RemoteForwardReq{ + RemoteAddr: fk.KeyAddr, + RemotePort: fk.KeyPort, + }) + } +} + +/* +func (cs *ClientSet) LocalForwardClose(ctx context.Context, epKey string, req LocalForwardReq) { + client := cs.getClient(epKey, ForwardKeyTypeL) + client.LocalForwardClose(ctx, req) +} + +func (cs *ClientSet) RemoteForwardClose(ctx context.Context, epKey string, req RemoteForwardReq) { + client := cs.getClient(epKey, ForwardKeyTypeR) + client.RemoteForwardClose(ctx, req) +} +*/ + +func (cs *ClientSet) getOrCreateClient(epKey string, typ string) (*Client, bool) { + return cs.getClient_(epKey, typ, true) +} + +func (cs *ClientSet) getClient(epKey string, typ string) *Client { + client, _ := cs.getClient_(epKey, typ, false) + return client +} + +func (cs *ClientSet) getClient_(epKey string, typ string, create bool) (*Client, bool) { + var client *Client + + clients, ok := cs.epClients[epKey] + if !ok || len(clients.clients) == 0 { + if !create { + return nil, false + } + client = NewClient(&clients.cc) + clients.clients = append(clients.clients, client) + cs.epClients[epKey] = clients + return client, true + } else { + client = clients.clients[0] + return client, false + } +} diff --git a/pkg/cloudproxy/agent/ssh/doc.go b/pkg/cloudproxy/agent/ssh/doc.go new file mode 100644 index 0000000000..676ee2b000 --- /dev/null +++ b/pkg/cloudproxy/agent/ssh/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 ssh // import "yunion.io/x/onecloud/pkg/proxyagent/agent/ssh" diff --git a/pkg/cloudproxy/agent/ssh/forward_key.go b/pkg/cloudproxy/agent/ssh/forward_key.go new file mode 100644 index 0000000000..8a1c3b7ee3 --- /dev/null +++ b/pkg/cloudproxy/agent/ssh/forward_key.go @@ -0,0 +1,69 @@ +// Copyright 2019 Yunion +// +// 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 ssh + +import ( + "fmt" +) + +const ( + ForwardKeyTypeL = "L" + ForwardKeyTypeR = "R" +) + +type ForwardKey struct { + EpKey string + Type string + KeyAddr string + KeyPort int + + Value interface{} +} + +func (fk *ForwardKey) Key() string { + return fmt.Sprintf("%s/%s/%s/%d", fk.EpKey, fk.Type, fk.KeyAddr, fk.KeyPort) +} + +type ForwardKeySet map[string]ForwardKey + +func (fks ForwardKeySet) addByPortMap(epKey, typ string, pm portMap) { + for port, addrMap := range pm { + for addr := range addrMap { + fk := ForwardKey{ + EpKey: epKey, + Type: typ, + KeyAddr: addr, + KeyPort: port, + } + fks.Add(fk) + } + } +} + +func (fks ForwardKeySet) Contains(fk ForwardKey) bool { + key := fk.Key() + if _, ok := fks[key]; ok { + return true + } + return false +} + +func (fks ForwardKeySet) Remove(fk ForwardKey) { + delete(fks, fk.Key()) +} + +func (fks ForwardKeySet) Add(fk ForwardKey) { + fks[fk.Key()] = fk +} diff --git a/pkg/cloudproxy/agent/ssh/forwarder.go b/pkg/cloudproxy/agent/ssh/forwarder.go new file mode 100644 index 0000000000..f7428a00a1 --- /dev/null +++ b/pkg/cloudproxy/agent/ssh/forwarder.go @@ -0,0 +1,147 @@ +// Copyright 2019 Yunion +// +// 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 ssh + +import ( + "context" + "fmt" + "io" + "net" + "time" + + "yunion.io/x/log" +) + +type TickFunc func(context.Context) + +type LocalForwardReq struct { + LocalAddr string + LocalPort int + RemoteAddr string + RemotePort int + + Tick time.Duration + TickCb TickFunc +} + +type RemoteForwardReq struct { + // LocalAddr is the address the forward will forward to + LocalAddr string + // LocalPort is the port the forward will forward to + LocalPort int + + // RemoteAddr is the address on the remote to listen on + RemoteAddr string + // RemotePort is the address on the remote to listen on + RemotePort int + + Tick time.Duration + TickCb TickFunc +} + +type dialFunc func(n, addr string) (net.Conn, error) +type doneFunc func(laddr string, lport int) + +type forwarder struct { + listener net.Listener + + dial dialFunc + dialAddr string + dialPort int + + done doneFunc + doneAddr string + donePort int + + tick time.Duration + tickCb TickFunc +} + +func (fwd *forwarder) Stop(ctx context.Context) { + fwd.listener.Close() +} + +func (fwd *forwarder) Start( + ctx context.Context, +) { + var ( + listener = fwd.listener + dial = fwd.dial + dialAddr = fwd.dialAddr + dialPort = fwd.dialPort + done = fwd.done + doneAddr = fwd.doneAddr + donePort = fwd.donePort + tick = fwd.tick + tickCb = fwd.tickCb + ) + + ctx, cancelFunc := context.WithCancel(ctx) + + if done != nil { + defer done(doneAddr, donePort) + } + + defer listener.Close() + + go func() { // accept local/remote connection + for { + conn, err := listener.Accept() + if err != nil { + log.Warningf("local forward: accept: %v", err) + cancelFunc() + break + } + go func(local net.Conn) { + defer local.Close() + + // dial remote/local + addr := net.JoinHostPort(dialAddr, fmt.Sprintf("%d", dialPort)) + remote, err := dial("tcp", addr) + if err != nil { + log.Warningf("local forward: dial remote: %v", err) + return + } + defer remote.Close() + + // forward + go io.Copy(local, remote) + go io.Copy(remote, local) + <-ctx.Done() + }(conn) + } + }() + + if tick > 0 && tickCb != nil { + go func() { + ticker := time.NewTicker(tick) + defer ticker.Stop() + for { + select { + case <-ticker.C: + tickCb(ctx) + case <-ctx.Done(): + return + } + } + }() + } + for { + select { + case <-ctx.Done(): + return + } + } +} diff --git a/pkg/cloudproxy/agent/worker/doc.go b/pkg/cloudproxy/agent/worker/doc.go new file mode 100644 index 0000000000..892cb3c9c8 --- /dev/null +++ b/pkg/cloudproxy/agent/worker/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 worker // import "yunion.io/x/onecloud/pkg/proxyagent/agent/worker" diff --git a/pkg/cloudproxy/agent/worker/tick.go b/pkg/cloudproxy/agent/worker/tick.go new file mode 100644 index 0000000000..a9f24c7496 --- /dev/null +++ b/pkg/cloudproxy/agent/worker/tick.go @@ -0,0 +1,44 @@ +// Copyright 2019 Yunion +// +// 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 worker + +import ( + "context" + "math/rand" + "time" + + "yunion.io/x/log" + + agentssh "yunion.io/x/onecloud/pkg/cloudproxy/agent/ssh" + "yunion.io/x/onecloud/pkg/mcclient/auth" + cloudproxy_modules "yunion.io/x/onecloud/pkg/mcclient/modules/cloudproxy" +) + +func tickDuration(timeout int) time.Duration { + if timeout > 30 { + return time.Duration(timeout-5-rand.Intn(10)) * time.Second + } + return (time.Duration(timeout) * time.Second / 3) * 2 +} + +func heartbeatFunc(fwdId string, sessionCache *auth.SessionCache) agentssh.TickFunc { + return func(ctx context.Context) { + s := sessionCache.Get(ctx) + _, err := cloudproxy_modules.Forwards.PerformAction(s, fwdId, "heartbeat", nil) + if err != nil { + log.Errorf("forwarder heartbeat: %s: %v", fwdId, err) + } + } +} diff --git a/pkg/cloudproxy/agent/worker/worker.go b/pkg/cloudproxy/agent/worker/worker.go new file mode 100644 index 0000000000..32effb402d --- /dev/null +++ b/pkg/cloudproxy/agent/worker/worker.go @@ -0,0 +1,275 @@ +// Copyright 2019 Yunion +// +// 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 worker + +import ( + "context" + "runtime" + "runtime/debug" + "sync" + "time" + + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/utils" + + "yunion.io/x/onecloud/pkg/apihelper" + api "yunion.io/x/onecloud/pkg/apis/cloudproxy" + common_options "yunion.io/x/onecloud/pkg/cloudcommon/options" + agentmodels "yunion.io/x/onecloud/pkg/cloudproxy/agent/models" + agentoptions "yunion.io/x/onecloud/pkg/cloudproxy/agent/options" + agentssh "yunion.io/x/onecloud/pkg/cloudproxy/agent/ssh" + "yunion.io/x/onecloud/pkg/mcclient/auth" + cloudproxy_modules "yunion.io/x/onecloud/pkg/mcclient/modules/cloudproxy" + "yunion.io/x/onecloud/pkg/util/netutils2" +) + +type Worker struct { + commonOpts *common_options.CommonOptions + opts *agentoptions.Options + + proxyAgentId string + bindAddr string + + apih *apihelper.APIHelper + clientSet *agentssh.ClientSet + sessionCache *auth.SessionCache +} + +func NewWorker(commonOpts *common_options.CommonOptions, opts *agentoptions.Options) *Worker { + modelSets := agentmodels.NewModelSets() + apiOpts := &apihelper.Options{ + CommonOptions: *commonOpts, + SyncInterval: opts.APISyncInterval, + ListBatchSize: opts.APIListBatchSize, + } + apih, err := apihelper.NewAPIHelper(apiOpts, modelSets) + if err != nil { + return nil + } + w := &Worker{ + commonOpts: commonOpts, + opts: opts, + proxyAgentId: opts.ProxyAgentId, + + apih: apih, + clientSet: agentssh.NewClientSet(), + sessionCache: &auth.SessionCache{ + Region: commonOpts.Region, + APIVersion: "v2", + UseAdminToken: true, + EarlyRefresh: time.Hour, + }, + } + return w +} + +func (w *Worker) initProxyAgent_(ctx context.Context) error { + s := w.sessionCache.Get(ctx) + + var agentDetail api.ProxyAgentDetails + { + j, err := cloudproxy_modules.ProxyAgents.Get(s, w.proxyAgentId, nil) + if err != nil { + return errors.Wrapf(err, "fetch proxy agent %s", w.proxyAgentId) + } + if err := j.Unmarshal(&agentDetail); err != nil { + return errors.Wrapf(err, "unmarshal proxy agent detail: %s", j.String()) + } + if agentDetail.Id == "" { + return errors.Error("proxy agent id is empty") + } + w.proxyAgentId = agentDetail.Id + } + + var bindAddr = agentDetail.BindAddr + if bindAddr == "" { + var err error + bindAddr, err = netutils2.MyIP() + if err != nil { + return errors.Wrap(err, "find bind Addr") + } + } + w.bindAddr = bindAddr + + if agentDetail.BindAddr == "" || agentDetail.AdvertiseAddr == "" { + var advertiseAddr = agentDetail.AdvertiseAddr + if advertiseAddr == "" { + if true { //TODO, fetch clusterIP from k8s environ + advertiseAddr = bindAddr + } + } + req := api.ProxyAgentUpdateInput{ + BindAddr: bindAddr, + AdvertiseAddr: advertiseAddr, + } + reqJ := req.JSON(req) + if _, err := cloudproxy_modules.ProxyAgents.Put(s, w.proxyAgentId, reqJ); err != nil { + return errors.Wrapf(err, "update proxy agent addr: %s", reqJ.String()) + } + } + + return nil +} + +func (w *Worker) initProxyAgent(ctx context.Context) error { + done, err := utils.NewFibonacciRetrierMaxElapse( + w.opts.GetProxyAgentInitWaitDuration(), + func(retrier utils.FibonacciRetrier) (bool, error) { + err := w.initProxyAgent_(ctx) + if err != nil { + return false, err + } + return true, nil + }).Start(ctx) + if done { + return nil + } + return err +} + +func (w *Worker) Start(ctx context.Context) { + wg := ctx.Value("wg").(*sync.WaitGroup) + wg.Add(1) + defer func() { + log.Infoln("agent: worker bye") + wg.Done() + }() + + if err := w.initProxyAgent(ctx); err != nil { + log.Errorf("init proxy agent: %v", err) + return + } + go w.apih.Start(ctx) + + var mss *agentmodels.ModelSets + for { + select { + case imss := <-w.apih.ModelSets(): + log.Infof("agent: got new data from api helper") + mss = imss.(*agentmodels.ModelSets) + if err := w.run(ctx, mss); err != nil { + log.Errorf("agent: %v", err) + } + case <-ctx.Done(): + return + } + } +} + +func (w *Worker) run(ctx context.Context, mss *agentmodels.ModelSets) (err error) { + defer func() { + if panicVal := recover(); panicVal != nil { + if panicErr, ok := panicVal.(runtime.Error); ok { + err = errors.Wrap(panicErr, string(debug.Stack())) + } else if panicErr, ok := panicVal.(error); ok { + err = panicErr + } else { + panic(panicVal) + } + } + }() + + w.clientSet.ClearAllMark() + for _, pep := range mss.ProxyEndpoints { + cc := agentssh.ClientConfig{ + User: pep.User, + Host: pep.Host, + Port: pep.Port, + Key: pep.PrivateKey, + } + if reset := w.clientSet.ResetIfChanged(ctx, pep.Id, cc); reset { + log.Warningf("proxy endpoint %s changed, connections reset", pep.Id) + } else if added := w.clientSet.AddIfNotExist(ctx, pep.Id, cc); added { + log.Infof("proxy endpoint %s added", pep.Id) + } + } + w.clientSet.ResetUnmarked(ctx) + + removes := w.clientSet.ForwardKeySet() + adds := agentssh.ForwardKeySet{} + for _, pep := range mss.ProxyEndpoints { + for _, forward := range pep.Forwards { + if forward.ProxyAgentId != w.proxyAgentId { + continue + } + if forward.ProxyEndpointId == "" { + continue + } + var ( + typ string + addr string + port int + ) + switch forward.Type { + case api.FORWARD_TYPE_LOCAL: + addr = w.bindAddr + port = forward.BindPort + typ = agentssh.ForwardKeyTypeL + case api.FORWARD_TYPE_REMOTE: + addr = forward.ProxyEndpoint.IntranetIpAddr + port = forward.BindPort + typ = agentssh.ForwardKeyTypeR + default: + log.Warningf("unknown forward type %s", forward.Type) + continue + } + fk := agentssh.ForwardKey{ + EpKey: forward.ProxyEndpointId, + Type: typ, + KeyAddr: addr, + KeyPort: port, + + Value: forward, + } + if removes.Contains(fk) { + removes.Remove(fk) + } else { + adds.Add(fk) + } + } + } + for _, fk := range removes { + log.Infof("close forward %s", fk.Key()) + w.clientSet.CloseForward(ctx, fk) + } + for _, fk := range adds { + log.Infof("open forward %s", fk.Key()) + forward := fk.Value.(*agentmodels.Forward) + tick := tickDuration(forward.LastSeenTimeout) + tickCb := heartbeatFunc(forward.Id, w.sessionCache) + switch fk.Type { + case agentssh.ForwardKeyTypeL: + w.clientSet.LocalForward(ctx, fk.EpKey, agentssh.LocalForwardReq{ + LocalAddr: fk.KeyAddr, + LocalPort: fk.KeyPort, + RemoteAddr: forward.RemoteAddr, + RemotePort: forward.RemotePort, + Tick: tick, + TickCb: tickCb, + }) + case agentssh.ForwardKeyTypeR: + w.clientSet.RemoteForward(ctx, fk.EpKey, agentssh.RemoteForwardReq{ + RemoteAddr: fk.KeyAddr, + RemotePort: fk.KeyPort, + LocalAddr: forward.RemoteAddr, + LocalPort: forward.RemotePort, + Tick: tick, + TickCb: tickCb, + }) + } + } + return nil +} diff --git a/pkg/cloudproxy/models/doc.go b/pkg/cloudproxy/models/doc.go new file mode 100644 index 0000000000..85669243c7 --- /dev/null +++ b/pkg/cloudproxy/models/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 models // import "yunion.io/x/onecloud/pkg/cloudproxy/models" diff --git a/pkg/cloudproxy/models/forwards.go b/pkg/cloudproxy/models/forwards.go new file mode 100644 index 0000000000..64cc8aa78b --- /dev/null +++ b/pkg/cloudproxy/models/forwards.go @@ -0,0 +1,470 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "context" + "fmt" + "math/rand" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/gotypes" + "yunion.io/x/sqlchemy" + + cloudproxy_api "yunion.io/x/onecloud/pkg/apis/cloudproxy" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/validators" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/rbacutils" + "yunion.io/x/onecloud/pkg/util/stringutils2" +) + +type SForward struct { + db.SVirtualResourceBase + + ProxyEndpointId string `width:"36" charset:"ascii" nullable:"false" list:"user" update:"user" create:"required"` + ProxyAgentId string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"` + + Type string `width:"16" charset:"ascii" nullable:"false" list:"user" create:"required"` + RemoteAddr string `width:"16" charset:"ascii" nullable:"false" list:"user" create:"required"` + RemotePort int `width:"16" charset:"ascii" nullable:"false" list:"user" create:"required"` + BindPortReq int `width:"16" charset:"ascii" nullable:"false" list:"user" update:"user" create:"optional"` + + Opaque string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"` + + BindPort int `width:"16" charset:"ascii" nullable:"false" list:"user" update:"user" create:"optional"` + LastSeen time.Time `nullable:"true" get:"user" list:"user"` + LastSeenTimeout int `width:"16" charset:"ascii" nullable:"false" list:"user" update:"user" create:"optional" default:"117"` +} + +type SForwardManager struct { + db.SVirtualResourceBaseManager +} + +var ForwardManager *SForwardManager + +func init() { + ForwardManager = &SForwardManager{ + SVirtualResourceBaseManager: db.NewVirtualResourceBaseManager( + SForward{}, + "forwards_tbl", + "forward", + "forwards", + ), + } + ForwardManager.SetVirtualObject(ForwardManager) +} + +func (man *SForwardManager) validateLocalSetPort(ctx context.Context, data *jsonutils.JSONDict, agentId string, portReq int) (*jsonutils.JSONDict, error) { + var ( + fwds []SForward + q = man.Query(). + Equals("proxy_agent_id", agentId). + Equals("bind_port_req", portReq) + ) + if err := db.FetchModelObjects(man, q, &fwds); err != nil { + return nil, httperrors.NewServerError("query forwards by agent failed: %v", err) + } + if len(fwds) > 0 { + return nil, httperrors.NewConflictError("port %d on agent %s was already occupied", + portReq, agentId) + } + data.Set("bind_port", jsonutils.NewInt(int64(portReq))) + return data, nil +} + +func (man *SForwardManager) validateRemoteSetPort(ctx context.Context, data *jsonutils.JSONDict, epId string, portReq int) (*jsonutils.JSONDict, error) { + var ( + fwds []SForward + q = man.Query(). + Equals("proxy_endpoint_id", epId). + Equals("bind_port_req", portReq) + ) + if err := db.FetchModelObjects(man, q, &fwds); err != nil { + return nil, httperrors.NewServerError("query forwards by endpoint id failed: %v", err) + } + if len(fwds) > 0 { + return nil, httperrors.NewConflictError("port %d on proxy endpoint %s was already occupied", + portReq, epId) + } + data.Set("bind_port", jsonutils.NewInt(int64(portReq))) + return data, nil +} + +func (man *SForwardManager) validateLocalSelectAgent(ctx context.Context, data *jsonutils.JSONDict, portReq int) (*jsonutils.JSONDict, error) { + agents, err := ProxyAgentManager.allAgents(ctx) + if err != nil { + return nil, httperrors.NewGeneralError(err) + } + agentsNum := len(agents) + if agentsNum == 0 { + return nil, httperrors.NewResourceNotFoundError("empty proxy agents set") + } + s := rand.Intn(agentsNum) + for i := s; ; { + agent := &agents[i] + + var err error + data, err = man.validateLocalSetPort(ctx, data, agent.Id, portReq) + if err == nil { + data.Set("proxy_agent_id", jsonutils.NewString(agent.Id)) + return data, nil + } + + i += 1 + if i == agentsNum { + i = 0 + } + if i == s { + break + } + } + return nil, httperrors.NewResourceNotFoundError("no proxy agent accepts request for port %d", portReq) +} + +func (man *SForwardManager) validateRemoteSelectAgent(ctx context.Context, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + agents, err := ProxyAgentManager.allAgents(ctx) + if err != nil { + return nil, httperrors.NewGeneralError(err) + } + agentsNum := len(agents) + if agentsNum == 0 { + return nil, httperrors.NewResourceNotFoundError("empty proxy agents set") + } + i := rand.Intn(agentsNum) + agent := agents[i] + data.Set("proxy_agent_id", jsonutils.NewString(agent.Id)) + return data, nil +} + +func (man *SForwardManager) validatePortReq( + ctx context.Context, + typ string, portReq int, agentId, epId string, + data *jsonutils.JSONDict, +) (*jsonutils.JSONDict, error) { + validateOne := func(portReq int) (*jsonutils.JSONDict, error) { + var err error + switch typ { + case cloudproxy_api.FORWARD_TYPE_LOCAL: + if agentId == "" { + data, err = man.validateLocalSelectAgent(ctx, data, portReq) + } else { + data, err = man.validateLocalSetPort(ctx, data, agentId, portReq) + } + case cloudproxy_api.FORWARD_TYPE_REMOTE: + data, err = man.validateRemoteSetPort(ctx, data, epId, portReq) + } + return data, err + } + + if typ == cloudproxy_api.FORWARD_TYPE_REMOTE && agentId == "" { + var err error + data, err = man.validateRemoteSelectAgent(ctx, data) + if err != nil { + return nil, httperrors.NewResourceNotFoundError("select proxy agent: %v", err) + } + } + + var err error + if portReq <= 0 { + portTotal := cloudproxy_api.BindPortMax - cloudproxy_api.BindPortMin + 1 + portReqStart := rand.Intn(portTotal) + for portInc := portReqStart; ; { + data, err = validateOne(cloudproxy_api.BindPortMin + portInc) + if err == nil { + break + } + portInc += 1 + if portInc == portTotal { + portInc = 0 + } + if portInc == portReqStart { + return nil, httperrors.NewOutOfResourceError("no available port for bind") + } + } + } else { + data, err = validateOne(portReq) + } + return data, err +} + +func (man *SForwardManager) AllowPerformCreateFromServer(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return db.IsAllowClassPerform(rbacutils.ScopeProject, userCred, man, "create-from-server") +} + +func (man *SForwardManager) PerformCreateFromServer(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input *cloudproxy_api.ForwardCreateFromServerInput) (jsonutils.JSONObject, error) { + data := jsonutils.Marshal(input).(*jsonutils.JSONDict) + + typeV := validators.NewStringChoicesValidator("type", cloudproxy_api.FORWARD_TYPES) + portReqV := validators.NewRangeValidator("bind_port_req", cloudproxy_api.BindPortMin, cloudproxy_api.BindPortMax) + remotePortV := validators.NewPortValidator("remote_port") + { + for _, v := range []validators.IValidator{ + typeV, + portReqV.Optional(true), + remotePortV, + + validators.NewNonNegativeValidator("last_seen_timeout").Optional(true), + } { + if err := v.Validate(data); err != nil { + return nil, err + } + } + } + + serverId := input.ServerId + if serverId == "" { + return nil, httperrors.NewBadRequestError("server_id is required") + } + serverInfo, err := getServerInfo(ctx, userCred, serverId) + if err != nil { + return nil, err + } + nic := serverInfo.GetNic() + if nic == nil { + return nil, httperrors.NewBadRequestError("cannot find network interface for this server") + } + proxymatch := ProxyMatchManager.findMatch(ctx, nic.NetworkId, nic.VpcId) + if proxymatch == nil { + return nil, httperrors.NewBadRequestError("cannot find an endpoint for this server") + } + data.Set("opaque", jsonutils.NewString(serverInfo.Server.Id)) + data.Set("remote_addr", jsonutils.NewString(nic.IpAddr)) + data.Set("proxy_endpoint_id", jsonutils.NewString(proxymatch.ProxyEndpointId)) + + typ := typeV.Value + agentId := "" + epId := proxymatch.ProxyEndpointId + if data.Contains("bind_port_req") { + portReq := int(portReqV.Value) + data, err = man.validatePortReq(ctx, typ, portReq, agentId, epId, data) + } else { + data, err = man.validatePortReq(ctx, typ, -1, agentId, epId, data) + } + + forward := &SForward{} + if err := data.Unmarshal(forward); err != nil { + return nil, httperrors.NewServerError("unmarshal create params: %v", err) + } + forward.Name = fmt.Sprintf("%s-%s-%d", serverInfo.Server.Name, typ, forward.RemotePort) + forward.DomainId = userCred.GetProjectDomainId() + forward.ProjectId = userCred.GetProjectId() + if err := man.TableSpec().Insert(ctx, forward); err != nil { + return nil, httperrors.NewServerError("database insertion error: %v", err) + } + + return jsonutils.Marshal(forward), nil +} + +func (man *SForwardManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + endpointV := validators.NewModelIdOrNameValidator("proxy_endpoint", ProxyEndpointManager.Keyword(), ownerId) + agentV := validators.NewModelIdOrNameValidator("proxy_agent", ProxyAgentManager.Keyword(), ownerId) + typeV := validators.NewStringChoicesValidator("type", cloudproxy_api.FORWARD_TYPES) + portReqV := validators.NewRangeValidator("bind_port_req", cloudproxy_api.BindPortMin, cloudproxy_api.BindPortMax) + for _, v := range []validators.IValidator{ + endpointV, + agentV.Optional(true), + + typeV, + validators.NewIPv4AddrValidator("remote_addr"), + validators.NewPortValidator("remote_port"), + portReqV.Optional(true), + + validators.NewNonNegativeValidator("last_seen_timeout").Optional(true), + } { + if err := v.Validate(data); err != nil { + return nil, err + } + } + + typ := typeV.Value + epId := endpointV.Model.GetId() + var agentId string + if agentV.Model != nil { + agentId = agentV.Model.GetId() + } + + var err error + if data.Contains("bind_port_req") { + portReq := int(portReqV.Value) + data, err = man.validatePortReq(ctx, typ, portReq, agentId, epId, data) + } else { + data, err = man.validatePortReq(ctx, typ, -1, agentId, epId, data) + } + return data, err +} + +func (fwd *SForward) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + endpointV := validators.NewModelIdOrNameValidator("proxy_endpoint", ProxyEndpointManager.Keyword(), userCred) + agentV := validators.NewModelIdOrNameValidator("proxy_agent", ProxyAgentManager.Keyword(), userCred) + portReqV := validators.NewRangeValidator("bind_port_req", cloudproxy_api.BindPortMin, cloudproxy_api.BindPortMax) + for _, v := range []validators.IValidator{ + endpointV, + agentV.Optional(true), + + validators.NewIPv4AddrValidator("remote_addr"), + validators.NewPortValidator("remote_port"), + portReqV, + + validators.NewNonNegativeValidator("last_seen_timeout"), + } { + v.Optional(true) + if err := v.Validate(data); err != nil { + return nil, err + } + } + + portReq := int(portReqV.Value) + if portReq != fwd.BindPortReq { + var err error + var agentId string + if agentV.Model == nil { + agentId = fwd.ProxyAgentId + } else { + agentId = agentV.Model.GetId() + } + switch typ := fwd.Type; typ { + case cloudproxy_api.FORWARD_TYPE_LOCAL: + data, err = ForwardManager.validateLocalSetPort(ctx, data, agentId, portReq) + case cloudproxy_api.FORWARD_TYPE_REMOTE: + data, err = ForwardManager.validateRemoteSetPort(ctx, data, agentId, portReq) + } + if err != nil { + return nil, err + } + } + return data, nil +} + +func (man *SForwardManager) ListItemFilter( + ctx context.Context, + q *sqlchemy.SQuery, + userCred mcclient.TokenCredential, + input cloudproxy_api.ForwardListInput, +) (*sqlchemy.SQuery, error) { + filters := [][2]string{ + [2]string{"type", input.Type}, + [2]string{"proxy_endpoint_id", input.ProxyEndpointId}, + [2]string{"proxy_agent_id", input.ProxyAgentId}, + [2]string{"opaque", input.Opaque}, + } + for _, filter := range filters { + if v := filter[1]; v != "" { + q = q.Equals(filter[0], v) + } + } + return q, nil +} + +func (man *SForwardManager) FetchCustomizeColumns( + ctx context.Context, + userCred mcclient.TokenCredential, + query jsonutils.JSONObject, + objs []interface{}, + fields stringutils2.SSortedStrings, + isList bool, +) []*jsonutils.JSONDict { + fwds := gotypes.ConvertSliceElemType(objs, (**SForward)(nil)).([]*SForward) + + paMap := map[string]*SProxyAgent{} + peMap := map[string]*SProxyEndpoint{} + { + var paIds []string + var peIds []string + { + paIdMap := map[string]string{} + peIdMap := map[string]string{} + for _, fwd := range fwds { + paIdMap[fwd.ProxyAgentId] = "" + peIdMap[fwd.ProxyEndpointId] = "" + } + for id := range paIdMap { + if id != "" { + paIds = append(paIds, id) + } + } + for id := range peIdMap { + if id != "" { + peIds = append(peIds, id) + } + } + } + + var pas []SProxyAgent + var pes []SProxyEndpoint + { + paQ := ProxyAgentManager.Query().In("id", paIds) + if err := db.FetchModelObjects(ProxyAgentManager, paQ, &pas); err != nil { + return nil + } + peQ := ProxyEndpointManager.Query().In("id", peIds) + if err := db.FetchModelObjects(ProxyEndpointManager, peQ, &pes); err != nil { + return nil + } + } + + for i := range pas { + pa := &pas[i] + paMap[pa.Id] = pa + } + for i := range pes { + pe := &pes[i] + peMap[pe.Id] = pe + } + } + + r := make([]*jsonutils.JSONDict, len(objs)) + for i, fwd := range fwds { + d := jsonutils.NewDict() + pa, paOK := paMap[fwd.ProxyAgentId] + pe, peOK := peMap[fwd.ProxyEndpointId] + if paOK || peOK { + if paOK { + d.Set("proxy_agent", jsonutils.NewString(pa.Name)) + } + if peOK { + d.Set("proxy_endpoint", jsonutils.NewString(pe.Name)) + } + switch fwd.Type { + case cloudproxy_api.FORWARD_TYPE_LOCAL: + if paOK { + d.Set("bind_addr", jsonutils.NewString(pa.AdvertiseAddr)) + } + case cloudproxy_api.FORWARD_TYPE_REMOTE: + if peOK { + d.Set("bind_addr", jsonutils.NewString(pe.IntranetIpAddr)) + } + } + r[i] = d + } + } + return r +} + +func (man *SForwardManager) AllowPerformHeartbeat(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return db.IsAllowClassPerform(rbacutils.ScopeProject, userCred, man, "heartbeat") +} + +func (fwd *SForward) PerformHeartbeat(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input *cloudproxy_api.ForwardHeartbeatInput) (jsonutils.JSONObject, error) { + if _, err := db.Update(fwd, func() error { + fwd.LastSeen = time.Now() + return nil + }); err != nil { + return nil, err + } + return nil, nil +} diff --git a/pkg/cloudproxy/models/initdb.go b/pkg/cloudproxy/models/initdb.go new file mode 100644 index 0000000000..173a348fb5 --- /dev/null +++ b/pkg/cloudproxy/models/initdb.go @@ -0,0 +1,19 @@ +// Copyright 2019 Yunion +// +// 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 models + +func InitDB() error { + return nil +} diff --git a/pkg/cloudproxy/models/proxy_agents.go b/pkg/cloudproxy/models/proxy_agents.go new file mode 100644 index 0000000000..cb1ffe7ead --- /dev/null +++ b/pkg/cloudproxy/models/proxy_agents.go @@ -0,0 +1,104 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "context" + + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/validators" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" +) + +// bind_addr, default 0.0.0.0 +// advertise_addr, default default route adddr, maybe k8s cluster ip +type SProxyAgent struct { + db.SStandaloneResourceBase + + BindAddr string `width:"16" charset:"ascii" nullable:"false" list:"user" create:"optional" update:"admin"` + AdvertiseAddr string `width:"16" charset:"ascii" nullable:"false" list:"user" create:"optional" update:"admin"` +} + +type SProxyAgentManager struct { + db.SStandaloneResourceBaseManager +} + +var ProxyAgentManager *SProxyAgentManager + +func init() { + ProxyAgentManager = &SProxyAgentManager{ + SStandaloneResourceBaseManager: db.NewStandaloneResourceBaseManager( + SProxyAgent{}, + "proxy_agents_tbl", + "proxy_agent", + "proxy_agents", + ), + } + ProxyAgentManager.SetVirtualObject(ProxyAgentManager) +} + +func (man *SProxyAgentManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + vs := []validators.IValidator{ + validators.NewIPv4AddrValidator("bind_addr").Optional(true), + validators.NewIPv4AddrValidator("advertise_addr").Optional(true), + } + for _, v := range vs { + if err := v.Validate(data); err != nil { + return nil, err + } + } + return data, nil +} + +func (proxyagent *SProxyAgent) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + vs := []validators.IValidator{ + validators.NewIPv4AddrValidator("bind_addr"), + validators.NewIPv4AddrValidator("advertise_addr"), + } + for _, v := range vs { + v.Optional(true) + if err := v.Validate(data); err != nil { + return nil, err + } + } + return data, nil +} + +func (proxyagent *SProxyAgent) ValidateDeleteCondition(ctx context.Context) error { + q := ForwardManager.Query().Equals("proxy_agent_id", proxyagent.Id) + if count, err := q.CountWithError(); err != nil { + return httperrors.NewServerError("count forwards using proxy endpoint %s(%s)", + proxyagent.Name, proxyagent.Id) + } else if count > 0 { + return httperrors.NewConflictError("proxy endpoint %s(%s) is still used by %d forward(s)", + proxyagent.Name, proxyagent.Id, count) + } else { + return nil + } +} + +func (man *SProxyAgentManager) allAgents(ctx context.Context) ([]SProxyAgent, error) { + var ( + agents []SProxyAgent + q = man.Query() + ) + if err := db.FetchModelObjects(man, q, &agents); err != nil { + return nil, httperrors.NewServerError("query forwards by agent failed: %v", err) + } + return agents, nil +} diff --git a/pkg/cloudproxy/models/proxy_endpoints.go b/pkg/cloudproxy/models/proxy_endpoints.go new file mode 100644 index 0000000000..2ab76710fa --- /dev/null +++ b/pkg/cloudproxy/models/proxy_endpoints.go @@ -0,0 +1,260 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/sqlchemy" + + cloudproxy_api "yunion.io/x/onecloud/pkg/apis/cloudproxy" + compute_apis "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/validators" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/rbacutils" +) + +// Add revision? +type SProxyEndpoint struct { + db.SVirtualResourceBase + + User string `nullable:"false" list:"user" update:"user" create:"optional"` + Host string `nullable:"false" list:"user" update:"user" create:"required"` + Port int `nullable:"false" list:"user" update:"user" create:"optional"` + PrivateKey string `nullable:"false" update:"user" list:"admin" get:"admin" create:"required"` // do not allow get, list + + IntranetIpAddr string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"required"` + + StatusDetail string `width:"128" charset:"ascii" nullable:"false" default:"init" list:"user" create:"optional" json:"status_detail"` +} + +type SProxyEndpointManager struct { + db.SVirtualResourceBaseManager +} + +var ProxyEndpointManager *SProxyEndpointManager + +func init() { + ProxyEndpointManager = &SProxyEndpointManager{ + SVirtualResourceBaseManager: db.NewVirtualResourceBaseManager( + SProxyEndpoint{}, + "proxy_endpoints_tbl", + "proxy_endpoint", + "proxy_endpoints", + ), + } + ProxyEndpointManager.SetVirtualObject(ProxyEndpointManager) +} + +func (man *SProxyEndpointManager) AllowPerformCreateFromServer(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return db.IsAllowClassPerform(rbacutils.ScopeProject, userCred, man, "create-from-server") +} + +func (man *SProxyEndpointManager) PerformCreateFromServer(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input *cloudproxy_api.ProxyEndpointCreateFromServerInput) (jsonutils.JSONObject, error) { + serverId := input.ServerId + if serverId == "" { + return nil, httperrors.NewBadRequestError("server_id is required") + } + + serverInfo, err := getServerInfo(ctx, userCred, serverId) + if err != nil { + return nil, err + } + if serverInfo.PrivateKey == "" { + return nil, httperrors.NewBadRequestError("cannot find ssh private key for this server") + } + + nic := serverInfo.GetNic() + if nic == nil { + return nil, httperrors.NewBadRequestError("cannot find usable network interface for this server") + } + host := serverInfo.Server.Eip + if host == "" && nic.VpcId == compute_apis.DEFAULT_VPC_ID { + host = nic.IpAddr + } + if host == "" { + return nil, httperrors.NewBadRequestError("cannot find ssh host ip address for this server") + } + + proxyendpoint := &SProxyEndpoint{ + User: "cloudroot", + Host: host, + Port: 22, + PrivateKey: serverInfo.PrivateKey, + + IntranetIpAddr: nic.IpAddr, + } + proxyendpoint.Name = serverInfo.Server.Name + proxyendpoint.DomainId = userCred.GetProjectDomainId() + proxyendpoint.ProjectId = userCred.GetProjectId() + if err := man.TableSpec().Insert(ctx, proxyendpoint); err != nil { + return nil, httperrors.NewServerError("database insertion error: %v", err) + } + + var proxymatches []*SProxyMatch + if nic.VpcId != "" { + pm := &SProxyMatch{ + ProxyEndpointId: proxyendpoint.Id, + MatchScope: cloudproxy_api.PM_SCOPE_VPC, + MatchValue: nic.VpcId, + } + pm.Name = "vpc-" + nic.VpcId + proxymatches = append(proxymatches, pm) + } + + if nic.NetworkId != "" { + pm := &SProxyMatch{ + ProxyEndpointId: proxyendpoint.Id, + MatchScope: cloudproxy_api.PM_SCOPE_NETWORK, + MatchValue: nic.NetworkId, + } + pm.Name = "network-" + nic.NetworkId + proxymatches = append(proxymatches, pm) + } + for _, proxymatch := range proxymatches { + proxymatch.DomainId = userCred.GetProjectDomainId() + proxymatch.ProjectId = userCred.GetProjectId() + if err := ProxyMatchManager.TableSpec().Insert(ctx, proxymatch); err != nil { + log.Errorf("failed insertion of proxy match %s: %v", proxymatch.Name, err) + } + } + + return jsonutils.Marshal(proxyendpoint), nil +} + +func (man *SProxyEndpointManager) ValidateCreateData( + ctx context.Context, + userCred mcclient.TokenCredential, + ownerId mcclient.IIdentityProvider, + query jsonutils.JSONObject, + input cloudproxy_api.ProxyEndpointCreateInput, +) (*jsonutils.JSONDict, error) { + data := jsonutils.Marshal(input).(*jsonutils.JSONDict) + + if input, err := man.SVirtualResourceBaseManager.ValidateCreateData(ctx, userCred, ownerId, query, input.VirtualResourceCreateInput); err != nil { + return nil, err + } else { + data.Update(jsonutils.Marshal(input)) + } + + vs := []validators.IValidator{ + validators.NewStringNonEmptyValidator("user").Default("cloudroot"), + validators.NewStringNonEmptyValidator("host"), + validators.NewPortValidator("port").Default(22), + validators.NewSSHKeyValidator("private_key").Optional(true), + + validators.NewIPv4AddrValidator("intranet_ip_addr"), + } + for _, v := range vs { + if err := v.Validate(data); err != nil { + return nil, err + } + } + // populate ssh credential through "cloudhost" + // + // if ! skip validation { + // ssh credential validation + // } + return data, nil +} + +func (man *SProxyEndpointManager) getById(id string) (*SProxyEndpoint, error) { + m, err := db.FetchById(man, id) + if err != nil { + return nil, err + } + proxyendpoint := m.(*SProxyEndpoint) + return proxyendpoint, err +} + +func (proxyendpoint *SProxyEndpoint) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input cloudproxy_api.ProxyEndpointUpdateInput) (cloudproxy_api.ProxyEndpointUpdateInput, error) { + var err error + input.VirtualResourceBaseUpdateInput, err = proxyendpoint.SVirtualResourceBase.ValidateUpdateData(ctx, userCred, query, input.VirtualResourceBaseUpdateInput) + if err != nil { + return input, errors.Wrap(err, "SVirtualResourceBase.ValidateUpdateData") + } + + data := jsonutils.Marshal(input).(*jsonutils.JSONDict) + vs := []validators.IValidator{ + validators.NewStringNonEmptyValidator("user"), + validators.NewStringNonEmptyValidator("host"), + validators.NewPortValidator("port"), + validators.NewSSHKeyValidator("private_key").Optional(true), + } + for _, v := range vs { + v.Optional(true) + if err := v.Validate(data); err != nil { + return input, err + } + } + return input, nil +} + +func (proxyendpoint *SProxyEndpoint) ValidateDeleteCondition(ctx context.Context) error { + q := ForwardManager.Query().Equals("proxy_endpoint_id", proxyendpoint.Id) + if count, err := q.CountWithError(); err != nil { + return httperrors.NewServerError("count forwards using proxy endpoint %s(%s)", + proxyendpoint.Name, proxyendpoint.Id) + } else if count > 0 { + return httperrors.NewConflictError("proxy endpoint %s(%s) is still used by %d forward(s)", + proxyendpoint.Name, proxyendpoint.Id, count) + } else { + return nil + } +} + +func (proxyendpoint *SProxyEndpoint) CustomizeDelete(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) error { + var pms []SProxyMatch + q := ProxyMatchManager.Query().Equals("proxy_endpoint_id", proxyendpoint.Id) + if err := db.FetchModelObjects(ProxyMatchManager, q, &pms); err != nil { + return httperrors.NewServerError("fetch proxy matches for endpoint %s(%s)", + proxyendpoint.Name, proxyendpoint.Id) + } + + for i := range pms { + pm := &pms[i] + err := db.DeleteModel(ctx, userCred, pm) + if err != nil { + return err + } + } + return nil +} + +func (man *SProxyEndpointManager) ListItemFilter( + ctx context.Context, + q *sqlchemy.SQuery, + userCred mcclient.TokenCredential, + input cloudproxy_api.ProxyEndpointListInput, +) (*sqlchemy.SQuery, error) { + filters := [][2]string{ + [2]string{cloudproxy_api.PM_SCOPE_VPC, input.VpcId}, + [2]string{cloudproxy_api.PM_SCOPE_NETWORK, input.NetworkId}, + } + for _, filter := range filters { + if v := filter[1]; v != "" { + pmQ := ProxyMatchManager.Query("proxy_endpoint_id"). + Equals("match_scope", filter[0]). + Equals("match_value", v) + q = q.In("id", pmQ.SubQuery()) + } + } + return q, nil +} diff --git a/pkg/cloudproxy/models/proxy_matches.go b/pkg/cloudproxy/models/proxy_matches.go new file mode 100644 index 0000000000..d43326350e --- /dev/null +++ b/pkg/cloudproxy/models/proxy_matches.go @@ -0,0 +1,113 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/sqlchemy" + + api "yunion.io/x/onecloud/pkg/apis/cloudproxy" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/validators" + "yunion.io/x/onecloud/pkg/mcclient" +) + +type SProxyMatch struct { + db.SVirtualResourceBase + + ProxyEndpointId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"` + MatchScope string `width:"16" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"` + MatchValue string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"` +} + +type SProxyMatchManager struct { + db.SVirtualResourceBaseManager +} + +var ProxyMatchManager *SProxyMatchManager + +func init() { + ProxyMatchManager = &SProxyMatchManager{ + SVirtualResourceBaseManager: db.NewVirtualResourceBaseManager( + SProxyMatch{}, + "proxy_matches_tbl", + "proxy_match", + "proxy_matches", + ), + } + ProxyMatchManager.SetVirtualObject(ProxyMatchManager) +} + +func (man *SProxyMatchManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + matchScopeV := validators.NewStringChoicesValidator("match_scope", api.PM_SCOPES) + endpointV := validators.NewModelIdOrNameValidator("proxy_endpoint", ProxyEndpointManager.Keyword(), ownerId) + for _, v := range []validators.IValidator{ + matchScopeV, + endpointV, + } { + if err := v.Validate(data); err != nil { + return nil, err + } + } + return data, nil +} + +func (pm *SProxyMatch) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + matchScopeV := validators.NewStringChoicesValidator("match_scope", api.PM_SCOPES) + endpointV := validators.NewModelIdOrNameValidator("proxy_endpoint", ProxyEndpointManager.Keyword(), userCred) + for _, v := range []validators.IValidator{ + matchScopeV, + endpointV, + } { + v.Optional(true) + if err := v.Validate(data); err != nil { + return nil, err + } + } + return data, nil +} + +func (man *SProxyMatchManager) findMatch(ctx context.Context, networkId, vpcId string) *SProxyMatch { + q := man.Query() + qfScope := q.Field("match_scope") + qfValue := q.Field("match_value") + q = q.Filter( + sqlchemy.OR( + sqlchemy.AND( + sqlchemy.Equals(qfScope, api.PM_SCOPE_VPC), + sqlchemy.Equals(qfValue, vpcId), + ), + sqlchemy.AND( + sqlchemy.Equals(qfScope, api.PM_SCOPE_NETWORK), + sqlchemy.Equals(qfValue, networkId), + ), + ), + ) + + var pms []SProxyMatch + if err := db.FetchModelObjects(man, q, &pms); err != nil { + return nil + } + var r *SProxyMatch + for _, pm := range pms { + if pm.MatchScope == api.PM_SCOPE_NETWORK { + return &pm + } + r = &pm + } + return r +} diff --git a/pkg/cloudproxy/models/server_util.go b/pkg/cloudproxy/models/server_util.go new file mode 100644 index 0000000000..7de28b73a0 --- /dev/null +++ b/pkg/cloudproxy/models/server_util.go @@ -0,0 +1,82 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "context" + + compute_apis "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + compute_models "yunion.io/x/onecloud/pkg/mcclient/models" + compute_modules "yunion.io/x/onecloud/pkg/mcclient/modules" +) + +type serverInfo struct { + Server *compute_models.Server + + // PrivateKey is the one corresponds to userCred when getting this + // serverInfo instance. It can be empty + PrivateKey string +} + +func (si *serverInfo) GetNic() *compute_models.ServerNic { + nic := si.getVPCNic() + if nic != nil { + return nic + } + + for _, nic := range si.Server.Nics { + if nic.IpAddr != "" && nic.VpcId == compute_apis.DEFAULT_VPC_ID { + return &nic + } + } + return nil +} + +func (si *serverInfo) getVPCNic() *compute_models.ServerNic { + for _, nic := range si.Server.Nics { + if nic.IpAddr != "" && nic.VpcId != compute_apis.DEFAULT_VPC_ID { + return &nic + } + } + return nil +} + +func getServerInfo( + ctx context.Context, + userCred mcclient.TokenCredential, + serverId string, +) (*serverInfo, error) { + sess := auth.GetSession(ctx, userCred, "", "") + serverJson, err := compute_modules.Servers.Get(sess, serverId, nil) + if err != nil { + return nil, httperrors.NewGeneralError(err) + } + + server := &compute_models.Server{} + if err := serverJson.Unmarshal(server); err != nil { + return nil, httperrors.NewServerError("unmarshal server %s: %v", serverId, err) + } + + privateKey, _ := compute_modules.Sshkeypairs.FetchPrivateKey(ctx, userCred) + + serverInfo := &serverInfo{ + Server: server, + PrivateKey: privateKey, + } + return serverInfo, nil +} diff --git a/pkg/cloudproxy/models/util_errors.go b/pkg/cloudproxy/models/util_errors.go new file mode 100644 index 0000000000..1dd6b78f01 --- /dev/null +++ b/pkg/cloudproxy/models/util_errors.go @@ -0,0 +1,30 @@ +// Copyright 2019 Yunion +// +// 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 models + +type ( + errNotFound error + errMoreThanOne error +) + +func IsNotFound(err error) bool { + _, ok := err.(errNotFound) + return ok +} + +func IsMoreThanOne(err error) bool { + _, ok := err.(errMoreThanOne) + return ok +} diff --git a/pkg/cloudproxy/models/util_subnets.go b/pkg/cloudproxy/models/util_subnets.go new file mode 100644 index 0000000000..9216e08233 --- /dev/null +++ b/pkg/cloudproxy/models/util_subnets.go @@ -0,0 +1,52 @@ +// Copyright 2019 Yunion +// +// 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 models + +import ( + "strings" + + "yunion.io/x/pkg/util/netutils" +) + +type Subnets []*netutils.IPV4Prefix + +func (nets Subnets) StrList() []string { + r := make([]string, 0, len(nets)) + for _, p := range nets { + r = append(r, p.String()) + } + return r +} + +func (nets Subnets) String() string { + r := nets.StrList() + return strings.Join(r, ",") +} + +func (nets Subnets) ContainsAny(nets1 Subnets) bool { + contains, _ := nets.ContainsAnyEx(nets1) + return contains +} + +func (nets Subnets) ContainsAnyEx(nets1 Subnets) (bool, *netutils.IPV4Prefix) { + for _, p0 := range nets { + for _, p1 := range nets1 { + if p0.Equals(p1) { + return true, p0 + } + } + } + return false, nil +} diff --git a/pkg/cloudproxy/options/doc.go b/pkg/cloudproxy/options/doc.go new file mode 100644 index 0000000000..0dcd449872 --- /dev/null +++ b/pkg/cloudproxy/options/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 options // import "yunion.io/x/onecloud/pkg/cloudproxy/options" diff --git a/pkg/cloudproxy/options/options.go b/pkg/cloudproxy/options/options.go new file mode 100644 index 0000000000..cb88a71965 --- /dev/null +++ b/pkg/cloudproxy/options/options.go @@ -0,0 +1,50 @@ +// Copyright 2019 Yunion +// +// 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 options + +import ( + "os" + "sync" + + common_options "yunion.io/x/onecloud/pkg/cloudcommon/options" + agent_options "yunion.io/x/onecloud/pkg/cloudproxy/agent/options" +) + +type Options struct { + EnableAPIServer bool + EnableProxyAgent bool + + common_options.CommonOptions + common_options.DBOptions + + agent_options.Options +} + +var ( + opts Options + optsOnce sync.Once +) + +func Get() *Options { + optsOnce.Do(func() { + common_options.ParseOptions( + &opts, + os.Args, + "cloudproxy.conf", + "cloudproxy", + ) + }) + return &opts +} diff --git a/pkg/cloudproxy/service/doc.go b/pkg/cloudproxy/service/doc.go new file mode 100644 index 0000000000..0a6cb5dc7c --- /dev/null +++ b/pkg/cloudproxy/service/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// 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 service // import "yunion.io/x/onecloud/pkg/cloudproxy/service" diff --git a/pkg/cloudproxy/service/handlers.go b/pkg/cloudproxy/service/handlers.go new file mode 100644 index 0000000000..e1ac99a3a9 --- /dev/null +++ b/pkg/cloudproxy/service/handlers.go @@ -0,0 +1,40 @@ +// Copyright 2019 Yunion +// +// 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 service + +import ( + "yunion.io/x/onecloud/pkg/appsrv" + "yunion.io/x/onecloud/pkg/appsrv/dispatcher" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudproxy/models" +) + +func InitHandlers(app *appsrv.Application) { + db.InitAllManagers() + + db.RegisterModelManager(db.OpsLog) + db.RegisterModelManager(db.TenantCacheManager) + db.RegisterModelManager(db.UserCacheManager) + for _, manager := range []db.IModelManager{ + models.ProxyEndpointManager, + models.ProxyMatchManager, + models.ProxyAgentManager, + models.ForwardManager, + } { + db.RegisterModelManager(manager) + handler := db.NewModelHandler(manager) + dispatcher.AddModelDispatcher("", app, handler) + } +} diff --git a/pkg/cloudproxy/service/service.go b/pkg/cloudproxy/service/service.go new file mode 100644 index 0000000000..5d4cfc6eb6 --- /dev/null +++ b/pkg/cloudproxy/service/service.go @@ -0,0 +1,41 @@ +// Copyright 2019 Yunion +// +// 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 service + +import ( + _ "github.com/go-sql-driver/mysql" + + "yunion.io/x/onecloud/pkg/cloudcommon" + common_app "yunion.io/x/onecloud/pkg/cloudcommon/app" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudproxy/models" + "yunion.io/x/onecloud/pkg/cloudproxy/options" +) + +func StartService() { + var ( + opts = options.Get() + dbOpts = &opts.DBOptions + baseOpts = &opts.BaseOptions + ) + + app := common_app.InitApp(baseOpts, false) + InitHandlers(app) + + db.EnsureAppInitSyncDB(app, dbOpts, models.InitDB) + defer cloudcommon.CloseDB() + + common_app.ServeForever(app, baseOpts) +}