concurrent

package
v0.3.68 Latest Latest
Warning

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

Go to latest
Published: Aug 14, 2026 License: LGPL-2.1 Imports: 10 Imported by: 2

Documentation

Overview

Package concurrent 提供轻量级并发控制辅助组件。

当前主要包含两类工具:

  • FutureController:用于管理带超时的异步请求和响应匹配
  • Listener/Listeners:用于维护监听者集合并做广播分发

这些组件主要服务于协议层、事件分发和异步协作场景。

Index

Constants

This section is empty.

Variables

View Source
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 返回用于匹配响应的唯一标识;零值不会被分配。

func (*FutureHandle) Resolve added in v0.3.68

func (h *FutureHandle) Resolve(ret async.Result) error

Resolve 以 ret 完成 Future;已完成或超时时返回 ErrFutureExceeded。

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

func NewListener[H, M any](handler H, size int) *Listener[H, M]

NewListener 创建携带 handler、Inbox 容量为 size 的监听者;size 为负数时 panic。

type Listeners added in v0.3.68

type Listeners[H, M any] atomic.Pointer[[]*Listener[H, M]]

Listeners 以原子写时复制快照维护监听者集合。 添加、删除、加载和广播可并发调用;集合开始使用后不得复制。

func NewListeners added in v0.3.68

func NewListeners[H, M any]() *Listeners[H, M]

NewListeners 创建空监听者集合;Listeners 的零值同样可直接使用。

func (*Listeners[H, M]) Add added in v0.3.68

func (ls *Listeners[H, M]) Add(handler H, size int) *Listener[H, M]

Add 创建并加入监听者,返回值用于接收消息和后续 Delete;相同 handler 可重复添加。

func (*Listeners[H, M]) Broadcast added in v0.3.68

func (ls *Listeners[H, M]) Broadcast(m M) (rejected int)

Broadcast 尝试以非阻塞方式向当前快照中的每个 Inbox 投递 m,并返回因通道已满而拒绝的数量。 调用方不得关闭仍在集合或并发广播旧快照中的 Inbox,否则发送会 panic。

func (*Listeners[H, M]) Delete added in v0.3.68

func (ls *Listeners[H, M]) Delete(l *Listener[H, M])

Delete 从后续快照中删除指定监听者,但不关闭其 Inbox。 已取得旧快照的并发广播仍可能向该监听者投递一次消息。

func (*Listeners[H, M]) Load added in v0.3.68

func (ls *Listeners[H, M]) Load() []*Listener[H, M]

Load 返回当前监听者快照;调用方必须将返回切片视为只读。

Jump to

Keyboard shortcuts

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