streamingutil

package
v0.59.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 17, 2026 License: Apache-2.0 Imports: 7 Imported by: 0

Documentation

Index

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

type AsyncJobManager[P any, R any] struct {
	// contains filtered or unexported fields
}

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.

func (*AsyncJobManager[P, R]) Poll

func (m *AsyncJobManager[P, R]) Poll(
	ctx context.Context,
	jobID string,
	fastWait time.Duration,
	runner JobRunner[P, R],
) (*PollStatus[P, R], error)

Poll handles starting a new job or retrieving the status of an existing job.

type Job

type Job[P any, R any] struct {
	ID string
	// contains filtered or unexported fields
}

Job represents an asynchronous in-flight or completed operation.

type JobRunner

type JobRunner[P any, R any] func(ctx context.Context, onProgress func(P) error) (R, error)

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

type SnapshotProvider[T any] func(ctx context.Context) (T, error)

SnapshotProvider returns the current state snapshot.

type SubscriptionProvider

type SubscriptionProvider[T any] func(ctx context.Context) (<-chan T, func())

SubscriptionProvider returns a channel for real-time events and a cleanup function.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL