unique

package module
v0.1.0 Latest Latest
Warning

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

Go to latest
Published: Aug 10, 2026 License: MIT Imports: 6 Imported by: 0

README

unique

单会话仲裁库:同一 (project, env, uid) 任一时刻恰有一个在线会话。新会话登录即踢掉旧会话(先踢后写),靠 gen fencing 在存储单键线性化点上把并发对撞压成「恰一 winner」。

纯 Go 公共库,零侵入:服务端核心是 Arbiter(你自挂 RPC/HTTP 框架),客户端核心网络层是 interface(不强制 HTTP)。你只补协议通信存储实现两块。

特性

  • 先踢后写四环:强读 → 重放短路 → 踢旧(gen 预告)→ 三形 CAS 围栏条件写
  • gen fencing:gen 单调不回退,CAS 恰一成功,预算耗尽 fail-closed
  • 重放幂等:同 claim_id 重试短路返回既有 gen,不误踢不写新
  • 双通道逐出:主动 kick(快)+ 心跳自逐(对),任一失效仍收敛
  • 可插拔:Store / Kicker / Transport / Observer 全是 interface
  • 参考实现:DynamoDB Store、内存 Store、HTTP Kicker、HTTP Transport、stdlog Observer

快速开始

// 服务端
arb := unique.New(memory.New(), nil) // 测试用内存 Store + fail-open 踢
out, _ := arb.Claim(ctx, unique.ClaimInput{
    Key:     unique.Key{Project: "game", Env: "prod", UID: 42},
    Owner:   unique.Owner{Endpoint: "https://node-a/kick"},
    ClaimID: "cid-abc",
})
// out.Gen == 1(首登)

// 客户端(网络层以参考 HTTP Transport 为例)
tr, _  := httptransport.New("https://unique-server")
cl     := client.NewClaimer(tr)
resp, _ := cl.Login(ctx, key, owner)   // 含 claim_id 纪律 + 退避重试

kh, _ := client.NewKickHandler(getLocalGen, destroy, client.KickOptions{}) // 被踢接收
hb := client.NewHeartbeat(tr, client.WithOnEvict(onEvict))                 // 心跳自逐
go hb.Run(ctx, key, resp.Gen)

端到端可跑示例见 example/;真实 HTTP 进程冒烟(S1–S6)见 scripts/smoke/

架构

        客户端                     服务端                     存储
  ┌──────────────────┐     ┌───────────────────┐     ┌────────────────┐
  │ Claimer          │     │ Arbiter           │     │ Store          │
  │  Login(claim_id) │────▶│  Claim 四环        │────▶│  ClaimCAS 恰一  │
  │  退避重试 ≤25s    │     │  重放短路/踢/CAS    │     │  gen fencing   │
  │ Heartbeat 自逐    │◀────│  Current          │────▶│  强一致读       │
  │ KickHandler 被踢  │◀────│  Kicker(先踢)      │     │                │
  └──────────────────┘     └───────────────────┘     └────────────────┘
        Transport iface           Kicker iface              Store iface
        (HTTP 为参考)           (HTTP 为参考)          (DynamoDB 参考)

仲裁时序:新会话 Claim → Arbiter 强读当前 owner → 同 claim_id 则短路返回 → 否则按旧 gen+1 预告踢旧端(fail-open)→ CAS 写「gen+1、新 owner、新 cid」,条件不满足(被人抢先)则重读重踢重试,预算 3 次耗尽返回 ErrUnavailable

三模块导航

你要接入的角色 读这份
服务端:挂 Arbiter 到你的框架、选 Store、配 Kicker/Observer SERVER.md
存储:实现新 Store 后端(正确性契约全在这) STORE.md
客户端:实现 Transport、接 Claimer/Heartbeat/KickHandler CLIENT.md

包布局

