MCPcopy Create free account
hub / github.com/PostHog/duckgres / openQueryLogDuckLakeDB

Function openQueryLogDuckLakeDB

server/querylog.go:135–206  ·  view source on GitHub ↗
(cfg Config)

Source from the content-addressed store, hash-verified

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
143func 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
153func (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
185func (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

Callers 2

NewQueryLoggerFunction · 0.85

Calls 12

BuildAttachStmtFunction · 0.92
MigrationNeededFunction · 0.92
queryLogDuckLakeDBDSNFunction · 0.85
setExtensionDirectoryFunction · 0.85
LoadExtensionsFunction · 0.85
createS3SecretFunction · 0.85
ensureQueryLogTableFunction · 0.85
CloseMethod · 0.65
ExecMethod · 0.65
OpenMethod · 0.45

Tested by

no test coverage detected