Documentation
¶
Overview ¶
Package memorybroker 是 Broker 的内存参考实现:单进程、单 sync.Mutex + sync.Cond。 它是三后端的"语义基准":所有状态流转都在同一个锁临界区内完成, 等价于 sqlite 的"同一个事务"——终态更新和子任务唤醒天然原子,不可能丢唤醒。 brokertest 的 18 条契约以它的行为为准。
Index ¶
- type Broker
- func (b *Broker) Ack(ctx context.Context, id, leaseToken string, result []byte) error
- func (b *Broker) Cancel(ctx context.Context, id string) error
- func (b *Broker) Close() error
- func (b *Broker) Counts(ctx context.Context) (map[string]map[taskgate.Status]int64, error)
- func (b *Broker) Dequeue(ctx context.Context, queues []string) (*taskgate.Task, error)
- func (b *Broker) Enqueue(ctx context.Context, t *taskgate.Task) error
- func (b *Broker) Fail(ctx context.Context, id, leaseToken, errMsg string, kind taskgate.FailKind, ...) error
- func (b *Broker) FinishCanceled(ctx context.Context, id, leaseToken string) error
- func (b *Broker) Get(ctx context.Context, id string) (*taskgate.Task, error)
- func (b *Broker) Heartbeat(ctx context.Context, id, leaseToken string) error
- func (b *Broker) Init(opts taskgate.BrokerOptions) error
- func (b *Broker) List(ctx context.Context, f taskgate.Filter) ([]*taskgate.Task, error)
- func (b *Broker) QueueLen(ctx context.Context, queue string) (int, error)
- func (b *Broker) QueueQuota(queue string, qc taskgate.QueueConfig) (taskgate.QuotaGate, error)
- func (b *Broker) ReapExpired(ctx context.Context) (int, error)
- func (b *Broker) Replay(ctx context.Context, req taskgate.ReplayRequest) (*taskgate.Task, error)
- func (b *Broker) Requeue(ctx context.Context, id, leaseToken string) error
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 (*Broker) Cancel ¶
Cancel 取消:排队类状态直接 canceled 并传播;running 只打标记(终态由 FinishCanceled 落); 终态报 ErrAlreadyFinal,不存在报 ErrTaskNotFound。
func (*Broker) Dequeue ¶
Dequeue 阻塞认领:直到某队列出现就绪任务,或 ctx 取消(返回 ctx.Err())。 认领本身原子:置 running、发新令牌、记租约、首次写 StartedAt。
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 ¶
FinishCanceled worker 响应取消后收尾:running → canceled 落库并传播。
func (*Broker) Heartbeat ¶
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 ¶
List 按 Filter 过滤,零值字段不过滤;先过滤 → 按 (CreatedAt, ID) 升序 → 跳过 Offset 再取 Limit(排序分页合同见 broker-contract.md,Offset 越界返回空)。
func (*Broker) QueueQuota ¶
QueueQuota 构造该队列的配额闸。介质是本进程内存,计数与任务共用同一把大锁。
func (*Broker) ReapExpired ¶
ReapExpired 回收过期租约:带取消标记的直接落 canceled(不占 LeaseLost,触发传播); 其余 LeaseLost+1,封顶进 failed(触发传播),否则回 pending。 顺带做防御性修复:blocked 但父实际全部终态的任务,按正常规则补唤醒/补取消 (这不是正常路径,是给"唤醒中途崩"这类事故兜底)。返回值只算租约回收条数。