unique/                核心类型(Key/Owner/OwnerState)、哨兵错、Arbiter、Kicker、Observer
store/memory/          内存 Store(测试/单进程)
store/dynamodb/        DynamoDB Store(参考生产实现)+ storecontract 契约
store/storecontract/   Store 可复用契约测试套件
kickhttp/              HTTP Kicker(+ kickercontract 契约)
client/                客户端核心:Claimer / Heartbeat / KickHandler / Transport iface
client/httptransport/  参考 HTTP Transport
client/kickrecv/       被踢 HTTP 接收适配(恒 200)
observer/stdlog/       参考 stdlog Observer
example/               端到端示例(服务端薄适配 + 双客户端对撞)
scripts/smoke/         真实 HTTP 进程冒烟 S1–S6

不变量(I1–I6)

单键线性化恰一 · gen 单调不回退 · 重放幂等 · CAS 耗尽 fail-closed · 踢失败 fail-open · 服务端 error 恒不误逐。详见 SERVER.md

License

LICENSE

Documentation

Overview

Package unique 提供「每用户仅一个活跃会话、新登录顶替旧登录」的单会话仲裁能力。

核心机制:先踢后写四环 + gen fencing

每个用户一条归属项,键为 Key{Project, Env, UID},含三个属性:

  • gen:单调递增世代号(fencing token),全生命周期严格单调、无 TTL 永久留存。
  • owner:当前会话归属(JSON),tombstone(登出)时为空串。
  • cid:最近一次成功 Claim 的幂等键(claim_id)。

Claim 走「先踢后写」四环(见 Arbiter.Claim):

  1. ReadOwner:强一致读旧像——重放判定、踢人决策与 CAS 围栏的共同基线。
  2. 重放判定:snap.HasCID ∧ snap.ClaimID==claimID ∧ snap.Owner!=nil → 短路返回既有 gen(幂等)。
  3. 先踢旧会话:按 gen 围栏踢更旧会话,阻塞到销毁完成或拟制下线(fail-open)。
  4. ClaimCAS:三形围栏条件写(按 snap.HasCID 选形),并发同键至多一个成功。

不变量

  • I1 单键串行化:同键操作按到达顺序线性化(由 Store 实现保证)。
  • I2 gen 单调:gen 全生命周期严格单调、永不重置。
  • I3 恰一 winner:并发 ClaimCAS 至多一个成功。
  • I4 重放幂等:同 claim_id 重试短路返回既有 gen。
  • I5a 严格互斥:2xx 确认顶替,旧确认下线后新才上线。
  • I5b 有界重叠:超时拟制放行,重叠 ≤ kickTimeout + 心跳间隔,自逐兜底收敛。
  • I6 gen 围栏:踢人与自逐均按 gen 判定,仅踢/逐更旧会话。

三模块

  • 服务端:Arbiter(本包)+ Store / Kicker /Observer 三个 interface。
  • 存储:store/memory(测试)、store/dynamodb(参考实现);语义契约见 doc/STORE.md。
  • 客户端:client 包(Transport interface / Claimer / Heartbeat / KickHandler), 网络层不强制 HTTP;client/httptransport 为参考实现。

Package unique 提供先踢后写 + gen fencing 单会话仲裁的核心库。

Index

Constants

View Source
const (
	// Arbiter 核心合成
	KickNoEndpoint         = "no-endpoint"          // 旧会话无 endpoint / owner 损坏,无从踢起(拟制下线)
	KickSelfReplaceSkipped = "self-replace-skipped" // 旧会话 endpoint 与本端相同,跳过自踢(不发起 HTTP)

	// Kicker 返回(三类传输结果)
	KickTimeout      = "timeout"       // 踢人整体超时(拟制下线)
	KickConnectError = "connect-error" // 连接失败(拟制下线)
	KickNon2xx       = "non-2xx"       // 被调方返回非 2xx(拟制下线)
)

kick_detail 脱敏类别。五类;前两类由 Arbiter 核心合成,后三类由 Kicker 返回。 Kicked=true 时 KickDetail == ""。

Variables

