fix: sync error (#19483)

Co-authored-by: Qiu Jian <qiujian@yunionyun.com>
This commit is contained in:
Jian Qiu
2024-02-10 18:00:28 +08:00
committed by GitHub
co-authored by Qiu Jian
parent 0f2e11591f
commit 97fc5258df
2 changed files with 23 additions and 46 deletions
+5 -44
View File
@@ -57,50 +57,11 @@ func GetModels(opts *GetModelsOptions) error {
tstr := timeutils.MysqlTime(time)
return fmt.Sprintf("updated_at.ge('%s')", tstr)
}
setNextListParams := func(params *jsonutils.JSONDict, lastUpdatedAt time.Time, lastResult *printutils.ListResult) (time.Time, error) {
setNextListParams := func(params *jsonutils.JSONDict, lastResult *printutils.ListResult) error {
// NOTE: the updated_at field has second-level resolution.
// If they all have the same date...
var max time.Time
nmax := 0
n := len(lastResult.Data)
// find out the max updated_at date in the result set, and how
// many in the set has this date
for i := n - 1; i >= 0; i-- {
j := lastResult.Data[i]
updatedAt, err := j.GetTime("updated_at")
if err != nil {
log.Warningf("%s: updated_at field: %s, %s",
manKeyPlural, err, j.String())
continue
}
if max.IsZero() {
max = updatedAt
}
if max.Equal(updatedAt) {
nmax += 1
}
}
// error if we do not have valid date
if max.IsZero() {
return time.Time{}, fmt.Errorf("%s: cannot find next updated_at after '%q'",
manKeyPlural, lastUpdatedAt)
}
var newTime time.Time
var newOffset int
// if not all updated_at date are the same, then we can
// continue to the next age.
if nmax < n || (!max.Equal(lastUpdatedAt) && !max.Equal(PseudoZeroTime)) {
newTime = max
newOffset = nmax
} else {
newTime = lastUpdatedAt
newOffset = lastResult.Offset + n
}
params.Set("filter.0", jsonutils.NewString(minUpdatedAtFilter(newTime)))
params.Set("offset", jsonutils.NewInt(int64(newOffset)))
return newTime, nil
params.Set("offset", jsonutils.NewInt(int64(lastResult.Offset+len(lastResult.Data))))
return nil
}
listOptions := options.BaseListOptions{
@@ -112,7 +73,7 @@ func GetModels(opts *GetModelsOptions) error {
Filter: []string{
minUpdatedAtFilter(minUpdatedAt), // order matters, filter.0
},
OrderBy: []string{"updated_at"},
OrderBy: []string{"updated_at", "created_at", "id"},
Order: "asc",
Limit: options.Int(opts.BatchListSize),
Offset: options.Int(0),
@@ -161,7 +122,7 @@ func GetModels(opts *GetModelsOptions) error {
if listResult.Offset+len(listResult.Data) >= listResult.Total {
break
}
minUpdatedAt, err = setNextListParams(params, minUpdatedAt, listResult)
err = setNextListParams(params, listResult)
if err != nil {
return fmt.Errorf("%s: %s", manKeyPlural, err)
}
+18 -2
View File
@@ -26,6 +26,7 @@ import (
"yunion.io/x/ovsdb/types"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/netutils"
"yunion.io/x/pkg/util/regutils"
"yunion.io/x/pkg/utils"
apis "yunion.io/x/onecloud/pkg/apis/compute"
@@ -436,6 +437,21 @@ func (keeper *OVNNorthboundKeeper) ClaimVpcEipgw(ctx context.Context, vpc *agent
return keeper.cli.Must(ctx, "ClaimVpcEipgw", args)
}
func formatNtpServers(srvs string) string {
srv := make([]string, 0)
for _, part := range strings.Split(srvs, ",") {
part = strings.TrimSpace(part)
if len(part) > 0 {
if regutils.MatchIPAddr(part) {
srv = append(srv, part)
} else {
srv = append(srv, "\""+part+"\"")
}
}
}
return strings.Join(srv, ",")
}
func generateDhcpOptions(ctx context.Context, guestnetwork *agentmodels.Guestnetwork, opts *options.Options) *ovn_nb.DHCPOptions {
var (
network = guestnetwork.Network
@@ -527,7 +543,7 @@ func generateDhcpOptions(ctx context.Context, guestnetwork *agentmodels.Guestnet
}
if len(ntpSrvs) > 0 {
// bug on OVN, should not use ntp server: QiuJian
// dhcpopts.Options["ntp_server"] = "{" + ntpSrvs + "}"
dhcpopts.Options["ntp_server"] = "{" + formatNtpServers(ntpSrvs) + "}"
}
}
return dhcpopts
@@ -964,7 +980,7 @@ func (keeper *OVNNorthboundKeeper) ClaimVpcGuestDnsRecords(ctx context.Context,
for _, guestnetwork := range network.Guestnetworks {
if guest := guestnetwork.Guest; guest != nil {
var (
name = guest.Name
name = guest.Hostname
ip = guestnetwork.IpAddr
)
grs[name] = append(grs[name], ip)