migrate_command.go 8.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392
  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. "github.com/coreos/etcd/etcdserver/api"
  28. pb "github.com/coreos/etcd/etcdserver/etcdserverpb"
  29. "github.com/coreos/etcd/etcdserver/membership"
  30. "github.com/coreos/etcd/mvcc"
  31. "github.com/coreos/etcd/mvcc/backend"
  32. "github.com/coreos/etcd/mvcc/mvccpb"
  33. "github.com/coreos/etcd/pkg/pbutil"
  34. "github.com/coreos/etcd/pkg/types"
  35. "github.com/coreos/etcd/raft/raftpb"
  36. "github.com/coreos/etcd/snap"
  37. "github.com/coreos/etcd/store"
  38. "github.com/coreos/etcd/wal"
  39. "github.com/coreos/etcd/wal/walpb"
  40. "github.com/gogo/protobuf/proto"
  41. "github.com/spf13/cobra"
  42. )
  43. var (
  44. migrateDatadir string
  45. migrateWALdir string
  46. migrateTransformer string
  47. )
  48. // NewMigrateCommand returns the cobra command for "migrate".
  49. func NewMigrateCommand() *cobra.Command {
  50. mc := &cobra.Command{
  51. Use: "migrate",
  52. Short: "Migrates keys in a v2 store to a mvcc store",
  53. Run: migrateCommandFunc,
  54. }
  55. mc.Flags().StringVar(&migrateDatadir, "data-dir", "", "Path to the data directory")
  56. mc.Flags().StringVar(&migrateWALdir, "wal-dir", "", "Path to the WAL directory")
  57. mc.Flags().StringVar(&migrateTransformer, "transformer", "", "Path to the user-provided transformer program")
  58. return mc
  59. }
  60. func migrateCommandFunc(cmd *cobra.Command, args []string) {
  61. var (
  62. writer io.WriteCloser
  63. reader io.ReadCloser
  64. errc chan error
  65. )
  66. if migrateTransformer != "" {
  67. writer, reader, errc = startTransformer()
  68. } else {
  69. fmt.Println("using default transformer")
  70. writer, reader, errc = defaultTransformer()
  71. }
  72. st, index := rebuildStoreV2()
  73. be := prepareBackend()
  74. defer be.Close()
  75. go func() {
  76. writeStore(writer, st)
  77. writer.Close()
  78. }()
  79. readKeys(reader, be)
  80. mvcc.UpdateConsistentIndex(be, index)
  81. err := <-errc
  82. if err != nil {
  83. fmt.Println("failed to transform keys")
  84. ExitWithError(ExitError, err)
  85. }
  86. fmt.Println("finished transforming keys")
  87. }
  88. func prepareBackend() backend.Backend {
  89. dbpath := path.Join(migrateDatadir, "member", "snap", "db")
  90. be := backend.New(dbpath, time.Second, 10000)
  91. tx := be.BatchTx()
  92. tx.Lock()
  93. tx.UnsafeCreateBucket([]byte("key"))
  94. tx.UnsafeCreateBucket([]byte("meta"))
  95. tx.Unlock()
  96. return be
  97. }
  98. func rebuildStoreV2() (store.Store, uint64) {
  99. var index uint64
  100. cl := membership.NewCluster("")
  101. waldir := migrateWALdir
  102. if len(waldir) == 0 {
  103. waldir = path.Join(migrateDatadir, "member", "wal")
  104. }
  105. snapdir := path.Join(migrateDatadir, "member", "snap")
  106. ss := snap.New(snapdir)
  107. snapshot, err := ss.Load()
  108. if err != nil && err != snap.ErrNoSnapshot {
  109. ExitWithError(ExitError, err)
  110. }
  111. var walsnap walpb.Snapshot
  112. if snapshot != nil {
  113. walsnap.Index, walsnap.Term = snapshot.Metadata.Index, snapshot.Metadata.Term
  114. index = snapshot.Metadata.Index
  115. }
  116. w, err := wal.OpenForRead(waldir, walsnap)
  117. if err != nil {
  118. ExitWithError(ExitError, err)
  119. }
  120. defer w.Close()
  121. _, _, ents, err := w.ReadAll()
  122. if err != nil {
  123. ExitWithError(ExitError, err)
  124. }
  125. st := store.New()
  126. if snapshot != nil {
  127. err := st.Recovery(snapshot.Data)
  128. if err != nil {
  129. ExitWithError(ExitError, err)
  130. }
  131. }
  132. cl.SetStore(st)
  133. cl.Recover(api.UpdateCapability)
  134. applier := etcdserver.NewApplierV2(st, cl)
  135. for _, ent := range ents {
  136. if ent.Type == raftpb.EntryConfChange {
  137. var cc raftpb.ConfChange
  138. pbutil.MustUnmarshal(&cc, ent.Data)
  139. applyConf(cc, cl)
  140. continue
  141. }
  142. var raftReq pb.InternalRaftRequest
  143. if !pbutil.MaybeUnmarshal(&raftReq, ent.Data) { // backward compatible
  144. var r pb.Request
  145. pbutil.MustUnmarshal(&r, ent.Data)
  146. applyRequest(&r, applier)
  147. } else {
  148. if raftReq.V2 != nil {
  149. req := raftReq.V2
  150. applyRequest(req, applier)
  151. }
  152. }
  153. if ent.Index > index {
  154. index = ent.Index
  155. }
  156. }
  157. return st, index
  158. }
  159. func applyConf(cc raftpb.ConfChange, cl *membership.RaftCluster) {
  160. if err := cl.ValidateConfigurationChange(cc); err != nil {
  161. return
  162. }
  163. switch cc.Type {
  164. case raftpb.ConfChangeAddNode:
  165. m := new(membership.Member)
  166. if err := json.Unmarshal(cc.Context, m); err != nil {
  167. panic(err)
  168. }
  169. cl.AddMember(m)
  170. case raftpb.ConfChangeRemoveNode:
  171. cl.RemoveMember(types.ID(cc.NodeID))
  172. case raftpb.ConfChangeUpdateNode:
  173. m := new(membership.Member)
  174. if err := json.Unmarshal(cc.Context, m); err != nil {
  175. panic(err)
  176. }
  177. cl.UpdateRaftAttributes(m.ID, m.RaftAttributes)
  178. }
  179. }
  180. func applyRequest(r *pb.Request, applyV2 etcdserver.ApplierV2) {
  181. toTTLOptions(r)
  182. switch r.Method {
  183. case "POST":
  184. applyV2.Post(r)
  185. case "PUT":
  186. applyV2.Put(r)
  187. case "DELETE":
  188. applyV2.Delete(r)
  189. case "QGET":
  190. applyV2.QGet(r)
  191. case "SYNC":
  192. applyV2.Sync(r)
  193. default:
  194. panic("unknown command")
  195. }
  196. }
  197. func toTTLOptions(r *pb.Request) store.TTLOptionSet {
  198. refresh, _ := pbutil.GetBool(r.Refresh)
  199. ttlOptions := store.TTLOptionSet{Refresh: refresh}
  200. if r.Expiration != 0 {
  201. ttlOptions.ExpireTime = time.Unix(0, r.Expiration)
  202. }
  203. return ttlOptions
  204. }
  205. func writeStore(w io.Writer, st store.Store) uint64 {
  206. all, err := st.Get("/1", true, true)
  207. if err != nil {
  208. if eerr, ok := err.(*etcdErr.Error); ok && eerr.ErrorCode == etcdErr.EcodeKeyNotFound {
  209. fmt.Println("no v2 keys to migrate")
  210. os.Exit(0)
  211. }
  212. ExitWithError(ExitError, err)
  213. }
  214. return writeKeys(w, all.Node)
  215. }
  216. func writeKeys(w io.Writer, n *store.NodeExtern) uint64 {
  217. maxIndex := n.ModifiedIndex
  218. nodes := n.Nodes
  219. // remove store v2 bucket prefix
  220. n.Key = n.Key[2:]
  221. if n.Key == "" {
  222. n.Key = "/"
  223. }
  224. if n.Dir {
  225. n.Nodes = nil
  226. }
  227. b, err := json.Marshal(n)
  228. if err != nil {
  229. ExitWithError(ExitError, err)
  230. }
  231. fmt.Fprint(w, string(b))
  232. for _, nn := range nodes {
  233. max := writeKeys(w, nn)
  234. if max > maxIndex {
  235. maxIndex = max
  236. }
  237. }
  238. return maxIndex
  239. }
  240. func readKeys(r io.Reader, be backend.Backend) error {
  241. for {
  242. length64, err := readInt64(r)
  243. if err != nil {
  244. if err == io.EOF {
  245. return nil
  246. }
  247. return err
  248. }
  249. buf := make([]byte, int(length64))
  250. if _, err = io.ReadFull(r, buf); err != nil {
  251. return err
  252. }
  253. var kv mvccpb.KeyValue
  254. err = proto.Unmarshal(buf, &kv)
  255. if err != nil {
  256. return err
  257. }
  258. mvcc.WriteKV(be, kv)
  259. }
  260. }
  261. func readInt64(r io.Reader) (int64, error) {
  262. var n int64
  263. err := binary.Read(r, binary.LittleEndian, &n)
  264. return n, err
  265. }
  266. func startTransformer() (io.WriteCloser, io.ReadCloser, chan error) {
  267. cmd := exec.Command(migrateTransformer)
  268. cmd.Stderr = os.Stderr
  269. writer, err := cmd.StdinPipe()
  270. if err != nil {
  271. ExitWithError(ExitError, err)
  272. }
  273. reader, rerr := cmd.StdoutPipe()
  274. if rerr != nil {
  275. ExitWithError(ExitError, rerr)
  276. }
  277. if err := cmd.Start(); err != nil {
  278. ExitWithError(ExitError, err)
  279. }
  280. errc := make(chan error, 1)
  281. go func() {
  282. errc <- cmd.Wait()
  283. }()
  284. return writer, reader, errc
  285. }
  286. func defaultTransformer() (io.WriteCloser, io.ReadCloser, chan error) {
  287. // transformer decodes v2 keys from sr
  288. sr, sw := io.Pipe()
  289. // transformer encodes v3 keys into dw
  290. dr, dw := io.Pipe()
  291. decoder := json.NewDecoder(sr)
  292. errc := make(chan error, 1)
  293. go func() {
  294. defer func() {
  295. sr.Close()
  296. dw.Close()
  297. }()
  298. for decoder.More() {
  299. node := &client.Node{}
  300. if err := decoder.Decode(node); err != nil {
  301. errc <- err
  302. return
  303. }
  304. kv := transform(node)
  305. if kv == nil {
  306. continue
  307. }
  308. data, err := proto.Marshal(kv)
  309. if err != nil {
  310. errc <- err
  311. return
  312. }
  313. buf := make([]byte, 8)
  314. binary.LittleEndian.PutUint64(buf, uint64(len(data)))
  315. if _, err := dw.Write(buf); err != nil {
  316. errc <- err
  317. return
  318. }
  319. if _, err := dw.Write(data); err != nil {
  320. errc <- err
  321. return
  322. }
  323. }
  324. errc <- nil
  325. }()
  326. return sw, dr, errc
  327. }
  328. func transform(n *client.Node) *mvccpb.KeyValue {
  329. const unKnownVersion = 1
  330. if n.Dir {
  331. return nil
  332. }
  333. kv := &mvccpb.KeyValue{
  334. Key: []byte(n.Key),
  335. Value: []byte(n.Value),
  336. CreateRevision: int64(n.CreatedIndex),
  337. ModRevision: int64(n.ModifiedIndex),
  338. Version: unKnownVersion,
  339. }
  340. return kv
  341. }