View Source
var (
	// ErrInvalidParams 入参非法(控制字符/超长/Project|Env 含 "#/零值/坏 endpoint scheme)。
	// 客户端立即失败不重试。
	ErrInvalidParams = errors.New("unique: invalid params")

	// ErrThrottled 存储节流。客户端以同一 claim_id 退避重试。
	ErrThrottled = errors.New("unique: throttled")

	// ErrUnavailable 存储故障或 CAS 预算耗尽。fail-closed:绝不放行登录,
	// 客户端以同一 claim_id 退避重试。
	ErrUnavailable = errors.New("unique: unavailable")

	// ErrCASFailed 围栏条件失败(并发竞争的正常结果,非存储故障)。
	// 服务端内部消化:重读、重判、重踢后重试(预算封顶)。
	ErrCASFailed = errors.New("unique: claim CAS condition failed (concurrent mutation)")
)

哨兵错。使用方适配层负责把哨兵错映射回自己的协议错误码。

AllKickDetails 全部五类,供校验/遍历。

Functions

func IsCASFailed

func IsCASFailed(err error) bool

IsCASFailed 判定错误是否为围栏条件失败(正常竞争,非故障)。

func IsInvalidParams

func IsInvalidParams(err error) bool

IsInvalidParams 判定错误是否为入参非法。

func IsPresumedOffline

func IsPresumedOffline(detail string) bool

IsPresumedOffline 判定 detail 是否属拟制下线(fail-open 放行路径)。 self-replace-skipped 不属拟制(同端会话已由新会话接管,无残留)。

func IsThrottled

func IsThrottled(err error) bool

IsThrottled 判定错误是否为存储节流。

func IsUnavailable

func IsUnavailable(err error) bool

IsUnavailable 判定错误是否为存储故障 / CAS 预算耗尽。

Types

type Arbiter

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

Arbiter 服务端核心仲裁器:先踢后写四环 + gen fencing 单会话仲裁。 零框架依赖:存储走 Store、踢人走 Kicker、观测走 Observer,全部注入。

func New

func New(store Store, kicker Kicker, opts ...Option) *Arbiter

New 构造 Arbiter。store 必填;kicker 为 nil 时踢人走 fail-open(视为已踢放行的兜底, 仅供测试/无回调场景)。依赖全注入、无进程级 init。

func (*Arbiter) Claim

func (a *Arbiter) Claim(ctx context.Context, in ClaimInput) (ClaimOutcome, error)

Claim 先踢后写四环:强一致读 → 重放判定 → 踢旧 → CAS 条件写。 CAS 冲突重读重踢重试,预算耗尽返回 ErrUnavailable(fail-closed)。

func (*Arbiter) Current

func (a *Arbiter) Current(ctx context.Context, key Key) (CurrentView, error)

Current 自逐基线读。委托 Store.Current;损坏 owner 硬失败映射 ErrUnavailable (调用方保留会话不误逐)。

func (*Arbiter) Release

func (a *Arbiter) Release(ctx context.Context, key Key, gen int64) (ReleaseResult, error)

Release 干净登出 tombstone。委托 Store.Release(gen 围栏)。

type ClaimInput

type ClaimInput struct {
	Key     Key
	Owner   Owner
	ClaimID string // 每次登录尝试一个新 id;该次尝试的网络重试复用同一 id
}

ClaimInput Claim 入参。

type ClaimOutcome

type ClaimOutcome struct {
	Gen        int64  // 新世代号 / fencing token
	PrevGen    int64  // 上一 gen
	HasPrev    bool   // 是否存在被顶替者
	PrevOwner  *Owner // 被顶替 owner
	Kicked     bool   // 旧会话确认下线
	KickDetail string // 脱敏类别(kick_detail 五类之一;Kicked=true 时为 "")
}

ClaimOutcome Claim 出参。

type CurrentView

type CurrentView struct {
	Gen    int64
	Owner  *Owner // tombstone 时 nil,须与「无记录」(Exists=false)可区分
	Exists bool
}

CurrentView Current 出参。自逐三件套最小完备集。

type EvictReason

type EvictReason int

EvictReason 自逐原因。

