(cfg Config)
| 133 | if ql.cancel != nil { |
| 134 | ql.cancel() |
| 135 | } |
| 136 | if ql.db != nil && ql.closeDB { |
| 137 | _ = ql.db.Close() |
| 138 | } |
| 139 | }) |
| 140 | return stopErr |
| 141 | } |
| 142 | |
| 143 | func queryLogStopContext(ctx context.Context, defaultTimeout time.Duration) (context.Context, context.CancelFunc) { |
| 144 | if ctx == nil { |
| 145 | ctx = context.Background() |
| 146 | } |
| 147 | if _, ok := ctx.Deadline(); ok { |
| 148 | return ctx, func() {} |
| 149 | } |
| 150 | return context.WithTimeout(ctx, defaultTimeout) |
| 151 | } |
| 152 | |
| 153 | func (ql *QueryLogger) flushLoop() { |
| 154 | defer close(ql.done) |
| 155 | defer observe.SetQueryLogBufferedEntries(0) |
| 156 | |
| 157 | batch := make([]QueryLogEntry, 0, ql.cfg.BatchSize) |
| 158 | flushTicker := time.NewTicker(ql.cfg.FlushInterval) |
| 159 | defer flushTicker.Stop() |
| 160 | |
| 161 | for { |
| 162 | select { |
| 163 | case entry, ok := <-ql.ch: |
| 164 | if !ok { |
| 165 | // Channel closed — drain and exit |
| 166 | if len(batch) > 0 { |
| 167 | ql.flushBatch(batch) |
| 168 | } |
| 169 | return |
| 170 | } |
| 171 | batch = append(batch, entry) |
| 172 | if len(batch) >= ql.cfg.BatchSize { |
| 173 | ql.flushBatch(batch) |
| 174 | batch = batch[:0] |
| 175 | } |
| 176 | case <-flushTicker.C: |
| 177 | if len(batch) > 0 { |
| 178 | ql.flushBatch(batch) |
| 179 | batch = batch[:0] |
| 180 | } |
| 181 | } |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | func (ql *QueryLogger) addBufferedEntries(delta int64) { |
| 186 | if ql == nil { |
| 187 | return |
| 188 | } |
| 189 | buffered := ql.buffered.Add(delta) |
| 190 | if buffered < 0 { |
| 191 | ql.buffered.Store(0) |
| 192 | buffered = 0 |
no test coverage detected