From ac1fd84413ef228ac55e7e2b9361a3a032ef7e5f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Toma=C5=BE=20Jerman?= Date: Sun, 1 Mar 2020 21:19:08 +0100 Subject: [PATCH] Add support for migration mapping This adds support for splitting a single source into multiple modules. --- compose/commands/migrator.go | 48 ++++++++- pkg/migrate/main.go | 131 +++++++++++++----------- pkg/migrate/stream.go | 164 +++++++++++++++++++++++++++++++ pkg/migrate/types/migrateable.go | 4 + 4 files changed, 284 insertions(+), 63 deletions(-) create mode 100644 pkg/migrate/stream.go diff --git a/compose/commands/migrator.go b/compose/commands/migrator.go index 9bd5def67..225184838 100644 --- a/compose/commands/migrator.go +++ b/compose/commands/migrator.go @@ -63,11 +63,28 @@ func Migrator() *cobra.Command { } ext := filepath.Ext(info.Name()) - mg = append(mg, mgt.Migrateable{ - Name: info.Name()[0 : len(info.Name())-len(ext)], - Path: path, - Source: file, - }) + name := info.Name()[0 : len(info.Name())-len(ext)] + mm := migrateableSource(mg, name) + mm.Name = name + mm.Path = path + mm.Source = file + + mg = migrateableAdd(mg, mm) + // @note yes Denis, we will support .yaml files + } else if strings.HasSuffix(info.Name(), ".map.json") { + file, err := os.Open(path) + if err != nil { + log.Fatal(err) + } + + ext := filepath.Ext(info.Name()) + // @todo improve this!! + name := info.Name()[0 : len(info.Name())-len(ext)-4] + mm := migrateableSource(mg, name) + mm.Name = name + mm.Map = file + + mg = migrateableAdd(mg, mm) } return nil }) @@ -89,3 +106,24 @@ func Migrator() *cobra.Command { return cmd } + +// small helper functions for migrateable node management +func migrateableSource(mg []mgt.Migrateable, name string) mgt.Migrateable { + for _, m := range mg { + if m.Name == name { + return m + } + } + + return mgt.Migrateable{} +} + +func migrateableAdd(mg []mgt.Migrateable, mm mgt.Migrateable) []mgt.Migrateable { + for i, m := range mg { + if m.Name == mm.Name { + mg[i] = mm + return mg + } + } + return append(mg, mm) +} diff --git a/pkg/migrate/main.go b/pkg/migrate/main.go index 5761ddf40..fcd5e8e57 100644 --- a/pkg/migrate/main.go +++ b/pkg/migrate/main.go @@ -1,6 +1,7 @@ package migrate import ( + "bytes" "context" "encoding/csv" "errors" @@ -51,71 +52,85 @@ func Migrate(mg []types.Migrateable, ns *cct.Namespace, ctx context.Context) err } // 2. prepare and link migration nodes - for _, m := range mg { - fmt.Printf("mg.processing > %s\n", m.Name) - - // 2.1 load module - mod, err := svcMod.FindByHandle(ns.ID, m.Name) + for _, mgR := range mg { + ss, err := splitStream(mgR) if err != nil { return err } - // 2.2 get header fields - r := csv.NewReader(m.Source) - header, err := r.Read() - if err == io.EOF { - break - } - if err != nil { - panic(err) - } + for _, m := range ss { + fmt.Printf("mg.processing > %s\n", m.Name) - // 2.3 create migration node - n := &types.Node{ - Name: m.Name, - Module: mod, - Namespace: ns, - Reader: r, - Header: header, - } - n = mig.AddNode(n) - - // 2.4 prepare additional migration nodes, to provide dep. constraints - for _, f := range mod.Fields { - if f.Kind == "Record" { - refMod := f.Options["moduleID"] - if refMod == nil { - return errors.New("moduleField.record.missingRef") - } - - modID, ok := refMod.(string) - if !ok { - return errors.New("moduleField.record.invalidRefFormat") - } - fmt.Printf("mg.node.link > %s [%s]\n", f.Name, modID) - - vv, err := strconv.ParseUint(modID, 10, 64) - if err != nil { - return err - } - - mm, err := svcMod.FindByID(ns.ID, vv) - if err != nil { - return err - } - - nn := &types.Node{ - Name: mm.Handle, - Module: mm, - Namespace: ns, - } - - nn = mig.AddNode(nn) - n.LinkAdd(nn) + // 2.1 load module + mod, err := svcMod.FindByHandle(ns.ID, m.Name) + if err != nil { + return err } - } - fmt.Printf("mg.processed > %s\n\n\n", m.Name) + // 2.2 get header fields + r := csv.NewReader(m.Source) + var header []string + if m.Header != nil { + header = *m.Header + } else { + header, err = r.Read() + if err == io.EOF { + break + } + if err != nil { + return err + } + } + + // 2.3 create migration node + n := &types.Node{ + Name: m.Name, + Module: mod, + Namespace: ns, + Reader: r, + Header: header, + Lock: &sync.Mutex{}, + } + n = mig.AddNode(n) + + // 2.4 prepare additional migration nodes, to provide dep. constraints + for _, f := range mod.Fields { + if f.Kind == "Record" { + refMod := f.Options["moduleID"] + if refMod == nil { + return errors.New("moduleField.record.missingRef") + } + + modID, ok := refMod.(string) + if !ok { + return errors.New("moduleField.record.invalidRefFormat") + } + fmt.Printf("mg.node.link > %s [%s]\n", f.Name, modID) + + vv, err := strconv.ParseUint(modID, 10, 64) + if err != nil { + return err + } + + mm, err := svcMod.FindByID(ns.ID, vv) + if err != nil { + return err + } + + nn := &types.Node{ + Name: mm.Handle, + Module: mm, + Namespace: ns, + Lock: &sync.Mutex{}, + } + + nn = mig.AddNode(nn) + n.LinkAdd(nn) + } + } + + fmt.Printf("mg.processed > %s\n\n\n", m.Name) + } } fmt.Printf("graph.remove.cycles\n") diff --git a/pkg/migrate/stream.go b/pkg/migrate/stream.go new file mode 100644 index 000000000..f932960bb --- /dev/null +++ b/pkg/migrate/stream.go @@ -0,0 +1,164 @@ +package migrate + +import ( + "bytes" + "encoding/csv" + "encoding/json" + "errors" + "io" + "io/ioutil" + "strings" + + "github.com/cortezaproject/corteza-server/pkg/migrate/types" +) + +type ( + SplitBuffer struct { + buffer *bytes.Buffer + name string + row []string + header []string + writer *csv.Writer + } +) + +// this function splits the stream of the given migrateable node. +// See readme for more info +func splitStream(m types.Migrateable) ([]types.Migrateable, error) { + var rr []types.Migrateable + rr = append(rr, m) + if m.Map == nil { + return rr, nil + } + + // unpack the map + // @todo provide a better structure!! + var streamMap []map[string]interface{} + src, _ := ioutil.ReadAll(m.Map) + err := json.Unmarshal(src, &streamMap) + if err != nil { + return nil, err + } + + // get header fields + r := csv.NewReader(m.Source) + header, err := r.Read() + if err == io.EOF { + return rr, nil + } + + if err != nil { + return nil, err + } + + // maps header field -> field index for a nicer lookup + hMap := make(map[string]int) + for i, h := range header { + hMap[h] = i + } + + bufs := make(map[string]*SplitBuffer) + + // splitting magic + + // @fix this hack will not always work. + // replace with a set or something similar + i := -1 + for { + i++ + + record, err := r.Read() + if err == io.EOF { + break + } + + if err != nil { + return nil, err + } + + // find first applicable map, that can be used for the given row. + // default maps should not inclide a where field + for _, strmp := range streamMap { + if checkWhere(strmp["where"], record, hMap) { + maps, ok := strmp["map"].([]interface{}) + if !ok { + return nil, errors.New("streamMap.invalidMap") + } + + // populate splitted streams + for _, mp := range maps { + mm, ok := mp.(map[string]interface{}) + if !ok { + return nil, errors.New("streamMap.map.invalidEntry") + } + + from, ok := mm["from"].(string) + if !ok { + return nil, errors.New("streamMap.map.entry.invalidFrom") + } + + to, ok := mm["to"].(string) + if !ok { + return nil, errors.New("streamMap.map.invalidTo") + } + + vv := strings.Split(to, ".") + nm := vv[0] + nmF := vv[1] + + if bufs[nm] == nil { + var bb bytes.Buffer + ww := csv.NewWriter(&bb) + defer ww.Flush() + bufs[nm] = &SplitBuffer{ + buffer: &bb, + writer: ww, + name: nm, + } + } + + bufs[nm].row = append(bufs[nm].row, record[hMap[from]]) + if i == 0 { + bufs[nm].header = append(bufs[nm].header, nmF) + } + } + } + + // write csv rows + for _, v := range bufs { + v.writer.Write(v.row) + var nn []string + v.row = nn + } + } + } + + // make migrateable nodes from the generated streams + for _, v := range bufs { + rr = append(rr, types.Migrateable{ + Name: v.name, + Source: v.buffer, + Header: &v.header, + }) + } + + return rr, nil +} + +// quick and dirty function to check the map's where condition. +// improve with our QL package +func checkWhere(where interface{}, row []string, hMap map[string]int) bool { + if where == nil { + return true + } + + ww, ok := where.(string) + if !ok { + return true + } + + pts := strings.Split(ww, "=") + org := pts[0] + val := pts[1] + return row[hMap[org]] == val +} diff --git a/pkg/migrate/types/migrateable.go b/pkg/migrate/types/migrateable.go index 1135ac652..81421e088 100644 --- a/pkg/migrate/types/migrateable.go +++ b/pkg/migrate/types/migrateable.go @@ -13,6 +13,10 @@ type ( Name string Path string + Header *[]string + Source io.Reader + // map is used for stream splitting + Map io.Reader } )