const (
	EvictNone        EvictReason = iota // 不自逐
	EvictNotExists                      // 远端无记录
	EvictGenMismatch                    // 远端 gen ≠ 本地 gen
	EvictOwnerEmpty                     // owner 为空(tombstone / 被顶替清场)
)

func (EvictReason) String

func (r EvictReason) String() string

type Key

type Key struct {
	Project string
	Env     string
	UID     int64
}

Key 归属项寻址三元组。同 UID 跨 Project/Env 互不踢(多租户/多环境隔离)。

func (Key) String

func (k Key) String() string

String 拼存储键 "project#env#uid"。Project/Env 禁含 "#"(见 validate),否则两不同三元组会拼出同键互踢。

type KickPayload

type KickPayload struct {
	Project string `json:"project"`
	Env     string `json:"env"`
	UID     int64  `json:"uid"`
	Gen     int64  `json:"gen"`
}

KickPayload / KickRequest 踢人回调载荷。 线上 JSON 字段名固化 {project, env, uid, gen}(与存量 kick 接收方互通红线)。 Gen 为预告值(旧像 gen+1,CAS 前回调),被调方按 gen 围栏仅踢更旧会话。

type KickRequest

type KickRequest = KickPayload

KickRequest 与 KickPayload 同形(被调方视角)。

type Kicker

type Kicker interface {
	// Kick 向 endpoint 指定的旧会话发踢出通知。实现须尊重 ctx 截止期。
	Kick(ctx context.Context, endpoint string, p KickPayload) (kicked bool, detail string)
}

Kicker 负责向旧会话发起踢出通知。实现方自行决定传输(HTTP/gRPC/MQ…)。

语义契约(Kicker 只上报三类传输结果,业务判定由 Arbiter 核心合成):

  • 踢成功(旧会话确认收到 / 已被移除): 返回 (true, "")
  • 传输失败: 返回 (false, detail),detail ∈ {KickTimeout, KickConnectError, KickNon2xx}

Kicker 不返回 KickNoEndpoint / KickSelfReplaceSkipped —— 这两类由 Arbiter 在 判定旧 owner 无 endpoint 或新旧 endpoint 相同(自顶替免疫)时自行合成,不经 Kicker。

endpoint 单独传参(不在 KickPayload 内):载荷仅固化 {project,env,uid,gen} 四字段, 寻址信息由 Arbiter 从旧 owner 取出注入。kicked=true 时 detail 必为 ""。

fail-open:kicked=false 且 detail 属拟制离线集合(见 IsPresumedOffline)时, Arbiter 视为旧会话已失联、拟制放行继续 CAS 抢占。实现不返回 Go error —— 传输类失败一律编码进 detail;只有参数非法等不应重试的致命情况才可在 实现内部记日志后归为 connect-error,接口层面不暴露 error。

type NoopObserver

type NoopObserver struct{}

NoopObserver 空实现:无分配、无格式化,零开销默认值。

func (NoopObserver) OnClaim

func (NoopObserver) OnKick

func (NoopObserver) OnKick(string, bool)

func (NoopObserver) OnRelease

func (NoopObserver) OnRelease(Key, bool)

type Observer

type Observer interface {
	// OnClaim 一次 Claim 编排结束(无论成败)触发。attempts=CAS 尝试次数,took=全程耗时。
	OnClaim(outcome ClaimOutcome, attempts int, took time.Duration)
	// OnRelease 一次 Release 触发。released=是否真正置 tombstone。
	OnRelease(key Key, released bool)
	// OnKick 按 Kick 调用次触发。detail=kick 脱敏类别(含核心合成的 no-endpoint/
	// self-replace-skipped),success=kicked。
	OnKick(detail string, success bool)
}

Observer 观测钩子:把仲裁关键事件暴露给使用方(指标/审计/日志)。

契约:

  • 实现须非阻塞、不得 panic(钩子绝不能反过来影响仲裁路径)。
  • 所有方法按调用次触发;OnKick 按 Kick 调用次触发(CAS 重试重踢每次各触发一次)。
  • Arbiter 在关键路径同步调用,实现内部若需 I/O 应自行异步化。

type Option

type Option func(*Arbiter)

Option 自定义 Arbiter。

func WithMaxClaimAttempts

func WithMaxClaimAttempts(n int) Option

WithMaxClaimAttempts 设 CAS 重试预算(含首次)。默认 3;<1 clamp 到 1。

func WithObserver

func WithObserver(o Observer) Option

WithObserver 注入观测钩子(默认 NoopObserver)。

func WithValidationLimits

func WithValidationLimits(l ValidationLimits) Option

WithValidationLimits 覆盖入参长度上限(默认 ValidationLimits{})。

type Owner

type Owner struct {
	Region   string `json:"region,omitempty"`
	Cluster  string `json:"cluster,omitempty"`
	Endpoint string `json:"endpoint,omitempty"`
}

Owner 会话归属。Endpoint 为踢人回调地址(必需),Region/Cluster 仅诊断。 JSON tag 用 proto3 camelCase 以兼容存量数据(否则存量记录 OwnerCorrupt)。

type OwnerState

type OwnerState struct {
	Exists       bool
	Gen          int64
	Owner        *Owner
	OwnerCorrupt bool   // owner 非空但不可解析(软路径,不算读错误)
	ClaimID      string // 最近一次成功 Claim 的幂等键
	HasCID       bool   // 区分「无 cid」与「cid 为空串」(选形与重放守卫)
}

OwnerState ReadOwner 读回的旧像快照。

type ReleaseResult

type ReleaseResult struct {
	Released bool
}

ReleaseResult Release 出参。

type Store

type Store interface {
	// ReadOwner 强一致读归属快照(重放判定、踢人决策与 CAS 围栏的共同基线,只读)。
	// 无记录返回 Exists=false 的非 nil 快照,非 error。
	// owner 非空但不可解析时不算读错误,置 OwnerCorrupt=true, Owner=nil(软路径)。
	// cid 与 HasCID 必须区分「无 cid 属性」与「cid 为空串」:本库写出的记录必有 cid(HasCID=true);
	// tombstone 与首登前无 cid(HasCID=false)。
	ReadOwner(ctx context.Context, key Key) (*OwnerState, error)

	// ClaimCAS 以 snap 为围栏基线原子抢占归属,单条原子写:
	// 条件校验 ∧ gen 自增 1 ∧ owner/cid 同写,三者在同一线性化点不可分割。
	//
	// 选形按 snap.HasCID(非按 snap.Owner,避免 OwnerCorrupt 记录 Owner=nil 时
	// 误选 tombstone 形恒假导致永久锁死):
	//   - !snap.Exists      → 首登形:记录必须仍不存在(attribute_not_exists(gen))。
	//   - snap.HasCID       → 双钉形:gen == snap.Gen 且 cid == snap.ClaimID
	//     (顶替在线会话与 OwnerCorrupt 记录,容许覆盖写防锁死)。
	//   - Exists 且无 cid   → tombstone 形:gen == snap.Gen 且无 cid。
	//
	// 条件失败返回 ErrCASFailed(可 errors.Is 判定),且无任何写入;
	// 其他错误不得伪装成 ErrCASFailed(节流/故障另行分流)。
	// 成功返回 newGen = 基线 gen + 1(首登为 1),返回值即落库值,调用方不再读回。
	// 并发同键至多一个成功。本原语不做踢人、不解释 claimID 语义。
	// owner 序列化失败(含需要时无 endpoint)必须在写之前显式报错,不静默写空串。
	ClaimCAS(ctx context.Context, key Key, owner Owner, claimID string, snap *OwnerState) (newGen int64, err error)

	// Current 强一致读当前权威 gen 与 owner(调用方自逐基线,只读)。
	// 无记录(含 gen==0)返回 exists=false。
	// tombstone(owner 为空)返回 exists=true, owner=nil —— 必须与「无记录」可区分。
	// owner 已存但损坏返回显式 error(调用方据此保留会话不误逐),不得静默降级。
	Current(ctx context.Context, key Key) (gen int64, owner *Owner, exists bool, err error)

	// Release 干净登出(tombstone):仅当当前 gen == 传入 gen 才
	// 原子地置空 owner、删除 cid,保留 gen(全生命周期单调,禁止重置/回收)。
	// gen 不匹配、无记录或 gen<=0:released=false, err=nil(无操作,非错误)。
	// 条件校验与 SET owner/REMOVE cid 必须在同一线性化点生效。
	Release(ctx context.Context, key Key, gen int64) (released bool, err error)
}

