async

package
v0.4.27 Latest Latest
Warning

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

Go to latest
Published: Aug 12, 2026 License: LGPL-2.1 Imports: 8 Imported by: 25

Documentation

Overview

Package async 提供一次性结果、完成信号、连续流和结构化异步作用域。

Package async 将不同异步语义拆分为独立类型:

  • Promise/Future:非泛型、一次性、可重放的 Result;
  • Completer/Signal:不携带 Result 的生命周期完成通知;
  • Emitter/Stream:连续 Result 的单消费流;
  • Scope/Spawn:绑定宿主生命周期的后台任务取消、汇合与统计;
  • Race、FirstSuccess、All、AllSettled、Zip2、Map、FlatMap 和 Timeout: 基于完成订阅的 Future 组合器。

Future 内部保存完成结果,并通过 OnComplete 在完成者 goroutine 中直接通知订阅者; 多个 Wait、TryGet 或 OnComplete 消费者读取的是同一个可重放结果,不会竞争消费, 也不会为每个订阅者启动等待 goroutine。

Stream 采用单消费语义;多个消费者会竞争元素。需要广播时应使用 event 或上层消息 设施。Scope.Close 只负责取消 Context 和等待协作式退出,不能强制终止 goroutine。

根包 core 的 Submit、Post、Spawn 和 ContinueOn 建立在这些能力之上。

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrAsync             = fmt.Errorf("%w: async", exception.ErrCore)       // ErrAsync 是异步模块错误的共同根错误。
	ErrScopeClosed       = fmt.Errorf("%w: scope closed", ErrAsync)         // ErrScopeClosed 表示异步作用域已经关闭。
	ErrNoCandidates      = fmt.Errorf("%w: no future candidates", ErrAsync) // ErrNoCandidates 表示组合器没有可用的候选 Future。
	ErrNoFutureSucceeded = fmt.Errorf("%w: no future succeeded", ErrAsync)  // ErrNoFutureSucceeded 表示所有候选 Future 均失败。
	ErrFutureTimeout     = fmt.Errorf("%w: future timeout", ErrAsync)       // ErrFutureTimeout 表示 Future 等待超时。
)

Functions

func NewPromise added in v0.4.27

func NewPromise(completionExecutorID ...ExecutorID) (Promise, Future)

NewPromise 创建一对生产者 Promise 和消费者 Future。

completionExecutorID 可选;仅使用第一个值。它应表示完成该 Future 必须运行的 Runtime 执行器,外部 I/O、后台 goroutine 或未知来源使用零值。

func NewSignal added in v0.4.27

func NewSignal() (Completer, Signal)

NewSignal 创建一对完成端 Completer 和等待端 Signal。

func NewStream added in v0.4.27

func NewStream(buffer ...int) (Emitter, Stream)

NewStream 创建单生产者、单消费语义的 Emitter/Stream。 buffer 省略或小于 1 时使用 1;仅使用第一个值。

Types

type Completer added in v0.4.27

type Completer struct {
	// contains filtered or unexported fields
}

Completer 是无结果完成信号的生产者端。

func (Completer) Complete added in v0.4.27

func (completer Completer) Complete() bool

Complete 完成 Signal。首次完成返回 true,后续调用返回 false。

func (Completer) IsNil added in v0.4.27

func (completer Completer) IsNil() bool

IsNil 报告 Completer 是否为零值。

func (Completer) Signal added in v0.4.27

func (completer Completer) Signal() Signal

Signal 返回共享状态的只读完成信号。

type Emitter added in v0.4.27

type Emitter struct {
	// contains filtered or unexported fields
}

Emitter 是 Stream 的生产者端。通常只应由一个 goroutine 使用。

func (Emitter) Close added in v0.4.27

func (emitter Emitter) Close() bool

Close 幂等关闭流。首次关闭返回 true。

func (Emitter) Emit added in v0.4.27

func (emitter Emitter) Emit(ctx context.Context, ret Result) bool

Emit 等待写入一项结果;流关闭或 ctx 取消时返回 false。

func (Emitter) IsNil added in v0.4.27

func (emitter Emitter) IsNil() bool

IsNil 报告 Emitter 是否为零值。

func (Emitter) Stream added in v0.4.27

func (emitter Emitter) Stream() Stream

Stream 返回共享状态的只读流。

func (Emitter) TryEmit added in v0.4.27

func (emitter Emitter) TryEmit(ret Result) bool

TryEmit 无阻塞地写入一项结果。

type ExecutorID added in v0.4.27

type ExecutorID uint64

ExecutorID 标识进程内的异步结果完成执行器。零值表示结果由外部执行器完成或归属未知。

ExecutorID 只用于运行时阻塞检测,不是持久化 ID,也不应跨进程传输。

func GenExecutorID added in v0.4.27

func GenExecutorID() ExecutorID

GenExecutorID 生成进程内唯一的非零异步执行器 ID。

type Future added in v0.4.27

type Future struct {
	// contains filtered or unexported fields
}

Future 是一次性异步结果的只读消费者视图。

完成结果会保存在共享状态中,Wait、TryGet 和 OnComplete 均具有重放语义;多个 消费者不会竞争同一个结果。

