Compare commits

..

109 Commits

Author SHA1 Message Date
Zexi Li
af3ee4af02 Merge pull request #9114 from yousong/automated-cherry-pick-of-#9109-upstream-release-3.2
Automated cherry pick of #9109: fix(host): make sure it won't match vpc guests by accident
2020-12-01 20:22:04 +08:00
Yousong Zhou
12a100f81b fix(host): make sure it won't match vpc guests by accident 2020-12-01 15:34:14 +08:00
yunion-ci-robot
e3509fcae8 Merge pull request #8529 from rainzm/automated-cherry-pick-of-#8525-upstream-release-3.2
Automated cherry pick of #8525: fix(esxi): no callback in uploadHandler
2020-10-28 00:58:40 +08:00
yunion-ci-robot
3cf538db43 Merge pull request #8516 from ioito/automated-cherry-pick-of-#8512-upstream-release-3.2
Automated cherry pick of #8512: fix: avoid sync project panic
2020-10-28 00:46:43 +08:00
rainzm
262d4d650b fix(esxi): filter vmdk item in lease.info 2020-10-27 21:34:58 +08:00
rainzm
fabcc221a9 fix(esxi): use correct device index
1. Sort devices via their Key.
2. There is no relationship between vdisk index and vnics length.
2020-10-27 21:34:58 +08:00
rainzm
a4c781b4d5 fix(esxi): no callback in uploadHandler 2020-10-27 21:34:58 +08:00
Qu Xuan
6256b3f88d fix: avoid sync project panic 2020-10-27 16:14:39 +08:00
yunion-ci-robot
4bd45eef3e Merge pull request #8503 from yousong/automated-cherry-pick-of-#8499-upstream-release-3.2
Automated cherry pick of #8499: lockman: test: add prefix setting
2020-10-26 21:24:42 +08:00
Yousong Zhou
21a0cad5e9 lockman: test: add prefix setting 2020-10-26 20:37:13 +08:00
yunion-ci-robot
311f4cd079 Merge pull request #8491 from ioito/automated-cherry-pick-of-#8487-upstream-release-3.2
Automated cherry pick of #8487: fix: avoid delete public ip when sync eip list
2020-10-26 17:02:41 +08:00
Qu Xuan
de817bd4a7 fix: avoid delete public ip when sync eip list 2020-10-26 16:19:49 +08:00
Zexi Li
fb5276aee4 Merge pull request #8376 from ioito/automated-cherry-pick-of-#8373-upstream-release-3.2
Automated cherry pick of #8373: fix: 新增google存储类型
2020-10-19 15:11:25 +08:00
Zexi Li
f8d689ba87 Merge pull request #8372 from ioito/automated-cherry-pick-of-#8369-upstream-release-3.2
Automated cherry pick of #8369: fix: use correct worker
2020-10-19 15:10:25 +08:00
Qu Xuan
698d13f3ef fix: 新增google存储类型 2020-10-19 14:23:52 +08:00
Qu Xuan
8943a16570 fix: use correct worker 2020-10-19 11:43:15 +08:00
yunion-ci-robot
82a3e65bc4 Merge pull request #8317 from rainzm/automated-cherry-pick-of-#8314-upstream-release-3.2
Automated cherry pick of #8314: feat(esxi): do not specify portkey in NewVNIC
2020-10-16 10:40:35 +08:00
rainzm
b272de37ec feat(esxi): do not specify portkey in NewVNIC 2020-10-16 08:25:09 +08:00
yunion-ci-robot
0447b4ade0 Merge pull request #8311 from ioito/automated-cherry-pick-of-#8308-upstream-release-3.2
Automated cherry pick of #8308: fix: 避免并发同步安全组导致规则混乱
2020-10-15 23:01:35 +08:00
Qu Xuan
64973bf486 fix: 避免并发同步安全组导致规则混乱 2020-10-15 22:27:49 +08:00
yunion-ci-robot
f335e270cf Merge pull request #8245 from zexi/automated-cherry-pick-of-#8242-upstream-release-3.2
Automated cherry pick of #8242: baremetal: fix pxe filter network by guest_dhcp
2020-10-13 10:30:29 +08:00
Zexi Li
5ba53e43b3 Merge pull request #8261 from rainzm/automated-cherry-pick-of-#8232-upstream-release-3.2
Automated cherry pick of #8232: fix(esxi): search template from whole vCenter
2020-10-13 10:06:49 +08:00
rainzm
d80f15af58 fix(esxi): use system disk of template only when cloning vm
1. Create a driver if necessary.
2. After cloning the virtual machine, remove all data disks.
3. Create new data disks.
2020-10-12 17:50:30 +08:00
rainzm
1fe71b8392 fix(esxi): search template from whole vCenter 2020-10-12 17:50:30 +08:00
Zexi Li
1e2e134ad0 baremetal: fix pxe filter network by guest_dhcp 2020-10-12 11:38:05 +08:00
yunion-ci-robot
d456364746 Merge pull request #8197 from rainzm/automated-cherry-pick-of-#8194-upstream-release-3.2
Automated cherry pick of #8194: fix(esxi): check nil value when NewVirtualMachine
2020-10-10 19:06:18 +08:00
yunion-ci-robot
1d82b2fe90 Merge pull request #8215 from tb365/automated-cherry-pick-of-#8214-upstream-release-3.2
Automated cherry pick of #8214: remove aws lbbg cache metedata
2020-10-10 18:55:20 +08:00
yunion-ci-robot
eb3d66750b Merge pull request #8221 from yousong/automated-cherry-pick-of-#8218-upstream-release-3.2
Automated cherry pick of #8218: webconsole: ssh: each argument on its own line
2020-10-10 15:34:49 +08:00
Yousong Zhou
88b752b35c webconsole: ssh: add keyboard-interactive as an option
For ESXi 6.0

  Authentications that can continue: publickey,keyboard-interactive
2020-10-10 14:44:26 +08:00
Yousong Zhou
225067eb88 webconsole: ssh: each argument on its own line 2020-10-10 14:44:26 +08:00
TangBin
d55429d4cd remove aws lbbg cache metedata 2020-10-10 14:27:25 +08:00
yunion-ci-robot
c3ec288464 Merge pull request #8186 from yousong/automated-cherry-pick-of-#8183-upstream-release-3.2
Automated cherry pick of #8183: vpcagent: ovn: increase dhcp lease, renew, rebind time
2020-10-10 11:56:12 +08:00
rainzm
af03b18155 fix(esxi): check nil value when NewVirtualMachine 2020-10-10 11:36:44 +08:00
Yousong Zhou
d40d6a6a74 vpcagent: ovn: increase dhcp lease, renew, rebind time 2020-10-09 10:26:59 +08:00
yunion-ci-robot
dd8629e9aa Merge pull request #8160 from tb365/automated-cherry-pick-of-#8159-upstream-release-3.2
Automated cherry pick of #8159: aliyun public elb region fix
2020-09-29 23:28:41 +08:00
TangBin
6d2f1d25c1 aliyun public elb region fix 2020-09-29 22:04:50 +08:00
yunion-ci-robot
1e95a83db0 Merge pull request #8155 from tb365/automated-cherry-pick-of-#8154-upstream-release-3.2
Automated cherry pick of #8154: aliyun public elb brand fix
2020-09-29 18:18:34 +08:00
TangBin
1e7776ba71 aliyun public elb brand fix 2020-09-29 16:47:25 +08:00
yunion-ci-robot
d4900cc6ac Merge pull request #8062 from ioito/automated-cherry-pick-of-#8059-upstream-release-3.2
Automated cherry pick of #8059: fix: avoid serial vm creation
2020-09-29 10:33:34 +08:00
Zexi Li
732120c796 Merge pull request #8130 from tb365/automated-cherry-pick-of-#8129-upstream-release-3.2
Automated cherry pick of #8129: aliyun create public lb fix
2020-09-28 16:24:23 +08:00
TangBin
70d22e6428 aliyun create public lb fix 2020-09-28 14:19:48 +08:00
yunion-ci-robot
62dfa77b87 Merge pull request #8119 from wanyaoqi/automated-cherry-pick-of-#8116-upstream-release-3.2
Automated cherry pick of #8116: umount nfs on detach storage
2020-09-27 21:03:36 +08:00
wanyaoqi
788e1bb225 increase rbd check timeout 2020-09-27 20:22:10 +08:00
wanyaoqi
02f7c5ac35 umount nfs on detach storage 2020-09-27 20:22:10 +08:00
Zexi Li
a675971d90 Merge pull request #8113 from tb365/automated-cherry-pick-of-#8087-upstream-release-3.2
Automated cherry pick of #8087: huawei snapshot delete fix
2020-09-27 17:29:49 +08:00
TangBin
38be1b0908 huawei util missing import fix 2020-09-27 15:14:37 +08:00
TangBin
e3242abf15 huawei client request add throttling control 2020-09-27 15:13:16 +08:00
TangBin
0171bf71d3 huawei elb get params log 2020-09-27 15:13:16 +08:00
TangBin
70170836fb qcloud elb create listener fix 2020-09-27 15:13:16 +08:00
TangBin
ec7707d397 huawei snapshot delete fix 2020-09-27 15:13:16 +08:00
Zexi Li
6e61174307 Merge pull request #8091 from wanyaoqi/automated-cherry-pick-of-#8088-upstream-release-3.2
Automated cherry pick of #8088: set stagefail on clean fail
2020-09-27 14:46:10 +08:00
wanyaoqi
0d96419b64 set stagefail on clean fail 2020-09-27 11:47:07 +08:00
Zexi Li
a7644951fe Merge pull request #8098 from ioito/automated-cherry-pick-of-#8095-upstream-release-3.2
Automated cherry pick of #8095: fix: avoid purge storages
2020-09-25 20:54:51 +08:00
Qu Xuan
96b6b3cf6f fix: avoid purge storages 2020-09-25 19:18:04 +08:00
Zexi Li
3cc7557964 Merge pull request #8074 from ioito/automated-cherry-pick-of-#8071-upstream-release-3.2
Automated cherry pick of #8071: fix: sql err
2020-09-24 21:19:38 +08:00
ioito
c433050318 fix: sql err 2020-09-24 20:23:33 +08:00
ioito
1505dc6e04 fix: avoid serial vm creation 2020-09-24 16:41:50 +08:00
Zexi Li
f14fbb4007 Merge pull request #8004 from tb365/automated-cherry-pick-of-#8003-upstream-release-3.2
Automated cherry pick of #8003: regional network filter bugfix
2020-09-23 21:09:14 +08:00
TangBin
396a99c4cb regional network filter bugfix 2020-09-22 20:02:33 +08:00
yunion-ci-robot
1356ae3124 Merge pull request #7995 from tb365/automated-cherry-pick-of-#7994-upstream-release-3.2
Automated cherry pick of #7994: ctyun sync secgroup fix & cloudaccount delete fix
2020-09-22 14:08:19 +08:00
TangBin
fc480ffa1a ctyun sync secgroup fix & cloudaccount delete fix 2020-09-22 13:53:20 +08:00
Zexi Li
4e7f02fa68 Merge pull request #7967 from tb365/automated-cherry-pick-of-#7966-upstream-release-3.2
Automated cherry pick of #7966: ctyun add region cn-bj1
2020-09-18 21:04:54 +08:00
TangBin
614ad609bf ctyun add region cn-bj1 2020-09-18 14:22:51 +08:00
Zexi Li
94a2645341 Merge pull request #7964 from zexi/automated-cherry-pick-of-#7961-upstream-release-3.2
Automated cherry pick of #7961: region: fix network schedtag not cleanup when network deleted
2020-09-18 12:45:07 +08:00
Zexi Li
3ad40222c2 region: fix network schedtag not cleanup when network deleted 2020-09-18 10:39:43 +08:00
Zexi Li
63e0f2ee76 Merge pull request #7958 from wanyaoqi/automated-cherry-pick-of-#7953-upstream-release-3.2
Automated cherry pick of #7953: fix set host image cache properties
2020-09-16 21:38:29 +08:00
Zexi Li
e14304a5af Merge pull request #7951 from wanyaoqi/automated-cherry-pick-of-#7948-upstream-release-3.2
Automated cherry pick of #7948: fix update image status after save failed
2020-09-16 21:35:47 +08:00
wanyaoqi
85c61ab14c fix set host image cache properties 2020-09-16 20:53:03 +08:00
wanyaoqi
f393424ca8 fix update image status after save failed 2020-09-16 19:15:05 +08:00
Zexi Li
5aebb04ce3 Merge pull request #7924 from ioito/automated-cherry-pick-of-#7921-upstream-release-3.2
Automated cherry pick of #7921: fix: bucket object cnt fix
2020-09-16 11:30:13 +08:00
Qu Xuan
0d99dc1981 fix: bucket object cnt fix 2020-09-15 19:15:32 +08:00
yunion-ci-robot
4334ccf537 Merge pull request #7908 from wanyaoqi/automated-cherry-pick-of-#7905-upstream-release-3.2
Automated cherry pick of #7905: webconsole: fix defer order to release zombie process
2020-09-15 10:26:10 +08:00
wanyaoqi
285c696e6d webconsole: fix defer order to release zombie process 2020-09-14 18:22:13 +08:00
Zexi Li
ce6e7542a5 Merge pull request #7898 from wanyaoqi/automated-cherry-pick-of-#7895-upstream-release-3.2
Automated cherry pick of #7895: disk: check disk is need renew on guest set renew
2020-09-14 13:02:17 +08:00
wanyaoqi
5ec462bc8e disk: check disk is need renew on guest set renew 2020-09-14 12:08:52 +08:00
Zexi Li
4868674fee Merge pull request #7882 from ioito/automated-cherry-pick-of-#7879-upstream-release-3.2
Automated cherry pick of #7879: fix: avoid sync natgateway eip panic
2020-09-11 21:03:57 +08:00
Zexi Li
2e560c1f04 Merge pull request #7873 from rainzm/automated-cherry-pick-of-#7870-upstream-release-3.2
Automated cherry pick of #7870: fix(esxiagent): add HostDelayTaskWorkerCount
2020-09-11 20:55:47 +08:00
Zexi Li
71b1244846 Merge pull request #7865 from wanyaoqi/automated-cherry-pick-of-#7836-upstream-release-3.2
Automated cherry pick of #7836: host: fix rbd storage cache iso image
2020-09-11 20:54:07 +08:00
Qu Xuan
545f9eebf1 fix: avoid sync natgateway eip panic 2020-09-11 19:18:15 +08:00
rainzm
e6ce176233 fix(esxiagent): add HostDelayTaskWorkerCount
之前,HostDelayWorker 是通过 hostutils.InitWorkerManager 来初始化,
worker的数量依赖于options.HostOptions.DefaultRequestWorkerCount, 因为
这个options没有经过初始化,所以就是0,导致worker count的数量变为1。

现在增加了 HostDelayTaskWorkerCount 来管理这个count,默认值为8。
2020-09-11 17:23:13 +08:00
wanyaoqi
f59d56e796 host: fix rbd storage cache iso image 2020-09-10 23:20:29 +08:00
yunion-ci-robot
37e36f5a08 Merge pull request #7848 from rainzm/automated-cherry-pick-of-#7845-upstream-release-3.2
Automated cherry pick of #7845: Feat & Fix & Refactor for ESXi
2020-09-10 09:41:57 +08:00
rainzm
05abebeea5 fix(region): add RequestSyncstatusOnHost for SESXiGuestDriver
esxi 平台 Guest 进行 syncstatus 和其他平台有所区别,Vcenter 可能会更改
 VM 所在的 Host,所以在 Host 上寻找 VM 可能会出现 ErrNotFound。此时,应该
