Documentation
¶
Index ¶
- func PollSnapshot[T any, Resp any](ctx context.Context, snapshot SnapshotProvider[T], mapResponse func(T) *Resp) (*connect.Response[Resp], error)
- func StreamJob[P any, R any, Resp any](ctx context.Context, stream *connect.ServerStream[Resp], ...) error
- func StreamWatch[T any, Resp any](ctx context.Context, stream *connect.ServerStream[Resp], ...) error
- type AsyncJobManager
- func (m *AsyncJobManager[P, R]) Cancel(jobID string) bool
- func (m *AsyncJobManager[P, R]) Close()
- func (m *AsyncJobManager[P, R]) Evict(key string) error
- func (m *AsyncJobManager[P, R]) Expirations() map[string]time.Time
- func (m *AsyncJobManager[P, R]) Poll(ctx context.Context, jobID string, fastWait time.Duration, ...) (*PollStatus[P, R], error)
- type Job
- type JobRunner
- type PollStatus
- type SnapshotProvider
- type SubscriptionProvider
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func PollSnapshot ¶
func PollSnapshot[T any, Resp any]( ctx context.Context, snapshot SnapshotProvider[T], mapResponse func(T) *Resp, ) (*connect.Response[Resp], error)
PollSnapshot resolves the current snapshot and constructs a Unary connect.Response.
func StreamJob ¶
func StreamJob[P any, R any, Resp any]( ctx context.Context, stream *connect.ServerStream[Resp], runner JobRunner[P, R], mapProgress func(P) *Resp, mapResult func(R) *Resp, ) error
StreamJob executes a JobRunner synchronously within a streaming RPC, forwarding progress and the final result.
func StreamWatch ¶
func StreamWatch[T any, Resp any]( ctx context.Context, stream *connect.ServerStream[Resp], cycleDuration time.Duration, pollInterval time.Duration, snapshot SnapshotProvider[T], subscribe SubscriptionProvider[T], mapResponse func(T) (*Resp, bool), ) error
StreamWatch executes a standard stream cycle (up to cycleDuration). It streams updates from either a subscription channel or periodic polling.
Types ¶
type AsyncJobManager ¶
AsyncJobManager manages execution, tracking, and polling for background jobs.
func NewAsyncJobManager ¶
func NewAsyncJobManager[P any, R any](abandonTimeout time.Duration, retentionTTL time.Duration) *AsyncJobManager[P, R]
NewAsyncJobManager creates a new initialized AsyncJobManager.
func (*AsyncJobManager[P, R]) Cancel ¶
func (m *AsyncJobManager[P, R]) Cancel(jobID string) bool
Cancel terminates a running job if it exists and returns true if canceled.
func (*AsyncJobManager[P, R]) Close ¶
func (m *AsyncJobManager[P, R]) Close()
Close stops the background cleaner and cancels all active jobs.
func (*AsyncJobManager[P, R]) Evict ¶
func (m *AsyncJobManager[P, R]) Evict(key string) error
Evict removes and cancels an expired job from the manager.
func (*AsyncJobManager[P, R]) Expirations ¶
func (m *AsyncJobManager[P, R]) Expirations() map[string]time.Time
Expirations returns a snapshot of active job expiration timestamps. If an in-progress job hasn't received any poll requests within abandonTimeout, its context is canceled.
type JobRunner ¶
JobRunner defines the signature of a long-running operation that reports progress and returns a result.
type PollStatus ¶
type PollStatus[P any, R any] struct { JobID string IsDone bool Progress P HasProgress bool Result R Err error }
PollStatus represents the output returned to a polling RPC caller.
type SnapshotProvider ¶
SnapshotProvider returns the current state snapshot.
type SubscriptionProvider ¶
SubscriptionProvider returns a channel for real-time events and a cleanup function.