Store 归属注册表存储契约。所有方法按 Key 单键寻址。

全局不变量(实现必须满足)

  • 单键线性化:同键操作按到达顺序串行生效,条件校验见到的是最近一次已提交的版本。
  • 强一致读:ReadOwner/Current 读最近一次已提交写,禁止过期快照/副本读。
  • 无 TTL:归属项永久留存,不得设过期回收 gen(gen 全生命周期严格单调,域 [1, MaxInt64))。
  • 单键设计:不得引入跨键事务/全局锁(水平扩展前提)。

语义契约完整版见 doc/STORE.md(能力清单 C1-C8 + 陷阱清单)。

type ValidationLimits

type ValidationLimits struct {
	Tenant   int // Project / Env 各自上限
	ClaimID  int
	Endpoint int
}

ValidationLimits 入参长度上限。Tenant 同时约束 Project 与 Env 各自长度。

Directories

Path Synopsis
Package client 提供先踢后写单会话仲裁的客户端核心。
Package client 提供先踢后写单会话仲裁的客户端核心。
httptransport
Package httptransport 提供 client.Transport 的 HTTP 参考实现(Wire Contract)。
Package httptransport 提供 client.Transport 的 HTTP 参考实现(Wire Contract)。
kickrecv
Package kickrecv 提供 KickHandler 的 net/http 参考适配器(被踢接收端)。
Package kickrecv 提供 KickHandler 的 net/http 参考适配器(被踢接收端)。
Package main 演示 unique 的端到端用法: 服务端薄适配(挂 Arbiter 到 HTTP)+ 双客户端对撞(顶替 + 被踢 + 心跳自逐)。
Package main 演示 unique 的端到端用法: 服务端薄适配(挂 Arbiter 到 HTTP)+ 双客户端对撞(顶替 + 被踢 + 心跳自逐)。
Package kickhttp 提供 Kicker 的 HTTP 参考实现。
Package kickhttp 提供 Kicker 的 HTTP 参考实现。
kickercontract
Package kickercontract 提供 Kicker 分类契约测试套件:任意 Kicker 实现可跑同一组 断言,验证其把传输结果正确归类为三类 detail(timeout/connect-error/non-2xx), 且 2xx → kicked=true 且 detail=""。
Package kickercontract 提供 Kicker 分类契约测试套件:任意 Kicker 实现可跑同一组 断言,验证其把传输结果正确归类为三类 detail(timeout/connect-error/non-2xx), 且 2xx → kicked=true 且 detail=""。
observer
stdlog
Package stdlog 提供 Observer 的标准库 log 参考实现。
Package stdlog 提供 Observer 的标准库 log 参考实现。
store
dynamodb
Package dynamodb 提供 Store 的 DynamoDB 参考实现。
Package dynamodb 提供 Store 的 DynamoDB 参考实现。
memory
Package memory 提供 Store 的内存实现,供测试与本地开发用。
Package memory 提供 Store 的内存实现,供测试与本地开发用。
storecontract
Package storecontract 提供 Store 契约测试套件:任意 Store 实现可跑同一组断言, 验证其满足「单键线性化 + 三形围栏 + 原子自增 + tombstone + 错误分类」语义契约。
Package storecontract 提供 Store 契约测试套件:任意 Store 实现可跑同一组断言, 验证其满足「单键线性化 + 三形围栏 + 原子自增 + tombstone + 错误分类」语义契约。

Jump to

Keyboard shortcuts

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