| 1 | package main |
| 2 | |
| 3 | import "context" |
| 4 | |
| 5 | // startCatalogMetadataRefresh owns at most one refresh and one pending wakeup. |
| 6 | // Requesting work never waits for the catalog writer. The caller cancels ctx |
| 7 | // and joins done before releasing the watcher lifecycle. |
| 8 | func startCatalogMetadataRefresh(ctx context.Context, sync func(context.Context)) (request func(), done <-chan struct{}) { |
| 9 | requests := make(chan struct{}, 1) |
| 10 | finished := make(chan struct{}) |
| 11 | go func() { |
| 12 | defer close(finished) |
| 13 | for { |
| 14 | select { |
| 15 | case <-ctx.Done(): |
| 16 | return |
| 17 | case <-requests: |
| 18 | if ctx.Err() != nil { |
| 19 | return |
| 20 | } |
| 21 | sync(ctx) |
| 22 | } |
| 23 | } |
| 24 | }() |
| 25 | return func() { |
| 26 | if ctx.Err() != nil { |
| 27 | return |
| 28 | } |
| 29 | select { |
| 30 | case requests <- struct{}{}: |
| 31 | default: |
| 32 | } |
| 33 | }, finished |
| 34 | } |
| 35 |