Documentation
¶
Overview ¶
Package concurrent 提供轻量级并发控制辅助组件。
当前主要包含两类工具:
- FutureController:用于管理带超时的异步请求和响应匹配
- Listener/Listeners:用于维护监听者集合并做广播分发
这些组件主要服务于协议层、事件分发和异步协作场景。
Index ¶
Constants ¶
This section is empty.
Variables ¶
var ( // ErrFutureControllerClosed 表示控制器的父 context 已结束,不再接受或等待任务。 ErrFutureControllerClosed = errors.New("future controller closed") // ErrFutureExceeded 表示 Future 已超过截止时间、已完成或不存在。 ErrFutureExceeded = errors.New("future exceeded deadline") )
Functions ¶
This section is empty.
Types ¶
type FutureController ¶ added in v0.3.68
type FutureController struct {
// contains filtered or unexported fields
}
FutureController 并发管理带唯一 ID 和统一超时时长的 Future,并将 Resolve 与等待方匹配。 控制器初始化后不得复制。
func NewFutureController ¶ added in v0.3.68
func NewFutureController(ctx context.Context, timeout time.Duration) *FutureController
NewFutureController 创建以 ctx 控制生命周期、以 timeout 设置每个 Future 截止时间的控制器。 ctx 为 nil 时使用 context.Background;调用方应取消 ctx 以终止内部 watcher。
func (*FutureController) New ¶ added in v0.3.68
func (fc *FutureController) New() (*FutureHandle, error)
New 注册一个待完成的 Future;控制器已结束时返回 ErrFutureControllerClosed。
func (*FutureController) Resolve ¶ added in v0.3.68
func (fc *FutureController) Resolve(id int64, ret async.Result) error
Resolve 以 ret 完成指定 Future。 ID 不存在、Future 已完成或已到截止时间时返回 ErrFutureExceeded。
func (*FutureController) Terminated ¶ added in v0.3.68
func (fc *FutureController) Terminated() async.Signal
Terminated 返回控制器结束信号;父 context 取消且所有待处理项完成收尾后该信号才会完成。
type FutureHandle ¶ added in v0.3.68
type FutureHandle struct {
// contains filtered or unexported fields
}
FutureHandle 标识由 FutureController 管理的一次异步结果等待。 Handle 初始化后不得复制,完成、取消与超时之间只会有一个结果生效。
func (*FutureHandle) Cancel ¶ added in v0.3.68
func (h *FutureHandle) Cancel(err error)
Cancel 尝试以 err 完成 Future;Future 已结束时调用无效果,失败不会返回给调用方。
func (*FutureHandle) Deadline ¶ added in v0.3.68
func (h *FutureHandle) Deadline() time.Time
Deadline 返回该 Future 的截止时间。
func (*FutureHandle) Future ¶ added in v0.3.68
func (h *FutureHandle) Future() async.Future
Future 返回用于等待结果的只读 Future。
func (*FutureHandle) Id ¶ added in v0.3.68
func (h *FutureHandle) Id() int64
Id 返回用于匹配响应的唯一标识;零值不会被分配。
type Listener ¶ added in v0.3.68
type Listener[H, M any] struct { Handler H // 与监听者关联的处理器或元数据。 Inbox chan M // 接收广播消息的有缓冲通道。 // contains filtered or unexported fields }
Listener 将调用方定义的 handler 与接收广播消息的 Inbox 关联起来。 Listener 初始化后不得复制;Inbox 的消费由调用方负责。
func NewListener ¶ added in v0.3.68
NewListener 创建携带 handler、Inbox 容量为 size 的监听者;size 为负数时 panic。
type Listeners ¶ added in v0.3.68
Listeners 以原子写时复制快照维护监听者集合。 添加、删除、加载和广播可并发调用;集合开始使用后不得复制。
func NewListeners ¶ added in v0.3.68
NewListeners 创建空监听者集合;Listeners 的零值同样可直接使用。
func (*Listeners[H, M]) Broadcast ¶ added in v0.3.68
Broadcast 尝试以非阻塞方式向当前快照中的每个 Inbox 投递 m,并返回因通道已满而拒绝的数量。 调用方不得关闭仍在集合或并发广播旧快照中的 Inbox,否则发送会 panic。