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

Method run

server/attach_timeout.go:123–160  ·  view source on GitHub ↗

run executes fn (one DuckDB statement) under a timeout. See the attachStepExecer doc comment for the two-layer deadline semantics.

(timeout time.Duration, query string, fn func(ctx context.Context) error)

Source from the content-addressed store, hash-verified

121// run executes fn (one DuckDB statement) under a timeout. See the
122// attachStepExecer doc comment for the two-layer deadline semantics.
123func (a *attachStepExecer) run(timeout time.Duration, query string, fn func(ctx context.Context) error) error {
124 ctx, cancel := context.WithTimeout(context.Background(), timeout)
125 done := make(chan struct{})
126 errCh := make(chan error, 1)
127 go func() {
128 defer close(done)
129 // Cancel only when fn returns: after the deadline fires the driver
130 // keeps re-asserting duckdb_interrupt until the call completes, which
131 // is exactly what we want for a statement that unblocks late.
132 defer cancel()
133 errCh <- fn(ctx)
134 }()
135
136 wall := time.NewTimer(timeout + attachInterruptGrace)
137 defer wall.Stop()
138 select {
139 case err := <-errCh:
140 if err != nil && ctx.Err() == context.DeadlineExceeded {
141 return fmt.Errorf("catalog attachment step timed out after %s (interrupted): %w", timeout, err)
142 }
143 return err
144 case <-wall.C:
145 a.mu.Lock()
146 a.abandoned = append(a.abandoned, abandonedAttachStmt{query: summarizeAttachStmt(query), done: done})
147 a.mu.Unlock()
148 slog.Error("Catalog attachment step exceeded its deadline and did not respond to interrupt; abandoning the call. "+
149 "The statement may still be executing on its DuckDB connection (likely blocked in non-interruptible network I/O); "+
150 "the attach semaphore will not be released until it returns.",
151 "timeout", timeout, "grace", attachInterruptGrace, "query", summarizeAttachStmt(query))
152 go func(q string, started time.Time) {
153 <-done
154 slog.Warn("Abandoned catalog attachment statement finally returned.",
155 "query", q, "overrun", time.Since(started))
156 }(summarizeAttachStmt(query), time.Now())
157 return fmt.Errorf("catalog attachment step timed out after %s and did not respond to interrupt (statement may still be executing): %s",
158 timeout, summarizeAttachStmt(query))
159 }
160}
161
162// releaseSem returns the attach semaphore slot acquired by the caller. If any
163// statement was abandoned past its deadline, the release is handed to a

Callers 5

execTimeoutMethod · 0.95
queryRowScanMethod · 0.95
onStartedLeadingMethod · 0.45
run_dbtFunction · 0.45

Calls 6

fnFunction · 0.85
summarizeAttachStmtFunction · 0.85
NowMethod · 0.80
ErrMethod · 0.65
StopMethod · 0.45
ErrorMethod · 0.45

Tested by 2

run_dbtFunction · 0.36