Add support for migration mapping
This adds support for splitting a single source into multiple modules.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
+73
-58
@@ -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")
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -13,6 +13,10 @@ type (
|
||||
Name string
|
||||
Path string
|
||||
|
||||
Header *[]string
|
||||
|
||||
Source io.Reader
|
||||
// map is used for stream splitting
|
||||
Map io.Reader
|
||||
}
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user