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)
| 121 | // run executes fn (one DuckDB statement) under a timeout. See the |
| 122 | // attachStepExecer doc comment for the two-layer deadline semantics. |
| 123 | func (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 |