memorybroker

package
v0.2.0 Latest Latest
Warning

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

Go to latest
Published: Jul 16, 2026 License: Apache-2.0 Imports: 8 Imported by: 0

Documentation

Overview

Package memorybroker 是 Broker 的内存参考实现:单进程、单 sync.Mutex + sync.Cond。 它是三后端的"语义基准":所有状态流转都在同一个锁临界区内完成, 等价于 sqlite 的"同一个事务"——终态更新和子任务唤醒天然原子,不可能丢唤醒。 brokertest 的 18 条契约以它的行为为准。

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Broker

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

Broker 内存后端。所有读写都拿同一把锁,Cond 用来唤醒阻塞中的 Dequeue。

func New

func New() *Broker

New 构造一个空的内存后端;用之前必须先 Init(由 taskgate.New(cfg) 统一调用)。

func (*Broker) Ack

func (b *Broker) Ack(ctx context.Context, id, leaseToken string, result []byte) error

Ack 成功完结:completed + Result + FinishedAt,并在同一临界区内唤醒子任务。

func (*Broker) Cancel

func (b *Broker) Cancel(ctx context.Context, id string) error

Cancel 取消:排队类状态直接 canceled 并传播;running 只打标记(终态由 FinishCanceled 落); 终态报 ErrAlreadyFinal,不存在报 ErrTaskNotFound。

func (*Broker) Close

func (b *Broker) Close() error

Close 关闭:标记 closed 并广播,让所有阻塞中的 Dequeue 尽快退出。内存后端没有别的资源要释放。

func (*Broker) Counts

func (b *Broker) Counts(ctx context.Context) (map[string]map[taskgate.Status]int64, error)

Counts 出现过的 Type×Status 稀疏矩阵,和逐个 Get 汇总必须一致(brokertest 验证)。

func (*Broker) Dequeue

func (b *Broker) Dequeue(ctx context.Context, queues []string) (*taskgate.Task, error)

Dequeue 阻塞认领:直到某队列出现就绪任务,或 ctx 取消(返回 ctx.Err())。 认领本身原子:置 running、发新令牌、记租约、首次写 StartedAt。

func (*Broker) Enqueue

func (b *Broker) Enqueue(ctx context.Context, t *taskgate.Task) error

Enqueue 入队。同一临界区内完成:查重、父存在性校验、初始状态判定、 登记依赖反向索引;生成的 ID 会回填到 t.ID。

func (*Broker) Fail

func (b *Broker) Fail(ctx context.Context, id, leaseToken, errMsg string, kind taskgate.FailKind, retryAt time.Time) error

Fail 失败路径:按 FailKind 动对应计数,封顶或耗尽进 failed(触发传播),否则 retrying。

func (*Broker) FinishCanceled

func (b *Broker) FinishCanceled(ctx context.Context, id, leaseToken string) error

FinishCanceled worker 响应取消后收尾:running → canceled 落库并传播。

func (*Broker) Get

func (b *Broker) Get(ctx context.Context, id string) (*taskgate.Task, error)

Get 取单个任务的副本。

func (*Broker) Heartbeat

func (b *Broker) Heartbeat(ctx context.Context, id, leaseToken string) error

Heartbeat 续租:lease_until = now + TTL。发现取消标记时续租照做, 但返回 ErrTaskCanceled 提醒 scheduler 去 cancel handler 的 ctx。

func (*Broker) Init

func (b *Broker) Init(opts taskgate.BrokerOptions) error

Init 装配运行参数,零值补默认(TTL 60s / LeaseLostMax 3 / ThrottledMax 100 / 真时钟)。

func (*Broker) List

func (b *Broker) List(ctx context.Context, f taskgate.Filter) ([]*taskgate.Task, error)

List 按 Filter 过滤,零值字段不过滤;先过滤 → 按 (CreatedAt, ID) 升序 → 跳过 Offset 再取 Limit(排序分页合同见 broker-contract.md,Offset 越界返回空)。

func (*Broker) QueueLen

func (b *Broker) QueueLen(ctx context.Context, queue string) (int, error)

QueueLen 队列积压:status∈{pending,retrying} 的数量(不看 RunAt 到没到点)。

func (*Broker) QueueQuota

func (b *Broker) QueueQuota(queue string, qc taskgate.QueueConfig) (taskgate.QuotaGate, error)

QueueQuota 构造该队列的配额闸。介质是本进程内存,计数与任务共用同一把大锁。

func (*Broker) ReapExpired

func (b *Broker) ReapExpired(ctx context.Context) (int, error)

ReapExpired 回收过期租约:带取消标记的直接落 canceled(不占 LeaseLost,触发传播); 其余 LeaseLost+1,封顶进 failed(触发传播),否则回 pending。 顺带做防御性修复:blocked 但父实际全部终态的任务,按正常规则补唤醒/补取消 (这不是正常路径,是给"唤醒中途崩"这类事故兜底)。返回值只算租约回收条数。

func (*Broker) Replay

func (b *Broker) Replay(ctx context.Context, req taskgate.ReplayRequest) (*taskgate.Task, error)

Replay 重放一次终态执行(spec 005):校验目标(终态/未被重放/completed 需显式允许) 与创建新执行在同一个锁临界区内原子完成——并发同目标重放恰好一个成功。 目标记录零改写;新执行沿用目标的 Type/Queue/MaxRetry/OnParentFailure 与 BusinessKey, 三计数清零、无依赖、pending 落库。

func (*Broker) Requeue

func (b *Broker) Requeue(ctx context.Context, id, leaseToken string) error

Requeue 优雅停机时归还任务:running → pending,三计数与 RunAt 全不动, 清租约和取消标记。这不算失败,一个计数都不许占(合同要求)。

Jump to

Keyboard shortcuts

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