migrate_command.go 7.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355
  1. // Copyright 2016 The etcd Authors
  2. //
  3. // Licensed under the Apache License, Version 2.0 (the "License");
  4. // you may not use this file except in compliance with the License.
  5. // You may obtain a copy of the License at
  6. //
  7. // http://www.apache.org/licenses/LICENSE-2.0
  8. //
  9. // Unless required by applicable law or agreed to in writing, software
  10. // distributed under the License is distributed on an "AS IS" BASIS,
  11. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  12. // See the License for the specific language governing permissions and
  13. // limitations under the License.
  14. package command
  15. import (
  16. "encoding/binary"
  17. "encoding/json"
  18. "fmt"
  19. "io"
  20. "os"
  21. "os/exec"
  22. "path"
  23. "time"
  24. "github.com/coreos/etcd/client"
  25. etcdErr "github.com/coreos/etcd/error"
  26. "github.com/coreos/etcd/etcdserver"
  27. pb "github.com/coreos/etcd/etcdserver/etcdserverpb"
  28. "github.com/coreos/etcd/mvcc"
  29. "github.com/coreos/etcd/mvcc/backend"
  30. "github.com/coreos/etcd/mvcc/mvccpb"
  31. "github.com/coreos/etcd/pkg/pbutil"
  32. "github.com/coreos/etcd/raft/raftpb"
  33. "github.com/coreos/etcd/snap"
  34. "github.com/coreos/etcd/store"
  35. "github.com/coreos/etcd/wal"
  36. "github.com/coreos/etcd/wal/walpb"
  37. "github.com/gogo/protobuf/proto"
  38. "github.com/spf13/cobra"
  39. )
  40. var (
  41. migrateDatadir string
  42. migrateWALdir string
  43. migrateTransformer string
  44. )
  45. // NewMigrateCommand returns the cobra command for "migrate".
  46. func NewMigrateCommand() *cobra.Command {
  47. mc := &cobra.Command{
  48. Use: "migrate",
  49. Short: "Migrates keys in a v2 store to a mvcc store",
  50. Run: migrateCommandFunc,
  51. }
  52. mc.Flags().StringVar(&migrateDatadir, "data-dir", "", "Path to the data directory")
  53. mc.Flags().StringVar(&migrateWALdir, "wal-dir", "", "Path to the WAL directory")
  54. mc.Flags().StringVar(&migrateTransformer, "transformer", "", "Path to the user-provided transformer program")
  55. return mc
  56. }
  57. func migrateCommandFunc(cmd *cobra.Command, args []string) {
  58. var (
  59. writer io.WriteCloser
  60. reader io.ReadCloser
  61. errc chan error
  62. )
  63. if migrateTransformer != "" {
  64. writer, reader, errc = startTransformer()
  65. } else {
  66. fmt.Println("using default transformer")
  67. writer, reader, errc = defaultTransformer()
  68. }
  69. st := rebuildStoreV2()
  70. be := prepareBackend()
  71. defer be.Close()
  72. maxIndexc := make(chan uint64, 1)
  73. go func() {
  74. maxIndexc <- writeStore(writer, st)
  75. writer.Close()
  76. }()
  77. readKeys(reader, be)
  78. mvcc.UpdateConsistentIndex(be, <-maxIndexc)
  79. err := <-errc
  80. if err != nil {
  81. fmt.Println("failed to transform keys")
  82. ExitWithError(ExitError, err)
  83. }
  84. fmt.Println("finished transforming keys")
  85. }
  86. func prepareBackend() backend.Backend {
  87. dbpath := path.Join(migrateDatadir, "member", "snap", "db")
  88. be := backend.New(dbpath, time.Second, 10000)
  89. tx := be.BatchTx()
  90. tx.Lock()
  91. tx.UnsafeCreateBucket([]byte("key"))
  92. tx.UnsafeCreateBucket([]byte("meta"))
  93. tx.Unlock()
  94. return be
  95. }
  96. func rebuildStoreV2() store.Store {
  97. waldir := migrateWALdir
  98. if len(waldir) == 0 {
  99. waldir = path.Join(migrateDatadir, "member", "wal")
  100. }
  101. snapdir := path.Join(migrateDatadir, "member", "snap")
  102. ss := snap.New(snapdir)
  103. snapshot, err := ss.Load()
  104. if err != nil && err != snap.ErrNoSnapshot {
  105. ExitWithError(ExitError, err)
  106. }
  107. var walsnap walpb.Snapshot
  108. if snapshot != nil {
  109. walsnap.Index, walsnap.Term = snapshot.Metadata.Index, snapshot.Metadata.Term
  110. }
  111. w, err := wal.OpenForRead(waldir, walsnap)
  112. if err != nil {
  113. ExitWithError(ExitError, err)
  114. }
  115. defer w.Close()
  116. _, _, ents, err := w.ReadAll()
  117. if err != nil {
  118. ExitWithError(ExitError, err)
  119. }
  120. st := store.New()
  121. if snapshot != nil {
  122. err := st.Recovery(snapshot.Data)
  123. if err != nil {
  124. ExitWithError(ExitError, err)
  125. }
  126. }
  127. applier := etcdserver.NewApplierV2(st, nil)
  128. for _, ent := range ents {
  129. if ent.Type != raftpb.EntryNormal {
  130. continue
  131. }
  132. var raftReq pb.InternalRaftRequest
  133. if !pbutil.MaybeUnmarshal(&raftReq, ent.Data) { // backward compatible
  134. var r pb.Request
  135. pbutil.MustUnmarshal(&r, ent.Data)
  136. applyRequest(&r, applier)
  137. } else {
  138. if raftReq.V2 != nil {
  139. req := raftReq.V2
  140. applyRequest(req, applier)
  141. }
  142. }
  143. }
  144. return st
  145. }
  146. func applyRequest(r *pb.Request, applyV2 etcdserver.ApplierV2) {
  147. toTTLOptions(r)
  148. switch r.Method {
  149. case "POST":
  150. applyV2.Post(r)
  151. case "PUT":
  152. applyV2.Put(r)
  153. case "DELETE":
  154. applyV2.Delete(r)
  155. case "QGET":
  156. applyV2.QGet(r)
  157. case "SYNC":
  158. applyV2.Sync(r)
  159. default:
  160. panic("unknown command")
  161. }
  162. }
  163. func toTTLOptions(r *pb.Request) store.TTLOptionSet {
  164. refresh, _ := pbutil.GetBool(r.Refresh)
  165. ttlOptions := store.TTLOptionSet{Refresh: refresh}
  166. if r.Expiration != 0 {
  167. ttlOptions.ExpireTime = time.Unix(0, r.Expiration)
  168. }
  169. return ttlOptions
  170. }
  171. func writeStore(w io.Writer, st store.Store) uint64 {
  172. all, err := st.Get("/1", true, true)
  173. if err != nil {
  174. if eerr, ok := err.(*etcdErr.Error); ok && eerr.ErrorCode == etcdErr.EcodeKeyNotFound {
  175. fmt.Println("no v2 keys to migrate")
  176. os.Exit(0)
  177. }
  178. ExitWithError(ExitError, err)
  179. }
  180. return writeKeys(w, all.Node)
  181. }
  182. func writeKeys(w io.Writer, n *store.NodeExtern) uint64 {
  183. maxIndex := n.ModifiedIndex
  184. nodes := n.Nodes
  185. // remove store v2 bucket prefix
  186. n.Key = n.Key[2:]
  187. if n.Key == "" {
  188. n.Key = "/"
  189. }
  190. if n.Dir {
  191. n.Nodes = nil
  192. }
  193. b, err := json.Marshal(n)
  194. if err != nil {
  195. ExitWithError(ExitError, err)
  196. }
  197. fmt.Fprint(w, string(b))
  198. for _, nn := range nodes {
  199. max := writeKeys(w, nn)
  200. if max > maxIndex {
  201. maxIndex = max
  202. }
  203. }
  204. return maxIndex
  205. }
  206. func readKeys(r io.Reader, be backend.Backend) error {
  207. for {
  208. length64, err := readInt64(r)
  209. if err != nil {
  210. if err == io.EOF {
  211. return nil
  212. }
  213. return err
  214. }
  215. buf := make([]byte, int(length64))
  216. if _, err = io.ReadFull(r, buf); err != nil {
  217. return err
  218. }
  219. var kv mvccpb.KeyValue
  220. err = proto.Unmarshal(buf, &kv)
  221. if err != nil {
  222. return err
  223. }
  224. mvcc.WriteKV(be, kv)
  225. }
  226. }
  227. func readInt64(r io.Reader) (int64, error) {
  228. var n int64
  229. err := binary.Read(r, binary.LittleEndian, &n)
  230. return n, err
  231. }
  232. func startTransformer() (io.WriteCloser, io.ReadCloser, chan error) {
  233. cmd := exec.Command(migrateTransformer)
  234. cmd.Stderr = os.Stderr
  235. writer, err := cmd.StdinPipe()
  236. if err != nil {
  237. ExitWithError(ExitError, err)
  238. }
  239. reader, rerr := cmd.StdoutPipe()
  240. if rerr != nil {
  241. ExitWithError(ExitError, rerr)
  242. }
  243. if err := cmd.Start(); err != nil {
  244. ExitWithError(ExitError, err)
  245. }
  246. errc := make(chan error, 1)
  247. go func() {
  248. errc <- cmd.Wait()
  249. }()
  250. return writer, reader, errc
  251. }
  252. func defaultTransformer() (io.WriteCloser, io.ReadCloser, chan error) {
  253. // transformer decodes v2 keys from sr
  254. sr, sw := io.Pipe()
  255. // transformer encodes v3 keys into dw
  256. dr, dw := io.Pipe()
  257. decoder := json.NewDecoder(sr)
  258. errc := make(chan error, 1)
  259. go func() {
  260. defer func() {
  261. sr.Close()
  262. dw.Close()
  263. }()
  264. for decoder.More() {
  265. node := &client.Node{}
  266. if err := decoder.Decode(node); err != nil {
  267. errc <- err
  268. return
  269. }
  270. kv := transform(node)
  271. if kv == nil {
  272. continue
  273. }
  274. data, err := proto.Marshal(kv)
  275. if err != nil {
  276. errc <- err
  277. return
  278. }
  279. buf := make([]byte, 8)
  280. binary.LittleEndian.PutUint64(buf, uint64(len(data)))
  281. if _, err := dw.Write(buf); err != nil {
  282. errc <- err
  283. return
  284. }
  285. if _, err := dw.Write(data); err != nil {
  286. errc <- err
  287. return
  288. }
  289. }
  290. errc <- nil
  291. }()
  292. return sw, dr, errc
  293. }
  294. func transform(n *client.Node) *mvccpb.KeyValue {
  295. const unKnownVersion = 1
  296. if n.Dir {
  297. return nil
  298. }
  299. kv := &mvccpb.KeyValue{
  300. Key: []byte(n.Key),
  301. Value: []byte(n.Value),
  302. CreateRevision: int64(n.CreatedIndex),
  303. ModRevision: int64(n.ModifiedIndex),
  304. Version: unKnownVersion,
  305. }
  306. return kv
  307. }