MCPcopy Create free account
hub / github.com/alibaba/open-code-review / triggerAsyncCompression

Method triggerAsyncCompression

internal/llmloop/compression.go:291–331  ·  view source on GitHub ↗

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)

Source from the content-addressed store, hash-verified

289// check-and-set happens under st.mu so concurrent callers cannot replace
290// (and thereby leak) an in-flight job.
291func (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

Calls 4

runCompressionMethod · 0.95
WriterFunction · 0.92
AddMethod · 0.80
copyMessagesFunction · 0.70