Concurrency Package
This package provides standardized concurrency primitives and patterns to ensure consistent, safe, and observable concurrent operations across the codebase.
Overview
The pkg/concurrency package addresses critical concurrency concerns:
- Deadlock prevention via timeout wrappers for all mutex operations
- Standard callback/notify pattern for operation lifecycle tracking
- Smart defaults for worker counts and concurrency parameters
- Observability via metrics and logging integration
Core Components
1. Run with max wait (RunWithMaxWait) โ Cap wait on fire-and-forget
Use when: You must run a function (e.g. writing to stdout/pipe) but must not block the current goroutine indefinitely if that operation blocks (e.g. broken pipe, full buffer). The function runs in a new goroutine; the caller returns when the function completes or when maxWait elapses, whichever comes first.
Note: The function is not cancelled when the timeout elapses; it may still complete in the background. This pattern only caps how long the caller waits.
concurrency.RunWithMaxWait(func() {
fmt.Fprintf(w, "Scheduler started (PID: %d).\n", pid)
}, 2*time.Second)
See run_with_max_wait.go for full documentation.
2. Transactional Critical Sections (RunInLock / RunInRLock) โ Preferred for shared state
Problem: Lock helpers that run the callback in a different goroutine (e.g. for timeout detection) leave the lock held by the caller while the callback runs without the lock, causing data races when the callback touches shared state.
Solution: Use RunInLock / RunInRLock for any critical section that reads or modifies shared state. The callback runs in the same goroutine that holds the lock, so the "transaction" (lock โ modify โ unlock) is atomic and no other goroutine can see half-updated state.
- Deadlock avoidance: establish a global lock order when multiple locks are needed (e.g. always A before B); prefer holding one lock at a time and copying data out when possible.
// Preferred: same-goroutine critical section
err := concurrency.RunInLock(&buf.mu, func() error {
buf.items[key] = value
return nil
})
// Read-only
err := concurrency.RunInRLock(&buf.mu, func() error {
v = buf.items[key]
return nil
})
With optional logging (when you want to see long-held locks):
err := concurrency.RunInLockWithLogger(&buf.mu, "audit_buffer_flush_copy", logging.GetLockLoggerFromProfile("system"), func() error {
// copy shared state
return nil
})
When to use RunInLock / RunInRLock:
- โ Any read or write to shared state protected by a mutex
- โ When you need a clear "transaction": read/copy under lock, then release before I/O or callbacks
When to use WithLockTimeout / WithLockCtxLogger instead:
- Only when you need timeout detection and your callback does not touch the same shared state (e.g. the callback only does logging). Note: WithLockTimeout runs the callback in a separate goroutine; the lock is held by the caller, so the callback does not run with the lock. Do not use it for critical sections that modify shared state.
2b. sync.Cond with context (WaitCondContext)
Problem: You cannot select on sync.Cond. Spawning a goroutine only to call cond.Wait so the parent can select on ctx.Done() is easy to get wrong (stray goroutines, misleading docs about locks) and ignores the same-goroutine lock rule above.
Solution: Wait in the calling goroutine with mu held, in a for !ready() { cond.Wait() } loop. Register context.AfterFunc so that when ctx is cancelled, a one-shot callback acquires mu, calls cond.Broadcast(), and releases mu. The waiter then wakes, observes ctx.Err(), and returns without an extra waiter goroutine.
Use WaitCondContext. Example call site: cmd/zqk/system.WaitProjectCacheBackgroundWork.
Architecture / anti-cycles: docs/architecture/concurrency-patterns-v1.0.md ยง Cache sidecar idle coordination documents wait vs callback vs coordinator layers and rules to avoid feedback loops when reacting to idle edges.
3. Timeout Mutex Wrappers (WithLock / WithRLock)
Problem: Direct mutex operations can deadlock indefinitely with no visibility.
Solution: Standardized convenience overloads provide deterministic timeouts and observability. Important: These run the callback in a separate goroutine for timeout monitoring; use RunInLock for critical sections that modify shared state.
Implementation note (WithLockTimeout / WithRLockTimeout): The inner operation runs via goroutinelabels.GoroutineBuilder (established pattern; see docs/architecture/concurrency-patterns-v1.0.md). The builder is used with WithCleanup so that a completion channel is always signaled when the inner goroutine exits (normal return, panic, or context cancel). That ensures the caller releases the lock promptly instead of blocking until the operation timeout. This enhancement keeps goroutine usage consistent with project patterns while fixing the case where a panic or early exit would otherwise leave the caller holding the lock until the timeout fired.
Standard Pattern (when timeout is required and callback does not touch shared state)
For most production code - use WithLockCtxLogger / WithRLockCtxLogger:
import (
"github.com/zqk-os/zqk/pkg/concurrency"
"github.com/zqk-os/zqk/pkg/logging"
)
// Standard pattern: context + logger (no metrics boilerplate)
err := concurrency.WithLockCtxLogger(
&myStruct.mu,
ctx,
"operation_name",
logging.GetLockLoggerFromProfile("system"),
func() error {
// Critical section - guaranteed to complete or timeout
myStruct.data = newValue
return nil
},
)
if err != nil {
// Handle timeout (deadlock prevented)
return fmt.Errorf("operation timed out: %w", err)
}
For read locks:
err := concurrency.WithRLockCtxLogger(
&myStruct.mu,
ctx,
"read_operation",
logging.GetLockLoggerFromProfile("system"),
func() error {
// Read-only critical section
value := myStruct.data
return nil
},
)
Import Cycle Avoidance Pattern
For packages that cannot import pkg/logging (e.g., pkg/context, pkg/logging itself) - use WithLockCtx / WithRLockCtx:
import "github.com/zqk-os/zqk/pkg/concurrency"
// No logger to avoid import cycles
err := concurrency.WithLockCtx(
&myStruct.mu,
ctx,
"operation_name",
func() error {
myStruct.data = newValue
return nil
},
)
Available Overloads
The package provides multiple convenience overloads to eliminate boilerplate:
WithLockCtxLogger/WithRLockCtxLoggerโญ Primary pattern - context + logger (most production code)WithLockCtx/WithRLockCtx- context only (import cycle avoidance)WithLockLogger/WithRLockLogger- logger only, usescontext.Background()WithLock/WithRLock- minimal, usescontext.Background()and no loggerWithLockTimeout/WithRLockTimeout- full control (legacy/internal use only)
Note: The verbose WithLockTimeout with all parameters is still available for special cases, but the convenience overloads should be preferred for cleaner, more readable code.
When to Use:
- โ High-contention operations (caches, shared state)
- โ Operations that might hold locks during I/O
- โ Operations in critical paths (validation, scheduling)
When NOT to Use:
- โ Simple getters/setters (just field assignments, < 1ms)
- โ Operations that already release locks before I/O
- โ Operations guaranteed to be very fast
4. Operation Callback Interface
Problem: No standard pattern for tracking long-running operation lifecycles.
Solution: OperationCallback interface provides consistent lifecycle tracking.
type OperationCallback interface {
OnStart(operationID string, metadata map[string]any)
OnProgress(operationID string, progress int, total int, message string)
OnComplete(operationID string, result any, duration time.Duration)
OnError(operationID string, err error)
OnCancel(operationID string, reason string)
}
Implementations:
NoOpOperationCallback: No-op implementation for testing/fallbackCoordinatorOperationCallback: Coordinator-integrated implementation (inpkg/coordination)
Usage:
import (
"github.com/zqk-os/zqk/pkg/concurrency"
"github.com/zqk-os/zqk/pkg/coordination"
)
// Create callback (coordinator-integrated)
callback := coordination.NewCoordinatorOperationCallback(
ctx,
projectRoot,
storageProvider,
"my_operation",
"human", // profile
)
// Use in operation
callback.OnStart(operationID, map[string]any{
"total_items": 100,
})
// ... do work ...
callback.OnComplete(operationID, result, duration)
Note: CoordinatorOperationCallback is in pkg/coordination to avoid import cycles. The interface is defined here in pkg/concurrency.
5. Goroutine Ceiling (Circuit Breaker)
Problem: Unbounded goroutine growth (e.g. many cron jobs or list workers) can make the process unresponsive; we want to block new work until existing work catches up.
Solution: Use WaitUnderGoroutineCeiling at "gate" points (e.g. before submitting to the triggered-job pool, or before acquiring a list slot). When runtime.NumGoroutine() is already at or above the ceiling (default 2000), the call blocks until the count drops or the context is cancelled. This prevents adding more work when the process is already overloaded.
import "github.com/zqk-os/zqk/pkg/concurrency"
// Before submitting work that would add goroutines
if err := concurrency.WaitUnderGoroutineCeiling(ctx, concurrency.DefaultGoroutineCeiling, 200*time.Millisecond); err != nil {
return err // e.g. context cancelled
}
// ... submit to channel, acquire slot, etc.
Constants:
DefaultGoroutineCeiling(2000): block new work when at or above this count.GoroutineCountWarningThreshold(1500): health monitor logs a warning when count exceeds this so operators see "things slowing down" before the ceiling blocks work.
When to use: At any point where you're about to add work that creates or triggers many goroutines (job submission, starting a list/count operation). Use together with bounded pools so each subsystem has a cap and a global ceiling acts as a safety net.
6. Smart Defaults Configuration
Problem: Hardcoded worker counts don't adapt to system resources.
Solution: CPU-aware defaults with global configuration.
import "github.com/zqk-os/zqk/pkg/concurrency"
// Get smart defaults (CPU-aware)
cfg := concurrency.GetGlobalConcurrencyConfig()
validatorWorkers := cfg.ValidatorMaxWorkers // NumCPU * 2, clamped [4, 32]
routerWorkers := cfg.AsyncRouterMaxWorkers // NumCPU, clamped [2, 16]
// Override defaults (for CLI, tests, etc.)
concurrency.SetGlobalConcurrencyConfig(&concurrency.ConcurrencyConfig{
ValidatorMaxWorkers: 16,
AsyncRouterMaxWorkers: 8,
})
Defaults:
ValidatorMaxWorkers:runtime.NumCPU() * 2, clamped to [4, 32]AsyncRouterMaxWorkers:runtime.NumCPU(), clamped to [2, 16]
Architecture Decisions
Import Cycle Resolution
Problem: pkg/concurrency needed pkg/storage for coordinator callbacks, but pkg/storage also needed pkg/concurrency for timeout wrappers.
Solution:
- Keep
OperationCallbackinterface inpkg/concurrency(no dependencies) - Move
CoordinatorOperationCallbacktopkg/coordination(can import both) - Result: Clean dependency graph with no cycles
pkg/concurrency โ (no storage/coordination deps)
pkg/coordination โ pkg/concurrency + pkg/storage
pkg/storage โ pkg/concurrency (no cycles!)
Timeout Calculation
Timeouts are calculated based on:
- Operation type (validation, scheduling, etc.)
- Optional metrics (if available, uses historical data)
- Default fallback (5 seconds for most operations)
This ensures timeouts are appropriate for the operation while preventing indefinite hangs.
LockMetrics Interface
Optional metrics interface for lock observability:
type LockMetrics interface {
GetObjectsPerSecond() float64
RecordLockWait(duration time.Duration)
RecordLockHold(duration time.Duration)
}
Implemented by pkg/validation/ValidationMetrics for validation operations.
Testing
RecordingOperationCallback
For tests that don't need coordinator/storage dependencies:
import "github.com/zqk-os/zqk/pkg/concurrency"
recorder := concurrency.NewRecordingOperationCallback()
// Use in operation
operation.SetCallback(recorder)
// ... execute operation ...
// Verify callbacks
recorder.AssertStartEventCount(t, 1)
recorder.AssertCompleteEventCount(t, 1)
events := recorder.GetStartEvents()
// ... verify event data ...
This avoids import cycles in tests while providing full verification capabilities.
Integration Examples
Async Validator
// Uses global config for worker count
cfg := concurrency.GetGlobalConcurrencyConfig()
validator := validation.NewAsyncValidator(projectRoot, cfg.ValidatorMaxWorkers, timeout)
Async Router
// Uses global config for worker count
cfg := concurrency.GetGlobalConcurrencyConfig()
router := transceiver.NewAsyncRouter(router, cfg.AsyncRouterMaxWorkers, queueSize, logger)
System Check Command
// High-level operation callback for lifecycle tracking
callback := coordination.NewCoordinatorOperationCallback(
ctx, projectRoot, storageProvider, "system_check", profile,
)
callback.OnStart(operationID, map[string]any{
"total_tasks": totalTasks,
})
// ... validation work ...
callback.OnComplete(operationID, result, duration)
Best Practices
- Always use timeout wrappers for high-contention or I/O-holding operations
- Prefer convenience overloads - use
WithLockCtxLogger/WithRLockCtxLoggerfor most production code - Avoid verbose patterns - don't use
WithLockTimeoutwith all parameters unless absolutely necessary - Use
WithLockCtx/WithRLockCtxwhen logger would cause import cycles (e.g., inpkg/contextorpkg/logging) - Use operation callbacks for long-running operations that need lifecycle tracking
- Leverage smart defaults - only override when necessary
- Follow import cycle patterns - keep interfaces in
pkg/concurrency, implementations inpkg/coordination - Use RecordingOperationCallback in tests to avoid coordinator dependencies
Standardization Status
โ Complete: All packages now use standardized concurrency patterns:
- Storage: 76 occurrences standardized
- Validation: ~60 occurrences standardized
- Objects: 4 occurrences standardized
- Logging: ~23 occurrences standardized (uses
WithLockCtx/WithRLockCtxfor import cycle avoidance) - Graph: ~39 occurrences standardized
- Scheduler: 83 occurrences standardized
- MCP: 102 occurrences standardized
- Runtime, Coordination, CLI, Config, and others: ~77 occurrences standardized
Total: ~464 occurrences across ~103 files now use the standardized convenience overloads, eliminating boilerplate and improving code readability.
Related Documentation
docs/architecture/LOCK_ORDERING.md: Global lock order when multiple locks are needed; enforce in code review (POL-ARCH-004).docs/architecture/GOROUTINE_ARCHITECTURE_POLICY.md: Goroutine patternsCONCURRENCY_STANDARDIZATION_PLAN.md: Standardization planCONCURRENCY_STANDARDIZATION_PROGRESS.md: Progress trackingpkg/coordination/operation_callback.go: Coordinator-based callback implementation