auth

package
v1.20.0 Latest Latest
Warning

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

Go to latest
Published: Sep 17, 2026 License: Apache-2.0 Imports: 28 Imported by: 0

Documentation

Overview

Package auth carries the canonical representation of the value stored under the token cache key (TokenCachePrefix+token) and provides versioned encoding helpers.

Historically the cache value was a `@`-joined string ("uid@name" or "uid@name@role") split ad-hoc at every call site. i18n 主方案 D10/D21 要求把 token cache 真相源收口,并在 payload 中带上用户语言偏好(D20 UserInfo), Encode 在 expand 阶段继续写 v2;Decode 提前兼容 v3,让所有 reader 先于 v3 writer 升级。灰度期老 token 不会因为升级失效。

Index

Constants

View Source
const ManagerRoleDashboardReader = "dashboardReader"

ManagerRoleDashboardReader is a temporary fixed role used before general manager RBAC exists. Keep it local to octo-server: adding it to octo-lib's CheckLoginRole would accidentally grant every admin endpoint.

View Source
const ManagerRoleMarketAdmin = "marketAdmin"

ManagerRoleMarketAdmin is a fixed role for staff who run the platform market — the MCP catalog, the Skill catalog and the Expert Market — without holding any other administrative power.

Same shape and same rationale as ManagerRoleDashboardReader: it is deliberately NOT known to octo-lib's CheckLoginRole, so an account holding it passes zero admin/superAdmin endpoint gates in octo-server. Its only effect here is the mcp.* / skill.* / expert.* capabilities in managerCapabilities.

