server.go 3.2 KB

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