Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,13 @@ func (c *ProductColl) GetCollectionName() string {

func (c *ProductColl) EnsureIndex(ctx context.Context) error {
mod := []mongo.IndexModel{
{
Keys: bson.D{
bson.E{Key: "product_name", Value: 1},
bson.E{Key: "env_name", Value: 1},
},
Options: options.Index().SetUnique(true),
},
{
Keys: bson.D{
bson.E{Key: "env_name", Value: 1},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -409,42 +409,7 @@ func (c *ProductionServiceColl) ListMaxRevisionsByProductWithFilter(productName
}

func (c *ProductionServiceColl) ListServicesWithSRevision(opt *SvcRevisionListOption) ([]*models.Service, error) {
productMatch := bson.M{}
productMatch["product_name"] = opt.ProductName

var serviceMatch bson.A
for _, sr := range opt.ServiceRevisions {
serviceMatch = append(serviceMatch, bson.M{
"service_name": sr.ServiceName,
"revision": sr.Revision,
})
}

pipeline := []bson.M{
{
"$match": productMatch,
},
}
if len(opt.ServiceRevisions) > 0 {
pipeline = append(pipeline, bson.M{
"$match": bson.M{
"$or": serviceMatch,
},
})
} else {
return []*models.Service{}, nil
}

cursor, err := c.Aggregate(context.TODO(), pipeline)
if err != nil {
return nil, err
}

res := make([]*models.Service, 0)
if err := cursor.All(context.TODO(), &res); err != nil {
return nil, err
}
return res, err
return listServicesWithSRevision(c.Collection, opt)
}

func (c *ProductionServiceColl) TransferServiceSource(productName, serviceName, source, newSource, username, yaml string) error {
Expand Down
61 changes: 37 additions & 24 deletions pkg/microservice/aslan/core/common/repository/mongodb/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -575,42 +575,55 @@ func (c *ServiceColl) ListAllRevisions() ([]*models.Service, error) {
}

func (c *ServiceColl) ListServicesWithSRevision(opt *SvcRevisionListOption) ([]*models.Service, error) {
productMatch := bson.M{}
productMatch["product_name"] = opt.ProductName

var serviceMatch bson.A
for _, sr := range opt.ServiceRevisions {
serviceMatch = append(serviceMatch, bson.M{
"service_name": sr.ServiceName,
"revision": sr.Revision,
})
}
return listServicesWithSRevision(c.Collection, opt)
}

pipeline := []bson.M{
{
"$match": productMatch,
},
}
if len(opt.ServiceRevisions) > 0 {
pipeline = append(pipeline, bson.M{
"$match": bson.M{
"$or": serviceMatch,
},
})
} else {
func listServicesWithSRevision(collection *mongo.Collection, opt *SvcRevisionListOption) ([]*models.Service, error) {
filter, ok := serviceRevisionListFilter(opt)
if !ok {
return []*models.Service{}, nil
}

cursor, err := c.Aggregate(context.TODO(), pipeline)
cursor, err := collection.Find(context.TODO(), filter)
if err != nil {
return nil, err
}
defer cursor.Close(context.TODO())

res := make([]*models.Service, 0)
if err := cursor.All(context.TODO(), &res); err != nil {
return nil, err
}
return res, err
return res, nil
}

func serviceRevisionListFilter(opt *SvcRevisionListOption) (bson.M, bool) {
if len(opt.ServiceRevisions) == 0 {
return nil, false
}
revisionsByService := make(map[string]map[int64]struct{})
for _, sr := range opt.ServiceRevisions {
if revisionsByService[sr.ServiceName] == nil {
revisionsByService[sr.ServiceName] = make(map[int64]struct{})
}
revisionsByService[sr.ServiceName][sr.Revision] = struct{}{}
}

serviceMatch := make(bson.A, 0, len(revisionsByService))
for serviceName, revisionSet := range revisionsByService {
revisions := make([]int64, 0, len(revisionSet))
for revision := range revisionSet {
revisions = append(revisions, revision)
}
serviceMatch = append(serviceMatch, bson.M{
"service_name": serviceName,
"revision": bson.M{"$in": revisions},
})
}
return bson.M{
"product_name": opt.ProductName,
"$or": serviceMatch,
}, true
}

func (c *ServiceColl) ListMaxRevisionsByProject(serviceName, serviceType string) ([]*models.Service, error) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,38 @@ func (c *ServiceModuleColl) ListByServiceRevision(ctx context.Context, projectNa
}))
}

// ListByServiceRevisions returns all manual records in a project together
// with the visible auto-discovered records for the requested service
// revisions. It is the batch equivalent of ListByServiceRevision and keeps
// the same deterministic ordering required by the merge logic.
func (c *ServiceModuleColl) ListByServiceRevisions(ctx context.Context, projectName string, serviceRevisions map[string][]int64) ([]*models.ServiceModule, error) {
autoMatches := make(bson.A, 0, len(serviceRevisions))
for serviceName, revisions := range serviceRevisions {
if serviceName == "" || len(revisions) == 0 {
continue
}
autoMatches = append(autoMatches, bson.M{
"service_name": serviceName,
"revision_bound": bson.M{"$in": revisions},
})
}

moduleMatches := bson.A{bson.M{"is_manual": true}}
if len(autoMatches) > 0 {
moduleMatches = append(moduleMatches, bson.M{
"is_manual": false,
"ignored": bson.M{"$ne": true},
"$or": autoMatches,
})
}

query := bson.M{
"project_name": projectName,
"$or": moduleMatches,
}
return c.findAll(ctx, query, options.Find().SetSort(bson.D{{Key: "create_time", Value: 1}, {Key: "_id", Value: 1}}))
}

// ListManual returns every manual record for a service. Used by the manual-
// module CRUD API to list user-declared modules independently of any revision.
func (c *ServiceModuleColl) ListManual(ctx context.Context, projectName, serviceName string) ([]*models.ServiceModule, error) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,63 @@ type ModuleConflict struct {
Shadowed []*models.ServiceModule
}

// ServiceModuleSnapshot holds modules preloaded for one request.
type ServiceModuleSnapshot struct {
Manual map[string][]*models.Container
Resolved map[string]map[int64][]*models.Container
RecordCount int
}

// LoadServiceModuleSnapshot loads all records required for serviceRevisions
// in one MongoDB find and prepares the same merge results as
// ResolveServiceModules.
func LoadServiceModuleSnapshot(ctx context.Context, projectName string, production bool, serviceRevisions map[string][]int64) (*ServiceModuleSnapshot, error) {
records, err := pickServiceModuleColl(production).ListByServiceRevisions(ctx, projectName, serviceRevisions)
if err != nil {
return nil, err
}
return newServiceModuleSnapshot(records, serviceRevisions), nil
}

func newServiceModuleSnapshot(records []*models.ServiceModule, serviceRevisions map[string][]int64) *ServiceModuleSnapshot {
snapshot := &ServiceModuleSnapshot{
Manual: make(map[string][]*models.Container),
Resolved: make(map[string]map[int64][]*models.Container),
RecordCount: len(records),
}
recordsByService := make(map[string][]*models.ServiceModule)
for _, record := range records {
if record == nil {
continue
}
recordsByService[record.ServiceName] = append(recordsByService[record.ServiceName], record)
if record.IsManual {
snapshot.Manual[record.ServiceName] = append(snapshot.Manual[record.ServiceName], &models.Container{
Name: record.Name, Type: record.Type, Image: record.Image, ImageName: record.ImageName, ImagePath: record.ImagePath,
})
}
}

for serviceName, revisions := range serviceRevisions {
snapshot.Resolved[serviceName] = make(map[int64][]*models.Container)
for _, revision := range revisions {
selected := make([]*models.ServiceModule, 0)
for _, record := range recordsByService[serviceName] {
// ListByServiceRevision historically excludes ignored records for
// both branches, while the standalone manual listing does not.
if record.Ignored {
continue
}
if record.IsManual || record.RevisionBound == revision {
selected = append(selected, record)
}
}
snapshot.Resolved[serviceName][revision] = mergeServiceModules(selected)
}
}
return snapshot
}

// ResolveServiceModules returns the merged module list for one (project,
// service, revision) plus any name conflicts. Production picks the right
// underlying collection (service_module vs production_service_module).
Expand Down
53 changes: 37 additions & 16 deletions pkg/microservice/aslan/core/common/service/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -1100,26 +1100,42 @@ func ListServicesInEnv(envName, productName string, newSvcKVsMap map[string][]*c
return BuildServiceInfoInEnv(env, latestSvcs, newSvcKVsMap, log)
}

// BuildServiceInfoOptions supplies data already loaded by a request-scoped caller.
type BuildServiceInfoOptions struct {
Project *template.Product
ProductTemplateServices []*commonmodels.Service
Modules *repository.ServiceModuleSnapshot
}

// @fixme newSvcKVsMap is old struct kv map, which are the kv are from deploy job config
// may need to be removed, or use new kv struct
// helm values need to be refactored
func BuildServiceInfoInEnv(productInfo *commonmodels.Product, templateSvcs []*commonmodels.Service, newSvcKVsMap map[string][]*commonmodels.ServiceKeyVal, log *zap.SugaredLogger) (*EnvServices, error) {
func BuildServiceInfoInEnv(productInfo *commonmodels.Product, templateSvcs []*commonmodels.Service, newSvcKVsMap map[string][]*commonmodels.ServiceKeyVal, log *zap.SugaredLogger, preload ...*BuildServiceInfoOptions) (*EnvServices, error) {
productName, envName := productInfo.ProductName, productInfo.EnvName
ret := &EnvServices{
ProductName: productName,
EnvName: envName,
Services: make([]*EnvService, 0),
}

project, err := templaterepo.NewProductColl().Find(productInfo.ProductName)
if err != nil {
return nil, e.ErrGetService.AddDesc(fmt.Sprintf("failed to find project %s, err: %v", productInfo.ProductName, err))
var project *template.Product
var productTemplateSvcs []*commonmodels.Service
var moduleSnapshot *repository.ServiceModuleSnapshot
var err error
if len(preload) > 0 && preload[0] != nil {
project = preload[0].Project
productTemplateSvcs = preload[0].ProductTemplateServices
moduleSnapshot = preload[0].Modules
} else {
project, err = templaterepo.NewProductColl().Find(productInfo.ProductName)
if err != nil {
return nil, e.ErrGetService.AddDesc(fmt.Sprintf("failed to find project %s, err: %v", productName, err))
}
productTemplateSvcs, err = commonutil.GetProductUsedTemplateSvcs(productInfo)
if err != nil {
return nil, e.ErrGetService.AddErr(errors.Wrapf(err, "failed to find product template services for env %s:%s", productName, envName))
}
}

productTemplateSvcs, err := commonutil.GetProductUsedTemplateSvcs(productInfo)
if err != nil {
return nil, e.ErrGetService.AddErr(errors.Wrapf(err, "failed to find product template services for env %s:%s", productName, envName))
}
productTemplateSvcMap := make(map[string]*commonmodels.Service)
for _, svc := range productTemplateSvcs {
productTemplateSvcMap[svc.ServiceName] = svc
Expand All @@ -1131,11 +1147,15 @@ func BuildServiceInfoInEnv(productInfo *commonmodels.Product, templateSvcs []*co
templateSvcMap[svc.ServiceName] = svc

svcModulesMap[svc.ServiceName] = make(map[string]*commonmodels.Container)
// Service.Containers is no longer persisted — read modules from the
// service_module collection.
resolved, _, rerr := repository.ResolveServiceModules(context.Background(), svc.ProductName, svc.ServiceName, productInfo.Production, svc.Revision)
if rerr != nil {
return nil, e.ErrGetService.AddErr(errors.Wrapf(rerr, "failed to resolve modules for %s/%s rev %d", svc.ProductName, svc.ServiceName, svc.Revision))
var resolved []*commonmodels.Container
if moduleSnapshot != nil {
resolved = moduleSnapshot.Resolved[svc.ServiceName][svc.Revision]
} else {
var rerr error
resolved, _, rerr = repository.ResolveServiceModules(context.Background(), svc.ProductName, svc.ServiceName, productInfo.Production, svc.Revision)
if rerr != nil {
return nil, e.ErrGetService.AddErr(errors.Wrapf(rerr, "failed to resolve modules for %s/%s rev %d", svc.ProductName, svc.ServiceName, svc.Revision))
}
}
for _, container := range resolved {
svcModulesMap[svc.ServiceName][container.Name] = container
Expand All @@ -1157,8 +1177,9 @@ func BuildServiceInfoInEnv(productInfo *commonmodels.Product, templateSvcs []*co
ret := make([]*commonmodels.Container, 0)
if modulesMap, ok := svcModulesMap[svcName]; ok {
for _, module := range modulesMap {
module.ImageName = commonutil.ExtractImageName(module.Image)
ret = append(ret, module)
copy := *module
copy.ImageName = commonutil.ExtractImageName(module.Image)
ret = append(ret, &copy)
}
}
return ret
Expand Down
Loading
Loading