主动去整个 datacenter 中寻找 VM,然后更改 Host,进行同步。
2020-09-09 21:58:48 +08:00
rainzm
747d70a747 fix(esxi): fetch templatevm from datacenter before clone vm
这里应该和image_cache那里保持一直,对template的获取应该从datacenter
中获取,而不是host。
2020-09-09 21:58:17 +08:00
rainzm
ebf3692ba0 refactor(esxi): replace GetTemplateVMById with FetchTemplateVMById
FetchTemplateVMById 更好,因为它更轻量,更快。
2020-09-09 21:58:16 +08:00
rainzm
7dfe68322a feat(esxicli): better vm operator
现在,可以通过Datacenter来获取vm,不一定要指定HostIp,
Datacenter和HostIp必须指定一个。
2020-09-09 21:58:16 +08:00
rainzm
59861adabc feat(esxi): support fetchVM form datacenter 2020-09-09 21:58:16 +08:00
rainzm
40b3e8e382 refactor(esxi): fetchVms and fetchHardwareInfo
1. fetchHardwareInfo 只有一种error,原因是moVM的某些字段为nil,这种情况下,完全可以打印日志直接返回。
2. fetchVms 现在只返回[]*SVirtualMachine, 进一步的过滤(是不是template)交给调用者。
2020-09-09 21:58:16 +08:00
Zexi Li
c81563b03d Merge pull request #7825 from rainzm/automated-cherry-pick-of-#7821-upstream-release-3.2
Automated cherry pick of #7821: fix(esxiagent): return image extid not id
2020-09-08 20:25:59 +08:00
rainzm
a18b34c365 fix(esxiagent): return image extid not id
调用 disk/image_cache 接口的Task StorageCacheImageTask 会根据回调
回来数据中的 image_id,来设置 storagecachedimage 中的 externalid,
所以这里的 image_id 应该是 externalid。
2020-09-08 15:34:35 +08:00
yunion-ci-robot
f64267897e Merge pull request #7801 from wanyaoqi/automated-cherry-pick-of-#7798-upstream-release-3.2
Automated cherry pick of #7798: guest short desc add backup host id
2020-09-07 21:01:51 +08:00
yunion-ci-robot
ad8d6e76b5 Merge pull request #7792 from wanyaoqi/automated-cherry-pick-of-#7789-upstream-release-3.2
Automated cherry pick of #7789: fix get ubuntu version
2020-09-07 20:59:44 +08:00
wanyaoqi
c266722c31 guest short desc add backup host id 2020-09-07 17:06:36 +08:00
wanyaoqi
4d0cb2b150 fix get ubuntu version 2020-09-07 16:01:59 +08:00
yunion-ci-robot
2171de57ee Merge pull request #7785 from rainzm/automated-cherry-pick-of-#7782-upstream-release-3.2
Automated cherry pick of #7782: fix(notify): Replace String() with GetString()
2020-09-07 15:26:42 +08:00
rainzm
021b4a1ae2 fix(notify): Replace String() with GetString()
The String() of jsonutils.JSONObject will add a pair of
double quotes around the content. The GetString() will
get a clean content.
2020-09-07 15:12:07 +08:00
yunion-ci-robot
b108fe81d2 Merge pull request #7775 from swordqiu/automated-cherry-pick-of-#7772-upstream-release-3.2
Automated cherry pick of #7772: fix: skup syncing lb, rds, redis instances in tasks
2020-09-05 15:47:42 +08:00
Qiu Jian
9f854385e5 fix: skup syncing lb, rds, redis instances in tasks 2020-09-05 14:31:06 +08:00
Zexi Li
b319be6d0a Merge pull request #7747 from wanyaoqi/automated-cherry-pick-of-#7744-upstream-release-3.2
Automated cherry pick of #7744: glance: check min disk size on update
2020-09-04 10:30:01 +08:00
wanyaoqi
fa28094906 glance: check min disk size on update 2020-09-03 20:02:17 +08:00
yunion-ci-robot
f797ec2fb9 Merge pull request #7726 from tb365/automated-cherry-pick-of-#7725-upstream-release-3.2
Automated cherry pick of #7725: create classic vpc fix
2020-09-03 10:32:41 +08:00
TangBin
5ca7179ce1 create classic vpc fix 2020-09-02 18:24:56 +08:00
wanyaoqi
c2b81cf638 hostinfo: try create network add is_on_premise (#7709) 2020-09-01 20:01:41 +08:00
wanyaoqi
ddb7981922 fix host get guests (#7685) 2020-09-01 01:08:36 +08:00
wanyaoqi
8020d70bab fix get storage capacity on init' (#7689) 2020-09-01 01:02:37 +08:00
Zexi Li
d78c549b4a Merge pull request #7625 from wanyaoqi/automated-cherry-pick-of-#7623-upstream-release-3.2
Automated cherry pick of #7623: fix gpfs check mountpoint
2020-08-25 17:41:23 +08:00
wanyaoqi
849a74c0a6 fix gpfs check mountpoint 2020-08-25 15:52:27 +08:00
屈轩
e016109d6e fix: avoid storage=nil, dis can not delete (#7610)
Co-authored-by: Qu Xuan <quxuan@yunionyun.com>
2020-08-22 12:33:36 +08:00
Jian Qiu
c67aa4991d fix: opslog filter by owner_project_ids and owner_domain_ids (#7602)
Co-authored-by: Qiu Jian <qiujian@yunionyun.com>
2020-08-21 10:57:23 +08:00
85 changed files with 1046 additions and 453 deletions

View File

@@ -35,7 +35,10 @@ type BaseEventListOptions struct {
Action []string `help:"Log action"`
User string `help:"filter by operator user"`
Project string `help:"filter by owner project"`
Project string `help:"filter by operator user's project"`
OwnerProjectIds []string `help:"filter by owner project ids"`
OwnerDomainIds []string `help:"filter by owner domain ids"`
PagingMarker string `help:"marker for pagination"`
}
@@ -107,6 +110,12 @@ func doEventList(man modulebase.ResourceManager, s *mcclient.ClientSession, args
if len(args.Scope) > 0 {
params.Add(jsonutils.NewString(args.Scope), "scope")
}
if len(args.OwnerProjectIds) > 0 {
params.Add(jsonutils.NewStringArray(args.OwnerProjectIds), "owner_project_ids")
}
if len(args.OwnerDomainIds) > 0 {
params.Add(jsonutils.NewStringArray(args.OwnerDomainIds), "owner_domain_ids")
}
if len(args.PagingMarker) > 0 {
params.Add(jsonutils.NewString(args.PagingMarker), "paging_marker")
}

View File

@@ -201,7 +201,10 @@ type LoadbalancerResourceInfo struct {
// 可用区ID
ZoneId string `json:"zone_id"`
ZoneResourceInfoBase
ZoneResourceInfo
// cloud provider info
ManagedResourceInfo
}
type LoadbalancerResourceInput struct {

View File

@@ -81,6 +81,7 @@ const (
STORAGE_GOOGLE_LOCAL_SSD = "local-ssd" //本地SSD暂存盘 (最多8个)
STORAGE_GOOGLE_PD_STANDARD = "pd-standard" //标准永久性磁盘
STORAGE_GOOGLE_PD_SSD = "pd-ssd" //SSD永久性磁盘
STORAGE_GOOGLE_PD_BALANCED = "pd-balanced" //平衡永久性磁盘
// ctyun storage type
STORAGE_CTYUN_SSD = "SSD" // 超高IO云硬盘

View File

@@ -88,3 +88,12 @@ type ImageCreateInput struct {
// 镜像属性
Properties map[string]string `json:"properties"`
}
type ImageUpdateStatusInput struct {
apis.Meta
// 镜像状态
Status string `json:"status"`
// 更新镜像状态原因
Reason string `json:"reason"`
}

View File

@@ -269,6 +269,9 @@ func (req *dhcpRequest) findNetworkConf(session *mcclient.ClientSession, filterU
params.Add(jsonutils.NewString(
fmt.Sprintf("guest_dhcp.endswith(',%s')", req.RelayAddr)),
"filter.3")
params.Add(jsonutils.NewString(
fmt.Sprintf("guest_dhcp.equals(%s)", req.RelayAddr)),
"filter.4")
params.Add(jsonutils.JSONTrue, "filter_any")
}
params.Add(jsonutils.JSONTrue, "is_on_premise")

View File

@@ -26,7 +26,8 @@ func TestEctdLockManager(t *testing.T) {
shared := newSharedObject()
for i := 0; i < 4; i++ {
lockman, err := NewEtcdLockManager(&SEtcdLockManagerConfig{
Endpoints: []string{"localhost:2379"},
LockPrefix: "test-etcd-lock-manager",
Endpoints: []string{"localhost:2379"},
})
if err != nil {
t.Skipf("new etcd lockman: %v", err)

View File

@@ -16,6 +16,7 @@ package db
import (
"context"
"database/sql"
"fmt"
"strconv"
"strings"
@@ -460,28 +461,7 @@ func (manager *SOpsLogManager) ListItemFilter(
userCred mcclient.TokenCredential,
query jsonutils.JSONObject,
) (*sqlchemy.SQuery, error) {
/*userStrs := jsonutils.GetQueryStringArray(query, "user")
if len(userStrs) > 0 {
for i := range userStrs {
usrObj, err := DefaultUserFetcher(ctx, userStrs[i])
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError2("user", userStrs[i])
} else if err == sqlchemy.ErrDuplicateEntry {
return nil, httperrors.NewDuplicateNameError("user", userStrs[i])
} else {
return nil, httperrors.NewGeneralError(err)
}
}
userStrs[i] = usrObj.GetId()
}
if len(userStrs) == 1 {
q = q.Filter(sqlchemy.Equals(q.Field("user_id"), userStrs[0]))
} else {
q = q.Filter(sqlchemy.In(q.Field("user_id"), userStrs))
}
}
projStrs := jsonutils.GetQueryStringArray(query, "project")
projStrs := jsonutils.GetQueryStringArray(query, "owner_project_ids")
if len(projStrs) > 0 {
for i := range projStrs {
projObj, err := DefaultProjectFetcher(ctx, projStrs[i])
@@ -494,12 +474,23 @@ func (manager *SOpsLogManager) ListItemFilter(
}
projStrs[i] = projObj.GetId()
}
if len(projStrs) == 1 {
q = q.Filter(sqlchemy.Equals(q.Field("owner_tenant_id"), projStrs[0]))
} else {
q = q.Filter(sqlchemy.In(q.Field("owner_tenant_id"), projStrs))
q = q.Filter(sqlchemy.In(q.Field("owner_tenant_id"), projStrs))
}
domainStrs := jsonutils.GetQueryStringArray(query, "owner_domain_ids")
if len(domainStrs) > 0 {
for i := range domainStrs {
domainObj, err := DefaultDomainFetcher(ctx, domainStrs[i])
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError2("domain", domainStrs[i])
} else {
return nil, httperrors.NewGeneralError(err)
}
}
domainStrs[i] = domainObj.GetId()
}
}*/
q = q.Filter(sqlchemy.In(q.Field("owner_domain_id"), domainStrs))
}
objTypes := jsonutils.GetQueryStringArray(query, "obj_type")
if len(objTypes) > 0 {
if len(objTypes) == 1 {

View File

@@ -126,7 +126,7 @@ func RawNotify(recipientId []string, isGroup bool, channel notify.TNotifyChannel
msg.Topic = topic
body, _ := getContent(event, "content", channel, data)
if len(body) == 0 {
body = data.String()
body, _ = data.GetString()
}
msg.Msg = body
// log.Debugf("send notification %s %s", topic, body)

View File

@@ -368,3 +368,42 @@ func (self *SESXiGuestDriver) RequestAssociateEip(ctx context.Context, userCred
func (self *SESXiGuestDriver) IsSupportCdrom(guest *models.SGuest) (bool, error) {
return false, nil
}
func (self *SESXiGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) {
ihost, err := host.GetIHost()
if err != nil {
return nil, err
}
ivm, err := ihost.GetIVMById(guest.GetExternalId())
if err != nil && errors.Cause(err) != errors.ErrNotFound {
return nil, err
}
// VM may be migrated by Vcenter, try to find VM from whole datacenter.
if err != nil {
ehost := ihost.(*esxi.SHost)
dc, err := ehost.GetDatacenter()
if err != nil {
return nil, errors.Wrapf(err, "ehost.GetDatacenter")
}
vm, err := dc.FetchVMById(guest.GetExternalId())
if err != nil {
log.Errorf("fail to find ivm by id %q in dc %q: %v", guest.GetExternalId(), dc.GetName(), err)
return nil, err
}
ihost = vm.GetIHost()
host = models.HostManager.FetchHostByExtId(ihost.GetGlobalId())
if host == nil {
return nil, errors.Wrapf(errors.ErrNotFound, "find ivm %q in ihost %q which is not existed here", guest.GetExternalId(), ihost.GetGlobalId())
}
ivm = vm
}
err = guest.SyncAllWithCloudVM(ctx, userCred, host, ivm)
if err != nil {
return nil, err
}
status := GetCloudVMStatus(ivm)
body := jsonutils.NewDict()
body.Add(jsonutils.NewString(status), "status")
return body, nil
}

View File

@@ -77,6 +77,7 @@ func (self *SGoogleGuestDriver) GetStorageTypes() []string {
return []string{
api.STORAGE_GOOGLE_PD_SSD,
api.STORAGE_GOOGLE_PD_STANDARD,
api.STORAGE_GOOGLE_PD_BALANCED,
api.STORAGE_GOOGLE_LOCAL_SSD,
}
}
@@ -113,7 +114,7 @@ func (self *SGoogleGuestDriver) ValidateResizeDisk(guest *models.SGuest, disk *m
if !utils.IsInStringArray(guest.Status, []string{api.VM_READY, api.VM_RUNNING}) {
return fmt.Errorf("Cannot resize disk when guest in status %s", guest.Status)
}
if !utils.IsInStringArray(storage.StorageType, []string{api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_PD_STANDARD}) {
if !utils.IsInStringArray(storage.StorageType, []string{api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_PD_STANDARD, api.STORAGE_GOOGLE_PD_BALANCED}) {
return fmt.Errorf("Cannot resize %s disk", storage.StorageType)
}
return nil
@@ -132,13 +133,15 @@ func (self *SGoogleGuestDriver) ValidateCreateData(ctx context.Context, userCred
minGB := -1
maxGB := -1
switch disk.Backend {
case api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_PD_STANDARD:
case api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_PD_STANDARD, api.STORAGE_GOOGLE_PD_BALANCED:
minGB = 10
maxGB = 65536
case api.STORAGE_GOOGLE_LOCAL_SSD:
minGB = 375
maxGB = 375
localDisk++
default:
return nil, httperrors.NewInputParameterError("Unknown google storage type %s", disk.Backend)
}
if i == 0 && disk.Backend == api.STORAGE_GOOGLE_LOCAL_SSD {
return nil, httperrors.NewInputParameterError("System disk does not support %s disk", disk.Backend)

View File

@@ -435,8 +435,8 @@ func (self *SManagedVirtualizedGuestDriver) RemoteDeployGuestForCreate(ctx conte
}
iVM, err := func() (cloudprovider.ICloudVM, error) {
lockman.LockObject(ctx, host)
defer lockman.ReleaseObject(ctx, host)
lockman.LockObject(ctx, guest)
defer lockman.ReleaseObject(ctx, guest)
iVM, err := ihost.CreateVM(&desc)
if err != nil {

View File

@@ -48,7 +48,7 @@ func (self *SGoogleHostDriver) ValidateDiskSize(storage *models.SStorage, sizeGb
minGB := 10
maxGB := -1
switch storage.StorageType {
case api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_PD_STANDARD:
case api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_PD_STANDARD, api.STORAGE_GOOGLE_PD_BALANCED:
maxGB = 65536
default:
return fmt.Errorf("Not support resize %s disk", storage.StorageType)

View File

@@ -199,7 +199,13 @@ func (manager *SBucketManager) newFromCloudBucket(
stats := extBucket.GetStats()
bucket.SizeBytes = stats.SizeBytes
if bucket.SizeBytes < 0 {
bucket.SizeBytes = 0
}
bucket.ObjectCnt = stats.ObjectCount
if bucket.ObjectCnt < 0 {
bucket.ObjectCnt = 0
}
limit := extBucket.GetLimit()
limitSupport := extBucket.LimitSupport()
@@ -259,7 +265,13 @@ func (bucket *SBucket) syncWithCloudBucket(
diff, err := db.UpdateWithLock(ctx, bucket, func() error {
stats := extBucket.GetStats()
bucket.SizeBytes = stats.SizeBytes
if bucket.SizeBytes < 0 {
bucket.SizeBytes = 0
}
bucket.ObjectCnt = stats.ObjectCount
if bucket.ObjectCnt < 0 {
bucket.ObjectCnt = 0
}
if !statsOnly {
limit := extBucket.GetLimit()

View File

@@ -117,6 +117,22 @@ func (self *SCloudregion) ValidateDeleteCondition(ctx context.Context) error {
return self.SEnabledStatusStandaloneResourceBase.ValidateDeleteCondition(ctx)
}
func (self *SCloudregion) GetElasticIps(managerId, eipMode string) ([]SElasticip, error) {
q := ElasticipManager.Query().Equals("cloudregion_id", self.Id)
if len(managerId) > 0 {
q = q.Equals("manager_id", managerId)
}
if len(eipMode) > 0 {
q = q.Equals("mode", eipMode)
}
eips := []SElasticip{}
err := db.FetchModelObjects(ElasticipManager, q, &eips)
if err != nil {
return nil, errors.Wrapf(err, "db.FetchModelObjects")
}
return eips, nil
}
func (self *SCloudregion) GetZoneQuery() *sqlchemy.SQuery {
zones := ZoneManager.Query()
if self.Id == api.DEFAULT_REGION_ID {

View File

@@ -1287,6 +1287,13 @@ func (manager *SDBInstanceManager) SyncDBInstances(ctx context.Context, userCred
return nil, nil, syncResult
}
for i := range dbInstances {
if taskman.TaskManager.IsInTask(&dbInstances[i]) {
syncResult.Error(fmt.Errorf("dbInstance %s(%s)in task", dbInstances[i].Name, dbInstances[i].Id))
return nil, nil, syncResult
}
}
removed := make([]SDBInstance, 0)
commondb := make([]SDBInstance, 0)
commonext := make([]cloudprovider.ICloudDBInstance, 0)

View File

@@ -448,6 +448,13 @@ func (manager *SElasticcacheManager) SyncElasticcaches(ctx context.Context, user
return nil, nil, syncResult
}
for i := range dbInstances {
if taskman.TaskManager.IsInTask(&dbInstances[i]) {
syncResult.Error(fmt.Errorf("ElasticCacheInstance %s(%s)in task", dbInstances[i].Name, dbInstances[i].Id))
return nil, nil, syncResult
}
}
removed := make([]SElasticcache, 0)
commondb := make([]SElasticcache, 0)
commonext := make([]cloudprovider.ICloudElasticcache, 0)

View File

@@ -242,19 +242,6 @@ func (manager *SElasticipManager) QueryDistinctExtraField(q *sqlchemy.SQuery, fi
return q, httperrors.ErrNotFound
}
func (manager *SElasticipManager) getEipsByRegion(region *SCloudregion, provider *SCloudprovider) ([]SElasticip, error) {
eips := make([]SElasticip, 0)
q := manager.Query().Equals("cloudregion_id", region.Id)
if provider != nil {
q = q.Equals("manager_id", provider.Id)
}
err := db.FetchModelObjects(manager, q, &eips)
if err != nil {
return nil, err
}
return eips, nil
}
func (self *SElasticip) GetRegion() *SCloudregion {
return CloudregionManager.FetchRegionById(self.CloudregionId)
}
@@ -323,7 +310,7 @@ func (manager *SElasticipManager) SyncEips(ctx context.Context, userCred mcclien
// remoteEips := make([]cloudprovider.ICloudEIP, 0)
syncResult := compare.SyncResult{}
dbEips, err := manager.getEipsByRegion(region, provider)
dbEips, err := region.GetElasticIps(provider.Id, api.EIP_MODE_STANDALONE_EIP)
if err != nil {
syncResult.Error(err)
return syncResult
@@ -421,8 +408,11 @@ func (self *SElasticip) SyncInstanceWithCloudEip(ctx context.Context, userCred m
case api.EIP_ASSOCIATE_TYPE_SERVER:
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), self.ManagerId))
case api.EIP_ASSOCIATE_TYPE_NAT_GATEWAY, api.EIP_ASSOCIATE_TYPE_LOADBALANCER:
case api.EIP_ASSOCIATE_TYPE_LOADBALANCER:
return q.Equals("manager_id", self.ManagerId)
case api.EIP_ASSOCIATE_TYPE_NAT_GATEWAY:
sq := VpcManager.Query("id").Equals("manager_id", self.ManagerId)
return q.In("vpc_id", sq.SubQuery())
}
return q
})

View File

@@ -3619,9 +3619,11 @@ func (self *SGuest) SaveRenewInfo(
guestdisks := self.GetDisks()
for i := 0; i < len(guestdisks); i += 1 {
disk := guestdisks[i].GetDisk()
err = disk.SaveRenewInfo(ctx, userCred, bc, expireAt, billingType)
if err != nil {
return err
if disk.AutoDelete {
err = disk.SaveRenewInfo(ctx, userCred, bc, expireAt, billingType)
if err != nil {
return err
}
}
}
return nil

View File

@@ -187,7 +187,7 @@ func (self *SGuestdisk) DoSave(driver string, cache string, mountpoint string) e
func (self *SGuestdisk) GetDisk() *SDisk {
disk, err := DiskManager.FetchById(self.DiskId)
if err != nil {
log.Errorf("GetDisk fail: %s", err)
log.Errorf("GetDisk %s fail: %s", self.DiskId, err)
return nil
}
return disk.(*SDisk)
@@ -287,6 +287,10 @@ func (self *SGuestdisk) GetDetailedJson() *jsonutils.JSONDict {
func (self *SGuestdisk) GetDetailedString() string {
disk := self.GetDisk()
if disk == nil {
return ""
}
var fs string
if len(disk.GetTemplateId()) > 0 {
fs = "root"

View File

@@ -1974,8 +1974,10 @@ func (self *SGuest) getNetworksDetails() string {
func (self *SGuest) getDisksDetails() string {
var buf bytes.Buffer
for _, disk := range self.GetDisks() {
buf.WriteString(disk.GetDetailedString())
buf.WriteString("\n")
if details := disk.GetDetailedString(); len(details) > 0 {
buf.WriteString(details)
buf.WriteString("\n")
}
}
return buf.String()
}
@@ -2730,7 +2732,7 @@ func getCloudNicNetwork(vnic cloudprovider.ICloudNic, host *SHost, ipList []stri
if vnet == nil {
if vnic.InClassicNetwork() {
region := host.GetRegion()
cloudprovider := region.GetCloudprovider()
cloudprovider := host.GetCloudprovider()
vpc, err := VpcManager.GetOrCreateVpcForClassicNetwork(cloudprovider, region)
if err != nil {
return nil, errors.Wrap(err, "NewVpcForClassicNetwork")
@@ -4220,6 +4222,14 @@ func (self *SGuest) GetShortDesc(ctx context.Context) *jsonutils.JSONDict {
billingInfo.SCloudProviderInfo = host.getCloudProviderInfo()
}
if len(self.BackupHostId) > 0 {
backupHost := HostManager.FetchHostById(self.BackupHostId)
if backupHost != nil {
desc.Set("backup_host", jsonutils.NewString(backupHost.Name))
desc.Set("backup_host_id", jsonutils.NewString(backupHost.Id))
}
}
if priceKey := self.GetMetadata("ext:price_key", nil); len(priceKey) > 0 {
billingInfo.PriceKey = priceKey
}

View File

@@ -1381,6 +1381,17 @@ func (self *SHost) GetGuests() []SGuest {
return guests
}
func (self *SHost) GetKvmGuests() []SGuest {
q := GuestManager.Query().Equals("host_id", self.Id).Equals("hypervisor", api.HYPERVISOR_KVM)
guests := make([]SGuest, 0)
err := db.FetchModelObjects(GuestManager, q, &guests)
if err != nil {
log.Errorf("GetGuests %s", err)
return nil
}
return guests
}
func (self *SHost) GetGuestCount() (int, error) {
q := self.GetGuestsQuery()
return q.CountWithError()
@@ -4809,7 +4820,7 @@ func (host *SHost) PerformHostMaintenance(ctx context.Context, userCred mcclient
preferHostId = host.Id
}
guests := host.GetGuests()
guests := host.GetKvmGuests()
for i := 0; i < len(guests); i++ {
lockman.LockObject(ctx, &guests[i])
defer lockman.ReleaseObject(ctx, &guests[i])
@@ -5229,3 +5240,15 @@ func (manager *SHostManager) ListItemExportKeys(ctx context.Context,
}
return q, nil
}
func (manager *SHostManager) FetchHostByExtId(extid string) *SHost {
host := SHost{}
host.SetModelManager(manager, &host)
err := manager.Query().Equals("external_id", extid).First(&host)
if err != nil {
log.Errorf("fetchHostByExtId fail %s", err)
return nil
} else {
return &host
}
}

View File

@@ -52,6 +52,7 @@ func InitDB() error {
LoadbalancerListenerRuleManager,
LoadbalancerBackendGroupManager,
LoadbalancerBackendManager,
AwsCachedLbbgManager,
LoadbalancerClusterManager,
SchedtagManager,
DynamicschedtagManager,

View File

@@ -159,16 +159,14 @@ func (man *SAwsCachedLbManager) SyncLoadbalancerBackends(ctx context.Context, us
if err != nil {
syncResult.UpdateError(err)
} else {
syncMetadata(ctx, userCred, &commondb[i], commonext[i])
syncResult.Update()
}
}
for i := 0; i < len(added); i++ {
local, err := man.newFromCloudLoadbalancerBackend(ctx, userCred, loadbalancerBackendgroup, added[i], syncOwnerId)
_, err := man.newFromCloudLoadbalancerBackend(ctx, userCred, loadbalancerBackendgroup, added[i], syncOwnerId)
if err != nil {
syncResult.AddError(err)
} else {
syncMetadata(ctx, userCred, local, added[i])
syncResult.Add()
}
}

View File

@@ -16,6 +16,7 @@ package models
import (
"context"
"database/sql"
"fmt"
"strconv"
"strings"
@@ -249,7 +250,6 @@ func (man *SAwsCachedLbbgManager) SyncLoadbalancerBackendgroups(ctx context.Cont
if err != nil {
syncResult.UpdateError(err)
} else {
syncMetadata(ctx, userCred, &commondb[i], commonext[i])
localLbgs = append(localLbgs, commondb[i])
remoteLbbgs = append(remoteLbbgs, commonext[i])
syncResult.Update()
@@ -277,7 +277,6 @@ func (man *SAwsCachedLbbgManager) SyncLoadbalancerBackendgroups(ctx context.Cont
if err != nil {
syncResult.AddError(err)
} else {
syncMetadata(ctx, userCred, new, added[i])
localLbgs = append(localLbgs, *new)
remoteLbbgs = append(remoteLbbgs, added[i])
syncResult.Add()
@@ -434,3 +433,30 @@ func (man *SAwsCachedLbbgManager) newFromCloudLoadbalancerBackendgroup(ctx conte
db.OpsLog.LogEvent(lbbg, db.ACT_CREATE, lbbg.GetShortDesc(ctx), userCred)
return lbbg, nil
}
func (man *SAwsCachedLbbgManager) InitializeData() error {
ret := []db.SMetadata{}
q := db.Metadata.Query().Equals("obj_type", man.Keyword())
err := db.FetchModelObjects(db.Metadata, q, &ret)
if err != nil {
if errors.Cause(err) != sql.ErrNoRows {
log.Debugf("SAwsCachedLbbgManager.InitializeData %s", err)
}
return nil
}
for i := range ret {
item := ret[i]
_, err := db.Update(&item, func() error {
return item.MarkDelete()
})
if err != nil {
log.Debugf("SAwsCachedLbbgManager.MarkDelete %s", err)
return nil
}
}
log.Debugf("SAwsCachedLbbgManager cleaned %d dirty data.", len(ret))
return nil
}

View File

@@ -162,16 +162,14 @@ func (man *SHuaweiCachedLbManager) SyncLoadbalancerBackends(ctx context.Context,
if err != nil {
syncResult.UpdateError(err)
} else {
syncMetadata(ctx, userCred, &commondb[i], commonext[i])
syncResult.Update()
}
}
for i := 0; i < len(added); i++ {
local, err := man.newFromCloudLoadbalancerBackend(ctx, userCred, loadbalancerBackendgroup, added[i], syncOwnerId)
_, err := man.newFromCloudLoadbalancerBackend(ctx, userCred, loadbalancerBackendgroup, added[i], syncOwnerId)
if err != nil {
syncResult.AddError(err)
} else {
syncMetadata(ctx, userCred, local, added[i])
syncResult.Add()
}
}

View File

@@ -222,7 +222,6 @@ func (man *SHuaweiCachedLbbgManager) SyncLoadbalancerBackendgroups(ctx context.C
if err != nil {
syncResult.UpdateError(err)
} else {
syncMetadata(ctx, userCred, &commondb[i], commonext[i])
localLbgs = append(localLbgs, commondb[i])
remoteLbbgs = append(remoteLbbgs, commonext[i])
syncResult.Update()
@@ -233,7 +232,6 @@ func (man *SHuaweiCachedLbbgManager) SyncLoadbalancerBackendgroups(ctx context.C
if err != nil {
syncResult.AddError(err)
} else {
syncMetadata(ctx, userCred, new, added[i])
localLbgs = append(localLbgs, *new)
remoteLbbgs = append(remoteLbbgs, added[i])
syncResult.Add()

View File

@@ -234,16 +234,14 @@ func (man *SQcloudCachedLbManager) SyncLoadbalancerBackends(ctx context.Context,
if err != nil {
syncResult.UpdateError(err)
} else {
syncMetadata(ctx, userCred, &commondb[i], commonext[i])
syncResult.Update()
}
}
for i := 0; i < len(added); i++ {
local, err := man.newFromCloudLoadbalancerBackend(ctx, userCred, loadbalancerBackendgroup, added[i], syncOwnerId)
_, err := man.newFromCloudLoadbalancerBackend(ctx, userCred, loadbalancerBackendgroup, added[i], syncOwnerId)
if err != nil {
syncResult.AddError(err)
} else {
syncMetadata(ctx, userCred, local, added[i])
syncResult.Add()
}
}

View File

@@ -294,7 +294,6 @@ func (man *SQcloudCachedLbbgManager) SyncLoadbalancerBackendgroups(ctx context.C
if err != nil {
syncResult.UpdateError(err)
} else {
syncMetadata(ctx, userCred, &commondb[i], commonext[i])
localLbgs = append(localLbgs, commondb[i])
remoteLbbgs = append(remoteLbbgs, commonext[i])
syncResult.Update()
@@ -305,7 +304,6 @@ func (man *SQcloudCachedLbbgManager) SyncLoadbalancerBackendgroups(ctx context.C
if err != nil {
syncResult.AddError(err)
} else {
syncMetadata(ctx, userCred, newlbbg, added[i])
localLbgs = append(localLbgs, *newlbbg)
remoteLbbgs = append(remoteLbbgs, added[i])
syncResult.Add()

View File

@@ -75,6 +75,12 @@ func (self *SLoadbalancerResourceBase) GetCloudprovider() *SCloudprovider {
if vpc != nil {
return vpc.GetCloudprovider()
}
lb := self.GetLoadbalancer()
if lb != nil {
return lb.GetCloudprovider()
}
return nil
}
@@ -161,23 +167,28 @@ func (manager *SLoadbalancerResourceBaseManager) FetchCustomizeColumns(
vpcList := make([]interface{}, len(rows))
zoneList := make([]interface{}, len(rows))
manList := make([]interface{}, len(rows))
for i := range rows {
rows[i] = api.LoadbalancerResourceInfo{}
if lb, ok := lbs[lbIds[i]]; ok {
rows[i].Loadbalancer = lb.Name
rows[i].VpcId = lb.VpcId
rows[i].ZoneId = lb.ZoneId
rows[i].ManagerId = lb.ManagerId
}
vpcList[i] = &SVpcResourceBase{rows[i].VpcId}
zoneList[i] = &SZoneResourceBase{rows[i].ZoneId}
manList[i] = &SManagedResourceBase{rows[i].ManagerId}
}
vpcRows := manager.SVpcResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, vpcList, fields, isList)
zoneRows := manager.SZoneResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, zoneList, fields, isList)
manRows := manager.SManagedResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, manList, fields, isList)
for i := range rows {
rows[i].VpcResourceInfo = vpcRows[i]
rows[i].ZoneResourceInfoBase = zoneRows[i].ZoneResourceInfoBase
rows[i].ZoneResourceInfo = zoneRows[i]
rows[i].ManagedResourceInfo = manRows[i]
}
return rows
}

View File

@@ -436,10 +436,7 @@ func (lb *SLoadbalancer) GetCreateLoadbalancerParams(iRegion cloudprovider.IClou
ChargeType: lb.ChargeType,
LoadbalancerSpec: lb.LoadbalancerSpec,
}
iRegion, err := lb.GetIRegion()
if err != nil {
return nil, err
}
if len(lb.ZoneId) > 0 {
zone := lb.GetZone()
if zone == nil {
@@ -447,7 +444,7 @@ func (lb *SLoadbalancer) GetCreateLoadbalancerParams(iRegion cloudprovider.IClou
}
iZone, err := iRegion.GetIZoneById(zone.ExternalId)
if err != nil {
return nil, err
return nil, errors.Wrap(err, "GetIZoneById")
}
params.ZoneID = iZone.GetId()
}
@@ -462,7 +459,7 @@ func (lb *SLoadbalancer) GetCreateLoadbalancerParams(iRegion cloudprovider.IClou
}
iVpc, err := iRegion.GetIVpcById(vpc.ExternalId)
if err != nil {
return nil, err
return nil, errors.Wrap(err, "GetIVpcById")
}
params.VpcID = iVpc.GetId()
}
@@ -476,7 +473,7 @@ func (lb *SLoadbalancer) GetCreateLoadbalancerParams(iRegion cloudprovider.IClou
for i := range networks {
iNetwork, err := networks[i].GetINetwork()
if err != nil {
return nil, err
return nil, errors.Wrap(err, "GetINetwork")
}
params.NetworkIDs = append(params.NetworkIDs, iNetwork.GetId())
}
@@ -706,6 +703,13 @@ func (man *SLoadbalancerManager) SyncLoadbalancers(ctx context.Context, userCred
return nil, nil, syncResult
}
for i := range dbLbs {
if taskman.TaskManager.IsInTask(&dbLbs[i]) {
syncResult.Error(fmt.Errorf("loadbalancer %s(%s)in task", dbLbs[i].Name, dbLbs[i].Id))
return nil, nil, syncResult
}
}
removed := []SLoadbalancer{}
commondb := []SLoadbalancer{}
commonext := []cloudprovider.ICloudLoadbalancer{}

View File

@@ -1691,6 +1691,7 @@ func (self *SNetwork) CustomizeDelete(ctx context.Context, userCred mcclient.Tok
}
func (self *SNetwork) RealDelete(ctx context.Context, userCred mcclient.TokenCredential) error {
DeleteResourceJointSchedtags(self, ctx, userCred)
db.OpsLog.LogEvent(self, db.ACT_DELOCATE, self.GetShortDesc(ctx), userCred)
self.SetStatus(userCred, api.NETWORK_STATUS_DELETED, "real delete")
self.ClearSchedDescCache()

View File

@@ -369,7 +369,7 @@ func (self *SSchedtag) GetObjectQuery() *sqlchemy.SQuery {
q := objs.Query()
q = q.Join(objschedtags, sqlchemy.AND(sqlchemy.Equals(objschedtags.Field(jointMan.GetMasterIdKey(jointMan)), objs.Field("id")),
sqlchemy.IsFalse(objschedtags.Field("deleted"))))
q = q.Filter(sqlchemy.IsTrue(objs.Field("enabled")))
// q = q.Filter(sqlchemy.IsTrue(objs.Field("enabled")))
q = q.Filter(sqlchemy.Equals(objschedtags.Field("schedtag_id"), self.Id))
return q
}
@@ -379,7 +379,8 @@ func (self *SSchedtag) GetJointManager() ISchedtagJointManager {
}
func (self *SSchedtag) GetObjectCount() (int, error) {
return self.GetJointManager().Query().Equals("schedtag_id", self.Id).CountWithError()
q := self.GetObjectQuery()
return q.CountWithError()
}
func (self *SSchedtag) getSchedPoliciesCount() (int, error) {

View File

@@ -26,9 +26,10 @@ import (
)
var (
syncAccountWorker *appsrv.SWorkerManager
syncWorkers []*appsrv.SWorkerManager
syncWorkerRing *hashring.HashRing
syncSecgroupWorker *appsrv.SWorkerManager
syncAccountWorker *appsrv.SWorkerManager
syncWorkers []*appsrv.SWorkerManager
syncWorkerRing *hashring.HashRing
)
func InitSyncWorkers(count int) {
@@ -50,6 +51,12 @@ func InitSyncWorkers(count int) {
2048,
true,
)
syncSecgroupWorker = appsrv.NewWorkerManager(
"syncSecgroupProbeWorkerManager",
1,
2048,
true,
)
}
func RunSyncCloudproviderRegionTask(key string, syncFunc func()) {
@@ -62,3 +69,7 @@ func RunSyncCloudproviderRegionTask(key string, syncFunc func()) {
func RunSyncCloudAccountTask(probeFunc func()) {
syncAccountWorker.Run(probeFunc, nil, nil)
}
func RunSyncSecgroupTask(syncFunc func()) {
syncSecgroupWorker.Run(syncFunc, nil, nil)
}

View File

@@ -183,7 +183,7 @@ func (manager *SVpcManager) getVpcExternalIdForClassicNetwork(regionId, cloudpro
func (manager *SVpcManager) GetOrCreateVpcForClassicNetwork(cloudprovider *SCloudprovider, region *SCloudregion) (*SVpc, error) {
externalId := manager.getVpcExternalIdForClassicNetwork(region.Id, cloudprovider.Id)
_vpc, err := db.FetchByExternalIdAndManagerId(manager, externalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", region.ManagerId)
return q.Equals("manager_id", cloudprovider.Id)
})
if err == nil {
return _vpc.(*SVpc), nil
@@ -200,7 +200,7 @@ func (manager *SVpcManager) GetOrCreateVpcForClassicNetwork(cloudprovider *SClou
vpc.SetEnabled(false)
vpc.Status = api.VPC_STATUS_UNAVAILABLE
vpc.ExternalId = externalId
vpc.ManagerId = region.ManagerId
vpc.ManagerId = cloudprovider.Id
err = manager.TableSpec().Insert(vpc)
if err != nil {
return nil, errors.Wrap(err, "Insert vpc for classic network")

View File

@@ -192,11 +192,11 @@ func (manager *SWireResourceBaseManager) ListItemFilter(
}
err = q.First(region)
if err != nil {
return nil, errors.Wrap(err, "q.First")
return nil, errors.Wrap(err, "regionQ.First")
}
if utils.IsInStringArray(region.Provider, api.REGIONAL_NETWORK_PROVIDERS) {
vpcQ := VpcManager.Query().SubQuery()
q = q.Join(vpcQ, sqlchemy.Equals(vpcQ.Field("id"), q.Field("vpc_id"))).
wireQ = wireQ.Join(vpcQ, sqlchemy.Equals(vpcQ.Field("id"), wireQ.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpcQ.Field("cloudregion_id"), region.Id))
} else {
zoneQuery := api.ZonalFilterListInput{

View File

@@ -1910,7 +1910,7 @@ func (self *SHuaWeiRegionDriver) RequestCreateLoadbalancer(ctx context.Context,
}
eip, err := db.FetchByExternalIdAndManagerId(models.ElasticipManager, ieip.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", lb.MarkUnDelete)
return q.Equals("manager_id", lb.ManagerId)
})
if err != nil {
return nil, errors.Wrap(err, "Huawei.RequestCreateLoadbalancer.FetchByExternalId")

View File

@@ -1476,95 +1476,109 @@ func (self *SManagedVirtualizationRegionDriver) RequestSyncSecurityGroup(ctx con
return "", errors.Wrap(err, "SSecurityGroupCache.Register")
}
iRegion, err := vpc.GetIRegion()
if err != nil {
return "", errors.Wrap(err, "vpc.GetIRegion")
}
waitChan := make(chan error)
var iSecgroup cloudprovider.ICloudSecurityGroup = nil
if len(cache.ExternalId) > 0 {
iSecgroup, err = iRegion.GetISecurityGroupById(cache.ExternalId)
if err != nil {
if errors.Cause(err) != cloudprovider.ErrNotFound {
return "", errors.Wrap(err, "iRegion.GetSecurityGroupById")
}
cache.ExternalId = ""
}
}
if len(cache.ExternalId) == 0 {
if strings.ToLower(secgroup.Name) == "default" { //避免有些云不支持default关键字
secgroup.Name = "DefaultGroup"
}
// 避免有的云不支持重名安全组
randomString := func(prefix string, length int) string {
return fmt.Sprintf("%s-%s", prefix, rand.String(length))
}
opts := &cloudprovider.SecurityGroupFilterOptions{
Name: randomString(secgroup.Name, 1),
VpcId: vpcId,
ProjectId: remoteProjectId,
}
for i := 2; i < 30; i++ {
_, err := iRegion.GetISecurityGroupByName(opts)
models.RunSyncSecgroupTask(func() {
err := func() error {
iRegion, err := vpc.GetIRegion()
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotFound {
break
}
if errors.Cause(err) != cloudprovider.ErrDuplicateId {
return "", err
return errors.Wrap(err, "vpc.GetIRegion")
}
var iSecgroup cloudprovider.ICloudSecurityGroup = nil
if len(cache.ExternalId) > 0 {
iSecgroup, err = iRegion.GetISecurityGroupById(cache.ExternalId)
if err != nil {
if errors.Cause(err) != cloudprovider.ErrNotFound {
return errors.Wrap(err, "iRegion.GetSecurityGroupById")
}
cache.ExternalId = ""
}
}
opts.Name = randomString(secgroup.Name, i)
}
conf := &cloudprovider.SecurityGroupCreateInput{
Name: opts.Name,
Desc: secgroup.Description,
VpcId: vpcId,
ProjectId: remoteProjectId,
Rules: secgroup.GetSecRules(""),
}
iSecgroup, err = iRegion.CreateISecurityGroup(conf)
if err != nil {
return "", errors.Wrap(err, "iRegion.CreateISecurityGroup")
}
}
_, err = db.Update(cache, func() error {
cache.ExternalId = iSecgroup.GetGlobalId()
cache.Name = iSecgroup.GetName()
cache.Status = api.SECGROUP_CACHE_STATUS_READY
return nil
if len(cache.ExternalId) == 0 {
if strings.ToLower(secgroup.Name) == "default" { //避免有些云不支持default关键字
secgroup.Name = "DefaultGroup"
}
// 避免有的云不支持重名安全组
randomString := func(prefix string, length int) string {
return fmt.Sprintf("%s-%s", prefix, rand.String(length))
}
opts := &cloudprovider.SecurityGroupFilterOptions{
Name: randomString(secgroup.Name, 1),
VpcId: vpcId,
ProjectId: remoteProjectId,
}
for i := 2; i < 30; i++ {
_, err := iRegion.GetISecurityGroupByName(opts)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotFound {
break
}
if errors.Cause(err) != cloudprovider.ErrDuplicateId {
return errors.Wrapf(err, "GetISecurityGroupByName")
}
}
opts.Name = randomString(secgroup.Name, i)
}
conf := &cloudprovider.SecurityGroupCreateInput{
Name: opts.Name,
Desc: secgroup.Description,
VpcId: vpcId,
ProjectId: remoteProjectId,
Rules: secgroup.GetSecRules(""),
}
iSecgroup, err = iRegion.CreateISecurityGroup(conf)
if err != nil {
return errors.Wrapf(err, "iRegion.CreateISecurityGroup")
}
}
_, err = db.Update(cache, func() error {
cache.ExternalId = iSecgroup.GetGlobalId()
cache.Name = iSecgroup.GetName()
cache.Status = api.SECGROUP_CACHE_STATUS_READY
return nil
})
if err != nil {
return errors.Wrapf(err, "db.Update")
}
rules, err := iSecgroup.GetRules()
if err != nil {
return errors.Wrapf(err, "iSecgroup.GetRules")
}
maxPriority := region.GetDriver().GetSecurityGroupRuleMaxPriority()
minPriority := region.GetDriver().GetSecurityGroupRuleMinPriority()
defaultInRule := region.GetDriver().GetDefaultSecurityGroupInRule()
defaultOutRule := region.GetDriver().GetDefaultSecurityGroupOutRule()
order := region.GetDriver().GetSecurityGroupRuleOrder()
onlyAllowRules := region.GetDriver().IsOnlySupportAllowRules()
localRules := secrules.SecurityRuleSet(secgroup.GetSecRules(""))
common, inAdds, outAdds, inDels, outDels := cloudprovider.CompareRules(minPriority, maxPriority, order, localRules, rules, defaultInRule, defaultOutRule, onlyAllowRules, false)
if len(inAdds) == 0 && len(inDels) == 0 && len(outAdds) == 0 && len(outDels) == 0 {
return nil
}
return iSecgroup.SyncRules(common, inAdds, outAdds, inDels, outDels)
}()
waitChan <- err
})
err = <-waitChan
if err != nil {
return "", errors.Wrap(err, "db.Update")
return "", err
}
rules, err := iSecgroup.GetRules()
cache, err = models.SecurityGroupCacheManager.Register(ctx, userCred, secgroup.Id, vpcId, region.Id, vpc.ManagerId, remoteProjectId)
if err != nil {
return "", errors.Wrap(err, "iSecgroup.GetRules")
}
maxPriority := region.GetDriver().GetSecurityGroupRuleMaxPriority()
minPriority := region.GetDriver().GetSecurityGroupRuleMinPriority()
defaultInRule := region.GetDriver().GetDefaultSecurityGroupInRule()
defaultOutRule := region.GetDriver().GetDefaultSecurityGroupOutRule()
order := region.GetDriver().GetSecurityGroupRuleOrder()
onlyAllowRules := region.GetDriver().IsOnlySupportAllowRules()
localRules := secrules.SecurityRuleSet(secgroup.GetSecRules(""))
common, inAdds, outAdds, inDels, outDels := cloudprovider.CompareRules(minPriority, maxPriority, order, localRules, rules, defaultInRule, defaultOutRule, onlyAllowRules, false)
if len(inAdds) == 0 && len(inDels) == 0 && len(outAdds) == 0 && len(outDels) == 0 {
return cache.ExternalId, nil
}
err = iSecgroup.SyncRules(common, inAdds, outAdds, inDels, outDels)
if err != nil {
return "", errors.Wrap(err, "iSecgroup.SyncRules")
return "", errors.Wrap(err, "SSecurityGroupCache.Register")
}
return cache.ExternalId, nil

View File

@@ -16,8 +16,11 @@ package tasks
import (
"context"
"database/sql"
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -62,6 +65,18 @@ func (self *CloudAccountDeleteTask) OnInit(ctx context.Context, obj db.IStandalo
func (self *CloudAccountDeleteTask) OnAllCloudProviderDeleteComplete(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) {
account := obj.(*models.SCloudaccount)
// check providers deleted success
count, err := account.GetProviderCount()
if err != nil && errors.Cause(err) != sql.ErrNoRows {
self.OnAllCloudProviderDeleteCompleteFailed(ctx, obj, jsonutils.NewString(err.Error()))
return
}
if count > 0 {
self.OnAllCloudProviderDeleteCompleteFailed(ctx, obj, jsonutils.NewString(fmt.Sprintf("cloudprovider deleted failed count %d", count)))
return
}
account.RealDelete(ctx, self.UserCred)
self.SetStageComplete(ctx, nil)

View File

@@ -157,5 +157,5 @@ func (self *SnapshotCleanupTask) OnDeleteSnapshot(ctx context.Context, obj db.IS
func (self *SnapshotCleanupTask) OnDeleteSnapshotFailed(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
log.Errorf("snapshot delete faield %s", data)
self.OnDeleteSnapshot(ctx, obj, data)
self.SetStageFailed(ctx, data.String())
}

View File

@@ -101,10 +101,13 @@ func (self *DiskDeleteTask) startDeleteDisk(ctx context.Context, disk *models.SD
)
storage = disk.GetStorage()
if storage != nil {
host = storage.GetMasterHost()
if storage == nil { // dirty data
self.OnGuestDiskDeleteComplete(ctx, disk, nil)
return
}
host = storage.GetMasterHost()
isPurge := false
if (host == nil || !host.GetEnabled()) && jsonutils.QueryBoolean(self.Params, "purge", false) {
isPurge = true
@@ -112,21 +115,25 @@ func (self *DiskDeleteTask) startDeleteDisk(ctx context.Context, disk *models.SD
disk.SetStatus(self.UserCred, api.DISK_DEALLOC, "")
if isPurge {
self.OnGuestDiskDeleteComplete(ctx, disk, nil)
return
}
if isNeed, _ := disk.IsNeedWaitSnapshotsDeleted(); isNeed { // for kvm rbd disk
self.OnGuestDiskDeleteComplete(ctx, disk, nil)
return
}
if len(disk.BackupStorageId) > 0 {
self.SetStage("OnMasterStorageDeleteDiskComplete", nil)
} else {
if isNeed, _ := disk.IsNeedWaitSnapshotsDeleted(); isNeed {
self.OnGuestDiskDeleteComplete(ctx, disk, nil)
return
}
if len(disk.BackupStorageId) > 0 {
self.SetStage("OnMasterStorageDeleteDiskComplete", nil)
} else {
self.SetStage("OnGuestDiskDeleteComplete", nil)
}
if host == nil {
self.OnGuestDiskDeleteCompleteFailed(ctx, disk, jsonutils.NewString("fail to find master host"))
} else if err := host.GetHostDriver().RequestDeallocateDiskOnHost(ctx, host, storage, disk, self); err != nil {
self.OnGuestDiskDeleteCompleteFailed(ctx, disk, jsonutils.NewString(err.Error()))
}
self.SetStage("OnGuestDiskDeleteComplete", nil)
}
if host == nil {
self.OnGuestDiskDeleteCompleteFailed(ctx, disk, jsonutils.NewString("fail to find master host"))
return
}
err := host.GetHostDriver().RequestDeallocateDiskOnHost(ctx, host, storage, disk, self)
if err != nil {
self.OnGuestDiskDeleteCompleteFailed(ctx, disk, jsonutils.NewString(err.Error()))
return
}
}

View File

@@ -133,7 +133,6 @@ func Start(app *appsrv.Application) error {
}
func (agent *SEsxiAgent) AddImageCacheHandler(prefix string, app *appsrv.Application) {
hostutils.InitWorkerManager()
app.AddHandler("POST",
fmt.Sprintf("%s/disks/image_cache", prefix),
auth.Authenticate(func(ctx context.Context, w http.ResponseWriter, r *http.Request) {

View File

@@ -75,7 +75,7 @@ func uploadHandler(ctx context.Context, w http.ResponseWriter, r *http.Request)
httperrors.MissingParameterError(w, "miss disk")
return
}
hostutils.DelayTask(ctx, esxi.EsxiAgent.AgentStorage.SaveToGlance, disk)
hostutils.DelayTaskWithoutReqctx(ctx, esxi.EsxiAgent.AgentStorage.SaveToGlance, disk)
hostutils.ResponseOk(ctx, w)
}

View File

@@ -19,18 +19,19 @@ import common_options "yunion.io/x/onecloud/pkg/cloudcommon/options"
type EsxiOptions struct {
common_options.CommonOptions
ListenInterface string `help:"Master address of host server" default:"br0"`
ListenAddress string `help:"Host serve IP address to select when multiple address bind to ListenInterface"`
EsxiAgentPath string `default:"/opt/cloud/workspace/esxi_agent" help:"Path for esxi agent configuration files"`
ImageCachePath string `help:"Path for storing image caches"`
ImageCacheLimit int `help:"Maximal storage space for image caching, in GB" default:"20"`
AgentTempPath string `help:"Path for ESXI Agent"`
AgentTempLimit int `help:"Maximal storage space for ESXi agent, in GB" default:"20"`
LinuxDefaultRootUser bool `help:"Default account for Linux system is root" default:"false"`
WindowsDefaultAdminUser bool `help:"Default account for Windows system is Administrator" default:"true"`
DefaultImageSaveFormat string `help:"Default image save format, default is vmdk, canbe qcow2" default:"vmdk"`
Zone string `help:"Zone where the agent locates"`
DeployServerSocketPath string `help:"Deploy server listen socket path" default:"/var/run/deploy.sock"`
ListenInterface string `help:"Master address of host server" default:"br0"`
ListenAddress string `help:"Host serve IP address to select when multiple address bind to ListenInterface"`
EsxiAgentPath string `default:"/opt/cloud/workspace/esxi_agent" help:"Path for esxi agent configuration files"`
ImageCachePath string `help:"Path for storing image caches"`
ImageCacheLimit int `help:"Maximal storage space for image caching, in GB" default:"20"`
AgentTempPath string `help:"Path for ESXI Agent"`
AgentTempLimit int `help:"Maximal storage space for ESXi agent, in GB" default:"20"`
LinuxDefaultRootUser bool `help:"Default account for Linux system is root" default:"false"`
WindowsDefaultAdminUser bool `help:"Default account for Windows system is Administrator" default:"true"`
DefaultImageSaveFormat string `help:"Default image save format, default is vmdk, canbe qcow2" default:"vmdk"`
Zone string `help:"Zone where the agent locates"`
DeployServerSocketPath string `help:"Deploy server listen socket path" default:"/var/run/deploy.sock"`
HostDelayTaskWorkerCount int `default:"8" help:"Host delay worker thread count, default is 8"`
}
var (

View File

@@ -30,6 +30,7 @@ import (
"yunion.io/x/onecloud/pkg/esxi/options"
"yunion.io/x/onecloud/pkg/hostman/guestfs/fsdriver"
"yunion.io/x/onecloud/pkg/hostman/hostdeployer/deployclient"
"yunion.io/x/onecloud/pkg/hostman/hostutils"
)
type SExsiAgentService struct {
@@ -70,6 +71,7 @@ func (s *SExsiAgentService) StartService() {
fsdriver.Init(nil)
deployclient.Init(options.Options.DeployServerSocketPath)
hostutils.InitWorkerManagerWithCount(options.Options.HostDelayTaskWorkerCount)
app := app_common.InitApp(&options.Options.BaseOptions, false)
handler.InitHandlers(app)

View File

@@ -764,7 +764,7 @@ func (d *SUbuntuRootFs) GetReleaseInfo(rootFs IDiskPartition) *deployapi.Release
lines := strings.Split(string(rel), "\n")
for _, l := range lines {
if strings.HasPrefix(l, distroKey) {
version = strings.TrimSpace(l[len(distroKey) : len(l)-1])
version = strings.TrimSpace(l[len(distroKey):])
}
}
return deployapi.NewReleaseInfo(d.GetName(), version, d.GetArch(rootFs))

View File

@@ -993,10 +993,13 @@ func (s *SKVMGuestInstance) ExecSuspendTask(ctx context.Context) {
func (s *SKVMGuestInstance) GetNicDescMatch(mac, ip, port, bridge string) jsonutils.JSONObject {
nics, _ := s.Desc.GetArray("nics")
for _, nic := range nics {
nicBridge, _ := nic.GetString("bridge")
if bridge == "" && nicBridge != "" && nicBridge == options.HostOptions.OvnIntegrationBridge {
continue
}
nicMac, _ := nic.GetString("mac")
nicIp, _ := nic.GetString("ip")
nicPort, _ := nic.GetString("ifname")
nicBridge, _ := nic.GetString("bridge")
if (len(mac) == 0 || netutils2.MacEqual(nicMac, mac)) &&
(len(ip) == 0 || nicIp == ip) &&
(len(port) == 0 || nicPort == port) &&

View File

@@ -757,6 +757,7 @@ func (h *SHostInfo) tryCreateNetworkOnWire() {
params.Set("mask", jsonutils.NewInt(int64(mask)))
params.Set("is_on_premise", jsonutils.JSONTrue)
params.Set("server_type", jsonutils.NewString(api.NETWORK_TYPE_BAREMETAL))
params.Set("is_on_premise", jsonutils.JSONTrue)
ret, err := modules.Networks.PerformClassAction(
hostutils.GetComputeSession(context.Background()),
"try-create-network", params)
@@ -1295,6 +1296,9 @@ func (h *SHostInfo) onGetStorageInfoSucc(hoststorages []jsonutils.JSONObject) {
func (h *SHostInfo) uploadStorageInfo() {
for _, s := range storageman.GetManager().Storages {
if err := s.SetStorageInfo(s.GetId(), s.GetStorageName(), s.GetStorageConf()); err != nil {
h.onFail(err)
}
res, err := s.SyncStorageInfo()
if err != nil {
h.onFail(err)

View File

@@ -189,7 +189,11 @@ func DelayTaskWithWorker(
}
func InitWorkerManager() {
wm = workmanager.NewWorkManger(TaskFailed, TaskComplete, options.HostOptions.DefaultRequestWorkerCount)
InitWorkerManagerWithCount(options.HostOptions.DefaultRequestWorkerCount)
}
func InitWorkerManagerWithCount(count int) {
wm = workmanager.NewWorkManger(TaskFailed, TaskComplete, count)
}
func InitK8sWorkerManager() {

View File

@@ -191,8 +191,8 @@ func (l *SLocalImageCache) prepare(ctx context.Context, zone, srcUrl, format str
}
func (l *SLocalImageCache) fetch(ctx context.Context, zone, srcUrl, format string) bool {
if (fileutils2.Exists(l.GetPath()) &&
l.remoteFile.VerifyIntegrity()) || l.remoteFile.Fetch() {
if (fileutils2.Exists(l.GetPath()) && l.remoteFile.VerifyIntegrity()) ||
l.remoteFile.Fetch() {
if len(l.Manager.GetId()) > 0 {
_, err := hostutils.RemoteStoragecacheCacheImage(ctx,
l.Manager.GetId(), l.imageId, "ready", l.GetPath())

View File

@@ -24,7 +24,6 @@ import (
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/hostman/hostutils"
"yunion.io/x/onecloud/pkg/hostman/options"
"yunion.io/x/onecloud/pkg/hostman/storageman/remotefile"
"yunion.io/x/onecloud/pkg/mcclient/modules"
"yunion.io/x/onecloud/pkg/util/procutils"
@@ -75,11 +74,11 @@ func (r *SRbdImageCache) Acquire(ctx context.Context, zone, srcUrl, format strin
}
r.imageName = localImageCache.GetName()
if !r.Load() {
log.Debugf("convert local image %s to rbd pool %s", r.imageId, r.Manager.GetPath())
log.Infof("convert local image %s to rbd pool %s", r.imageId, r.Manager.GetPath())
err := procutils.NewRemoteCommandAsFarAsPossible(qemutils.GetQemuImg(),
"convert", "-O", "raw", localImageCache.GetPath(), r.GetPath()).Run()
if err != nil {
log.Errorf("failed to convert image %s", options.HostOptions.ServersPath)
log.Errorf("failed to convert image %s", err)
return false
}
}

View File

@@ -211,26 +211,16 @@ func (c *SAgentImageCacheManager) prefetchImageCacheByUpload(ctx context.Context
}
func (c *SAgentImageCacheManager) perfetchTemplateVMImageCache(ctx context.Context, data *sImageCacheData) (jsonutils.JSONObject, error) {
client, err := esxi.NewESXiClientFromAccessInfo(ctx, &data.Datastore)
if err != nil {
return nil, errors.Wrap(err, "esxi.NewESXiClientFromJson")
}
// data.StorageCacheExternalId is Host Ip associated with the StorageCache in where CachedImage stored
host, err := client.FindHostByIp(data.StorageCacheHostIp)
_, err = client.SearchTemplateVM(data.ImageExternalId)
if err != nil {
return nil, err
}
dc, err := host.GetDatacenter()
if err != nil {
return nil, errors.Wrap(err, "host.GetDatacenter")
}
_, err = dc.GetTemplateVMById(data.ImageExternalId)
if err != nil {
return nil, err
return nil, errors.Wrapf(err, "SEsxiClient.SearchTemplateVM for image %q", data.ImageExternalId)
}
res := jsonutils.NewDict()
res.Add(jsonutils.NewString(data.ImageId), "image_id")
res.Add(jsonutils.NewString(data.ImageExternalId), "image_id")
return res, nil
}

View File

@@ -111,7 +111,7 @@ func (c *SRbdImageCacheManager) PrefetchImageCache(ctx context.Context, data int
if err != nil {
return nil, err
}
format := "qcow2"
format, _ := body.GetString("format")
srcUrl, _ := body.GetString("src_url")
zone, _ := body.GetString("zone")

View File

@@ -259,7 +259,13 @@ func (r *SRemoteFile) downloadInternal(getData bool, preChksum string) bool {
}
func (r *SRemoteFile) setProperties(header http.Header) {
r.chksum = header.Get("X-Image-Meta-Checksum")
r.format = header.Get("X-Image-Meta-Disk_format")
r.name = header.Get("X-Image-Meta-Name")
if chksum := header.Get("X-Image-Meta-Checksum"); len(chksum) > 0 {
r.chksum = chksum
}
if format := header.Get("X-Image-Meta-Disk_format"); len(format) > 0 {
r.format = format
}
if name := header.Get("X-Image-Meta-Name"); len(name) > 0 {
r.name = name
}
}

View File

@@ -125,6 +125,7 @@ type IStorage interface {
disksBackingFile, srcSnapshots jsonutils.JSONObject, rebaseDisks bool, diskDesc jsonutils.JSONObject) error
Accessible() error
Detach() error
}
type SBaseStorage struct {
@@ -262,6 +263,7 @@ func (s *SBaseStorage) bindMountTo(sPath string) error {
return errors.Errorf("bind mount temp path to local image path failed %s", out)
}
}
log.Infof("bind mount %s -> %s", tempPath, sPath)
return nil
}

View File

@@ -174,6 +174,10 @@ func (s *SLocalStorage) Accessible() error {
}
func (s *SLocalStorage) Detach() error {
return nil
}
func (s *SLocalStorage) DeleteDiskfile(diskpath string) error {
log.Infof("Start Delete %s", diskpath)
if options.HostOptions.RecycleDiskfile {
@@ -219,7 +223,7 @@ func (s *SLocalStorage) SaveToGlance(ctx context.Context, params interface{}) (j
if err := s.saveToGlance(ctx, imageId, imagePath, compress, format); err != nil {
log.Errorf("Save to glance failed: %s", err)
s.onSaveToGlanceFailed(ctx, imageId)
s.onSaveToGlanceFailed(ctx, imageId, err.Error())
}
imagecacheManager := s.Manager.LocalStorageImagecacheManager
@@ -305,11 +309,14 @@ func (s *SLocalStorage) saveToGlance(ctx context.Context, imageId, imagePath str
return err
}
func (s *SLocalStorage) onSaveToGlanceFailed(ctx context.Context, imageId string) {
func (s *SLocalStorage) onSaveToGlanceFailed(ctx context.Context, imageId string, reason string) {
params := jsonutils.NewDict()
params.Set("status", jsonutils.NewString("killed"))
_, err := modules.Images.Update(hostutils.GetImageSession(ctx, s.GetZoneName()),
imageId, params)
params.Set("reason", jsonutils.NewString(reason))
_, err := modules.Images.PerformAction(
hostutils.GetImageSession(ctx, s.GetZoneName()),
imageId, "update-status", params,
)
if err != nil {
log.Errorln(err)
}

View File

@@ -17,6 +17,7 @@ package storageman
import (
"context"
"fmt"
"path"
"strings"
"time"
@@ -129,3 +130,22 @@ func (s *SNFSStorage) checkAndMount() error {
}
return nil
}
func (s *SNFSStorage) Detach() error {
if !strings.HasPrefix(s.Path, "/opt/cloud") {
tmpPath := path.Join(TempBindMountPath, s.Path)
out, err := procutils.NewCommand("umount", s.Path).Output()
if err != nil {
return errors.Wrapf(err, "1. umount %s failed %s", s.Path, out)
}
out, err = procutils.NewRemoteCommandAsFarAsPossible("umount", tmpPath).Output()
if err != nil {
return errors.Wrapf(err, "2. umount %s failed %s", tmpPath, out)
}
}
out, err := procutils.NewRemoteCommandAsFarAsPossible("umount", s.Path).Output()
if err != nil {
return errors.Wrapf(err, "3. umount %s failed %s", s.Path, out)
}
return nil
}

View File

@@ -634,12 +634,16 @@ func (s *SRbdStorage) Accessible() error {
select {
case err = <-c:
break
case <-time.After(time.Second * 10):
case <-time.After(time.Second * 30):
err = ErrStorageTimeout
}
return err
}
func (s *SRbdStorage) Detach() error {
return nil
}
func (s *SRbdStorage) SaveToGlance(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) {
data, ok := params.(*jsonutils.JSONDict)
if !ok {

View File

@@ -66,7 +66,7 @@ func storageVerifyMountPoint(ctx context.Context, w http.ResponseWriter, r *http
hostutils.Response(ctx, w, httperrors.NewMissingParameterError("mount_point"))
return
}
output, err := procutils.NewCommand("mountpoint", mountPoint).Output()
output, err := procutils.NewRemoteCommandAsFarAsPossible("mountpoint", mountPoint).Output()
if err == nil {
appsrv.SendStruct(w, map[string]interface{}{"is_mount_point": true})
} else {
@@ -136,6 +136,9 @@ func storageDetach(ctx context.Context, body jsonutils.JSONObject) (interface{},
if storage == nil {
return nil, httperrors.NewBadRequestError("ShareStorage[%s] Has detach from host ...", name)
}
if err := storage.Detach(); err != nil {
log.Errorf("detach storage %s failed: %s", storage.GetPath(), err)
}
storageman.GetManager().Remove(storage)
return nil, nil
}

View File

@@ -603,6 +603,16 @@ func (self *SImage) ValidateUpdateData(ctx context.Context, userCred mcclient.To
if appParams != nil && appParams.Request.ContentLength > 0 {
return nil, httperrors.NewInvalidStatusError("cannot upload in status %s", self.Status)
}
if minDiskSize, err := data.Int("min_disk"); err == nil {
img, err := qemuimg.NewQemuImage(self.getLocalLocation())
if err != nil {
return nil, errors.Wrap(err, "open image")
}
virtualSizeMB := img.SizeBytes / 1024 / 1024
if virtualSizeMB > 0 && minDiskSize < virtualSizeMB {
return nil, httperrors.NewBadRequestError("min disk size must >= %v", virtualSizeMB)
}
}
} else {
appParams := appsrv.AppContextGetParams(ctx)
if appParams != nil {
@@ -654,7 +664,6 @@ func (self *SImage) ValidateUpdateData(ctx context.Context, userCred mcclient.To
return nil, errors.Wrap(err, "SSharableVirtualResourceBase.ValidateUpdateData")
}
data.Update(jsonutils.Marshal(input))
return data, nil
}
@@ -1443,6 +1452,26 @@ func (img *SImage) GetUsages() []db.IUsage {
}
}
func (self *SImage) AllowPerformUpdateStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return db.IsAdminAllowPerform(userCred, self, "update-status")
}
func (img *SImage) PerformUpdateStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ImageUpdateStatusInput) (jsonutils.JSONObject, error) {
if !utils.IsInStringArray(input.Status, api.ImageDeadStatus) {
return nil, httperrors.NewBadRequestError("can't udpate image to status %s, must in %v", input.Status, api.ImageDeadStatus)
}
_, err := db.Update(img, func() error {
img.Status = input.Status
return nil
})
if err != nil {
return nil, errors.Wrap(err, "udpate image status")
}
db.OpsLog.LogEvent(img, db.ACT_UPDATE_STATUS, input.Reason, userCred)
logclient.AddSimpleActionLog(img, logclient.ACT_UPDATE_STATUS, input.Reason, userCred, true)
return nil, nil
}
func (img *SImage) PerformPublic(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input apis.PerformPublicProjectInput) (jsonutils.JSONObject, error) {
if img.IsStandard.IsTrue() {
return nil, errors.Wrap(httperrors.ErrForbidden, "cannot perform public for standard image")

View File

@@ -19,6 +19,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
api "yunion.io/x/onecloud/pkg/apis/compute"
@@ -202,14 +203,20 @@ func (self *SZone) getStorageByCategory(category string) (*SStorage, error) {
func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
return self.istorages, nil
}
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
for i := 0; i < len(self.istorages); i += 1 {
if self.istorages[i].GetGlobalId() == id {

View File

@@ -19,6 +19,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudprovider"
@@ -145,21 +146,27 @@ func (self *SZone) GetIHostById(id string) (cloudprovider.ICloudHost, error) {
func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
return self.istorages, nil
}
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
for i := 0; i < len(self.istorages); i += 1 {
if self.istorages[i].GetGlobalId() == id {
return self.istorages[i], nil
}
}
return nil, ErrorNotFound()
return nil, errors.Wrapf(cloudprovider.ErrNotFound, "not found %s", id)
}
func (self *SZone) getStorageByCategory(category string) (*SStorage, error) {

View File

@@ -19,6 +19,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/cloudprovider"
)
@@ -117,16 +118,22 @@ func (self *SZone) fetchStorages() error {
func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
err := self.fetchStorages()
if err != nil {
return nil, err
return nil, errors.Wrapf(err, "fetchStorages")
}
err = self.fetchClassicStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchClassicStorages")
}
self.fetchClassicStorages()
istorages := append(self.istorages, self.iclassicStorages...)
return istorages, nil
}
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
for i := 0; i < len(self.istorages); i += 1 {
if self.istorages[i].GetGlobalId() == id {

View File

@@ -365,7 +365,7 @@ func (self *SInstance) GetInstanceType() string {
func (self *SInstance) GetSecurityGroupIds() ([]string, error) {
if len(self.SecurityGroups) == 0 {
return nil, nil
return []string{}, nil
}
if len(self.MasterOrderId) > 0 {
@@ -423,8 +423,34 @@ func (self *SInstance) AssignSecurityGroup(secgroupId string) error {
}
func (self *SInstance) SetSecurityGroups(secgroupIds []string) error {
for i := 0; i < len(secgroupIds); i++ {
err := self.host.zone.region.AssignSecurityGroup(self.GetId(), secgroupIds[i])
currentIds, err := self.GetSecurityGroupIds()
if err != nil {
return errors.Wrap(err, "GetSecurityGroupIds")
}
adds := []string{}
for i := range secgroupIds {
if !utils.IsInStringArray(secgroupIds[i], currentIds) {
adds = append(adds, secgroupIds[i])
}
}
for i := range adds {
err := self.host.zone.region.AssignSecurityGroup(self.GetId(), adds[i])
if err != nil {
return errors.Wrap(err, "Instance.SetSecurityGroups")
}
}
removes := []string{}
for i := range currentIds {
if !utils.IsInStringArray(currentIds[i], secgroupIds) {
removes = append(removes, currentIds[i])
}
}
for i := range removes {
err := self.host.zone.region.UnsignSecurityGroup(self.GetId(), removes[i])
if err != nil {
return errors.Wrap(err, "Instance.SetSecurityGroups")
}

View File

@@ -53,6 +53,7 @@ var LatitudeAndLongitude = map[string]cloudprovider.SGeographicInfo{
"cn-baoding1": {Latitude: 38.8739745619, Longitude: 115.4646082830, City: api.CITY_BAO_DING, CountryCode: api.COUNTRY_CODE_CN},
"cn-nj2": {Latitude: 32.0584065670, Longitude: 118.7964897811, City: api.CITY_NAN_JING, CountryCode: api.COUNTRY_CODE_CN},
"cn-gdgz1": {Latitude: 23.12911, Longitude: 113.264385, City: api.CITY_GUANG_ZHOU, CountryCode: api.COUNTRY_CODE_CN},
"cn-bj1": {Latitude: 39.904202, Longitude: 116.407394, City: api.CITY_BEI_JING, CountryCode: api.COUNTRY_CODE_CN},
"cn-bj4": {Latitude: 39.904202, Longitude: 116.407394, City: api.CITY_BEI_JING, CountryCode: api.COUNTRY_CODE_CN},
"cn-neimeng5": {Latitude: 40.842358, Longitude: 111.749992, City: api.CITY_HU_HE_HAO_TE, CountryCode: api.COUNTRY_CODE_CN},
"cn-shanghai5": {Latitude: 31.210344, Longitude: 121.455364, City: api.CITY_SHANG_HAI, CountryCode: api.COUNTRY_CODE_CN},

View File

@@ -18,6 +18,7 @@ import (
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/cloudprovider"
)
@@ -94,14 +95,20 @@ func (self *SZone) GetIHostById(id string) (cloudprovider.ICloudHost, error) {
func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
return self.istorages, nil
}
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
for i := 0; i < len(self.istorages); i += 1 {
if self.istorages[i].GetGlobalId() == id {

View File

@@ -18,6 +18,8 @@ import (
"strings"
"github.com/vmware/govmomi/object"
"github.com/vmware/govmomi/property"
"github.com/vmware/govmomi/view"
"github.com/vmware/govmomi/vim25/mo"
"github.com/vmware/govmomi/vim25/types"
@@ -237,32 +239,90 @@ func (dc *SDatacenter) getDcObj() *object.Datacenter {
return object.NewDatacenter(dc.manager.client.Client, dc.object.Reference())
}
// fetchVms will identify if VM is a template and return two different arrays; the latter contains all template vms.
func (dc *SDatacenter) fetchVms(vmRefs []types.ManagedObjectReference, all bool) ([]cloudprovider.ICloudVM, []*SVirtualMachine, error) {
var vms []mo.VirtualMachine
func (dc *SDatacenter) fetchVms(vmRefs []types.ManagedObjectReference, all bool) ([]*SVirtualMachine, error) {
var movms []mo.VirtualMachine
if vmRefs != nil {
err := dc.manager.references2Objects(vmRefs, VIRTUAL_MACHINE_PROPS, &vms)
err := dc.manager.references2Objects(vmRefs, VIRTUAL_MACHINE_PROPS, &movms)
if err != nil {
return nil, nil, errors.Wrap(err, "dc.manager.references2Objects")
return nil, errors.Wrap(err, "dc.manager.references2Objects")
}
}
// avoid applying new memory and copying
retVms := make([]cloudprovider.ICloudVM, 0, len(vms)/2)
templateVMs := make([]*SVirtualMachine, 0, 2)
for i := 0; i < len(vms); i += 1 {
if all || !strings.HasPrefix(vms[i].Entity().Name, api.ESXI_IMAGE_CACHE_TMP_PREFIX) {
vmObj := NewVirtualMachine(dc.manager, &vms[i], dc)
if vms[i].Config != nil && vms[i].Config.Template {
templateVMs = append(templateVMs, vmObj)
vms := make([]*SVirtualMachine, 0, len(movms))
for i := range movms {
if all || !strings.HasPrefix(movms[i].Entity().Name, api.ESXI_IMAGE_CACHE_TMP_PREFIX) {
vm := NewVirtualMachine(dc.manager, &movms[i], dc)
// must
if vm == nil {
continue
}
if vmObj != nil {
retVms = append(retVms, vmObj)
}
vms = append(vms, vm)
}
}
return retVms, templateVMs, nil
return vms, nil
}
func (dc *SDatacenter) FetchVMs() ([]*SVirtualMachine, error) {
return dc.fetchVMs(property.Filter{})
}
func (dc *SDatacenter) FetchNoTemplateVMs() ([]*SVirtualMachine, error) {
filter := property.Filter{}
filter["config.template"] = false
return dc.fetchVMs(filter)
}
func (dc *SDatacenter) fetchVMs(filter property.Filter) ([]*SVirtualMachine, error) {
odc := dc.getObjectDatacenter()
root := odc.Reference()
m := view.NewManager(dc.manager.client.Client)
v, err := m.CreateContainerView(dc.manager.context, root, []string{"VirtualMachine"}, true)
if err != nil {
return nil, err
}
defer func() {
_ = v.Destroy(dc.manager.context)
}()
objs, err := v.Find(dc.manager.context, []string{"VirtualMachine"}, filter)
if err != nil {
return nil, err
}
vms, err := dc.fetchVms(objs, false)
return vms, err
}
func (dc *SDatacenter) FetchTemplateVMs() ([]*SVirtualMachine, error) {
filter := property.Filter{}
filter["config.template"] = true
return dc.fetchVMs(filter)
}
func (dc *SDatacenter) FetchTemplateVMById(id string) (*SVirtualMachine, error) {
filter := property.Filter{}
filter["config.template"] = true
filter["summary.config.uuid"] = id
vms, err := dc.fetchVMs(filter)
if err != nil {
return nil, err
}
if len(vms) == 0 {
return nil, errors.ErrNotFound
}
return vms[0], nil
}
func (dc *SDatacenter) FetchVMById(id string) (*SVirtualMachine, error) {
filter := property.Filter{}
filter["summary.config.uuid"] = id
vms, err := dc.fetchVMs(filter)
if err != nil {
return nil, err
}
if len(vms) == 0 {
return nil, errors.ErrNotFound
}
return vms[0], nil
}
func (dc *SDatacenter) fetchDatastores(datastoreRefs []types.ManagedObjectReference) ([]cloudprovider.ICloudStorage, error) {
@@ -381,24 +441,3 @@ func (dc *SDatacenter) GetTemplateVMs() ([]*SVirtualMachine, error) {
}
return templateVms, nil
}
func (dc *SDatacenter) GetTemplateVMById(id string) (*SVirtualMachine, error) {
id = dc.manager.getPrivateId(id)
hosts, err := dc.GetIHosts()
if err != nil {
return nil, errors.Wrap(err, "SDatacenter.GetIHosts")
}
for _, ihost := range hosts {
host := ihost.(*SHost)
tvms, err := host.GetTemplateVMs()
if err != nil {
return nil, errors.Wrap(err, "host.GetTemplateVMs")
}
for i := range tvms {
if tvms[i].GetGlobalId() == id {
return tvms[i], nil
}
}
}
return nil, cloudprovider.ErrNotFound
}

View File

@@ -158,18 +158,23 @@ func NewVNICDev(host *SHost, mac, driver string, vlanId int32, key, ctlKey, inde
var backing types.BaseVirtualDeviceBackingInfo
switch inet.(type) {
case *SDistributedVirtualPortgroup:
// net := inet.(*SDistributedVirtualPortgroup)
// port, err := net.FindPort()
//if err != nil {
// return nil, errors.Wrap(err, "net.FindPort")
// }
// if port == nil {
// return nil, fmt.Errorf("no active port for dvportgroup %q", net.GetName())
// }
net := inet.(*SDistributedVirtualPortgroup)
port, err := net.FindPort()
dvpg := net.getMODVPortgroup()
uuid, err := net.GetDVSUuid()
if err != nil {
return nil, errors.Wrap(err, "net.FindPort")
}
if port == nil {
return nil, errors.Error("no valid port on DVS, exhausted")
return nil, errors.Wrap(err, "GetDVSUuid")
}
portCon := types.DistributedVirtualSwitchPortConnection{
PortgroupKey: port.PortgroupKey,
SwitchUuid: port.DvsUuid,
PortKey: port.Key,
PortgroupKey: dvpg.Key,
SwitchUuid: uuid,
}
backing = &types.VirtualEthernetCardDistributedVirtualPortBackingInfo{Port: portCon}
case *SNetwork:

View File

@@ -18,7 +18,6 @@ import (
"context"
"fmt"
"regexp"
"sort"
"strings"
"time"
@@ -179,6 +178,7 @@ func (self *SHost) fetchVMs(all bool) error {
}
MAX_TRIES := 3
var vms []*SVirtualMachine
for tried := 0; tried < MAX_TRIES; tried += 1 {
hostVms := self.getHostSystem().Vm
if len(hostVms) == 0 {
@@ -186,15 +186,20 @@ func (self *SHost) fetchVMs(all bool) error {
return nil
}
vms, templatevms, err := dc.fetchVms(hostVms, all)
vms, err = dc.fetchVms(hostVms, all)
if err != nil {
log.Errorf("dc.fetchVms fail %s", err)
time.Sleep(time.Second)
self.Refresh()
continue
}
self.vms = vms
self.tempalteVMs = templatevms
}
for _, vm := range vms {
if vm.IsTemplate() {
self.tempalteVMs = append(self.tempalteVMs, vm)
} else {
self.vms = append(self.vms, vm)
}
}
return nil
}
@@ -679,14 +684,9 @@ func (self *SHost) CreateVM2(ctx context.Context, ds *SDatastore, params SCreate
if imageInfo.ImageType != cloudprovider.CachedImageTypeSystem {
return self.DoCreateVM(ctx, ds, params)
}
// get host
imgHost, err := self.manager.FindHostByIp(imageInfo.StorageCacheHostIp)
temvm, err := self.manager.SearchTemplateVM(imageInfo.ImageExternalId)
if err != nil {
return nil, errors.Wrap(err, "SEsxiClient.FindHostByIp")
}
temvm, err := imgHost.GetTemplateVMById(imageInfo.ImageExternalId)
if err != nil {
return nil, errors.Wrap(err, "SHost.GetTemplateVMById")
return nil, errors.Wrapf(err, "SEsxiClient.SearchTemplateVM for image %q", imageInfo.ImageExternalId)
}
return self.CloneVM(ctx, temvm, ds, params)
}
@@ -868,7 +868,11 @@ func (self *SHost) DoCreateVM(ctx context.Context, ds *SDatastore, params SCreat
return nil, errors.Wrap(err, "fail to fetch virtual machine just created")
}
return NewVirtualMachine(self.manager, &moVM, self.datacenter), nil
evm := NewVirtualMachine(self.manager, &moVM, self.datacenter)
if evm == nil {
return nil, errors.Error("create successfully but unable to NewVirtualMachine")
}
return evm, nil
}
func (host *SHost) CloneVM(ctx context.Context, from *SVirtualMachine, ds *SDatastore, params SCreateVMParam) (*SVirtualMachine, error) {
@@ -917,69 +921,59 @@ func (host *SHost) CloneVM(ctx context.Context, from *SVirtualMachine, ds *SData
}
}
// check scsi controller
var ctlKey int32
scsiDevs, err := from.FindController(ctx, "scsi")
if err != nil {
return nil, errors.Wrap(err, "SVirtualMachine.FindController")
}
if len(scsiDevs) == 0 {
key := from.FindMinDiffKey(1000)
driver := "pvscsi"
if host.isVersion50() {
driver = "scsi"
}
deviceChange = append(deviceChange, addDevSpec(NewSCSIDev(key, 100, driver)))
ctlKey = key
} else {
ctlKey = minDevKey(scsiDevs)
}
// change disk if set
newSizes := make([]int64, 0, len(from.vdisks))
if params.Disks != nil && len(params.Disks) > 0 {
var (
i int
disk SDiskInfo
)
// resize existed disk
for i, disk = range params.Disks {
if i == len(from.vdisks) {
break
if len(params.Disks) > 0 {
driver := params.Disks[0].Driver
if driver == "scsi" || driver == "pvscsi" {
scsiDevs, err := from.FindController(ctx, "scsi")
if err != nil {
return nil, errors.Wrap(err, "SVirtualMachine.FindController")
}
size := disk.Size
if size == 0 {
size = 30 * 1024
}
newSizes = append(newSizes, size)
}
// create new disk
if i == len(from.vdisks) {
// find same disk
var index int32
var key int32 = 2000
sameDisk := from.FindDiskByDriver("scsi", "pvscsi")
index += int32(len(sameDisk))
if index >= 7 {
index++
}
if len(sameDisk) > 0 {
key = minDiskKey(sameDisk)
}
for ; i < len(params.Disks); i++ {
size := params.Disks[i].Size
if size == 0 {
size = 30 * 1024
if len(scsiDevs) == 0 {
key := from.FindMinDiffKey(1000)
driver := "pvscsi"
if host.isVersion50() {
driver = "scsi"
}
uuid := params.Disks[i].DiskId
spec := addDevSpec(NewDiskDev(size, "", uuid, index, key, ctlKey, 0))
spec.FileOperation = "create"
deviceChange = append(deviceChange, spec)
deviceChange = append(deviceChange, addDevSpec(NewSCSIDev(key, 100, driver)))
}
} else {
ideDevs, err := from.FindController(ctx, "ide")
if err != nil {
return nil, errors.Wrap(err, "SVirtualMachine.FindController")
}
if len(ideDevs) == 0 {
// add ide driver
deviceChange = append(deviceChange, addDevSpec(NewIDEDev(200, 0)))
deviceChange = append(deviceChange, addDevSpec(NewIDEDev(200, 1)))
}
}
// resize system disk
sysDiskSize := params.Disks[0].Size
if sysDiskSize == 0 {
sysDiskSize = 30 * 1024
}
if int64(from.vdisks[0].GetDiskSizeMB()) != sysDiskSize {
vdisk := from.vdisks[0].getVirtualDisk()
vdisk.CapacityInKB = sysDiskSize * 1024
spec := &types.VirtualDeviceConfigSpec{}
spec.Operation = types.VirtualDeviceConfigSpecOperationEdit
spec.Device = vdisk
deviceChange = append(deviceChange, spec)
log.Infof("resize system disk: %dGB => %dGB", from.vdisks[0].GetDiskSizeMB()/1024, vdisk.CapacityInKB/1024/1024)
}
// remove extra disk
for i := 1; i < len(from.vdisks); i++ {
dev := from.vdisks[i].dev
spec := &types.VirtualDeviceConfigSpec{}
spec.Operation = types.VirtualDeviceConfigSpecOperationRemove
spec.Device = dev
spec.FileOperation = types.VirtualDeviceConfigSpecFileOperationDestroy
deviceChange = append(deviceChange, spec)
log.Debugf("remove disk, index: %d", i)
}
}
dc, err := host.GetDatacenter()
if err != nil {
return nil, errors.Wrapf(err, "SHost.GetDatacenter for host '%s'", host.GetId())
@@ -1038,13 +1032,22 @@ func (host *SHost) CloneVM(ctx context.Context, from *SVirtualMachine, ds *SData
return nil, errors.Wrap(err, "fail to fetch virtual machine just created")
}
// resize the disk
vm := NewVirtualMachine(host.manager, &moVM, host.datacenter)
sort.Sort(byDiskType(vm.vdisks))
for i, s := range newSizes {
err := vm.vdisks[i].Resize(ctx, s)
if vm == nil {
return nil, errors.Error("clone successfully but unable to NewVirtualMachine")
}
// add data disk
for i := 1; i < len(params.Disks); i++ {
size := params.Disks[i].Size
if size == 0 {
size = 30 * 1024
}
uuid := params.Disks[i].DiskId
driver := params.Disks[i].Driver
err := vm.CreateDisk(ctx, int(size), uuid, driver)
if err != nil {
log.Errorf("no.%d vdisk.Resize failed: %s", i, err.Error())
log.Errorf("unable to add No.%d disk for vm %s", i, vm.GetId())
return vm, nil
}
}
return vm, nil

View File

@@ -345,6 +345,54 @@ func (cli *SESXiClient) scanAllMObjects(props []string, dst interface{}) error {
return cli.scanMObjects(cli.client.ServiceContent.RootFolder, props, dst)
}
func (cli *SESXiClient) SearchTemplateVM(id string) (*SVirtualMachine, error) {
filter := property.Filter{}
filter["config.template"] = true
filter["summary.config.uuid"] = id
var movms []mo.VirtualMachine
err := cli.scanMObjectsWithFilter(cli.client.ServiceContent.RootFolder, VIRTUAL_MACHINE_PROPS, &movms, filter)
if err != nil {
return nil, err
}
if len(movms) == 0 {
return nil, errors.ErrNotFound
}
vm := NewVirtualMachine(cli, &movms[0], nil)
dc, err := vm.fetchDatacenter()
if err != nil {
return nil, errors.Wrap(err, "fetchDatacenter")
}
vm.datacenter = dc
return vm, nil
}
func (cli *SESXiClient) scanMObjectsWithFilter(folder types.ManagedObjectReference, props []string, dst interface{}, filter property.Filter) error {
dstValue := reflect.Indirect(reflect.ValueOf(dst))
dstType := dstValue.Type()
dstEleType := dstType.Elem()
resType := dstEleType.Name()
m := view.NewManager(cli.client.Client)
v, err := m.CreateContainerView(cli.context, folder, []string{resType}, true)
if err != nil {
return errors.Wrapf(err, "m.CreateContainerView %s", resType)
}
defer v.Destroy(cli.context)
err = v.RetrieveWithFilter(cli.context, []string{resType}, props, dst, filter)
if err != nil {
// hack
if strings.Contains(err.Error(), "object references is empty") {
return nil
}
return errors.Wrapf(err, "v.RetrieveWithFilter %s", resType)
}
return nil
}
func (cli *SESXiClient) scanMObjects(folder types.ManagedObjectReference, props []string, dst interface{}) error {
dstValue := reflect.Indirect(reflect.ValueOf(dst))
dstType := dstValue.Type()
@@ -536,7 +584,11 @@ func (cli *SESXiClient) FindVMByPrivateID(idstr string) (*SVirtualMachine, error
return nil, errors.Wrap(err, "reference2Object fail")
}
return NewVirtualMachine(cli, &vm, nil), nil
ret := NewVirtualMachine(cli, &vm, nil)
if ret == nil {
return nil, errors.Error("invalid vm")
}
return ret, nil
}
func (cli *SESXiClient) DoExtendDiskOnline(_vm *SVirtualMachine, _disk *SVirtualDisk, newSizeMb int64) error {

View File

@@ -43,7 +43,7 @@ const (
)
var NETWORK_PROPS = []string{"name", "parent", "summary", "host", "vm"}
var DVPORTGROUP_PROPS = []string{"name", "parent", "summary", "host", "vm", "config"}
var DVPORTGROUP_PROPS = []string{"name", "parent", "summary", "host", "vm", "config", "key"}
type SNetwork struct {
SManagedObject
@@ -175,6 +175,16 @@ func (net *SDistributedVirtualPortgroup) Uplink() bool {
return *dvpg.Config.Uplink
}
func (net *SDistributedVirtualPortgroup) GetDVSUuid() (string, error) {
dvgp := net.getMODVPortgroup()
var dvs mo.DistributedVirtualSwitch
err := net.manager.reference2Object(*dvgp.Config.DistributedVirtualSwitch, []string{"uuid"}, &dvs)
if err != nil {
return "", errors.Wrap(err, "reference2Object")
}
return dvs.Uuid, nil
}
func (net *SDistributedVirtualPortgroup) FindPort() (*types.DistributedVirtualPort, error) {
dvgp := net.getMODVPortgroup()
odvs := object.NewDistributedVirtualSwitch(net.manager.client.Client, *dvgp.Config.DistributedVirtualSwitch)

View File

@@ -27,28 +27,53 @@ import (
func init() {
type VirtualMachineListOptions struct {
HOSTIP string `help:"Host IP"`
Template bool `help:"Whether it is tempalte virtual machine"`
Datacenter string `help:"Datacenter"`
HostIP string `help:"HostIP"`
Template bool `help:"Whether it is tempalte virtual machine, default:false"`
}
shellutils.R(&VirtualMachineListOptions{}, "vm-list", "List vms of a host", func(cli *esxi.SESXiClient, args *VirtualMachineListOptions) error {
host, err := cli.FindHostByIp(args.HOSTIP)
if err != nil {
return err
}
if args.Template {
vms, err := host.GetTemplateVMs()
switch {
case len(args.HostIP) > 0:
host, err := cli.FindHostByIp(args.HostIP)
if err != nil {
return err
}
if args.Template {
vms, err := host.GetTemplateVMs()
if err != nil {
return err
}
printList(vms, []string{})
return nil
}
vms, err := host.GetIVMs2()
if err != nil {
return err
}
printList(vms, []string{})
return nil
case len(args.Datacenter) > 0:
dc, err := cli.FindDatacenterByMoId(args.Datacenter)
if err != nil {
return errors.Wrap(err, "FindDatacenterByMoId")
}
var vms []*esxi.SVirtualMachine
if args.Template {
vms, err = dc.FetchTemplateVMs()
if err != nil {
return errors.Wrap(err, "FetchTemplateVMs")
}
} else {
vms, err = dc.FetchNoTemplateVMs()
if err != nil {
return errors.Wrap(err, "FetchNoTemplateVMs")
}
}
printList(vms, []string{})
return nil
default:
return fmt.Errorf("Both Datacenter and HostIP cannot be empty")
}
vms, err := host.GetIVMs2()
if err != nil {
return err
}
printList(vms, []string{})
return nil
})
type VirtualMachineCloneOptions struct {
@@ -91,24 +116,45 @@ func init() {
})
type VirtualMachineShowOptions struct {
HOSTIP string `help:"Host IP"`
VMID string `help:"VM ID"`
Template bool
Datacenter string `help:"Datacenter"`
HostIP string `help:"Host IP"`
VMID string `help:"VM ID"`
}
getVM := func(cli *esxi.SESXiClient, args *VirtualMachineShowOptions) (*esxi.SVirtualMachine, error) {
var vm *esxi.SVirtualMachine
switch {
case len(args.HostIP) > 0:
host, err := cli.FindHostByIp(args.HostIP)
if err != nil {
return nil, errors.Wrap(err, "FindHostByIp")
}
ivm, err := host.GetIVMById(args.VMID)
if err != nil && errors.Cause(err) != errors.ErrNotFound {
return nil, err
}
if err != nil {
vm, err = host.GetTemplateVMById(args.VMID)
if err != nil {
return nil, errors.Wrap(err, "GetTemplateVMById")
}
}
vm = ivm.(*esxi.SVirtualMachine)
case len(args.Datacenter) > 0:
dc, err := cli.FindDatacenterByMoId(args.Datacenter)
if err != nil {
return nil, errors.Wrap(err, "FindDatacenterByMoId")
}
vm, err = dc.FetchVMById(args.VMID)
if err != nil {
return nil, errors.Wrap(err, "FetchVMById")
}
default:
return nil, fmt.Errorf("Both Datacenter and HostIP cannot be empty")
}
return vm, nil
}
shellutils.R(&VirtualMachineShowOptions{}, "vm-show", "Show vm details", func(cli *esxi.SESXiClient, args *VirtualMachineShowOptions) error {
host, err := cli.FindHostByIp(args.HOSTIP)
if err != nil {
return err
}
if args.Template {
vm, err := host.GetTemplateVMById(args.VMID)
if err != nil {
return err
}
printObject(vm)
return nil
}
vm, err := host.GetIVMById(args.VMID)
vm, err := getVM(cli, args)
if err != nil {
return err
}
@@ -117,11 +163,7 @@ func init() {
})
shellutils.R(&VirtualMachineShowOptions{}, "vm-nics", "Show vm nics details", func(cli *esxi.SESXiClient, args *VirtualMachineShowOptions) error {
host, err := cli.FindHostByIp(args.HOSTIP)
if err != nil {
return err
}
vm, err := host.GetIVMById(args.VMID)
vm, err := getVM(cli, args)
if err != nil {
return err
}
@@ -134,11 +176,7 @@ func init() {
})
shellutils.R(&VirtualMachineShowOptions{}, "vm-disks", "Show vm disks details", func(cli *esxi.SESXiClient, args *VirtualMachineShowOptions) error {
host, err := cli.FindHostByIp(args.HOSTIP)
if err != nil {
return err
}
vm, err := host.GetIVMById(args.VMID)
vm, err := getVM(cli, args)
if err != nil {
return err
}
@@ -151,17 +189,12 @@ func init() {
})
type VirtualMachineDiskResizeOptions struct {
HOSTIP string `help:"host ip"`
VMID string `help:"virtual machine UUID"`
DISKIDX int `help:"disk index"`
SIZEGB int64 `help:"new size of disk"`
VirtualMachineShowOptions
DISKIDX int `help:"disk index"`
SIZEGB int64 `help:"new size of disk"`
}
shellutils.R(&VirtualMachineDiskResizeOptions{}, "vm-disk-resize", "Resize a vm disk", func(cli *esxi.SESXiClient, args *VirtualMachineDiskResizeOptions) error {
host, err := cli.FindHostByIp(args.HOSTIP)
if err != nil {
return err
}
vm, err := host.GetIVMById(args.VMID)
vm, err := getVM(cli, &args.VirtualMachineShowOptions)
if err != nil {
return err
}
@@ -178,11 +211,7 @@ func init() {
})
shellutils.R(&VirtualMachineShowOptions{}, "vm-vnc", "Show vm VNC details", func(cli *esxi.SESXiClient, args *VirtualMachineShowOptions) error {
host, err := cli.FindHostByIp(args.HOSTIP)
if err != nil {
return err
}
vm, err := host.GetIVMById(args.VMID)
vm, err := getVM(cli, args)
if err != nil {
return err
}
@@ -195,15 +224,11 @@ func init() {
})
shellutils.R(&VirtualMachineShowOptions{}, "vm-file-status", "Show vm files details", func(cli *esxi.SESXiClient, args *VirtualMachineShowOptions) error {
host, err := cli.FindHostByIp(args.HOSTIP)
vm, err := getVM(cli, args)
if err != nil {
return err
}
vm, err := host.GetIVMById(args.VMID)
if err != nil {
return err
}
err = vm.(*esxi.SVirtualMachine).CheckFileInfo(context.Background())
err = vm.CheckFileInfo(context.Background())
if err != nil {
return err
}

View File

@@ -259,7 +259,14 @@ func (self *SDatastore) getVMs() ([]cloudprovider.ICloudVM, error) {
if len(vms) == 0 {
return nil, nil
}
ret, _, err := dc.fetchVms(vms, false)
svms, err := dc.fetchVms(vms, false)
if err != nil {
return nil, err
}
ret := make([]cloudprovider.ICloudVM, len(svms))
for i := range svms {
ret[i] = svms[i]
}
return ret, err
}

View File

@@ -22,6 +22,7 @@ import (
"strings"
"time"
"github.com/vmware/govmomi/nfc"
"github.com/vmware/govmomi/object"
"github.com/vmware/govmomi/vim25/mo"
"github.com/vmware/govmomi/vim25/soap"
@@ -74,6 +75,7 @@ func NewVirtualMachine(manager *SESXiClient, vm *mo.VirtualMachine, dc *SDatacen
svm := &SVirtualMachine{SManagedObject: newManagedObject(manager, vm, dc)}
err := svm.fetchHardwareInfo()
if err != nil {
log.Errorf("NewVirtualMachine: %v", err)
return nil
}
return svm
@@ -722,11 +724,16 @@ func (self *SVirtualMachine) fetchHardwareInfo() error {
}
if moVM == nil || moVM.Config == nil || moVM.Config.Hardware.Device == nil {
return errors.Error("invalid vm config")
return fmt.Errorf("invalid vm")
}
for i := 0; i < len(moVM.Config.Hardware.Device); i += 1 {
dev := moVM.Config.Hardware.Device[i]
// sort devices via their Key
devices := moVM.Config.Hardware.Device
sort.Slice(devices, func(i, j int) bool {
return devices[i].GetVirtualDevice().Key < devices[j].GetVirtualDevice().Key
})
for i := 0; i < len(devices); i += 1 {
dev := devices[i]
devType := reflect.Indirect(reflect.ValueOf(dev)).Type()
etherType := reflect.TypeOf((*types.VirtualEthernetCard)(nil)).Elem()
@@ -737,7 +744,7 @@ func (self *SVirtualMachine) fetchHardwareInfo() error {
if reflectutils.StructContains(devType, etherType) {
self.vnics = append(self.vnics, NewVirtualNIC(self, dev, len(self.vnics)))
} else if reflectutils.StructContains(devType, diskType) {
self.vdisks = append(self.vdisks, NewVirtualDisk(self, dev, len(self.vnics)))
self.vdisks = append(self.vdisks, NewVirtualDisk(self, dev, len(self.vdisks)))
} else if reflectutils.StructContains(devType, vgaType) {
self.vga = NewVirtualVGA(self, dev, 0)
} else if reflectutils.StructContains(devType, cdromType) {
@@ -1166,8 +1173,19 @@ func (self *SVirtualMachine) ExportTemplate(ctx context.Context, idx int, diskPa
lr := newLeaseLogger("download vmdk", 5)
lr.Log()
defer lr.End()
// filter vmdk item
vmdkItems := make([]nfc.FileItem, 0, len(info.Items)/2)
for i := range info.Items {
if strings.HasSuffix(info.Items[i].Path, ".vmdk") {
vmdkItems = append(vmdkItems, info.Items[i])
} else {
log.Infof("item.Path does not end in '.vmdk': %#v", info.Items[i])
}
}
log.Debugf("download to %s start...", diskPath)
err = lease.DownloadFile(ctx, diskPath, info.Items[idx], soap.Download{Progress: lr})
err = lease.DownloadFile(ctx, diskPath, vmdkItems[idx], soap.Download{Progress: lr})
if err != nil {
return errors.Wrap(err, "lease.DownloadFile")
}
@@ -1218,3 +1236,8 @@ func (self *SVirtualMachine) FindMinDiffKey(limit int32) int32 {
}
return limit
}
func (self *SVirtualMachine) IsTemplate() bool {
movm := self.getVirtualMachine()
return movm.Config != nil && movm.Config.Template
}

View File

@@ -529,7 +529,7 @@ func getDiskInfo(disk string) (cloudprovider.SDiskInfo, error) {
result := cloudprovider.SDiskInfo{}
diskInfo := strings.Split(disk, ":")
for _, d := range diskInfo {
if utils.IsInStringArray(d, []string{api.STORAGE_GOOGLE_PD_STANDARD, api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_LOCAL_SSD}) {
if utils.IsInStringArray(d, []string{api.STORAGE_GOOGLE_PD_STANDARD, api.STORAGE_GOOGLE_PD_SSD, api.STORAGE_GOOGLE_LOCAL_SSD, api.STORAGE_GOOGLE_PD_BALANCED}) {
result.StorageType = d
} else if memSize, err := fileutils.GetSizeMb(d, 'M', 1024); err == nil {
result.SizeGB = memSize >> 10

View File

@@ -46,6 +46,33 @@ type SBaseManager struct {
debug bool
}
type sThrottlingThreshold struct {
locked bool
lockTime time.Time
}
func (t *sThrottlingThreshold) CheckingLock() {
if !t.locked {
return
}
for {
if t.lockTime.Sub(time.Now()).Seconds() < 0 {
return
}
log.Debugf("throttling threshold has been reached. release at %s", t.lockTime)
time.Sleep(5 * time.Second)
}
}
func (t *sThrottlingThreshold) Lock() {
// 锁定至少15秒
t.locked = true
t.lockTime = time.Now().Add(15 * time.Second)
}
var ThrottlingLock = sThrottlingThreshold{locked: false, lockTime: time.Time{}}
func NewBaseManager2(signer auth.Signer, debug bool, requesthk IRequestHook) SBaseManager {
return SBaseManager{
signer: signer,
@@ -154,6 +181,7 @@ func (ce *HuaweiClientError) ParseErrorFromJsonResponse(statusCode int, body jso
}
func (self *SBaseManager) jsonRequest(request requests.IRequest) (http.Header, jsonutils.JSONObject, error) {
ThrottlingLock.CheckingLock()
ctx := context.Background()
// hook request
if self.requestHook != nil {
@@ -182,7 +210,6 @@ func (self *SBaseManager) jsonRequest(request requests.IRequest) (http.Header, j
req := httputils.NewJsonRequest(httputils.THttpMethod(request.GetMethod()), request.BuildUrl(), jsonBody)
req.SetHeader(header)
resp := &HuaweiClientError{}
// 发送 request。todo: 支持debug
const MAX_RETRY = 3
retry := MAX_RETRY
for {
@@ -195,11 +222,16 @@ func (self *SBaseManager) jsonRequest(request requests.IRequest) (http.Header, j
switch err := e.(type) {
case *HuaweiClientError:
if (err.Code == 499 || err.Code == 429) && retry > 0 && request.GetMethod() == "GET" {
if err.Code == 499 && retry > 0 && request.GetMethod() == "GET" {
retry -= 1
time.Sleep(3 * time.Second * time.Duration(MAX_RETRY-retry))
} else if (err.Code == 404 || strings.Contains(err.Details, "could not be found") || strings.Contains(err.Details, "does not exist")) && request.GetMethod() != "POST" {
return h, b, errors.Wrap(cloudprovider.ErrNotFound, err.Error())
} else if err.Code == 429 && retry > 0 {
// 当前请求过多。
ThrottlingLock.Lock()
retry -= 1
time.Sleep(15 * time.Second)
} else {
return h, b, e
}

View File

@@ -21,6 +21,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/cloudprovider"
@@ -180,7 +181,7 @@ func doListPart(doList listFunc, queries map[string]string, result interface{})
func DoGet(doGet getFunc, id string, queries map[string]string, result interface{}) error {
if len(id) == 0 {
resultType := reflect.Indirect(reflect.ValueOf(result)).Type()
return fmt.Errorf(" Get %s id should not be empty", resultType.Name())
return errors.Wrap(cloudprovider.ErrNotFound, fmt.Sprintf(" Get %s id should not be empty", resultType.Name()))
}
ret, err := doGet(id, queries)

View File

@@ -18,6 +18,7 @@ import (
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudprovider"
@@ -125,14 +126,20 @@ func (self *SZone) GetIHostById(id string) (cloudprovider.ICloudHost, error) {
func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
return self.istorages, nil
}
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
for i := 0; i < len(self.istorages); i += 1 {
if self.istorages[i].GetGlobalId() == id {

View File

@@ -23,6 +23,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
api "yunion.io/x/onecloud/pkg/apis/compute"
@@ -187,13 +188,28 @@ func (self *SLoadbalancer) CreateILoadBalancerListener(ctx context.Context, list
hc,
cert)
}
if err != nil {
return nil, err
}
time.Sleep(3 * time.Second)
return self.GetILoadBalancerListenerById(listenId)
var lblis cloudprovider.ICloudLoadbalancerListener
err = cloudprovider.Wait(3*time.Second, 30*time.Second, func() (bool, error) {
lblis, err = self.GetILoadBalancerListenerById(listenId)
if err != nil {
if errors.Cause(err) != cloudprovider.ErrNotFound {
return false, err
} else {
return false, nil
}
}
return true, nil
})
if err != nil {
return nil, errors.Wrap(err, "GetILoadBalancerListenerById.Wait")
}
return lblis, nil
}
func (self *SLoadbalancer) GetILoadBalancerListenerById(listenerId string) (cloudprovider.ICloudLoadbalancerListener, error) {

View File

@@ -79,6 +79,9 @@ func (client *SQcloudClient) GetProjects() ([]SProject, error) {
projects := []SProject{}
params := map[string]string{"allList": "1"}
resp, err := client.accountRequestRequest("DescribeProject", params)
if err != nil {
return nil, errors.Wrapf(err, "DescribeProject")
}
err = resp.Unmarshal(&projects)
if err != nil {
return nil, errors.Wrap(err, "resp.Unmarshal")

View File

@@ -170,7 +170,10 @@ func (self *SZone) fetchStorages() error {
func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
return self.istorages, nil
}
@@ -221,7 +224,10 @@ func (self *SZone) getStorageByCategory(category string) (*SStorage, error) {
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
for i := 0; i < len(self.istorages); i += 1 {
if self.istorages[i].GetGlobalId() == id {

View File

@@ -18,6 +18,7 @@ import (
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudprovider"
@@ -127,14 +128,20 @@ func (self *SZone) GetIHostById(id string) (cloudprovider.ICloudHost, error) {
func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
return self.istorages, nil
}
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
if self.istorages == nil {
self.fetchStorages()
err := self.fetchStorages()
if err != nil {
return nil, errors.Wrapf(err, "fetchStorages")
}
}
for i := 0; i < len(self.istorages); i += 1 {
if self.istorages[i].GetGlobalId() == id {

View File

@@ -147,15 +147,22 @@ func (keeper *OVNNorthboundKeeper) ClaimNetwork(ctx context.Context, network *ag
"0.0.0.0/0", network.GuestGateway,
}
mtu -= 58
const (
leaseTime = 86400 * 365 * 3
renewTime = 86400
rebindTime = 86400 * 3
)
dhcpopts := &ovnutil.DHCPOptions{
Cidr: fmt.Sprintf("%s/%d", network.GuestIpStart, network.GuestIpMask),
Options: map[string]string{
"server_id": network.GuestGateway,
"server_mac": dhcpMac,
"lease_time": fmt.Sprintf("%d", 86400),
"router": network.GuestGateway,
"classless_static_route": fmt.Sprintf("{%s}", strings.Join(routes, ",")),
"mtu": fmt.Sprintf("%d", mtu),
"lease_time": fmt.Sprintf("%d", leaseTime),
"T1": fmt.Sprintf("%d", renewTime),
"T2": fmt.Sprintf("%d", rebindTime),
},
ExternalIds: map[string]string{
externalKeyOcRef: network.Id,

View File

@@ -129,8 +129,11 @@ func (c *SSHtoolSol) GetCommand() *exec.Cmd {
args := []string{
o.Options.SshpassToolPath, "-p", c.password,
o.Options.SshToolPath, "-p", fmt.Sprintf("%d", c.Port), fmt.Sprintf("%s@%s", c.username, c.IP),
"-oGlobalKnownHostsFile=/dev/null", "-oUserKnownHostsFile=/dev/null", "-oStrictHostKeyChecking=no",
"-oPreferredAuthentications=password", "-oPubkeyAuthentication=no", //密码登录时,避免搜寻秘钥登录
"-oGlobalKnownHostsFile=/dev/null",
"-oUserKnownHostsFile=/dev/null",
"-oStrictHostKeyChecking=no",
"-oPreferredAuthentications=password,keyboard-interactive",
"-oPubkeyAuthentication=no", // 密码登录时,避免搜寻秘钥登录
"-oNumberOfPasswordPrompts=1",
}
cmd := exec.Command(args[0], args[1:]...)

View File

@@ -26,6 +26,7 @@ import (
"github.com/gorilla/mux"
"yunion.io/x/log"
"yunion.io/x/pkg/util/signalutils"
api "yunion.io/x/onecloud/pkg/apis/webconsole"
"yunion.io/x/onecloud/pkg/appsrv"
@@ -66,9 +67,15 @@ func StartService() {
common_options.StartOptionManager(opts, opts.ConfigSyncPeriodSeconds, api.SERVICE_TYPE, api.SERVICE_VERSION, o.OnOptionsChange)
registerSigTraps()
start()
}
func registerSigTraps() {
signalutils.SetDumpStackSignal()
signalutils.StartTrap()
}
func start() {
baseOpts := &o.Options.BaseOptions
// commonOpts := &o.Options.CommonOptions

View File

@@ -102,6 +102,10 @@ func (p *Pty) Resize(size *pty.Winsize) {
func (p *Pty) Stop() (err error) {
var errs []error
defer func() {
p.Cmd, p.Pty = nil, nil
}()
defer func() {
err = errors.NewAggregate(errs)
}()
@@ -139,8 +143,5 @@ func (p *Pty) Stop() (err error) {
}
}()
defer func() {
p.Cmd, p.Pty = nil, nil
}()
return
}