func All added in v0.4.27

func All(futures ...Future) Future

All 按输入顺序收集所有值;任一 Future 失败时立即失败。 空输入成功返回空 []any。

func AllSettled added in v0.4.27

func AllSettled(futures ...Future) Future

AllSettled 按输入顺序返回全部 Result,无论各项成功或失败。 空输入成功返回空 []Result。

func FirstSuccess added in v0.4.27

func FirstSuccess(futures ...Future) Future

FirstSuccess 返回第一个成功结果;所有候选均失败时返回包装 ErrNoFutureSucceeded 的错误。

func FlatMap added in v0.4.27

func FlatMap(future Future, fn func(Result) Future) Future

FlatMap 使用源结果选择下一个 Future,并将其结果展平到返回 Future。

func Map added in v0.4.27

func Map(future Future, fn func(Result) Result) Future

Map 在源 Future 完成的 goroutine 中转换结果。fn 必须快速返回。

func Race added in v0.4.27

func Race(futures ...Future) Future

Race 返回第一个完成的 Future 结果。候选为空时返回 ErrNoCandidates。

Race 只解除失败者订阅,不取消可能被其他调用者共享的生产任务。

func Rejected added in v0.4.27

func Rejected(err error) Future

Rejected 创建已经以 err 失败的 Future。

func Resolved added in v0.4.27

func Resolved(ret Result) Future

Resolved 创建已经以 ret 完成的 Future。

func Spawn added in v0.4.27

func Spawn(scope *Scope, task func(context.Context) Result) Future

Spawn 在 Scope 中启动后台任务并返回其一次性结果 Future。 panic 会转换为带堆栈的 Result.Error。

func SpawnVoid added in v0.4.27

func SpawnVoid(scope *Scope, task func(context.Context)) Future

SpawnVoid 在 Scope 中启动无业务返回值的后台任务,并返回可等待错误的 Future。

func Timeout added in v0.4.27

func Timeout(ctx context.Context, future Future, duration time.Duration) Future

Timeout 返回在源 Future、ctx 取消或 duration 到期三者中最先完成的结果。 它只解除源 Future 的订阅,不取消源任务。

func Zip2 added in v0.4.27

func Zip2(first, second Future) Future

Zip2 等待两个 Future 成功,并按参数顺序返回 Pair;任一失败时立即失败。

func (Future) CompletionExecutorID added in v0.4.27

func (future Future) CompletionExecutorID() ExecutorID

CompletionExecutorID 返回完成 Future 所依赖的执行器 ID。

func (Future) Context added in v0.4.27

func (future Future) Context(ctx context.Context) context.Context

Context 返回由 ctx 派生、并在 Future 完成时取消的上下文。

func (Future) Done added in v0.4.27

func (future Future) Done() <-chan struct{}

Done 返回 Future 完成时关闭的共享频道;零值 Future 会导致 panic。

func (Future) ID added in v0.4.27

func (future Future) ID() uint64

ID 返回 Future 的进程内诊断 ID;零值 Future 返回 0。

func (Future) IsNil added in v0.4.27

func (future Future) IsNil() bool

IsNil 报告 Future 是否为零值。

func (Future) OnComplete added in v0.4.27

func (future Future) OnComplete(callback func(Result)) Subscription

OnComplete 订阅 Future 完成。Future 已完成时 callback 会在调用者 goroutine 中立即执行。 callback 为 nil 时会导致 panic。

func (Future) TryGet added in v0.4.27

func (future Future) TryGet() (ret Result, ok bool)

TryGet 无阻塞地读取完成结果。Future 尚未完成时 ok 为 false。

func (Future) Wait added in v0.4.27

func (future Future) Wait(ctx context.Context) Result

Wait 等待 Future 完成或 ctx 取消。

nil ctx 按 context.Background 处理。结果具有重放语义。若 ctx 实现 WaitGuard, Future 会在真正阻塞前调用它,以阻止 Runtime 自等待等非法调度。

type Pair added in v0.4.27

type Pair struct {
	First  any
	Second any
}

Pair 是 Zip2 成功时保存在 Result.Value 中的二元结果。

type Promise added in v0.4.27

type Promise struct {
	// contains filtered or unexported fields
}

Promise 是一次性异步结果的生产者端。

Promise 可安全地由多个 goroutine 竞争完成;只有第一次 Resolve 成功。

func (Promise) Future added in v0.4.27

func (promise Promise) Future() Future

Future 返回与 Promise 共享状态的只读 Future。

func (Promise) IsNil added in v0.4.27

func (promise Promise) IsNil() bool

IsNil 报告 Promise 是否为零值。

func (Promise) Resolve added in v0.4.27

func (promise Promise) Resolve(ret Result) bool

Resolve 以 ret 完成 Future。首次完成返回 true,后续调用返回 false。

回调在完成者 goroutine 中、状态锁之外执行,因此回调必须快速返回;需要修改 Runtime 状态时应通过 core.ContinueOn 投递续体。

type Result added in v0.4.27

type Result struct {
	Value any   // Value 是本次产出携带的返回值。
	Error error // Error 非 nil 时表示本次产出失败。
}

