server.go 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163
  1. package etcdserver
  2. import (
  3. "errors"
  4. "sync"
  5. "time"
  6. "code.google.com/p/go.net/context"
  7. "github.com/coreos/etcd/raft"
  8. "github.com/coreos/etcd/raft/raftpb"
  9. "github.com/coreos/etcd/store"
  10. "github.com/coreos/etcd/wait"
  11. pb "github.com/coreos/etcd/etcdserver2/etcdserverpb"
  12. )
  13. var ErrUnknownMethod = errors.New("etcdserver: unknown method")
  14. type SendFunc func(m []raftpb.Message)
  15. type Response struct {
  16. // The last seen term raft was at when this request was built.
  17. Term int64
  18. // The last seen index raft was at when this request was built.
  19. Commit int64
  20. *store.Event
  21. *store.Watcher
  22. err error
  23. }
  24. type Server struct {
  25. once sync.Once
  26. w wait.List
  27. Node raft.Node
  28. Store store.Store
  29. // Send specifies the send function for sending msgs to peers. Send
  30. // MUST NOT block. It is okay to drop messages, since clients should
  31. // timeout and reissue their messages. If Send is nil, Server will
  32. // panic.
  33. Send SendFunc
  34. // Save specifies the save function for saving ents to stable storage.
  35. // Save MUST block until st and ents are on stable storage. If Send is
  36. // nil, Server will panic.
  37. Save func(st raftpb.State, ents []raftpb.Entry)
  38. }
  39. func (s *Server) init() { s.w = wait.New() }
  40. func (s *Server) Run(ctx context.Context) {
  41. s.once.Do(s.init)
  42. for {
  43. select {
  44. case rd := <-s.Node.Ready():
  45. s.Save(rd.State, rd.Entries)
  46. s.Send(rd.Messages)
  47. // TODO(bmizerany): do this in the background, but take
  48. // care to apply entries in a single goroutine, and not
  49. // race them.
  50. for _, e := range rd.CommittedEntries {
  51. var resp Response
  52. resp.Event, resp.err = s.apply(context.TODO(), e)
  53. resp.Term = rd.Term
  54. resp.Commit = rd.Commit
  55. s.w.Trigger(e.Id, resp)
  56. }
  57. case <-ctx.Done():
  58. return
  59. }
  60. }
  61. }
  62. func (s *Server) Do(ctx context.Context, r pb.Request) (Response, error) {
  63. s.once.Do(s.init)
  64. if r.Id == 0 {
  65. panic("r.Id cannot be 0")
  66. }
  67. switch r.Method {
  68. case "POST", "PUT", "DELETE":
  69. data, err := r.Marshal()
  70. if err != nil {
  71. return Response{}, err
  72. }
  73. ch := s.w.Register(r.Id)
  74. s.Node.Propose(ctx, r.Id, data)
  75. select {
  76. case x := <-ch:
  77. resp := x.(Response)
  78. return resp, resp.err
  79. case <-ctx.Done():
  80. s.w.Trigger(r.Id, nil) // GC wait
  81. return Response{}, ctx.Err()
  82. }
  83. case "GET":
  84. switch {
  85. case r.Wait:
  86. wc, err := s.Store.Watch(r.Path, r.Recursive, false, r.Since)
  87. if err != nil {
  88. return Response{}, err
  89. }
  90. return Response{Watcher: wc}, nil
  91. default:
  92. ev, err := s.Store.Get(r.Path, r.Recursive, r.Sorted)
  93. if err != nil {
  94. return Response{}, err
  95. }
  96. return Response{Event: ev}, nil
  97. }
  98. default:
  99. return Response{}, ErrUnknownMethod
  100. }
  101. }
  102. // apply interprets r as a call to store.X and returns an Response interpreted from store.Event
  103. func (s *Server) apply(ctx context.Context, e raftpb.Entry) (*store.Event, error) {
  104. var r pb.Request
  105. if err := r.Unmarshal(e.Data); err != nil {
  106. return nil, err
  107. }
  108. expr := time.Unix(0, r.Expiration)
  109. switch r.Method {
  110. case "POST":
  111. return s.Store.Create(r.Path, r.Dir, r.Val, true, expr)
  112. case "PUT":
  113. exists, existsSet := getBool(r.PrevExists)
  114. switch {
  115. case existsSet:
  116. if exists {
  117. return s.Store.Update(r.Path, r.Val, expr)
  118. } else {
  119. return s.Store.Create(r.Path, r.Dir, r.Val, false, expr)
  120. }
  121. case r.PrevIndex > 0 || r.PrevValue != "":
  122. return s.Store.CompareAndSwap(r.Path, r.PrevValue, r.PrevIndex, r.Val, expr)
  123. default:
  124. return s.Store.Set(r.Path, r.Dir, r.Val, expr)
  125. }
  126. case "DELETE":
  127. switch {
  128. case r.PrevIndex > 0 || r.PrevValue != "":
  129. return s.Store.CompareAndDelete(r.Path, r.PrevValue, r.PrevIndex)
  130. default:
  131. return s.Store.Delete(r.Path, r.Recursive, r.Dir)
  132. }
  133. default:
  134. // This should never be reached, but just in case:
  135. return nil, ErrUnknownMethod
  136. }
  137. }
  138. func getBool(v *bool) (vv bool, set bool) {
  139. if v == nil {
  140. return false, false
  141. }
  142. return *v, true
  143. }