triggerAsyncCompression kicks off a background compression job for the conversation owning st. A no-op when a job is already pending — the check-and-set happens under st.mu so concurrent callers cannot replace (and thereby leak) an in-flight job.
(ctx context.Context, st *compressionState, messages []llm.Message, filePath string)
| 289 | // check-and-set happens under st.mu so concurrent callers cannot replace |
| 290 | // (and thereby leak) an in-flight job. |
| 291 | func (r *Runner) triggerAsyncCompression(ctx context.Context, st *compressionState, messages []llm.Message, filePath string) { |
| 292 | st.mu.Lock() |
| 293 | if st.pendingJob != nil { |
| 294 | st.mu.Unlock() |
| 295 | return |
| 296 | } |
| 297 | msgSnapshot := copyMessages(messages) |
| 298 | asyncCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Minute) |
| 299 | job := &compressionJob{done: make(chan struct{}), cancel: cancel, snapshotLen: len(messages)} |
| 300 | st.pendingJob = job |
| 301 | st.mu.Unlock() |
| 302 | |
| 303 | // Registered before the goroutine starts so WaitBackground can never miss |
| 304 | // a job that was launched but has not run yet. |
| 305 | r.bg.Add(1) |
| 306 | go func() { |
| 307 | defer r.bg.Done() |
| 308 | defer cancel() |
| 309 | rebuilt, err := r.runCompression(asyncCtx, msgSnapshot, filePath) |
| 310 | |
| 311 | st.mu.Lock() |
| 312 | defer st.mu.Unlock() |
| 313 | |
| 314 | if st.pendingJob != job { |
| 315 | return // cancelled or superseded |
| 316 | } |
| 317 | if err != nil { |
| 318 | // Still the owner, so this is a genuine failure rather than a |
| 319 | // deliberate cancel (cancelPendingCompression cancels and clears |
| 320 | // pendingJob under the lock, so cancelled jobs fail the ownership |
| 321 | // check above and die silently). Abandon the job rather than |
| 322 | // applying a truncated/unmodified snapshot over live messages. |
| 323 | fmt.Fprintf(stdout.Writer(), "[ocr] Memory compression failed: %v\n", err) |
| 324 | st.pendingJob = nil |
| 325 | close(job.done) |
| 326 | return |
| 327 | } |
| 328 | job.rebuilt = rebuilt |
| 329 | close(job.done) |
| 330 | }() |
| 331 | } |
| 332 | |
| 333 | // tryApplyPendingCompression checks whether a background compression has |
| 334 | // completed and swaps the rebuilt messages into place. Returns true if |