Result 保存 Future 的一次产出值或错误。

func NewResult added in v0.4.27

func NewResult(value any, err error) Result

NewResult 将 value 与 err 组合为异步结果。

func (Result) OK added in v0.4.27

func (ret Result) OK() bool

OK 报告结果是否不含错误。

func (Result) String added in v0.4.27

func (ret Result) String() string

String 返回错误文本;无错误时返回 Value 的默认格式。

type Scope added in v0.4.27

type Scope struct {
	// contains filtered or unexported fields
}

Scope 将一组后台任务绑定到同一生命周期。

Scope 组合 Context 取消、关闭后拒绝注册、活动任务计数和完成信号。Close 不会 强制终止 goroutine;任务必须观察传入的 Context 才能及时退出。

func NewScope added in v0.4.27

func NewScope(parent context.Context) *Scope

NewScope 创建由 parent 控制生命周期的异步作用域。nil parent 使用 Background。

func (*Scope) AsyncScope added in v0.4.27

func (scope *Scope) AsyncScope() *Scope

AsyncScope 返回自身,使 Scope 可直接作为生命周期 Scope 提供者使用。

func (*Scope) Close added in v0.4.27

func (scope *Scope) Close() bool

Close 幂等关闭 Scope、取消 Context 并禁止注册新任务。 返回值表示本次调用是否首次关闭 Scope。

func (*Scope) Context added in v0.4.27

func (scope *Scope) Context() context.Context

Context 返回传递给所属异步任务的取消上下文。

func (*Scope) Done added in v0.4.27

func (scope *Scope) Done() Signal

Done 返回 Scope 关闭且所有已注册任务退出后完成的 Signal。

func (*Scope) Err added in v0.4.27

func (scope *Scope) Err() error

Err 返回 Scope 的取消原因;尚未关闭时返回 nil。

func (*Scope) Stats added in v0.4.27

func (scope *Scope) Stats() ScopeStats

Stats 返回 Scope 的并发安全统计快照。

type ScopeStats added in v0.4.27

type ScopeStats struct {
	Spawned   int64 // 成功注册过的任务总数。
	Active    int64 // 当前仍未退出的任务数。
	Completed int64 // 正常返回的任务总数。
	Canceled  int64 // 在 Context 已取消状态下退出的任务总数。
	Rejected  int64 // Scope 关闭后拒绝的任务总数。
	Closed    bool  // Scope 是否已经关闭。
}

ScopeStats 是异步作用域的瞬时统计快照。

type Signal added in v0.4.27

type Signal struct {
	// contains filtered or unexported fields
}

Signal 是可重放、无结果的完成通知。

func CompletedSignal added in v0.4.27

func CompletedSignal() Signal

CompletedSignal 创建已经完成的 Signal。

func (Signal) Completed added in v0.4.27

func (signal Signal) Completed() bool

Completed 报告 Signal 是否已经完成。

func (Signal) Done added in v0.4.27

func (signal Signal) Done() <-chan struct{}

Done 返回 Signal 完成时关闭的频道。

func (Signal) IsNil added in v0.4.27

func (signal Signal) IsNil() bool

IsNil 报告 Signal 是否为零值。

func (Signal) Wait added in v0.4.27

func (signal Signal) Wait(ctx context.Context) error

Wait 等待 Signal 完成或 ctx 取消,并返回取消错误。

type Stream added in v0.4.27

type Stream struct {
	// contains filtered or unexported fields
}

Stream 是连续 Result 的只读消费者视图。

Stream 采用单消费语义;多个消费者读取 Chan 会竞争元素。需要广播时应在更高层 使用事件或消息总线,而不是共享 Stream。

func (Stream) Chan added in v0.4.27

func (stream Stream) Chan() <-chan Result

Chan 返回结果频道。

func (Stream) Done added in v0.4.27

func (stream Stream) Done() <-chan struct{}

Done 返回 Stream 关闭时关闭的频道。

func (Stream) IsNil added in v0.4.27

func (stream Stream) IsNil() bool

IsNil 报告 Stream 是否为零值。

func (Stream) Next added in v0.4.27

func (stream Stream) Next(ctx context.Context) (ret Result, ok bool)

Next 等待下一项结果。流关闭时 ok 为 false;ctx 取消时返回包含 ctx.Err 的结果。

type Subscription added in v0.4.27

type Subscription struct {
	// contains filtered or unexported fields
}

Subscription 表示一个可取消的 Future 完成订阅。

func (Subscription) Cancel added in v0.4.27

func (subscription Subscription) Cancel() bool

Cancel 取消尚未执行的订阅。成功移除返回 true;已完成、已取消或零值返回 false。

type WaitGuard added in v0.4.27

type WaitGuard interface {
	BeforeFutureWait(futureID uint64, completionExecutorID ExecutorID) error
	AfterFutureWait(futureID uint64)
}

WaitGuard 允许执行上下文在 Future 进入阻塞等待前实施调度约束。 Runtime Context 使用该接口阻止 Actor 执行协程等待自身队列产生的结果。

Jump to

Keyboard shortcuts

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