Enforcement lives in octo-marketplace, whose /api/v1/admin/* groups are gated per resource; each group that admits this role opts into it explicitly, and a group registered without it stays superAdmin-only. Both sides must agree: a capability advertised here that marketplace does not admit renders the page and then 403s every call behind it. The marketplace half is octo-marketplace#55 (per-resource gating) and #56 (admitting this role on the Expert Market groups).

Before granting this to anyone, five things are worth knowing:

  • It is a publishing authority, not a read-mostly editor. A holder can create, edit and delete the public Skills, system MCPs and experts that every user on the platform installs and runs locally, and restructure the catalog taxonomy. It is genuinely narrower than superAdmin — no system settings, backups, user/group writes or space destruction — but pick people at a supply-chain bar, not at a "content editor" one.
  • Do not DEPLOY this service before octo-marketplace has the matching gate live. The capabilities are computed per request from CanAdminMarketplace, so every account already holding the role picks up the advertised surface at deploy — withholding new grants does not narrow the blast radius, and the population that gains it at deploy is fixed at release time. (Later grants do of course extend it; the point is that deploy order, not grant policy, is the lever for the accounts that already hold the role.)
  • A GRANT DOES NOT REACH MARKETPLACE UNTIL THE HOLDER RE-AUTHENTICATES. Mirror image of the revoke trap below, same cause: /v1/manager/me reads the live role through the parser's RoleResolver, while marketplace reads the token snapshot through /v1/auth/verify. So granting this to someone with an open console session makes the menu appear immediately and every call behind it 403 until they log out and back in. Have them re-login.
  • REVOKING THE ROLE DOES NOT CUT OFF MARKET ACCESS. Revoke the session too. Marketplace resolves callers through /v1/auth/verify, which answers from tokenValidator.Validate — the role snapshotted into the session token at login — not from the live user.role column. (The RoleResolver that keeps octo-server's own console fresh within RoleCacheTTL is wired into CacheTokenParser only; see main.go.) So for an existing session, clearing user.role never lands: catalog access survives for the remaining token lifetime, up to Cache.TokenExpire — 30 days by default. Revoking the session does work, bounded by marketplace's own token-keyed identity cache (AUTH_CACHE_TTL, 30s default), which has no invalidation entry point. So: role + session, and expect up to ~30s of residual access.
  • See the fixed-role section in modules/user/api_manager.go for the two lifecycle traps both fixed roles share (one-way downgrade, and accounts that cannot be deleted until the role is revoked).
View Source
const (

	// SessionMaxPerUIDLimit is shared by the persisted authority and the operator
	// command so they cannot drift on the accepted control range.
	SessionMaxPerUIDLimit = 10_000
)

Variables

View Source
var (
	ErrMigrationLockHeld                          = errors.New("auth: token migration lock is held")
	ErrMigrationLockLost                          = errors.New("auth: token migration lock was lost")
	ErrMigrationElapsedCutoffConfirmationRequired = errors.New("auth: migration apply requires confirm elapsed cutoff")
)
View Source
var (
	ErrRolloutControlChanged = errors.New("auth: session rollout control changed concurrently")
	ErrRolloutFloorNotNext   = errors.New("auth: rollout floor must advance exactly one phase")
)
View Source
var (
	ErrTokenCollision        = errors.New("auth: generated token already exists")
	ErrTokenVersionDowngrade = errors.New("auth: token payload version downgrade rejected")
)
View Source
var (
	ErrV3SessionsDisabled        = errors.New("auth: v3 session writer is disabled")
	ErrSessionLimitReached       = errors.New("auth: session limit reached")
	ErrIssueFenceChanged         = errors.New("auth: session issue fence changed")
	ErrRevocationAlreadyApplied  = errors.New("auth: revocation event already applied")
	ErrSessionGenerationInactive = errors.New("auth: session generation is not active")
	ErrLegacySessionDenied       = errors.New("auth: legacy session denied")
)
View Source
var ErrEmptyToken = errors.New("auth: empty token payload")

ErrEmptyToken indicates an empty cache value (missing or evicted token).

View Source
var ErrInvalidToken = errors.New("auth: invalid token payload")

ErrInvalidToken indicates a token payload that matches neither v2 JSON nor the legacy "uid@name[@role]" string.

View Source
var ErrRolloutScanLeaseHeld = errors.New("auth: session rollout scan lease is held")
View Source
var ErrRolloutStateUninitialized = errors.New("auth: session rollout state is not initialized")
View Source
var ErrSessionCapUnavailable = errors.New("auth: bounded sessions per UID unavailable; refusing to create a new session")

ErrSessionCapUnavailable is returned when a v3-writing mode is active with no usable per-UID cap. Reader strictness still applies; only new credentials are refused, so the failure is closed in both directions.

View Source
var ErrWriterLeaseLost = errors.New("auth: writer lease lost; refusing to create a new session")

ErrWriterLeaseLost is returned when this process may not create credentials because its write lease has expired. It is a degradation, not a crash: the caller surfaces the existing unauthenticated envelope and existing sessions keep working.

Functions

func CanAdminMarketplace added in v1.15.0

func CanAdminMarketplace(role string) bool

CanAdminMarketplace is the server-authoritative policy for the platform market admin surface — MCP catalog, Skill catalog and Expert Market. It is the octo-server half of a contract whose enforcement lives in octo-marketplace (internal/middleware/admin.go): this function decides what /v1/manager/me advertises, marketplace decides what the /api/v1/admin/* routes actually admit. Keep the two in sync.

func CanReadManagerDashboard added in v1.14.0

func CanReadManagerDashboard(role string) bool

CanReadManagerDashboard is the server-authoritative policy for the global operations dashboard read surface. Mutating operations keep their narrower SuperAdmin checks at the handler.

func Encode

func Encode(info TokenInfo) (string, error)

Encode serializes a TokenInfo as the versioned JSON envelope. The UID is the only required field; callers should populate Name/Role exactly as they did for the legacy "uid@name@role" string, and Language with the value resolved at write time (may be empty if unknown).

func EncodeV3 added in v1.14.0

func EncodeV3(info TokenInfo) (string, error)

EncodeV3 is intentionally separate from Encode. Release A readers can decode and validate v3 while production writers still emit v2.

func InitializeSessionRollout added in v1.15.0

func InitializeSessionRollout(ctx *config.Context) (RolloutBoot, *RolloutControlStore, error)

InitializeSessionRollout runs after module migrations and before HTTP serve. The #725 Redis floor and deprecated MODE are consulted only if the MySQL singleton does not exist yet; every later boot reads MySQL exclusively.

func IsManagerConsoleRole added in v1.14.0

func IsManagerConsoleRole(role string) bool

IsManagerConsoleRole reports whether a role may establish a manager-console session and read its own /v1/manager/me capability map.

func SplitUndeterminableBlockers added in v1.15.0

func SplitUndeterminableBlockers(blockedBy []string) (undeterminable, remaining []string)

SplitUndeterminableBlockers separates blockers a one-shot caller cannot evaluate from ones it can. Only the convergence window qualifies: it asserts that a set held still across a lease TTL, and a single invocation has no window to have observed. Reporting it next to real obstacles made `status` contradict the reconciler it exists to diagnose.

Types

type CacheTokenParser

type CacheTokenParser struct {
	Cache  cache.Cache
	Prefix string
	// contains filtered or unexported fields
}

CacheTokenParser implements octo-lib's wkhttp.TokenParser using the shared pkg/auth codec. It supersedes octo-lib's legacyTokenParser so that octo-server can write v2 JSON envelopes while still decoding any legacy uid@name[@role] values left in cache from older binaries.

When a LanguageResolver is injected via WithLanguageResolver, Parse hits the resolver after Decode to upgrade the token-cache language snapshot to the authoritative value before octo-lib's AuthMiddleware stores UserInfo on the request context. Resolver failures are non-fatal — the decoded snapshot is preserved so a Redis/DB outage degrades to "stale language" rather than "authentication failure".

Construct once at boot and register with WKHttp.SetTokenParser; the parser is safe for concurrent use as long as the underlying cache + resolver are.

func NewCacheTokenParser

func NewCacheTokenParser(c cache.Cache, prefix string, opts ...ParserOption) *CacheTokenParser

NewCacheTokenParser is a convenience constructor; nil cache is a programmer error and panics rather than silently degrading to a parser that fails every request.

func (*CacheTokenParser) Parse

func (p *CacheTokenParser) Parse(ctx context.Context, token string) (wkhttp.UserInfo, error)

Parse implements wkhttp.TokenParser. ctx is propagated to the optional LanguageResolver so resolver implementations can honour deadlines / cancellation set by the surrounding request.

type IssueFence added in v1.15.0

type IssueFence struct {
	Generation string
	Revision   uint64
}

type LanguageResolver

type LanguageResolver interface {
	Resolve(ctx context.Context, uid string) (string, error)
}

LanguageResolver hydrates UserInfo.Language with the freshest user-language preference (Redis cache → DB → ""). It is intentionally a tiny interface shaped at the consumer side so pkg/auth does not need to import the i18n package or know about Redis / DB primitives. The concrete implementation lives in modules/user.

type LegacyFinitePolicy added in v1.15.0

type LegacyFinitePolicy string
const (
	// LegacyFinitePolicyNatural preserves finite legacy deadlines that already
	// fit within maxTTL. Persistent and over-max records are still bounded.
	LegacyFinitePolicyNatural LegacyFinitePolicy = "natural"
	// LegacyFinitePolicyCap also compresses finite legacy deadlines to the
	// campaign cutoff. This may log users out earlier and requires approval.
	LegacyFinitePolicyCap LegacyFinitePolicy = "cap"
)

type LegacyMigrationOptions added in v1.15.0

type LegacyMigrationOptions struct {
	CampaignID           string
	CutoffAt             time.Time
	FinitePolicy         LegacyFinitePolicy
	BatchSize            int64
	Interval             time.Duration
	Apply                bool
	ConfirmElapsedCutoff bool
	Lease                time.Duration
}

type LegacyMigrationResult added in v1.15.0

type LegacyMigrationResult struct {
	Complete   bool   `json:"complete"`
	CampaignID string `json:"campaign_id"`
	Scanned    int64  `json:"scanned"`
	// InvalidPayload counts records that cannot be decoded as any token
	// version. It mirrors SessionObservation.DecodeInvalid so observe and
	// migrate report the same number for the same keyspace.
	InvalidPayload int64  `json:"invalid_payload"`
	Missing        int64  `json:"missing"`
	V1             int64  `json:"v1"`
	V2             int64  `json:"v2"`
	V3             int64  `json:"v3"`
	Shortened      int64  `json:"shortened"`
	WouldDelete    int64  `json:"would_delete"`
	Deleted        int64  `json:"deleted"`
	Unchanged      int64  `json:"unchanged"`
	Invalid        int64  `json:"invalid"`
	V3NonFinite    int64  `json:"v3_non_finite"`
	LastCursor     uint64 `json:"last_cursor"`
	LockLost       bool   `json:"lock_lost"`
}

type LegacySessionPolicy added in v1.15.0

type LegacySessionPolicy interface {
	ValidateLegacySession(ctx context.Context, info TokenInfo, record TokenRecord) error
}

type ParserOption

type ParserOption func(*CacheTokenParser)

ParserOption configures optional CacheTokenParser behaviour.

func WithLanguageResolver

func WithLanguageResolver(r LanguageResolver) ParserOption

WithLanguageResolver wires a LanguageResolver into the parser; nil resolver is a no-op so callers can pass an interface value that may be unset in test environments without an extra guard.

func WithRoleResolver added in v1.7.0

func WithRoleResolver(r RoleResolver) ParserOption

WithRoleResolver wires a RoleResolver into the parser; nil resolver is a no-op so callers (and tests) can pass an interface value that may be unset without an extra guard. When unset, Parse falls back to the role snapshot decoded from the token — i.e. legacy behaviour.

func WithTokenValidator added in v1.14.0

func WithTokenValidator(v *TokenValidator) ParserOption

WithTokenValidator makes Parse use the canonical payload+PTTL validator.

type ReconcilerOptions added in v1.15.0

type ReconcilerOptions struct {
	Registry      *WriterRegistry
	Control       RolloutStateController
	AutoAdvance   bool
	CanaryAhead   bool
	ExpectWriters int
	ScanBatchSize int64
	// ScanInterval throttles the keyspace scan. This runs unattended, on the
	// shared session pool, from every replica, so it must not be zero in
	// production — ObserveRateLimited reserves that for tests.
	ScanInterval time.Duration
	Log          func(format string, args ...interface{})
}

type RedisSessionStore added in v1.14.0

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

RedisSessionStore is the sole v2 credential writer in Release A. It owns token and compatibility UIDToken key construction so handlers cannot accidentally drop or extend a bearer deadline.

func NewRedisSessionStore added in v1.14.0

func NewRedisSessionStore(client *rd.Client, tokenPrefix, uidTokenPrefix string, maxTTL time.Duration, opts ...SessionStoreOption) *RedisSessionStore

func SessionStoreAndClientForContext added in v1.14.0

func SessionStoreAndClientForContext(ctx *config.Context) (*RedisSessionStore, *rd.Client)

SessionStoreAndClientForContext also exposes the shared client to adjacent auth stores that require Lua. Callers must not close it independently; its lifetime is the same as the config.Context, matching existing Redis pools.

func SessionStoreForContext added in v1.14.0

func SessionStoreForContext(ctx *config.Context) *RedisSessionStore

SessionStoreForContext returns one bounded Redis pool per server context. All modules in a replica share it; security hardening must not multiply connection pools with the number of token-consuming modules.

func (*RedisSessionStore) ApplyAndPublishRolloutState added in v1.15.0

func (s *RedisSessionStore) ApplyAndPublishRolloutState(
	registry *WriterRegistry,
	state RolloutState,
	mode SessionMode,
) error

ApplyAndPublishRolloutState is the sole runtime mode transition. Issuance is fenced before the local reader changes; the registry state and lease are then published atomically. Any failure keeps the fence closed until a later poll successfully republishes the same applied state.

func (*RedisSessionStore) ApplyRolloutState added in v1.15.0

func (s *RedisSessionStore) ApplyRolloutState(state RolloutState, mode SessionMode) error

ApplyRolloutState raises the derived state to match a newly observed MySQL floor. It only ever moves upward: a stale snapshot or operator rollback must not re-admit sessions this replica has already stopped accepting.

A missing or unusable cap does NOT block reader strictness. The cap is enforced on the write path (see writableState), so a malformed control row cannot loosen validation and credential issuance still fails closed.

func (*RedisSessionStore) BeginIssue added in v1.15.0

func (s *RedisSessionStore) BeginIssue(ctx context.Context, uid string) (fence IssueFence, err error)

func (*RedisSessionStore) CanIssue added in v1.15.0

func (s *RedisSessionStore) CanIssue() error

CanIssue reports whether this replica may create a credential right now. It exists so a caller can check BEFORE a destructive step: paths that revoke an old token and then issue a replacement would otherwise turn a lease loss into a logout rather than a refused login.

func (*RedisSessionStore) Client added in v1.15.0

func (s *RedisSessionStore) Client() *rd.Client

Client exposes the shared Redis client for adjacent rollout components (the writer registry, the operator subcommands) that must not open a second pool.

func (*RedisSessionStore) CurrentGeneration added in v1.15.0

func (s *RedisSessionStore) CurrentGeneration(ctx context.Context, uid string) (string, error)

func (*RedisSessionStore) DeleteToken added in v1.14.0

func (s *RedisSessionStore) DeleteToken(ctx context.Context, token string) error

func (*RedisSessionStore) DeviceToken added in v1.14.0

func (s *RedisSessionStore) DeviceToken(ctx context.Context, uid string, deviceFlag int) (string, error)

func (*RedisSessionStore) EvaluateRolloutAdvance added in v1.15.0

func (s *RedisSessionStore) EvaluateRolloutAdvance(ctx context.Context, in RolloutAdvanceInput) (RolloutAdvanceDecision, error)

EvaluateRolloutAdvance decides whether the floor may move one phase.

func (*RedisSessionStore) ForceAdvanceRollout added in v1.15.0

func (s *RedisSessionStore) ForceAdvanceRollout(
	ctx context.Context,
	decision RolloutAdvanceDecision,
	control RolloutStateController,
) error

ForceAdvanceRollout is the operator fault channel, for when the reconciler itself is broken. The caller must have already evaluated the predicate, which --force does not skip.

The MySQL controller performs the evidence insert and versioned floor CAS in one transaction, so this path cannot bypass pause, ordering, or audit rules.

func (*RedisSessionStore) InvalidateCurrentToken added in v1.15.0

func (s *RedisSessionStore) InvalidateCurrentToken(ctx context.Context, uid, token string) error

InvalidateCurrentToken is the module-facing current-logout primitive. It uses v3 metadata to clean the bounded index; legacy credentials are deleted directly because their payload has no trustworthy device scope.

func (*RedisSessionStore) IssueNew added in v1.14.0

func (s *RedisSessionStore) IssueNew(ctx context.Context, token, payload, uid string, deviceFlag int) (err error)

IssueNew creates a finite bearer and compatibility device index. The stored payload carries an internal ownership marker until another request rewrites it through ReuseExisting or UpdatePayloadKeepDeadline. If the index write fails, only the still-owned credential is deleted as compensation.

func (*RedisSessionStore) IssueNewSession added in v1.15.0

func (s *RedisSessionStore) IssueNewSession(ctx context.Context, token string, info TokenInfo, fence IssueFence) (err error)

func (*RedisSessionStore) MigrateLegacySessions added in v1.15.0

func (s *RedisSessionStore) MigrateLegacySessions(ctx context.Context, options LegacyMigrationOptions) (result LegacyMigrationResult, err error)

func (*RedisSessionStore) Mode added in v1.15.0

func (s *RedisSessionStore) Mode() SessionMode

Mode is the phase this replica is currently running at. It can change underneath a caller when the floor advances, so read it once per decision rather than caching it.

func (*RedisSessionStore) Observe added in v1.14.0

func (s *RedisSessionStore) Observe(ctx context.Context, batchSize int64) (stats SessionObservation, err error)

Observe performs an explicit, read-only cursor scan for migration planning. It is never called at process startup. The result is aggregate-only, and each batch checks context cancellation so operators can stop it safely.

func (*RedisSessionStore) ObserveRateLimited added in v1.14.0

func (s *RedisSessionStore) ObserveRateLimited(ctx context.Context, batchSize int64, interval time.Duration) (stats SessionObservation, err error)

ObserveRateLimited is Observe with an optional minimum interval between token reads. Production tooling must pass a positive interval to cap Redis load; zero is intended for bounded tests and offline environments.

func (*RedisSessionStore) Probe added in v1.14.0

func (s *RedisSessionStore) Probe(ctx context.Context) error

Probe executes the same read-only, single-key Lua used by the authentication hot path. It is intended for startup compatibility checks before accepting traffic; it never creates or mutates a Redis key.

func (*RedisSessionStore) ReadToken added in v1.14.0

func (s *RedisSessionStore) ReadToken(ctx context.Context, key string) (record TokenRecord, err error)

func (*RedisSessionStore) RedisInstanceFingerprint added in v1.15.0

func (s *RedisSessionStore) RedisInstanceFingerprint() (string, error)

RedisInstanceFingerprint identifies the Redis process this store is talking to. It is printed by every operator subcommand: identity derived purely from config cannot tell two endpoints apart, which is how a misplaced config key once pointed a tool at the wrong Redis without a word of complaint.

func (*RedisSessionStore) ReuseExisting added in v1.14.0

func (s *RedisSessionStore) ReuseExisting(ctx context.Context, token, payload, uid string, deviceFlag int) (ok bool, err error)

ReuseExisting updates a bearer without extending its deadline and aligns the compatibility index to the same TTL. A missing bearer is never recreated.

func (*RedisSessionStore) ReuseSession added in v1.15.0

func (s *RedisSessionStore) ReuseSession(ctx context.Context, token string, snapshot TokenInfo, fence IssueFence) (ok bool, err error)

ReuseSession refreshes an existing credential without extending its deadline. In v3 writer modes it also promotes a finite legacy credential to v3 under the supplied issuance fence; a missing credential returns false so the caller can issue a new random token.

func (*RedisSessionStore) RevokeAll added in v1.15.0

func (s *RedisSessionStore) RevokeAll(ctx context.Context, uid string, event RevocationEvent) (err error)

func (*RedisSessionStore) RevokeCurrent added in v1.15.0

func (s *RedisSessionStore) RevokeCurrent(ctx context.Context, token, uid string, deviceFlag int) (err error)

RevokeCurrent removes one bearer with compare-delete semantics and cleans both compatibility and v3 indexes only after that exact payload was removed. A concurrent payload adoption is retried instead of being deleted based on a stale read.

func (*RedisSessionStore) RevokeIssued added in v1.14.0

func (s *RedisSessionStore) RevokeIssued(ctx context.Context, token, uid string, deviceFlag int) error

RevokeIssued compensates a credential only while it is still owned by its original issue attempt. Adoption by another request removes the ownership marker atomically with its payload update, making this method a no-op.

func (*RedisSessionStore) RolloutControl added in v1.15.0

func (s *RedisSessionStore) RolloutControl(ctx context.Context) (*SessionRolloutControl, error)

func (*RedisSessionStore) UIDTokenPrefix added in v1.15.0

func (s *RedisSessionStore) UIDTokenPrefix() string

UIDTokenPrefix is the namespace all rollout safety state lives under.

func (*RedisSessionStore) UpdatePayloadKeepDeadline added in v1.14.0

func (s *RedisSessionStore) UpdatePayloadKeepDeadline(ctx context.Context, token, payload string) (ok bool, err error)

func (*RedisSessionStore) UpdateSessionSnapshot added in v1.15.0

func (s *RedisSessionStore) UpdateSessionSnapshot(ctx context.Context, token string, snapshot TokenInfo) (ok bool, err error)

UpdateSessionSnapshot changes only display and authorization snapshot fields. v3 security claims are copied from the current payload and the CAS prevents a concurrent writer from being overwritten with stale claims.

func (*RedisSessionStore) UseWriterLease added in v1.15.0

func (s *RedisSessionStore) UseWriterLease(registry *WriterRegistry)

UseWriterLease binds the write lease. Once set, creating a new credential requires a live lease: that is what makes "absent from the registry" prove "not writing", which the advance gate depends on. Validation reads are deliberately unaffected — fencing new logins during a Redis outage is the intended degradation, logging everyone out is not.

func (*RedisSessionStore) ValidateLegacySession added in v1.15.0

func (s *RedisSessionStore) ValidateLegacySession(ctx context.Context, info TokenInfo, record TokenRecord) error

type RevocationEvent added in v1.15.0

type RevocationEvent struct {
	Version uint64
	ID      string
}

type RoleResolver added in v1.7.0

type RoleResolver interface {
	ResolveRole(ctx context.Context, uid string) (string, error)
}

RoleResolver hydrates UserInfo.Role with the user's *current* system role (Redis cache → DB → "") instead of the value snapshotted into the token at issuance. Without it, a system role baked into the token at login keeps granting admin / superAdmin access until the token expires — a demotion or admin-account removal cannot be honoured promptly. Resolving per request bounds that staleness to the resolver's cache TTL.

Like LanguageResolver it is shaped at the consumer side so pkg/auth stays free of DB / Redis imports; the concrete implementation lives in modules/user (RoleService).

type RolloutAdvanceDecision added in v1.15.0

type RolloutAdvanceDecision struct {
	Current      SessionMode `json:"current_floor"`
	Target       SessionMode `json:"target_floor"`
	StateVersion int64       `json:"state_version"`
	Allowed      bool        `json:"allowed"`
	BlockedBy    []string    `json:"blocked_by,omitempty"`
	MaxPerUID    int         `json:"max_per_uid,omitempty"`
	// Scanned reports whether this evaluation ran a keyspace scan, so a caller
	// can back off rather than rescanning every cycle while waiting out a
	// deadline measured in days.
	Scanned bool `json:"scanned"`
	// Observation is the scan that justified the decision. It is carried here so
	// the advance snapshot records the counts it claims to audit — a snapshot
	// reading v1=0 v2=0 because nobody filled it in is indistinguishable from
	// one that actually looked.
	Observation *SessionObservation `json:"observation,omitempty"`
	// Options are operator-actionable next steps for the blockers above.
	Options []string `json:"options,omitempty"`
	// These bind the decision to the exact fleet and Redis process observed.
	// They are persisted in the same MySQL transaction as the floor CAS.
	WriterFingerprint string `json:"writer_fingerprint,omitempty"`
	RedisInstanceID   string `json:"redis_instance_id,omitempty"`
}

RolloutAdvanceDecision is what the reconciler and `advance --force` both act on. BlockedBy entries are operator-facing and low cardinality.

func (RolloutAdvanceDecision) BlockedOnlyOnConvergence added in v1.15.0

func (d RolloutAdvanceDecision) BlockedOnlyOnConvergence() bool

BlockedOnlyOnConvergence is the exported form, for callers outside this package that poll the predicate.

func (RolloutAdvanceDecision) BlockedSummary added in v1.15.0

func (d RolloutAdvanceDecision) BlockedSummary() string

BlockedSummary renders BlockedBy for a log line or status output.

type RolloutAdvanceInput added in v1.15.0

type RolloutAdvanceInput struct {
	// State is the MySQL-authoritative snapshot whose version will be used by
	// the advance CAS. Redis never supplies a floor after takeover.
	State    RolloutState
	Registry *WriterRegistry
	// ExpectWriters is the replica count the deployment intends to run. Required
	// for the first v3 floor and only that one, because it is the single
	// transition where a non-participating pre-#725 build could still be writing
	// v2 and the registry structurally cannot see it. Every later transition is
	// fully machine-gated.
	ExpectWriters int
	// Convergence carries the caller's observation window across evaluations.
	// Required for the first v3 floor, where a count taken at one instant is not
	// enough — see WriterConvergence. A caller that cannot observe over time
	// (a one-shot command) supplies one and evaluates repeatedly.
	Convergence *WriterConvergence
	// ScanBatchSize and ScanInterval throttle the keyspace scan. A zero interval
	// is reserved for tests; production callers pass a positive one, because
	// this scan runs unattended and on the shared session pool.
	ScanBatchSize int64
	ScanInterval  time.Duration
}

RolloutAdvanceInput carries what the predicate cannot read for itself.

type RolloutAdvanceRecord added in v1.15.0

type RolloutAdvanceRecord struct {
	From              SessionMode         `json:"from"`
	To                SessionMode         `json:"to"`
	Actor             string              `json:"actor"`
	AtMS              int64               `json:"at_unix_ms"`
	Kind              string              `json:"transition_kind"`
	RedisID           string              `json:"redis_instance_id,omitempty"`
	WriterFingerprint string              `json:"writer_fingerprint,omitempty"`
	Observation       *SessionObservation `json:"observation,omitempty"`
	CapChange         *RolloutCapChange   `json:"cap_change,omitempty"`
}

RolloutAdvanceRecord is the audit payload inserted in the same MySQL transaction as the versioned floor or cap CAS.

type RolloutBoot added in v1.15.0

type RolloutBoot struct {
	Outcome       RolloutBootOutcome
	Floor         SessionMode
	Mode          SessionMode
	MaxPerUID     int
	Version       int64
	Warning       string
	AutoAdvance   bool
	CanaryAhead   bool
	ExpectWriters int
}

func SessionBootForContext added in v1.15.0

func SessionBootForContext(ctx *config.Context) (RolloutBoot, []string)

SessionBootForContext exposes what boot resolved, for the startup log line and the rollout status subcommand.

type RolloutBootOutcome added in v1.15.0

type RolloutBootOutcome string

RolloutBootOutcome is a low-cardinality startup classification. Redis is consulted only during the one-time legacy takeover; after the singleton exists every boot is normal regardless of Redis contents.

const (
	RolloutBootFresh   RolloutBootOutcome = "fresh"
	RolloutBootAdopted RolloutBootOutcome = "adopted"
	RolloutBootNormal  RolloutBootOutcome = "normal"
)

type RolloutCapChange added in v1.15.0

type RolloutCapChange struct {
	FromMaxPerUID int `json:"from_max_per_uid"`
	ToMaxPerUID   int `json:"to_max_per_uid"`
}

RolloutCapChange is stored in the append-only audit row for a cap-only control transition. The floor remains unchanged for this transition.

type RolloutControlStore added in v1.15.0

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

RolloutControlStore owns the singleton state and its append-only advance evidence. All irreversible changes happen in one MySQL transaction.

func NewRolloutControlStore added in v1.15.0

func NewRolloutControlStore(db *dbr.Session) *RolloutControlStore

func (*RolloutControlStore) Advance added in v1.15.0

Advance atomically records the evidence and performs a monotonic, one-phase CAS. A loser rolls the evidence row back with the state update.

func (*RolloutControlStore) Initialize added in v1.15.0

func (s *RolloutControlStore) Initialize(ctx context.Context, seed RolloutSeed) (RolloutState, error)

Initialize creates or monotonically adopts a takeover seed under a row lock. Concurrent starters serialize here; the strictest seed eventually wins.

func (*RolloutControlStore) LastAdvance added in v1.15.0

func (*RolloutControlStore) Load added in v1.15.0

func (*RolloutControlStore) SetMaxPerUID added in v1.15.0

func (s *RolloutControlStore) SetMaxPerUID(
	ctx context.Context,
	current RolloutState,
	maxPerUID int,
	actor string,
) (RolloutState, error)

SetMaxPerUID atomically updates the durable cap and appends the old/new value to the same audit stream. It does not move the rollout floor and it does not evict existing sessions; a lower cap gates subsequent index additions.

func (*RolloutControlStore) SetPaused added in v1.15.0

func (s *RolloutControlStore) SetPaused(ctx context.Context, paused bool) error

type RolloutReconciler added in v1.15.0

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

RolloutReconciler polls the floor and, when enabled, advances it.

func NewRolloutReconciler added in v1.15.0

func NewRolloutReconciler(store *RedisSessionStore, opts ReconcilerOptions) *RolloutReconciler

func (*RolloutReconciler) Run added in v1.15.0

func (r *RolloutReconciler) Run(ctx context.Context)

Run polls the floor and, when enabled, reconciles — in two goroutines.

One loop was not enough: a rate-limited full-keyspace scan inside reconcileOnce can run for an hour on a large keyspace, and while it did, the select never reached the poll tick. A floor advance published by another replica would not be applied here for that whole hour, defeating the five-second propagation the poller exists for, and the registry would keep advertising a stale applied state — which is exactly what the convergence gate reads.

It is started only from the server wiring, never from NewRedisSessionStore: a background goroutine that advances floors, started by a constructor, would fire in every test that happens to build a store.

type RolloutSeed added in v1.15.0

type RolloutSeed struct {
	Floor     SessionMode
	MaxPerUID int
	Actor     string
	Source    string
	RedisID   string
}

RolloutSeed is used only while taking authority over from the #725 Redis record and deprecated environment. Once the singleton exists, Initialize may raise it but can never lower it.

func ResolveRolloutSeed added in v1.15.0

func ResolveRolloutSeed(
	redisFloor, legacyMode SessionMode,
	redisMaxPerUID, legacyMaxPerUID int,
) (RolloutSeed, error)

ResolveRolloutSeed persists the strictest floor already in effect during takeover. A cap already carried by the #725 Redis record remains authoritative; the deprecated environment cap is only a fallback for older records without one.

type RolloutState added in v1.15.0

type RolloutState struct {
	Floor     SessionMode `db:"floor" json:"floor"`
	MaxPerUID int         `db:"max_per_uid" json:"max_per_uid"`
	Version   int64       `db:"version" json:"version"`
	Paused    bool        `db:"paused" json:"paused"`
	UpdatedAt time.Time   `db:"updated_at" json:"updated_at"`
}

RolloutState is the sole durable authority for the session rollout. Redis contains session data and scan leases only; losing or restoring Redis cannot move this value in either direction.

type RolloutStateController added in v1.15.0

type RolloutStateController interface {
	Load(context.Context) (RolloutState, error)
	Advance(context.Context, RolloutState, RolloutAdvanceRecord) (RolloutState, error)
}

type SessionGenerationResolver added in v1.14.0

type SessionGenerationResolver interface {
	CurrentGeneration(ctx context.Context, uid string) (string, error)
}

type SessionMode added in v1.15.0

type SessionMode string

SessionMode is the deployment phase for user HTTP sessions. Its zero value is deliberately invalid.

const (
	SessionModeExpand  SessionMode = "expand"
	SessionModeV3Write SessionMode = "v3-write"
	SessionModeRevoke  SessionMode = "revoke"
	SessionModeBounded SessionMode = "bounded"
	SessionModeEnforce SessionMode = "enforce"
)

func (SessionMode) RevokesSessions added in v1.15.0

func (m SessionMode) RevokesSessions() bool

RevokesSessions reports whether durable security events may rotate generations and create legacy deny markers in this rollout phase.

func (SessionMode) WritesV3 added in v1.15.0

func (m SessionMode) WritesV3() bool

WritesV3 reports whether the rollout phase permits creating or promoting user HTTP sessions with the v3 envelope.

type SessionObservation added in v1.14.0

type SessionObservation struct {
	ScanID           string `json:"scan_id"`
	ScopeFingerprint string `json:"scope_fingerprint"`
	RedisInstanceID  string `json:"redis_instance_id"`
	Complete         bool   `json:"complete"`
	Total            int64  `json:"total"`
	Missing          int64  `json:"missing"`
	Persistent       int64  `json:"persistent"`
	Finite           int64  `json:"finite"`
	OverMax          int64  `json:"over_max"`
	InvalidTTL       int64  `json:"invalid_ttl"`
	DecodeInvalid    int64  `json:"decode_invalid"`
	ReadErrors       int64  `json:"read_errors"`
	V1               int64  `json:"v1"`
	V2               int64  `json:"v2"`
	V3               int64  `json:"v3"`
}

SessionObservation contains low-cardinality migration facts only. It never includes a token, Redis key, UID, or payload.

type SessionRolloutControl added in v1.15.0

type SessionRolloutControl struct {
	ModeFloor           SessionMode `json:"mode_floor"`
	WriterVersion       int         `json:"writer_version"`
	ObservationMinGapMS int64       `json:"observation_min_gap_ms"`
	MaxPerUID           int         `json:"max_per_uid,omitempty"`
}

type SessionStoreOption added in v1.15.0

type SessionStoreOption func(*RedisSessionStore)

func WithSessionClock added in v1.15.0

func WithSessionClock(now func() time.Time) SessionStoreOption

func WithSessionMaxPerUID added in v1.15.0

func WithSessionMaxPerUID(max int) SessionStoreOption

func WithSessionMode added in v1.15.0

func WithSessionMode(mode SessionMode) SessionStoreOption

type TokenInfo

type TokenInfo struct {
	UID               string
	Name              string
	Role              string
	Language          string
	IssuedAt          int64
	ExpiresAt         int64
	DeviceFlag        int
	DeviceID          string
	SessionGeneration string
	SessionRevision   uint64
}

TokenInfo is the structured payload stored under TokenCachePrefix+token. UID is required; Role/Language may be empty.

func Decode

func Decode(raw string) (TokenInfo, error)

Decode reverses Encode and tolerates the legacy "uid@name[@role]" string so that tokens written by older binaries (and cached before the upgrade) keep working until they expire. UID emptiness is the only structural check applied here — language validity is enforced at consumption sites via i18n.MatchSupportedLanguage.

func (TokenInfo) IsV3 added in v1.14.0

func (i TokenInfo) IsV3() bool

IsV3 reports whether Decode returned the v3 envelope.

type TokenRecord added in v1.14.0

type TokenRecord struct {
	Payload string
	TTL     time.Duration
}

TokenRecord is an atomic snapshot of a token payload and its Redis PTTL. TTL follows Redis semantics: -2 means missing and -1 means persistent.

type TokenRecordReader added in v1.14.0

type TokenRecordReader interface {
	ReadToken(ctx context.Context, key string) (TokenRecord, error)
}

type TokenValidator added in v1.14.0

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

func NewTokenValidator added in v1.14.0

func NewTokenValidator(reader TokenRecordReader, prefix string, opts ...ValidatorOption) *TokenValidator

func (*TokenValidator) Validate added in v1.14.0

func (v *TokenValidator) Validate(ctx context.Context, token string) (TokenInfo, error)

Validate is the canonical token-read policy. Legacy v1/v2 persistent keys remain readable in Release A so migration can be observed before apply. Any v3 key is strict immediately: finite Redis TTL plus absolute payload expiry.

type ValidatorOption added in v1.14.0

type ValidatorOption func(*TokenValidator)

func WithSessionGenerationResolver added in v1.14.0

func WithSessionGenerationResolver(resolver SessionGenerationResolver) ValidatorOption

func WithValidatorClock added in v1.14.0

func WithValidatorClock(now func() time.Time) ValidatorOption

type WriterConvergence added in v1.15.0

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

WriterConvergence answers "how long has the live writer set looked exactly like this?".

It exists because a count is not a convergence proof. `live == ExpectWriters` holds at the single instant during a maxSurge:1 rollout when the new replicas are all up and the last old one — which predates the registry and therefore registers nowhere — has not yet terminated. That replica can still mint v2 credentials on the next login, which is the one thing the first v3 floor must rule out.

Stability of the SET, held across at least one lease TTL, is a much stronger statement than the count, and it is sound for a specific reason: writer identities are per-incarnation. An entry present at both ends of a window longer than the lease TTL cannot have expired and been recreated in between — a restart mints a new id — and any pod joining or leaving changes the set. So an unchanged set across that window proves no membership change occurred, i.e. the rollout has stopped moving.

It is a mitigation and not a proof of the absent replica itself, which is unobservable by construction; that residual risk is why EXPECT_WRITERS is a deployment-supplied number rather than something derived.

func NewWriterConvergence added in v1.15.0

func NewWriterConvergence() *WriterConvergence

func (*WriterConvergence) Observe added in v1.15.0

func (c *WriterConvergence) Observe(entries []WriterEntry, now time.Time) time.Duration

Observe records the roster and returns how long the CURRENT set has held. A set that differs from the previous observation restarts the window and returns zero.

type WriterEntry added in v1.15.0

type WriterEntry struct {
	ID           string `json:"id"`
	Build        string `json:"build"`
	AppliedState string `json:"applied_state"`
	Pod          string `json:"pod,omitempty"`
	StartedAtMS  int64  `json:"started_at_unix_ms"`
}

WriterEntry is one live participant. It carries nothing that could identify a user: build, pod name and applied state only.

type WriterRegistry added in v1.15.0

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

WriterRegistry tracks which processes currently hold a write lease.

func NewWriterRegistry added in v1.15.0

func NewWriterRegistry(client *rd.Client, keyPrefix string) *WriterRegistry

NewWriterRegistry builds a registry under the given key prefix. The prefix is supplied rather than derived so the type stays independent of auth config.

func (*WriterRegistry) Join added in v1.15.0

func (r *WriterRegistry) Join(ctx context.Context, build, pod, appliedState string, now func() time.Time) error

Join registers this process and starts refreshing its lease until ctx ends. The identity is a fresh random value per process incarnation, not the pod UID: private deployments are not always Kubernetes, and an in-place container restart keeps its pod UID, so keying on it would let a new process silently overwrite the old entry and hide the restart. The pod name rides along as a label so a crash-looping pod shows up as several live registrations rather than as confusion.

func (*WriterRegistry) Live added in v1.15.0

func (r *WriterRegistry) Live() ([]WriterEntry, error)

Live enumerates writers whose lease has not expired, and conditionally prunes roster members whose entry is still missing or malformed. Cleanup is a compare-and-remove operation: a writer that refreshes after MGET must stay registered because that refresh also makes MayWrite true.

SMEMBERS + MGET rather than SCAN: the runbook's Redis preflight already flags proxies with incomplete cursor semantics, and there is no reason to take that dependency for a set this small. Liveness comes from each entry key's own TTL, so Redis's clock is the only clock involved.

func (*WriterRegistry) MayWrite added in v1.15.0

func (r *WriterRegistry) MayWrite() bool

MayWrite reports whether this process still holds its lease. A writer without one must refuse to create new tokens, or "absent from the registry" stops meaning "not writing" and the gate becomes fail-open.

func (*WriterRegistry) PublishAppliedState added in v1.15.0

func (r *WriterRegistry) PublishAppliedState(state string) error

PublishAppliedState atomically renews this process's write lease and changes its advertised state. The in-memory entry is committed only after Redis accepts the candidate, so a failed publication cannot be refreshed later by the heartbeat as though it succeeded.

func (*WriterRegistry) SetAppliedState added in v1.15.0

func (r *WriterRegistry) SetAppliedState(state string) error

func (*WriterRegistry) SetAppliedStateIfChanged added in v1.15.0

func (r *WriterRegistry) SetAppliedStateIfChanged(state string) error

SetAppliedState records that this process has applied a new rollout state. The gate reads it to prove the fleet has converged before advancing. SetAppliedStateIfChanged skips the Redis round trip when nothing moved. The applied mode changes a handful of times across an entire rollout, while the poller runs every five seconds.

Jump to

Keyboard shortcuts

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