ring.go 4.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210
  1. package logger
  2. import (
  3. "errors"
  4. "sync"
  5. "sync/atomic"
  6. "github.com/Sirupsen/logrus"
  7. )
  8. const (
  9. defaultRingMaxSize = 1e6 // 1MB
  10. )
  11. // RingLogger is a ring buffer that implements the Logger interface.
  12. // This is used when lossy logging is OK.
  13. type RingLogger struct {
  14. buffer *messageRing
  15. l Logger
  16. logInfo Info
  17. closeFlag int32
  18. }
  19. type ringWithReader struct {
  20. *RingLogger
  21. }
  22. func (r *ringWithReader) ReadLogs(cfg ReadConfig) *LogWatcher {
  23. reader, ok := r.l.(LogReader)
  24. if !ok {
  25. // something is wrong if we get here
  26. panic("expected log reader")
  27. }
  28. return reader.ReadLogs(cfg)
  29. }
  30. func newRingLogger(driver Logger, logInfo Info, maxSize int64) *RingLogger {
  31. l := &RingLogger{
  32. buffer: newRing(maxSize),
  33. l: driver,
  34. logInfo: logInfo,
  35. }
  36. go l.run()
  37. return l
  38. }
  39. // NewRingLogger creates a new Logger that is implemented as a RingBuffer wrapping
  40. // the passed in logger.
  41. func NewRingLogger(driver Logger, logInfo Info, maxSize int64) Logger {
  42. if maxSize < 0 {
  43. maxSize = defaultRingMaxSize
  44. }
  45. l := newRingLogger(driver, logInfo, maxSize)
  46. if _, ok := driver.(LogReader); ok {
  47. return &ringWithReader{l}
  48. }
  49. return l
  50. }
  51. // Log queues messages into the ring buffer
  52. func (r *RingLogger) Log(msg *Message) error {
  53. if r.closed() {
  54. return errClosed
  55. }
  56. return r.buffer.Enqueue(msg)
  57. }
  58. // Name returns the name of the underlying logger
  59. func (r *RingLogger) Name() string {
  60. return r.l.Name()
  61. }
  62. func (r *RingLogger) closed() bool {
  63. return atomic.LoadInt32(&r.closeFlag) == 1
  64. }
  65. func (r *RingLogger) setClosed() {
  66. atomic.StoreInt32(&r.closeFlag, 1)
  67. }
  68. // Close closes the logger
  69. func (r *RingLogger) Close() error {
  70. r.setClosed()
  71. r.buffer.Close()
  72. // empty out the queue
  73. for _, msg := range r.buffer.Drain() {
  74. if err := r.l.Log(msg); err != nil {
  75. logrus.WithField("driver", r.l.Name()).WithField("container", r.logInfo.ContainerID).Errorf("Error writing log message: %v", r.l)
  76. break
  77. }
  78. }
  79. return r.l.Close()
  80. }
  81. // run consumes messages from the ring buffer and forwards them to the underling
  82. // logger.
  83. // This is run in a goroutine when the RingLogger is created
  84. func (r *RingLogger) run() {
  85. for {
  86. if r.closed() {
  87. return
  88. }
  89. msg, err := r.buffer.Dequeue()
  90. if err != nil {
  91. // buffer is closed
  92. return
  93. }
  94. if err := r.l.Log(msg); err != nil {
  95. logrus.WithField("driver", r.l.Name()).WithField("container", r.logInfo.ContainerID).Errorf("Error writing log message: %v", r.l)
  96. }
  97. }
  98. }
  99. type messageRing struct {
  100. mu sync.Mutex
  101. // singals callers of `Dequeue` to wake up either on `Close` or when a new `Message` is added
  102. wait *sync.Cond
  103. sizeBytes int64 // current buffer size
  104. maxBytes int64 // max buffer size size
  105. queue []*Message
  106. closed bool
  107. }
  108. func newRing(maxBytes int64) *messageRing {
  109. queueSize := 1000
  110. if maxBytes == 0 || maxBytes == 1 {
  111. // With 0 or 1 max byte size, the maximum size of the queue would only ever be 1
  112. // message long.
  113. queueSize = 1
  114. }
  115. r := &messageRing{queue: make([]*Message, 0, queueSize), maxBytes: maxBytes}
  116. r.wait = sync.NewCond(&r.mu)
  117. return r
  118. }
  119. // Enqueue adds a message to the buffer queue
  120. // If the message is too big for the buffer it drops the oldest messages to make room
  121. // If there are no messages in the queue and the message is still too big, it adds the message anyway.
  122. func (r *messageRing) Enqueue(m *Message) error {
  123. mSize := int64(len(m.Line))
  124. r.mu.Lock()
  125. if r.closed {
  126. r.mu.Unlock()
  127. return errClosed
  128. }
  129. if mSize+r.sizeBytes > r.maxBytes && len(r.queue) > 0 {
  130. r.wait.Signal()
  131. r.mu.Unlock()
  132. return nil
  133. }
  134. r.queue = append(r.queue, m)
  135. r.sizeBytes += mSize
  136. r.wait.Signal()
  137. r.mu.Unlock()
  138. return nil
  139. }
  140. // Dequeue pulls a message off the queue
  141. // If there are no messages, it waits for one.
  142. // If the buffer is closed, it will return immediately.
  143. func (r *messageRing) Dequeue() (*Message, error) {
  144. r.mu.Lock()
  145. for len(r.queue) == 0 && !r.closed {
  146. r.wait.Wait()
  147. }
  148. if r.closed {
  149. r.mu.Unlock()
  150. return nil, errClosed
  151. }
  152. msg := r.queue[0]
  153. r.queue = r.queue[1:]
  154. r.sizeBytes -= int64(len(msg.Line))
  155. r.mu.Unlock()
  156. return msg, nil
  157. }
  158. var errClosed = errors.New("closed")
  159. // Close closes the buffer ensuring no new messages can be added.
  160. // Any callers waiting to dequeue a message will be woken up.
  161. func (r *messageRing) Close() {
  162. r.mu.Lock()
  163. if r.closed {
  164. r.mu.Unlock()
  165. return
  166. }
  167. r.closed = true
  168. r.wait.Broadcast()
  169. r.mu.Unlock()
  170. return
  171. }
  172. // Drain drains all messages from the queue.
  173. // This can be used after `Close()` to get any remaining messages that were in queue.
  174. func (r *messageRing) Drain() []*Message {
  175. r.mu.Lock()
  176. ls := make([]*Message, 0, len(r.queue))
  177. ls = append(ls, r.queue...)
  178. r.sizeBytes = 0
  179. r.queue = r.queue[:0]
  180. r.mu.Unlock()
  181. return ls
  182. }