Documentation
¶
Overview ¶
This package contains shared types that are used across multiple packages in the Parti library. By keeping these types in a separate package, we avoid import cycles between the main parti package and its internal implementations.
Key types:
- State: Worker lifecycle state
- Partition: Logical work partition
- Assignment: Versioned partition assignment
- Logger: Structured logging interface
- MetricsCollector: Metrics recording interface
Index ¶
- Variables
- func IsNoKeysFoundError(err error) bool
- type Assignment
- type AssignmentMetrics
- type AssignmentStrategy
- type CalculatorMetrics
- type CalculatorState
- type ElectionAgent
- type HandoffState
- type Hooks
- type Logger
- type ManagerMetrics
- type MetricsCollector
- type OwnershipResolver
- type Partition
- type PartitionSource
- type PartitionSourceUpdater
- type PartitionUpdater
- type State
- type StateProvider
- type WatchablePartitionSource
- type WorkerConsumerMetrics
- type WorkerMetrics
Constants ¶
This section is empty.
Variables ¶
var ( // ErrInvalidConfig is returned when the configuration is invalid. ErrInvalidConfig = errors.New("invalid configuration") // ErrNATSConnectionRequired is returned when NATS connection is nil. ErrNATSConnectionRequired = errors.New("NATS connection is required") // ErrPartitionSourceRequired is returned when partition source is nil. ErrPartitionSourceRequired = errors.New("partition source is required") // ErrAssignmentStrategyRequired is returned when assignment strategy is nil. ErrAssignmentStrategyRequired = errors.New("assignment strategy is required") // ErrAlreadyStarted is returned when Start is called on an already running manager. ErrAlreadyStarted = errors.New("manager already started") // ErrNotStarted is returned when operations require a started manager. ErrNotStarted = errors.New("manager not started") // ErrNotImplemented is returned for functionality not yet implemented. ErrNotImplemented = errors.New("not implemented") // ErrNoWorkersAvailable is returned when trying to assign partitions with no workers. ErrNoWorkersAvailable = errors.New("no workers available") // ErrInvalidWorkerID is returned when a worker ID is invalid. ErrInvalidWorkerID = errors.New("invalid worker ID") // ErrElectionFailed is returned when leader election fails. ErrElectionFailed = errors.New("leader election failed") // ErrConnectivity indicates a NATS/KV connectivity issue. // This is used to distinguish network failures from application errors // and triggers degraded mode operation with cached data. ErrConnectivity = errors.New("connectivity issue") // ErrDegraded indicates the system is operating in degraded mode. // Operations may use cached data due to unavailable NATS connection. ErrDegraded = errors.New("degraded operation: using cached data") // ErrIDClaimFailed is returned when stable ID claiming fails. ErrIDClaimFailed = errors.New("failed to claim stable worker ID") // ErrAssignmentFailed is returned when assignment calculation or distribution fails. ErrAssignmentFailed = errors.New("assignment failed") )
Manager errors - Public API errors returned by Manager component.
var ( // ErrCalculatorAlreadyStarted is returned when Start is called on an already running calculator. ErrCalculatorAlreadyStarted = errors.New("calculator already started") // ErrCalculatorNotStarted is returned when operations require a started calculator. ErrCalculatorNotStarted = errors.New("calculator not started") )
Calculator errors - Internal assignment calculator component errors.
var ( // ErrWorkerMonitorAlreadyStarted is returned when Start is called on an already running monitor. ErrWorkerMonitorAlreadyStarted = errors.New("worker monitor already started") // ErrWorkerMonitorAlreadyStopped is returned when Stop is called on an already stopped monitor. ErrWorkerMonitorAlreadyStopped = errors.New("worker monitor already stopped") // ErrWorkerMonitorNotStarted is returned when Stop is called before Start. ErrWorkerMonitorNotStarted = errors.New("worker monitor not started") // ErrWatcherFailed is returned when NATS KV watcher operations fail. ErrWatcherFailed = errors.New("watcher operation failed") )
WorkerMonitor errors - Internal worker monitoring component errors.
var ( // ErrPublishFailed is returned when publishing assignment to NATS KV fails. ErrPublishFailed = errors.New("failed to publish assignment") // ErrDeleteFailed is returned when deleting assignment from NATS KV fails. ErrDeleteFailed = errors.New("failed to delete assignment") )
AssignmentPublisher errors - Internal assignment publishing component errors.
var ( // ErrContextCanceled is returned when an operation is canceled by context. ErrContextCanceled = errors.New("operation canceled by context") // ErrNoKeysFound is returned when NATS KV returns no keys (expected condition). ErrNoKeysFound = errors.New("no keys found") )
Common errors - Shared errors used across multiple components.
Functions ¶
func IsNoKeysFoundError ¶
IsNoKeysFoundError checks if an error indicates that no keys were found in NATS KV.
This function handles NATS-specific "no keys found" errors which may come as:
- Direct error: "nats: no keys found"
- Wrapped error: "failed to list KV keys: nats: no keys found"
Parameters:
- err: The error to check
Returns:
- bool: true if the error indicates no keys were found, false otherwise
Types ¶
type Assignment ¶
type Assignment struct {
// Version is a monotonically increasing assignment version.
// Used to detect stale assignments and coordinate updates.
Version int64 `json:"version"`
// Lifecycle indicates the assignment phase (e.g., "stable", "scaling", "rebalancing").
Lifecycle string `json:"lifecycle"`
// Partitions is the list of partitions assigned to this worker.
Partitions []Partition `json:"partitions"`
}
Assignment contains the current partition assignment for a worker.
Assignments are versioned and include lifecycle metadata for coordination.
type AssignmentMetrics ¶
type AssignmentMetrics interface {
// RecordAssignmentChange records partition assignment changes.
RecordAssignmentChange(added, removed int, version int64)
}
AssignmentMetrics defines metrics for partition assignment operations.
type AssignmentStrategy ¶
type AssignmentStrategy interface {
// Assign calculates partition assignments for the given workers.
//
// The strategy should distribute partitions across workers considering:
// - Partition weights (if supported)
// - Load balancing (even distribution)
// - Cache affinity (minimize reassignment)
//
// Parameters:
// - workers: List of worker IDs to assign partitions to
// - partitions: List of partitions to assign
//
// Returns:
// - map[string][]Partition: Map from workerID to assigned partitions
// - error: Assignment error (e.g., ErrNoWorkersAvailable)
Assign(workers []string, partitions []Partition) (map[string][]Partition, error)
}
AssignmentStrategy calculates partition assignments for a set of workers.
Strategies implement different assignment algorithms:
- ConsistentHash: Weighted consistent hashing with virtual nodes (>80% cache affinity)
- RoundRobin: Simple round-robin distribution (no cache affinity)
- Custom: User-defined algorithms
The leader worker calls Assign during:
- Initial assignment calculation
- Worker scaling (join/leave)
- Partition changes (add/remove)
- Manual rebalancing
Strategy implementations should:
- Be deterministic (same input → same output)
- Handle edge cases (no workers, no partitions, weights)
- Run quickly (called on hot path)
- Be stateless (no side effects)
type CalculatorMetrics ¶
type CalculatorMetrics interface {
// RecordRebalanceDuration records the time taken for a rebalance operation.
//
// Parameters:
// - duration: Time taken in seconds
// - reason: Rebalance reason ("cold_start", "planned_scale", "emergency", "restart")
RecordRebalanceDuration(duration float64, reason string)
// RecordRebalanceAttempt records a rebalance attempt (success or failure).
//
// Parameters:
// - reason: Rebalance reason
// - success: true if rebalance succeeded, false otherwise
RecordRebalanceAttempt(reason string, success bool)
// RecordPartitionCount sets the current partition count (gauge metric).
//
// Parameters:
// - count: Current number of partitions being managed
RecordPartitionCount(count int)
// RecordKVOperationDuration records NATS KV operation latency.
//
// Parameters:
// - operation: Operation type ("get", "put", "delete", "watch")
// - duration: Time taken in seconds
RecordKVOperationDuration(operation string, duration float64)
// RecordStateChangeDropped records when state change notifications are dropped due to slow subscribers.
RecordStateChangeDropped()
// RecordEmergencyRebalance records an emergency rebalance trigger.
//
// Parameters:
// - disappearedWorkers: Number of workers that disappeared suddenly
RecordEmergencyRebalance(disappearedWorkers int)
// RecordWorkerChange records worker topology changes detected by the calculator.
//
// Parameters:
// - added: Number of workers added (0 if none)
// - removed: Number of workers removed (0 if none)
RecordWorkerChange(added, removed int)
// RecordOrphanedPartitions records the number of partitions that were not assigned to any worker.
//
// Parameters:
// - count: Number of orphaned partitions
RecordOrphanedPartitions(count int)
// RecordActiveWorkers sets the current active worker count (gauge metric).
//
// Parameters:
// - count: Current number of active workers
RecordActiveWorkers(count int)
// RecordCacheUsage records when cached data is used instead of fresh KV data.
//
// Parameters:
// - cacheType: Type of cache used ("workers", "assignments")
// - age: Age of cached data in seconds
RecordCacheUsage(cacheType string, age float64)
// IncrementCacheFallback increments the counter for cache fallback events.
//
// Parameters:
// - reason: Reason for fallback ("connectivity_error", "timeout", "unknown")
IncrementCacheFallback(reason string)
}
CalculatorMetrics defines metrics for calculator operations.
type CalculatorState ¶
type CalculatorState int
CalculatorState represents the state of the assignment calculator.
The calculator transitions through these states during partition assignment:
Idle → Scaling → Rebalancing → Idle (normal scaling) Idle → Emergency → Idle (worker crash)
This provides type-safe state management for the assignment calculator, preventing string comparison errors in the Manager.
const ( // CalcStateIdle indicates the calculator is idle (stable operation). // No active rebalancing or scaling in progress. CalcStateIdle CalculatorState = iota // CalcStateScaling indicates the calculator is in a stabilization window. // Waiting for worker topology to stabilize before rebalancing. // Used during planned scaling (workers joining gradually). CalcStateScaling // CalcStateRebalancing indicates active rebalancing is in progress. // Assignments are being calculated and published to workers. CalcStateRebalancing // CalcStateEmergency indicates emergency rebalancing (worker crash). // No stabilization window - immediate rebalancing required. CalcStateEmergency )
func (CalculatorState) String ¶
func (s CalculatorState) String() string
String returns the string representation of calculator state.
Returns:
- string: Human-readable state name
type ElectionAgent ¶
type ElectionAgent interface {
// RequestLeadership attempts to acquire leadership.
//
// Should use a lease-based mechanism with the specified duration.
// If already leader, should extend the lease.
//
// Parameters:
// - ctx: Context for cancellation and timeout
// - workerID: The worker ID requesting leadership
// - leaseDuration: Lease duration in seconds
//
// Returns:
// - bool: true if leadership acquired/held, false otherwise
// - error: Election error (nil on success)
RequestLeadership(ctx context.Context, workerID string, leaseDuration int64) (bool, error)
// RenewLeadership renews the current leadership lease.
//
// Called periodically by the leader to maintain leadership.
// Should fail if leadership was lost (another worker became leader).
//
// Parameters:
// - ctx: Context for cancellation and timeout
//
// Returns:
// - error: Renewal error (nil on success, indicates leadership lost)
RenewLeadership(ctx context.Context) error
// ReleaseLeadership voluntarily releases leadership.
//
// Called during graceful shutdown to allow fast leader failover.
//
// Parameters:
// - ctx: Context for cancellation and timeout
//
// Returns:
// - error: Release error (nil on success)
ReleaseLeadership(ctx context.Context) error
// IsLeader checks if this worker is currently the leader.
//
// Used for state verification and metrics.
//
// Parameters:
// - ctx: Context for cancellation and timeout
//
// Returns:
// - bool: true if this worker is the leader
// - error: Check error (nil on success)
IsLeader(ctx context.Context) (bool, error)
}
ElectionAgent handles leader election for assignment coordination.
Leader election ensures exactly one worker calculates and distributes partition assignments. The leader is responsible for:
- Monitoring worker join/leave events
- Calculating new assignments
- Publishing assignments to all workers
Implementations can use:
- NATS KV (built-in, recommended)
- External agents (Consul, etcd, Zookeeper)
- Custom coordination services
The Manager calls ElectionAgent methods during:
- Startup (request leadership)
- Background loop (renew leadership)
- Shutdown (release leadership)
type HandoffState ¶
type HandoffState int
HandoffState represents the ownership handoff state for a partition during two-phase handoff.
States generally follow the sequence: Stable → Prepare → Commit → Stable. The Processing Gate uses these states to determine whether a partition message should be processed by the current owner.
const ( // HandoffStateUnknown indicates the state is not known (partition not found in claims). HandoffStateUnknown HandoffState = iota // HandoffStateStable indicates normal operation with stable ownership. HandoffStateStable // HandoffStatePrepare indicates the first phase of handoff (old owner still processing). HandoffStatePrepare // HandoffStateCommit indicates the second phase of handoff (new owner can start processing). HandoffStateCommit )
type Hooks ¶
type Hooks struct {
// OnAssignmentChanged is called when this worker's partition assignment changes.
//
// Parameters:
// - oldPartitions: previous complete assignment set (may be empty on first assignment)
// - newPartitions: new complete assignment set (may be empty if worker loses all partitions)
//
// Invocation behavior:
// - Triggered ONLY when the assignment KV key is updated with a new version
// - NOT triggered on manager startup if no assignment exists for the worker
// - If a worker joins but receives no partitions (more workers than partitions),
// this hook will NOT be called since no assignment key is created
// - Receiving an empty newPartitions slice indicates the worker should release
// all partition-related resources (caches, connections, etc.)
//
// Example scenarios:
// - First assignment: oldPartitions=[], newPartitions=[p1,p2,p3]
// - Rebalance: oldPartitions=[p1,p2], newPartitions=[p1,p3]
// - All partitions removed: oldPartitions=[p1,p2], newPartitions=[]
OnAssignmentChanged func(ctx context.Context, oldPartitions, newPartitions []Partition) error
// OnStateChanged is called when worker state transitions.
//
// Invocation behavior:
// - Called asynchronously in a background goroutine
// - Triggered on all valid state transitions (see Manager state machine)
// - Always invoked when entering or exiting degraded mode
//
// Parameters:
// - from: previous state
// - to: new state
OnStateChanged func(ctx context.Context, from, to State) error
// OnError is called when a recoverable error occurs.
//
// Invocation behavior:
// - Called for transient errors that don't cause manager shutdown
// - May be called multiple times in rapid succession during issues
//
// Parameters:
// - err: the error that occurred
OnError func(ctx context.Context, err error) error
// OnLeadershipChanged is called when the worker acquires or loses leadership.
//
// Invocation behavior:
// - Called immediately after leadership status changes
// - Called with isLeader=true when initially elected as leader
// - Called with isLeader=false when leadership is lost (lease expiry, etc.)
// - Called with isLeader=true when a follower successfully claims vacant leadership
// - NOT called if worker remains a follower on startup
//
// Parameters:
// - isLeader: true if the worker is now the leader, false if leadership was lost
OnLeadershipChanged func(ctx context.Context, isLeader bool) error
// OnPartitionsAssigned is called when new partitions are assigned to this worker.
//
// This is a convenience hook derived from OnAssignmentChanged that provides
// only the newly added partitions (not the full assignment set).
//
// Invocation behavior:
// - Called ONLY when there are partitions being added (len(added) > 0)
// - NOT called if assignment changes but no new partitions are added
// - NOT called for workers that never receive an assignment
//
// Parameters:
// - partitions: slice of newly assigned partitions (never empty when called)
OnPartitionsAssigned func(ctx context.Context, partitions []Partition) error
// OnPartitionsRevoked is called when partitions are removed from this worker.
//
// This is a convenience hook derived from OnAssignmentChanged that provides
// only the removed partitions (not the full assignment set).
//
// Invocation behavior:
// - Called ONLY when there are partitions being removed (len(removed) > 0)
// - Called when a worker loses all its partitions (e.g., scaling down)
// - NOT called if assignment changes but no partitions are removed
// - Use this hook to clean up partition-specific resources (caches, connections)
//
// Parameters:
// - partitions: slice of revoked partitions (never empty when called)
OnPartitionsRevoked func(ctx context.Context, partitions []Partition) error
// OnDegraded is called when the manager enters degraded mode.
//
// Invocation behavior:
// - Called once when transitioning into degraded mode
// - NOT called repeatedly while in degraded mode
// - OnStateChanged is also called with the degraded state transition
//
// Parameters:
// - reason: description of the cause (e.g., "NATS connection down",
// "KV error threshold exceeded")
OnDegraded func(ctx context.Context, reason string) error
}
Hooks defines callbacks for Manager lifecycle events.
All hooks are optional and called asynchronously in background goroutines to avoid blocking the state machine. Hooks receive the manager's lifecycle context which will be cancelled during shutdown.
IMPORTANT: Hook execution behavior:
- Hooks run concurrently and may not complete before Stop() returns
- The context passed to hooks is cancelled when manager stops
- Hook errors are logged but don't fail manager operations
Best practices for hook implementation:
- Complete quickly (< 1 second recommended)
- Respect context cancellation
- Don't block on long I/O operations
- Make hooks idempotent (may be called multiple times)
- Handle errors gracefully (return error for logging)
Example:
hooks := &parti.Hooks{
OnStateChanged: func(ctx context.Context, from, to parti.State) error {
select {
case <-ctx.Done():
return ctx.Err() // Manager is shutting down
case metricsChan <- StateMetric{from, to}:
return nil
case <-time.After(500 * time.Millisecond):
return errors.New("metric send timeout")
}
},
}
type Logger ¶
type Logger interface {
// Debug logs a message at DebugLevel.
// The message includes any fields passed at the log site, as well as any fields accumulated on the logger.
Debug(msg string, keysAndValues ...any)
// Info logs a message at InfoLevel.
// The message includes any fields passed at the log site, as well as any fields accumulated on the logger.
Info(msg string, keysAndValues ...any)
// Warn logs a message at WarnLevel.
// The message includes any fields passed at the log site, as well as any fields accumulated on the logger.
Warn(msg string, keysAndValues ...any)
// Error logs a message at ErrorLevel.
// The message includes any fields passed at the log site, as well as any fields accumulated on the logger.
Error(msg string, keysAndValues ...any)
// Fatal logs a message at FatalLevel and calls os.Exit(1).
// The message includes any fields passed at the log site, as well as any fields accumulated on the logger.
//
// The logger then calls os.Exit(1), even if logging at FatalLevel is disabled.
Fatal(msg string, keysAndValues ...any)
}
Logger defines methods for structured logging.
Compatible with zap.SugaredLogger and other structured loggers. All methods accept key-value pairs for structured fields.
type ManagerMetrics ¶
type ManagerMetrics interface {
// RecordStateTransition records a manager state transition event.
RecordStateTransition(from, to State, duration float64)
// RecordLeadershipChange records a leadership change.
RecordLeadershipChange(newLeader string)
// RecordDegradedDuration records the duration spent in degraded mode.
//
// Parameters:
// - duration: Time spent in degraded mode
RecordDegradedDuration(duration float64)
// SetDegradedMode sets the current degraded mode status (0 or 1).
//
// Parameters:
// - degraded: 1.0 if in degraded mode, 0.0 otherwise
SetDegradedMode(degraded float64)
// SetCacheAge sets the age of cached assignment data in seconds.
//
// Parameters:
// - age: Age of cached data (0 if no cache)
SetCacheAge(age float64)
// SetAlertLevel sets the current degraded mode alert level (0-3).
//
// Parameters:
// - level: Alert level (0=none, 1=info, 2=warn, 3=error, 4=critical)
SetAlertLevel(level int)
// IncrementAlertEmitted tracks alert emission by level for spam detection.
//
// Parameters:
// - level: Alert level name ("info", "warn", "error", "critical")
IncrementAlertEmitted(level string)
}
ManagerMetrics defines metrics for manager-level operations.
type MetricsCollector ¶
type MetricsCollector interface {
ManagerMetrics
CalculatorMetrics
WorkerMetrics
AssignmentMetrics
WorkerConsumerMetrics
}
MetricsCollector defines methods for recording operational metrics.
Implementations should be non-blocking and handle failures gracefully. All methods are called from internal goroutines and must be thread-safe.
This interface composes smaller, domain-focused interfaces for better modularity.
type OwnershipResolver ¶
type OwnershipResolver interface {
// GetOwner returns ownership information for a partition.
//
// Parameters:
// - partitionID: Partition identifier to look up
//
// Returns:
// - string: Owner worker ID (empty if not found)
// - HandoffState: Current handoff state
// - int64: Claim epoch number
// - bool: true if partition found, false otherwise
GetOwner(partitionID string) (string, HandoffState, int64, bool)
// ForceRefreshPartition performs a best-effort on-demand refresh for a specific partition.
//
// Parameters:
// - ctx: Context for cancellation and timeouts
// - partitionID: Partition identifier to refresh
//
// Returns:
// - error: Non-nil if the refresh failed
ForceRefreshPartition(ctx context.Context, partitionID string) error
}
OwnershipResolver provides the current owner and handoff state for a partition.
Implementations should provide O(1) reads without lock contention for use in the message processing hot path. The returned `ok` value indicates whether the partition was found in the resolver's cache.
type Partition ¶
type Partition struct {
// Keys uniquely identify this partition.
// For Kafka: ["topic", "partition_id"]
Keys []string `json:"keys"`
// Weight represents the relative processing cost (default: 100).
// Used by weighted assignment strategies for load balancing.
Weight int64 `json:"weight"`
}
Partition represents a logical work partition.
A partition is the unit of work assignment. Each partition can contain multiple keys and has an associated weight for load balancing.
func (Partition) Compare ¶
Compare performs a lexicographic comparison of partition key sequences.
Ordering rules:
- Compare Keys element-wise using string order
- If all shared elements are equal, the shorter Keys slice sorts first
- Returns 0 when both key sequences are identical (weight is not considered)
Returns:
- int: -1 if p < q, 0 if equal, +1 if p > q
func (Partition) HashID ¶
HashID returns a stable 64-bit hash of the partition's key sequence using chained XXH3 hashing.
This method computes the hash without allocating by folding each key into the hash of the previous key (seed chaining). Boundary ambiguity is avoided because each key boundary starts a fresh seeded hash step instead of raw concatenation.
Algorithm:
- If the partition has no keys, returns 0
- For the first key: xxh3.HashString(key)
- For subsequent keys: xxh3.HashStringSeed(key, previousHash)
- Final hash value is returned as the partition's HashID
Returns:
- uint64: 64-bit hash (0 if no keys)
Example:
p := Partition{Keys: []string{"topic", "42"}}
h := p.HashID()
_ = h // use in consistent hashing or caching structures
func (Partition) HashIDSeed ¶
HashIDSeed returns a stable 64-bit hash of the partition's key sequence using chained XXH3 hashing incorporating an explicit seed value. When seed == 0 this is equivalent to HashID().
The first key is hashed with the provided seed (if non-zero) or unseeded; subsequent keys are hashed using the previous hash value as the seed, preserving boundary separation without concatenation allocations.
Returns:
- uint64: 64-bit hash (0 if no keys)
func (Partition) ID ¶
ID returns the canonical durable name fragment for the partition by joining the Keys with a dash ("-"). This is suitable for durable consumer names and hashing contexts requiring a stable, human-readable identifier.
Returns:
- string: Dash-joined key sequence ("" if no keys)
func (Partition) SubjectKey ¶
SubjectKey returns the canonical subject identifier for the partition formed by joining the Keys with a dot ("."). This is used for subject templating and JetStream FilterSubjects construction.
Returns:
- string: Dot-joined key sequence ("" if no keys)
type PartitionSource ¶
type PartitionSource interface {
// Start initializes the source (e.g., starts watchers).
//
// Parameters:
// - ctx: Context for initialization
//
// Returns:
// - error: Initialization error
Start(ctx context.Context) error
// List returns all available partitions.
//
// Implementations should:
// - Return consistent results for the same backend state
// - Handle context cancellation gracefully
// - Return errors for transient failures (will be retried)
//
// Parameters:
// - ctx: Context for cancellation and timeout
//
// Returns:
// - []Partition: List of discovered partitions
// - error: Discovery error (nil on success)
List(ctx context.Context) ([]Partition, error)
// Stop cleans up resources (e.g., stops watchers).
//
// Parameters:
// - ctx: Context for cancellation and timeout
//
// Returns:
// - error: Cleanup error
Stop(ctx context.Context) error
}
PartitionSource discovers and provides the list of available partitions.
Implementations can query various backends:
- Cassandra: the signal sources from system tables
- Static: fixed list for testing
- Custom: any dynamic partition discovery logic
The Manager calls List during:
- Startup (initial discovery)
- RefreshPartitions() (manual refresh)
- Periodic refresh (if configured)
type PartitionSourceUpdater ¶ added in v1.3.0
type PartitionSourceUpdater interface {
PartitionSource
PartitionUpdater
}
PartitionSourceUpdater combines PartitionSource and PartitionUpdater.
type PartitionUpdater ¶ added in v1.1.0
type PartitionUpdater interface {
// Update updates the partition list in the backend.
//
// Parameters:
// - ctx: Context for the operation
// - partitions: New list of partitions
//
// Returns:
// - error: Update error
Update(ctx context.Context, partitions []Partition) error
}
PartitionUpdater allows updating the partition definition in the backend.
This interface is optional and typically used by admin tools or tests to modify the source of truth (e.g., NATS KV, file).
type State ¶
type State int
State represents the worker lifecycle state.
States follow a defined progression during normal operation:
StateInit → StateClaimingID → StateElection → StateWaitingAssignment → StateStable
During scaling or rebalancing:
StateStable → StateScaling/StateRebalancing → StateStable
Emergency and shutdown are terminal states.
const ( // StateInit is the initial state before any operations. StateInit State = iota // StateClaimingID indicates the worker is claiming a stable ID. StateClaimingID // StateElection indicates the worker is participating in leader election. StateElection // StateWaitingAssignment indicates waiting for initial partition assignment. StateWaitingAssignment // StateStable indicates normal operation with stable assignment. StateStable // StateScaling indicates dynamic scaling is in progress. StateScaling // StateRebalancing indicates partition rebalancing is in progress. StateRebalancing // StateEmergency indicates an error condition requiring intervention. StateEmergency // StateDegraded indicates operating with cached data due to NATS connectivity issues. // The system continues processing with last known good assignments while NATS is unavailable. StateDegraded // StateShutdown indicates graceful shutdown is in progress. StateShutdown )
type StateProvider ¶
type StateProvider interface {
// State returns the current manager state.
State() State
// IsInRecoveryGrace returns true if in post-recovery grace period.
// This prevents split-brain scenarios during recovery.
IsInRecoveryGrace() bool
}
StateProvider allows components to check manager state without circular dependencies.
This interface enables the Calculator to check for degraded mode and recovery grace period without requiring a direct dependency on the Manager type.
type WatchablePartitionSource ¶ added in v1.3.0
type WatchablePartitionSource interface {
PartitionSource
// Watch returns a channel that emits a signal when the partition list changes.
//
// The channel should be closed when the context is canceled or the source is stopped.
// The caller should drain the channel to prevent blocking the source.
//
// Parameters:
// - ctx: Context for the watcher
//
// Returns:
// - <-chan struct{}: Signal channel
Watch(ctx context.Context) <-chan struct{}
}
WatchablePartitionSource is an extension for sources that can signal updates.
Implementations of this interface allow the Calculator to react immediately to partition list changes without polling.
type WorkerConsumerMetrics ¶
type WorkerConsumerMetrics interface {
// IncrementWorkerConsumerControlRetry increments retry attempts by operation (create_update, iterate, info).
//
// Parameters:
// - op: Operation name ("create_update", "iterate", "info")
IncrementWorkerConsumerControlRetry(op string)
// RecordWorkerConsumerRetryBackoff records backoff delay durations (seconds) by operation.
// Typically emitted alongside retry increments to populate histogram buckets.
//
// Parameters:
// - op: Operation name
// - seconds: Backoff duration in seconds
RecordWorkerConsumerRetryBackoff(op string, seconds float64)
// SetWorkerConsumerSubjectsCurrent sets the current number of subjects (gauge).
SetWorkerConsumerSubjectsCurrent(count int)
// IncrementWorkerConsumerSubjectChange increments add/remove counts on subject diffs.
IncrementWorkerConsumerSubjectChange(kind string, count int)
// IncrementWorkerConsumerGuardrailViolation increments violations (max_subjects, workerid_mutation).
IncrementWorkerConsumerGuardrailViolation(kind string)
// IncrementWorkerConsumerSubjectThresholdWarning increments threshold warning events.
IncrementWorkerConsumerSubjectThresholdWarning()
// RecordWorkerConsumerUpdate increments update results (success|failure|noop).
RecordWorkerConsumerUpdate(result string)
// ObserveWorkerConsumerUpdateLatency records update latency in seconds.
ObserveWorkerConsumerUpdateLatency(seconds float64)
// IncrementWorkerConsumerIteratorRestart increments iterator restart counts by reason (transient|heartbeat).
IncrementWorkerConsumerIteratorRestart(reason string)
// IncrementWorkerConsumerIteratorEscalation increments the counter when iterator failures
// escalate to a consumer refresh action. This captures bursts of failures that
// warrant proactive intervention beyond simple iterator recreation.
IncrementWorkerConsumerIteratorEscalation(reason string)
// SetWorkerConsumerConsecutiveIteratorFailures sets the current consecutive iterator failures gauge.
SetWorkerConsumerConsecutiveIteratorFailures(count int)
// SetWorkerConsumerHealthStatus sets worker consumer health status gauge (1 healthy, 0 unhealthy).
SetWorkerConsumerHealthStatus(healthy bool)
// IncrementWorkerConsumerRecreationAttempt increments recreation attempts by reason.
//
// Parameters:
// - reason: Reason category ("not_found"|"iterator_error"|"unknown")
IncrementWorkerConsumerRecreationAttempt(reason string)
// RecordWorkerConsumerRecreation records recreation outcomes by result and reason.
//
// Parameters:
// - result: "success"|"failure"
// - reason: "not_found"|"iterator_error"|"unknown"
RecordWorkerConsumerRecreation(result string, reason string)
// ObserveWorkerConsumerRecreationDuration observes total recreation latency in seconds.
//
// Parameters:
// - seconds: Total duration of a recreation attempt sequence (success or terminal failure)
ObserveWorkerConsumerRecreationDuration(seconds float64)
// IncrementWorkerConsumerPullSuppressed increments the counter when a pull is suppressed
// due to ownership/state gating (v2 per-subject mode). Reason examples: "not_owner", "state_blocked".
IncrementWorkerConsumerPullSuppressed(reason string)
}
WorkerConsumerMetrics defines metrics for the single durable consumer helper.
type WorkerMetrics ¶
type WorkerMetrics interface {
// RecordHeartbeat records a heartbeat event from an individual worker.
//
// Parameters:
// - workerID: The ID of the worker publishing the heartbeat
// - success: true if heartbeat was successfully published, false otherwise
RecordHeartbeat(workerID string, success bool)
}
WorkerMetrics defines metrics for individual worker heartbeat operations.
These metrics are recorded by individual workers publishing their heartbeats, not by the calculator monitoring workers.