Enjoying Blockend?

Give us a star on GitHub

Star on GitHub
blockend

Blockend

02 blocks

Token Service

New

Issue and verify JWT access tokens, rotate opaque refresh tokens, and connect key and token storage to your production infrastructure.

The Token Service block issues short lived JWT access tokens and one time opaque refresh tokens. It validates issuer, audience, algorithm, token type, and revocation state before returning trusted claims.

Your application still owns login, authorization policy, browser cookie settings, CSRF protection, and account recovery.


Features

  • Production configuration requires an HTTPS or URN issuer, an audience, and asymmetric JWT algorithms.
  • Access token lifetime defaults to 15 minutes and is capped at one hour.
  • Refresh tokens contain a versioned family identifier and a 256-bit random secret. Stores keep only their SHA-256 hashes.
  • Refresh families have a 90-day absolute lifetime and a 10,000-rotation ceiling by default. Both limits are configurable and stored with each family, bounding retained replay state.
  • Refresh rotation is atomic in the supplied Redis and PostgreSQL stores. Reuse revokes the token family.
  • Key IDs (kid) support signing key rotation while old public keys remain available for verification.
  • Local PEM, remote JWKS, and vendor neutral KMS signing providers are included.
  • Express, Fastify, and Hono adapters expose refresh by default. Issue and revoke routes require an authorization callback.
  • Claims are size limited and reject reserved token fields and common secret field names.

MemoryTokenStore is process local and loses state on restart. Use PostgreSQL or Redis for refresh tokens in a multi instance production service. Use StatelessTokenStore only when you issue access tokens without refresh or server side revocation.

This release changes the refresh-token format. Existing refresh tokens must be invalidated and clients must authenticate again. Redis also moves to the blockend:v2: namespace. Roll issuer instances together; do not run old and new token-service versions against one active session set.

Run Redis token state with maxmemory-policy noeviction in a dedicated deployment. Redis must reject writes rather than evict family or JTI state. Redis Cluster replication is asynchronous and can lose acknowledged writes during failover; choose PostgreSQL if your threat model requires strict revocation durability across failover.

File structure

token-service
├── src
│   ├── adapters
│   │   ├── express.ts
│   │   ├── fastify.ts
│   │   ├── hono.ts
│   │   └── http.ts
│   ├── core
│   │   ├── crypto.ts
│   │   ├── errors.ts
│   │   ├── token-service.ts
│   │   ├── types.ts
│   │   └── validation.ts
│   ├── providers
│   │   ├── kms-rsa-key-provider.ts
│   │   ├── local-key-provider.ts
│   │   └── remote-jwks-provider.ts
│   ├── stores
│   │   ├── memory-token-store.ts
│   │   ├── postgres-token-store.ts
│   │   ├── redis-token-store.ts
│   │   └── stateless-token-store.ts
│   └── index.ts
├── sql
│   └── postgres.sql
└── tests
  • src/core/ validates input, signs and verifies access tokens, rotates refresh tokens, and defines the storage and provider contracts.
  • src/providers/ adapts PEM keys, remote JWKS, and an external RS256 signer.
  • src/stores/ provides process local, stateless, Redis, and PostgreSQL implementations.
  • src/adapters/ maps the token service to framework routes and HTTP errors.
  • sql/postgres.sql creates the refresh token and revoked JTI tables.
  • tests/ covers core behavior, adapters, and store contracts.

Installation

pnpm dlx blockend-cli add token-service

Detect project

Blockend detects your project and selects the framework adapter.

Install dependencies

The CLI installs jose and zod, plus the selected framework package.

Generate files

The service, selected adapter, tests, and PostgreSQL schema are copied into your blocks directory.

Copy the source files into blocks/token-service and install the dependencies used by your chosen adapter.

Runtime dependencies

PackageRequired for
joseJWT signing, verification, and remote JWKS
zodStartup and input validation
expressExpress adapter only
fastifyFastify adapter only
honoHono adapter only

Core files

index.ts exports the service, provider, store, and public types.

export { createTokenService } from "./core/token-service.js";
export { TokenError, isTokenError } from "./core/errors.js";
export type * from "./core/types.js";
export { LocalKeyProvider } from "./providers/local-key-provider.js";
export type { LocalKey } from "./providers/local-key-provider.js";
export { RemoteJwksKeyProvider } from "./providers/remote-jwks-provider.js";
export { KmsRsaKeyProvider } from "./providers/kms-rsa-key-provider.js";
export type { KmsRsaSigner } from "./providers/kms-rsa-key-provider.js";
export { MemoryTokenStore } from "./stores/memory-token-store.js";
export { StatelessTokenStore } from "./stores/stateless-token-store.js";
export { RedisTokenStore } from "./stores/redis-token-store.js";
export type { RedisLike, RedisTokenStoreOptions } from "./stores/redis-token-store.js";
export { PostgresTokenStore } from "./stores/postgres-token-store.js";
export type { PgLike, PgTransaction, QueryResult } from "./stores/postgres-token-store.js";

core/crypto.ts creates opaque tokens, IDs, hashes, and claim size measurements.

import { createHash, randomBytes, timingSafeEqual } from "node:crypto";
import { TokenError } from "./errors.js";

export function newOpaqueToken(familyId: string): string {
  return `rt1.${Buffer.from(familyId, "utf8").toString("base64url")}.${randomBytes(32).toString("base64url")}`;
}
export function familyIdFromOpaqueToken(token: string): string | undefined {
  const parts = token.split(".");
  if (parts.length !== 3 || parts[0] !== "rt1" || !/^[A-Za-z0-9_-]{43}$/.test(parts[2]!))
    return undefined;
  try {
    const secret = Buffer.from(parts[2]!, "base64url");
    if (secret.byteLength !== 32 || secret.toString("base64url") !== parts[2]) return undefined;
    const encoded = parts[1]!;
    const familyId = Buffer.from(encoded, "base64url").toString("utf8");
    if (Buffer.from(familyId, "utf8").toString("base64url") !== encoded) return undefined;
    if (
      !/^(?:[A-Za-z0-9_-]{43}\.)?[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(
        familyId
      )
    )
      return undefined;
    return familyId;
  } catch {
    return undefined;
  }
}
export function newId(): string {
  return crypto.randomUUID();
}
export function hashToken(token: string): string {
  return createHash("sha256").update(token, "utf8").digest("base64url");
}
export function equalHash(left: string, right: string): boolean {
  const a = Buffer.from(left);
  const b = Buffer.from(right);
  return a.length === b.length && timingSafeEqual(a, b);
}

export function jsonByteLength(value: unknown): number {
  try {
    return Buffer.byteLength(JSON.stringify(value), "utf8");
  } catch (cause) {
    throw new TokenError("INVALID_INPUT", "Claims must be JSON serializable", {
      cause
    });
  }
}

export function clamp(
  value: number | undefined,
  fallback: number,
  min: number,
  max: number
): number {
  return Math.min(max, Math.max(min, value ?? fallback));
}

core/errors.ts defines stable token error codes.

import type { TokenErrorCode } from "./types.js";

export class TokenError extends Error {
  override readonly name = "TokenError";
  constructor(
    public readonly code: TokenErrorCode,
    message: string,
    options?: ErrorOptions
  ) {
    super(message, options);
  }
}

export function isTokenError(value: unknown): value is TokenError {
  return value instanceof TokenError;
}

core/token-service.ts implements issue, verify, refresh, and revoke operations.

import {
  decodeProtectedHeader,
  errors as JoseErrors,
  jwtVerify,
  SignJWT,
  type JWTPayload
} from "jose";
import { TokenError, isTokenError } from "./errors.js";
import {
  clamp,
  familyIdFromOpaqueToken,
  hashToken,
  jsonByteLength,
  newId,
  newOpaqueToken
} from "./crypto.js";
import {
  assertAllowedAlgorithm,
  parseIssue,
  parseOptions,
  parseRevoke,
  type IssueInput
} from "./validation.js";
import type {
  AccessTokenClaims,
  RefreshTokenRecord,
  RevokeInput,
  SigningAlgorithm,
  TokenEvent,
  TokenPair,
  TokenService,
  TokenServiceOptions,
  VerifyResult
} from "./types.js";

const ACCESS_DEFAULT = 900,
  ACCESS_MIN = 60,
  ACCESS_MAX = 3600;
const REFRESH_DEFAULT = 604_800,
  REFRESH_MIN = 3600,
  REFRESH_MAX = 2_592_000;
const FAMILY_MAX_DEFAULT = 7_776_000,
  FAMILY_MAX_MIN = 3600,
  FAMILY_MAX_MAX = 31_536_000,
  FAMILY_ROTATIONS_DEFAULT = 10_000,
  FAMILY_ROTATIONS_MAX = 1_000_000;
const DEFAULT_FORBIDDEN = [
  "password",
  "password_hash",
  "secret",
  "private_key",
  "privatekey",
  "access_token",
  "refresh_token",
  "token",
  "__proto__",
  "prototype",
  "constructor"
];
const RESERVED = new Set(["iss", "sub", "aud", "exp", "nbf", "iat", "jti", "type"]);

class DefaultTokenService implements TokenService {
  private readonly options: TokenServiceOptions;
  private readonly allowedAlgorithms: SigningAlgorithm[];
  private readonly forbidden: Set<string>;
  private readonly accessTtl: number;
  private readonly refreshTtl: number;
  private readonly familyMaxLifetime: number;
  private readonly familyMaxRotations: number;

  constructor(options: TokenServiceOptions) {
    this.options = parseOptions(options);
    this.allowedAlgorithms = options.algorithms ?? ["RS256", "ES256"];
    this.forbidden = new Set(
      (options.forbiddenClaimKeys ?? DEFAULT_FORBIDDEN).map((key) => key.toLowerCase())
    );
    this.accessTtl = clamp(options.accessTokenTtlSeconds, ACCESS_DEFAULT, ACCESS_MIN, ACCESS_MAX);
    this.refreshTtl = clamp(
      options.refreshTokenTtlSeconds,
      REFRESH_DEFAULT,
      REFRESH_MIN,
      REFRESH_MAX
    );
    this.familyMaxLifetime = clamp(
      options.maxRefreshFamilyLifetimeSeconds,
      FAMILY_MAX_DEFAULT,
      FAMILY_MAX_MIN,
      FAMILY_MAX_MAX
    );
    this.familyMaxRotations = clamp(
      options.maxRefreshRotations,
      FAMILY_ROTATIONS_DEFAULT,
      1,
      FAMILY_ROTATIONS_MAX
    );
  }

  async issue(raw: IssueInput): Promise<TokenPair> {
    try {
      const input = parseIssue(raw);
      this.assertClaims(input.claims ?? {});
      this.assertAudience(input.audience);
      const access = input.tokens?.access ?? true;
      const refresh = input.tokens?.refresh ?? true;
      if (!access && !refresh)
        throw new TokenError("INVALID_INPUT", "At least one token must be requested");
      const now = this.nowSeconds();
      const pair: TokenPair = { tokenType: "Bearer", issuedAt: now };
      if (access) {
        const ttl = clamp(input.accessTokenTtlSeconds, this.accessTtl, ACCESS_MIN, ACCESS_MAX);
        pair.accessToken = await this.signAccess(
          input.sub,
          input.claims ?? {},
          input.audience ?? this.options.audience,
          ttl,
          now
        );
        pair.expiresIn = ttl;
      }
      if (refresh) {
        const ttl = clamp(input.refreshTokenTtlSeconds, this.refreshTtl, REFRESH_MIN, REFRESH_MAX);
        const familyId = this.options.tokenStore.createFamilyId?.(input.sub) ?? newId();
        const familyExpiresAt = new Date((now + this.familyMaxLifetime) * 1000);
        const token = newOpaqueToken(familyId);
        const record = this.refreshRecord(
          token,
          input.sub,
          familyId,
          ttl,
          now,
          familyExpiresAt,
          0,
          this.familyMaxRotations
        );
        await this.store(() => this.options.tokenStore.saveRefreshToken(record));
        pair.refreshToken = token;
        pair.refreshExpiresIn = Math.max(
          0,
          Math.floor((record.expiresAt.getTime() - now * 1000) / 1000)
        );
      }
      this.emit({
        name: "token.issued",
        at: this.now().toISOString(),
        sub: input.sub
      });
      return pair;
    } catch (error) {
      throw this.failure(error);
    }
  }

  async verify(accessToken: string): Promise<VerifyResult> {
    try {
      if (typeof accessToken !== "string" || accessToken.length < 20 || accessToken.length > 16_384)
        throw new TokenError("INVALID_TOKEN", "Invalid access token");
      const header = decodeProtectedHeader(accessToken);
      if (!header.alg || header.alg === "none")
        throw new TokenError("INVALID_TOKEN", "Missing or invalid algorithm");
      assertAllowedAlgorithm(header.alg, this.allowedAlgorithms);
      const verification = await this.options.keyProvider.getVerificationKey({
        ...(header.kid ? { kid: header.kid } : {}),
        alg: header.alg
      });
      if (verification.alg !== header.alg)
        throw new TokenError("INVALID_TOKEN", "Key algorithm mismatch");
      const result = await jwtVerify(accessToken, verification.key, {
        algorithms: this.allowedAlgorithms,
        issuer: this.options.issuer,
        ...(this.options.audience ? { audience: this.options.audience } : {}),
        clockTolerance: this.options.clockSkewSeconds ?? 30,
        currentDate: this.now(),
        requiredClaims: ["sub", "exp", "iat", "jti", "type"]
      });
      const claims = this.toAccessClaims(result.payload);
      if (
        this.options.tokenStore.isJtiRevoked &&
        (await this.store(() => this.options.tokenStore.isJtiRevoked!(claims.jti)))
      ) {
        throw new TokenError("TOKEN_REVOKED", "Access token has been revoked");
      }
      this.emit({
        name: "token.verified",
        at: this.now().toISOString(),
        sub: claims.sub,
        jti: claims.jti
      });
      return { valid: true, claims };
    } catch (error) {
      throw this.failure(error);
    }
  }

  async refresh(refreshToken: string): Promise<TokenPair> {
    try {
      if (typeof refreshToken !== "string" || refreshToken.length < 32 || refreshToken.length > 512)
        throw new TokenError("INVALID_TOKEN", "Invalid refresh token");
      const tokenFamilyId = familyIdFromOpaqueToken(refreshToken);
      if (!tokenFamilyId) throw new TokenError("INVALID_TOKEN", "Invalid refresh token");
      const nowDate = this.now();
      const now = Math.floor(nowDate.getTime() / 1000);
      const currentHash = hashToken(refreshToken);
      const current = await this.store(() =>
        this.options.tokenStore.getRefreshToken(currentHash, tokenFamilyId)
      );
      if (!current) throw new TokenError("INVALID_TOKEN", "Unknown refresh token");
      if (current.familyId !== tokenFamilyId)
        throw new TokenError("INVALID_TOKEN", "Invalid refresh token");
      if (current.status === "revoked")
        throw new TokenError("TOKEN_REVOKED", "Refresh token revoked");
      if (current.status === "active" && current.expiresAt.getTime() <= nowDate.getTime())
        throw new TokenError("TOKEN_EXPIRED", "Refresh token expired");
      if (current.status === "active" && current.familyExpiresAt.getTime() <= nowDate.getTime())
        throw new TokenError("TOKEN_EXPIRED", "Refresh family expired");
      if (
        current.status === "active" &&
        current.rotationCount >= Math.min(current.rotationLimit, this.familyMaxRotations)
      )
        throw new TokenError("REFRESH_LIMIT", "Refresh family reached its rotation limit");

      // Complete fallible identity lookups and signing before consuming the current token.
      // The store still performs the final status check and rotation atomically.
      let accessToken: string | undefined;
      if (current.status === "active") {
        const context = (await this.options.resolveRefreshContext?.(current.sub)) ?? {};
        this.assertClaims(context.claims ?? {});
        const audience = context.audience ?? this.options.audience;
        this.assertAudience(audience);
        accessToken = await this.signAccess(
          current.sub,
          context.claims ?? {},
          audience,
          this.accessTtl,
          now
        );
      }
      // Context resolution and key signing can be slow. Recheck expiration at
      // the rotation boundary instead of letting an earlier timestamp extend
      // an expired credential's authority.
      const rotationDate = this.now();
      const rotationNow = Math.floor(rotationDate.getTime() / 1000);
      if (current.status === "active" && current.expiresAt.getTime() <= rotationDate.getTime())
        throw new TokenError("TOKEN_EXPIRED", "Refresh token expired");
      const replacementToken = newOpaqueToken(current.familyId);
      const placeholder = this.refreshRecord(
        replacementToken,
        current.sub,
        current.familyId,
        this.refreshTtl,
        rotationNow,
        current.familyExpiresAt,
        current.rotationCount + 1,
        current.rotationLimit
      );
      const result = await this.store(() =>
        this.options.tokenStore.rotateRefreshToken({
          currentHash,
          replacement: placeholder,
          now: rotationDate,
          maxRotations: this.familyMaxRotations
        })
      );
      if (result.status === "missing")
        throw new TokenError("INVALID_TOKEN", "Unknown refresh token");
      if (result.status === "expired")
        throw new TokenError("TOKEN_EXPIRED", "Refresh token expired");
      if (result.status === "limit")
        throw new TokenError("REFRESH_LIMIT", "Refresh family reached its rotation limit");
      if (result.status === "revoked")
        throw new TokenError("TOKEN_REVOKED", "Refresh token revoked");
      if (result.status === "reused") {
        await this.store(() => this.options.tokenStore.revokeByFamily(result.record.familyId));
        this.emit({
          name: "token.reuse_detected",
          at: nowDate.toISOString(),
          sub: result.record.sub,
          familyId: result.record.familyId
        });
        throw new TokenError(
          "REFRESH_REUSED",
          "Refresh token reuse detected; token family revoked"
        );
      }
      if (!accessToken)
        throw new TokenError("KEY_UNAVAILABLE", "Could not prepare the access token");
      const { sub, familyId } = result.record;
      this.emit({
        name: "token.refreshed",
        at: nowDate.toISOString(),
        sub,
        familyId
      });
      return {
        accessToken,
        refreshToken: replacementToken,
        tokenType: "Bearer",
        expiresIn: this.accessTtl,
        refreshExpiresIn: Math.max(
          0,
          Math.floor((result.record.expiresAt.getTime() - rotationDate.getTime()) / 1000)
        ),
        issuedAt: now
      };
    } catch (error) {
      throw this.failure(error);
    }
  }

  async revoke(raw: RevokeInput): Promise<void> {
    try {
      const input = parseRevoke(raw);
      const nowDate = this.now();
      if ("refreshToken" in input) {
        const familyId = familyIdFromOpaqueToken(input.refreshToken);
        if (!familyId) throw new TokenError("INVALID_INPUT", "Invalid refresh token");
        await this.store(() =>
          this.options.tokenStore.revokeByHash(hashToken(input.refreshToken), familyId)
        );
      } else if ("familyId" in input)
        await this.store(() => this.options.tokenStore.revokeByFamily(input.familyId));
      else if ("sub" in input)
        await this.store(() => this.options.tokenStore.revokeBySub(input.sub));
      else {
        if (!this.options.tokenStore.revokeByJti)
          throw new TokenError(
            "UNSUPPORTED_OPERATION",
            "This store does not support access-token revocation"
          );
        const clockTolerance = this.options.clockSkewSeconds ?? 30;
        const expiresAt = new Date(nowDate.getTime() + (ACCESS_MAX + clockTolerance) * 1000);
        await this.store(() => this.options.tokenStore.revokeByJti!(input.jti, expiresAt));
      }
      this.emit({
        name: "token.revoked",
        at: this.now().toISOString(),
        ...("jti" in input ? { jti: input.jti } : {}),
        ...("sub" in input ? { sub: input.sub } : {}),
        ...("familyId" in input ? { familyId: input.familyId } : {})
      });
    } catch (error) {
      throw this.failure(error);
    }
  }

  private async signAccess(
    sub: string,
    claims: Record<string, unknown>,
    audience: string | string[] | undefined,
    ttl: number,
    now: number
  ): Promise<string> {
    const customSigner = this.options.keyProvider.signJwt;
    if (customSigner) {
      const signing = await this.options.keyProvider.getSigningKey?.();
      if (!signing)
        throw new TokenError(
          "KEY_UNAVAILABLE",
          "Custom signer must expose its active kid and algorithm"
        );
      if (!this.allowedAlgorithms.includes(signing.alg))
        throw new TokenError("KEY_UNAVAILABLE", "Signing key algorithm is not allowed");
      return customSigner.call(this.options.keyProvider, {
        protectedHeader: { alg: signing.alg, kid: signing.kid, typ: "JWT" },
        payload: {
          ...claims,
          type: "access",
          iss: this.options.issuer,
          sub,
          ...(audience ? { aud: audience } : {}),
          iat: now,
          exp: now + ttl,
          jti: newId()
        }
      });
    }
    const signing = await this.options.keyProvider.getSigningKey?.();
    if (!signing) throw new TokenError("KEY_UNAVAILABLE", "No signing capability is configured");
    if (!this.allowedAlgorithms.includes(signing.alg))
      throw new TokenError("KEY_UNAVAILABLE", "Signing key algorithm is not allowed");
    let jwt = new SignJWT({ ...claims, type: "access" })
      .setProtectedHeader({ alg: signing.alg, kid: signing.kid, typ: "JWT" })
      .setIssuer(this.options.issuer)
      .setSubject(sub)
      .setIssuedAt(now)
      .setExpirationTime(now + ttl)
      .setJti(newId());
    if (audience) jwt = jwt.setAudience(audience);
    return jwt.sign(signing.key);
  }

  private assertClaims(claims: Record<string, unknown>): void {
    if (jsonByteLength(claims) > (this.options.maxClaimsBytes ?? 4096))
      throw new TokenError("CLAIMS_TOO_LARGE", "Custom claims exceed the configured byte limit");
    const inspect = (value: unknown, root = true): void => {
      if (!value || typeof value !== "object") return;
      for (const [key, nested] of Object.entries(value)) {
        const normalized = key.toLowerCase();
        if (root && RESERVED.has(normalized))
          throw new TokenError("FORBIDDEN_CLAIM", `Reserved claim: ${key}`);
        if (this.forbidden.has(normalized))
          throw new TokenError("FORBIDDEN_CLAIM", `Forbidden claim: ${key}`);
        inspect(nested, false);
      }
    };
    inspect(claims);
  }

  private assertAudience(requested: string | string[] | undefined): void {
    if (!requested || !this.options.audience) return;
    const allowed = new Set(
      Array.isArray(this.options.audience) ? this.options.audience : [this.options.audience]
    );
    const values = Array.isArray(requested) ? requested : [requested];
    if (values.some((audience) => !allowed.has(audience)))
      throw new TokenError("INVALID_INPUT", "Requested audience is not allowed");
  }

  private refreshRecord(
    token: string,
    sub: string,
    familyId: string,
    ttl: number,
    now: number,
    familyExpiresAt: Date,
    rotationCount: number,
    rotationLimit: number
  ): RefreshTokenRecord {
    const expiresAt = new Date(Math.min((now + ttl) * 1000, familyExpiresAt.getTime()));
    return {
      tokenHash: hashToken(token),
      sub,
      familyId,
      status: "active",
      createdAt: new Date(now * 1000),
      expiresAt,
      familyExpiresAt,
      rotationCount,
      rotationLimit
    };
  }
  private toAccessClaims(payload: JWTPayload): AccessTokenClaims {
    if (
      payload.type !== "access" ||
      typeof payload.sub !== "string" ||
      typeof payload.iss !== "string" ||
      typeof payload.exp !== "number" ||
      typeof payload.iat !== "number" ||
      typeof payload.jti !== "string"
    )
      throw new TokenError("INVALID_TOKEN", "Invalid access-token claims");
    return payload as AccessTokenClaims;
  }
  private now(): Date {
    return this.options.now?.() ?? new Date();
  }
  private nowSeconds(): number {
    return Math.floor(this.now().getTime() / 1000);
  }
  private emit(event: TokenEvent): void {
    try {
      this.options.onEvent?.(event);
    } catch {
      /* telemetry must not change auth behavior */
    }
  }
  private async store<T>(operation: () => Promise<T>): Promise<T> {
    try {
      return await operation();
    } catch (cause) {
      if (isTokenError(cause)) throw cause;
      throw new TokenError("STORE_UNAVAILABLE", "Token store operation failed", { cause });
    }
  }
  private failure(error: unknown): TokenError {
    const mapped = isTokenError(error)
      ? error
      : error instanceof JoseErrors.JWTExpired
        ? new TokenError("TOKEN_EXPIRED", "Access token expired", {
            cause: error
          })
        : error instanceof JoseErrors.JOSEError
          ? new TokenError("INVALID_TOKEN", "Access token validation failed", {
              cause: error
            })
          : new TokenError("KEY_UNAVAILABLE", "Token operation failed", {
              cause: error
            });
    this.emit({
      name: "token.failed",
      at: this.now().toISOString(),
      code: mapped.code
    });
    return mapped;
  }
}

/** Creates an isolated token service. Configuration is validated immediately. */
export function createTokenService(options: TokenServiceOptions): TokenService {
  return new DefaultTokenService(options);
}

core/types.ts defines service, event, key provider, and token store contracts.

import type { JWK } from "jose";
import type { IssueInput } from "./validation.js";

export type SigningAlgorithm =
  | "RS256"
  | "RS384"
  | "RS512"
  | "PS256"
  | "PS384"
  | "PS512"
  | "ES256"
  | "ES384"
  | "ES512"
  | "EdDSA"
  | "HS256"
  | "HS384"
  | "HS512";

export interface TokenPair {
  accessToken?: string;
  refreshToken?: string;
  tokenType: "Bearer";
  expiresIn?: number;
  refreshExpiresIn?: number;
  issuedAt: number;
}

export interface AccessTokenClaims {
  sub: string;
  iss: string;
  aud?: string | string[];
  exp: number;
  iat: number;
  jti: string;
  type: "access";
  [key: string]: unknown;
}

export interface VerifyResult {
  valid: true;
  claims: AccessTokenClaims;
}

export type RevokeInput =
  | { refreshToken: string }
  | { jti: string }
  | { sub: string }
  | { familyId: string };

export type TokenErrorCode =
  | "INVALID_INPUT"
  | "INVALID_TOKEN"
  | "TOKEN_EXPIRED"
  | "TOKEN_REVOKED"
  | "REFRESH_REUSED"
  | "REFRESH_LIMIT"
  | "CLAIMS_TOO_LARGE"
  | "FORBIDDEN_CLAIM"
  | "KEY_UNAVAILABLE"
  | "STORE_UNAVAILABLE"
  | "UNSUPPORTED_OPERATION";

export type TokenEvent = {
  name:
    | "token.issued"
    | "token.verified"
    | "token.refreshed"
    | "token.revoked"
    | "token.reuse_detected"
    | "token.failed";
  at: string;
  sub?: string;
  jti?: string;
  familyId?: string;
  code?: TokenErrorCode;
};

export interface SigningKey {
  key: CryptoKey | Uint8Array;
  kid: string;
  alg: SigningAlgorithm;
}
export interface VerificationKey {
  key: CryptoKey | Uint8Array;
  kid?: string;
  alg: SigningAlgorithm;
}

/** Supplies private signing keys and public verification keys. Providers may rotate by changing the returned `kid`. */
export interface ProviderSignInput {
  protectedHeader: { alg: SigningAlgorithm; kid: string; typ: "JWT" };
  payload: Record<string, unknown>;
}

export interface KeyProvider {
  getSigningKey?(): Promise<SigningKey>;
  signJwt?(input: ProviderSignInput): Promise<string>;
  getVerificationKey(input: { kid?: string; alg: SigningAlgorithm }): Promise<VerificationKey>;
  getPublicJwks?(): Promise<{ keys: JWK[] }>;
}

export interface RefreshTokenRecord {
  tokenHash: string;
  familyId: string;
  sub: string;
  expiresAt: Date;
  /** Fixed absolute deadline shared by every descendant in the family. */
  familyExpiresAt: Date;
  /** Number of successful rotations already performed in this family. */
  rotationCount: number;
  /** Rotation ceiling fixed when the family was issued. */
  rotationLimit: number;
  /** Redis-internal generation used to invalidate every family for a subject. */
  subjectGeneration?: string;
  status: "active" | "used" | "revoked";
  createdAt: Date;
  usedAt?: Date;
  metadata?: Record<string, unknown>;
}

export type RotateResult =
  | { status: "rotated"; record: RefreshTokenRecord }
  | { status: "missing" }
  | { status: "expired"; record: RefreshTokenRecord }
  | { status: "limit"; record: RefreshTokenRecord }
  | { status: "reused"; record: RefreshTokenRecord }
  | { status: "revoked"; record: RefreshTokenRecord };

/** Persistence boundary. Rotation and reuse-triggered family revocation must each be atomic across all instances. */
export interface TokenStore {
  /**
   * Optional store-specific family ID. Return a UUID or `<43-char-base64url-route>.<UUID>`;
   * Redis uses the latter to route a subject's state to one cluster slot.
   */
  createFamilyId?(sub: string): string;
  saveRefreshToken(record: RefreshTokenRecord): Promise<void>;
  getRefreshToken(tokenHash: string, familyId?: string): Promise<RefreshTokenRecord | undefined>;
  rotateRefreshToken(input: {
    currentHash: string;
    replacement: RefreshTokenRecord;
    now: Date;
    maxRotations: number;
  }): Promise<RotateResult>;
  revokeByHash(tokenHash: string, familyId?: string): Promise<void>;
  revokeByFamily(familyId: string): Promise<void>;
  revokeBySub(sub: string): Promise<void>;
  revokeByJti?(jti: string, expiresAt?: Date): Promise<void>;
  isJtiRevoked?(jti: string): Promise<boolean>;
}

export interface TokenServiceOptions {
  /** Defaults to production; development permits localhost HTTP and an omitted audience. */
  deploymentMode?: "development" | "production";
  issuer: string;
  audience?: string | string[];
  accessTokenTtlSeconds?: number;
  refreshTokenTtlSeconds?: number;
  /** Absolute lifetime of one refresh family; defaults to 90 days. */
  maxRefreshFamilyLifetimeSeconds?: number;
  /** Maximum successful rotations per family; defaults to 10,000. */
  maxRefreshRotations?: number;
  clockSkewSeconds?: number;
  algorithms?: SigningAlgorithm[];
  keyProvider: KeyProvider;
  tokenStore: TokenStore;
  maxClaimsBytes?: number;
  forbiddenClaimKeys?: string[];
  onEvent?: (event: TokenEvent) => void;
  /** Reload current authorization claims instead of copying stale data from refresh tokens. */
  resolveRefreshContext?: (sub: string) => Promise<{
    claims?: Record<string, unknown>;
    audience?: string | string[];
  }>;
  now?: () => Date;
}

export interface TokenService {
  issue(input: IssueInput): Promise<TokenPair>;
  verify(accessToken: string): Promise<VerifyResult>;
  refresh(refreshToken: string): Promise<TokenPair>;
  revoke(input: RevokeInput): Promise<void>;
}

core/validation.ts validates options, issue input, revocation targets, and allowed algorithms.

import { z } from "zod";
import { TokenError } from "./errors.js";
import type { RevokeInput, SigningAlgorithm, TokenServiceOptions } from "./types.js";

const algorithms = [
  "RS256",
  "RS384",
  "RS512",
  "PS256",
  "PS384",
  "PS512",
  "ES256",
  "ES384",
  "ES512",
  "EdDSA",
  "HS256",
  "HS384",
  "HS512"
] as const;
const issueSchema = z
  .object({
    sub: z.string().trim().min(1).max(512),
    claims: z.record(z.string(), z.unknown()).optional(),
    tokens: z
      .object({
        access: z.boolean().optional(),
        refresh: z.boolean().optional()
      })
      .optional(),
    accessTokenTtlSeconds: z.number().int().positive().optional(),
    refreshTokenTtlSeconds: z.number().int().positive().optional(),
    audience: z.union([z.string().min(1), z.array(z.string().min(1)).min(1)]).optional()
  })
  .strict();
export type IssueInput = z.infer<typeof issueSchema>;
const optionsSchema = z
  .object({
    deploymentMode: z.enum(["development", "production"]).optional(),
    issuer: z.string().trim().min(1).max(2048),
    audience: z.union([z.string().min(1), z.array(z.string().min(1)).min(1)]).optional(),
    accessTokenTtlSeconds: z.number().int().positive().optional(),
    refreshTokenTtlSeconds: z.number().int().positive().optional(),
    maxRefreshFamilyLifetimeSeconds: z.number().int().min(3600).max(31_536_000).optional(),
    maxRefreshRotations: z.number().int().min(1).max(1_000_000).optional(),
    clockSkewSeconds: z.number().int().min(0).max(300).optional(),
    algorithms: z.array(z.enum(algorithms)).min(1).optional(),
    keyProvider: z.custom<object>((v) => typeof v === "object" && v !== null),
    tokenStore: z.custom<object>((v) => typeof v === "object" && v !== null),
    maxClaimsBytes: z.number().int().min(256).max(16_384).optional(),
    forbiddenClaimKeys: z.array(z.string().min(1)).optional(),
    onEvent: z.function().optional(),
    resolveRefreshContext: z.function().optional(),
    now: z.function().optional()
  })
  .strict()
  .superRefine((value, ctx) => {
    const validateMethods = (
      candidate: object,
      path: string,
      required: string[],
      optional: string[]
    ) => {
      const methods = candidate as Record<string, unknown>;
      for (const method of required) {
        if (typeof methods[method] !== "function")
          ctx.addIssue({
            code: "custom",
            path: [path, method],
            message: `${path}.${method} must be a function`
          });
      }
      for (const method of optional) {
        if (methods[method] !== undefined && typeof methods[method] !== "function")
          ctx.addIssue({
            code: "custom",
            path: [path, method],
            message: `${path}.${method} must be a function when provided`
          });
      }
    };
    validateMethods(
      value.keyProvider,
      "keyProvider",
      ["getVerificationKey"],
      ["getSigningKey", "signJwt", "getPublicJwks"]
    );
    const provider = value.keyProvider as Record<string, unknown>;
    if (typeof provider.signJwt === "function" && typeof provider.getSigningKey !== "function")
      ctx.addIssue({
        code: "custom",
        path: ["keyProvider", "getSigningKey"],
        message: "keyProvider.getSigningKey is required when keyProvider.signJwt is provided"
      });
    validateMethods(
      value.tokenStore,
      "tokenStore",
      [
        "saveRefreshToken",
        "getRefreshToken",
        "rotateRefreshToken",
        "revokeByHash",
        "revokeByFamily",
        "revokeBySub"
      ],
      ["revokeByJti", "isJtiRevoked"]
    );
    const store = value.tokenStore as Record<string, unknown>;
    if ((typeof store.revokeByJti === "function") !== (typeof store.isJtiRevoked === "function"))
      ctx.addIssue({
        code: "custom",
        path: ["tokenStore"],
        message: "tokenStore must implement revokeByJti and isJtiRevoked together"
      });
    if ((value.deploymentMode ?? "production") === "production") {
      if (!value.audience)
        ctx.addIssue({
          code: "custom",
          path: ["audience"],
          message: "Audience is required in production"
        });
      if (!(value.issuer.startsWith("https://") || value.issuer.startsWith("urn:")))
        ctx.addIssue({
          code: "custom",
          path: ["issuer"],
          message: "Production issuer must be HTTPS or a URN"
        });
      if (value.algorithms?.some((alg) => alg.startsWith("HS")))
        ctx.addIssue({
          code: "custom",
          path: ["algorithms"],
          message: "Symmetric JWT algorithms are not allowed in production mode"
        });
    }
  });

export function parseIssue(input: IssueInput): IssueInput {
  const result = issueSchema.safeParse(input);
  if (!result.success)
    throw new TokenError("INVALID_INPUT", "Invalid token issue input", {
      cause: result.error
    });
  return result.data as IssueInput;
}

export function parseOptions(input: TokenServiceOptions): TokenServiceOptions {
  const result = optionsSchema.safeParse(input);
  if (!result.success)
    throw new TokenError("INVALID_INPUT", "Invalid token service options", {
      cause: result.error
    });
  return input;
}

export function parseRevoke(input: RevokeInput): RevokeInput {
  if (!input || typeof input !== "object")
    throw new TokenError("INVALID_INPUT", "Invalid revocation target");
  const entries = Object.entries(input).filter(
    ([, value]) => typeof value === "string" && value.length > 0
  );
  if (
    entries.length !== 1 ||
    Object.keys(input).length !== 1 ||
    !["refreshToken", "jti", "sub", "familyId"].includes(entries[0]![0])
  ) {
    throw new TokenError("INVALID_INPUT", "Exactly one valid revocation target is required");
  }
  return input;
}

export function assertAllowedAlgorithm(
  alg: string,
  allowed: SigningAlgorithm[]
): asserts alg is SigningAlgorithm {
  if (!allowed.includes(alg as SigningAlgorithm))
    throw new TokenError("INVALID_TOKEN", "Token algorithm is not allowed");
}

Key providers

providers/kms-rsa-key-provider.ts adapts an external RS256 signing operation.

import type { KeyProvider, ProviderSignInput, SigningKey, VerificationKey } from "../core/types.js";
import { TokenError } from "../core/errors.js";

export interface KmsRsaSigner {
  /** Return an RSASSA-PKCS1-v1_5 SHA-256 signature over `data`. */
  sign(data: Uint8Array): Promise<Uint8Array>;
}

/** RS256 provider for AWS/GCP/Azure/HSM clients without coupling the block to a cloud SDK. */
export class KmsRsaKeyProvider implements KeyProvider {
  constructor(
    private readonly input: { kid: string; publicKey: CryptoKey; signer: KmsRsaSigner }
  ) {}
  async getSigningKey(): Promise<SigningKey> {
    // Public key is returned only as signing metadata; signJwt performs the private operation.
    return { key: this.input.publicKey, kid: this.input.kid, alg: "RS256" };
  }
  async signJwt(input: ProviderSignInput): Promise<string> {
    const encoder = new TextEncoder();
    const header = Buffer.from(JSON.stringify(input.protectedHeader)).toString("base64url");
    const payload = Buffer.from(JSON.stringify(input.payload)).toString("base64url");
    const signingInput = `${header}.${payload}`;
    const signature = await this.input.signer.sign(encoder.encode(signingInput));
    return `${signingInput}.${Buffer.from(signature).toString("base64url")}`;
  }
  async getVerificationKey(input: {
    kid?: string;
    alg: import("../core/types.js").SigningAlgorithm;
  }): Promise<VerificationKey> {
    if (input.kid !== this.input.kid || input.alg !== "RS256")
      throw new TokenError("KEY_UNAVAILABLE", "No matching KMS verification key");
    return { key: this.input.publicKey, kid: this.input.kid, alg: "RS256" };
  }
}

providers/local-key-provider.ts loads PEM signing and verification keys and publishes their public JWKS.

import { exportJWK, importPKCS8, importSPKI, type JWK } from "jose";
import { TokenError } from "../core/errors.js";
import type { KeyProvider, SigningAlgorithm, SigningKey, VerificationKey } from "../core/types.js";

export interface LocalKey {
  kid: string;
  alg: SigningAlgorithm;
  privateKeyPem: string;
  publicKeyPem: string;
}

/** PEM-backed provider supporting a current signing key and old verification keys during rotation. */
export class LocalKeyProvider implements KeyProvider {
  private readonly keys: LocalKey[];
  private readonly signingKid: string;
  private readonly privateCache = new Map<string, CryptoKey>();
  private readonly publicCache = new Map<string, CryptoKey>();
  constructor(input: { keys: LocalKey[]; signingKid: string }) {
    if (!input.keys.length || new Set(input.keys.map((key) => key.kid)).size !== input.keys.length)
      throw new TokenError("INVALID_INPUT", "Local keys require unique kid values");
    if (!input.keys.some((key) => key.kid === input.signingKid))
      throw new TokenError("INVALID_INPUT", "signingKid was not found");
    this.keys = input.keys;
    this.signingKid = input.signingKid;
  }
  async getSigningKey(): Promise<SigningKey> {
    const entry = this.keys.find((key) => key.kid === this.signingKid)!;
    let key = this.privateCache.get(entry.kid);
    if (!key) {
      key = await importPKCS8(entry.privateKeyPem, entry.alg);
      this.privateCache.set(entry.kid, key);
    }
    return { key, kid: entry.kid, alg: entry.alg };
  }
  async getVerificationKey(input: {
    kid?: string;
    alg: SigningAlgorithm;
  }): Promise<VerificationKey> {
    if (!input.kid) throw new TokenError("INVALID_TOKEN", "Token is missing kid");
    const entry = this.keys.find((key) => key.kid === input.kid && key.alg === input.alg);
    if (!entry) throw new TokenError("KEY_UNAVAILABLE", "No matching verification key");
    let key = this.publicCache.get(entry.kid);
    if (!key) {
      key = await importSPKI(entry.publicKeyPem, entry.alg);
      this.publicCache.set(entry.kid, key);
    }
    return { key, kid: entry.kid, alg: entry.alg };
  }
  async getPublicJwks(): Promise<{ keys: JWK[] }> {
    const keys = await Promise.all(
      this.keys.map(async (entry) => ({
        ...(await exportJWK(
          (await this.getVerificationKey({ kid: entry.kid, alg: entry.alg })).key as CryptoKey
        )),
        kid: entry.kid,
        alg: entry.alg,
        use: "sig"
      }))
    );
    return { keys };
  }
}

providers/remote-jwks-provider.ts resolves verification keys from an HTTPS JWKS endpoint.

import { createRemoteJWKSet, type JWTHeaderParameters } from "jose";
import { TokenError } from "../core/errors.js";
import type { KeyProvider, SigningAlgorithm, SigningKey, VerificationKey } from "../core/types.js";

/** Verification-only provider for resource services. Issuing through this provider intentionally fails closed. */
export class RemoteJwksKeyProvider implements KeyProvider {
  private readonly resolver: ReturnType<typeof createRemoteJWKSet>;
  constructor(input: {
    jwksUri: string;
    cooldownDurationMs?: number;
    cacheMaxAgeMs?: number;
    timeoutDurationMs?: number;
  }) {
    const url = new URL(input.jwksUri);
    if (url.protocol !== "https:")
      throw new TokenError("INVALID_INPUT", "Remote JWKS URL must use HTTPS");
    this.resolver = createRemoteJWKSet(url, {
      cooldownDuration: input.cooldownDurationMs ?? 30_000,
      cacheMaxAge: input.cacheMaxAgeMs ?? 600_000,
      timeoutDuration: input.timeoutDurationMs ?? 5_000
    });
  }
  async getSigningKey(): Promise<SigningKey> {
    throw new TokenError(
      "UNSUPPORTED_OPERATION",
      "Remote JWKS is verification-only; inject a signing/KMS provider in the issuer"
    );
  }
  async getVerificationKey(input: {
    kid?: string;
    alg: SigningAlgorithm;
  }): Promise<VerificationKey> {
    const header: JWTHeaderParameters = {
      alg: input.alg,
      ...(input.kid ? { kid: input.kid } : {})
    };
    const resolved = await this.resolver(header);
    return { key: resolved, ...(input.kid ? { kid: input.kid } : {}), alg: input.alg };
  }
}

Token stores

stores/memory-token-store.ts is intended for tests and single process development.

import type { RefreshTokenRecord, RotateResult, TokenStore } from "../core/types.js";

/** Development/test store. It is process-local and must not be used across multiple instances. */
export class MemoryTokenStore implements TokenStore {
  private readonly refresh = new Map<string, RefreshTokenRecord>();
  private readonly revokedJtis = new Map<string, number>();
  async saveRefreshToken(record: RefreshTokenRecord): Promise<void> {
    if (this.refresh.has(record.tokenHash)) throw new Error("Duplicate refresh hash");
    this.refresh.set(record.tokenHash, { ...record });
  }
  async getRefreshToken(tokenHash: string): Promise<RefreshTokenRecord | undefined> {
    const record = this.refresh.get(tokenHash);
    return record ? { ...record } : undefined;
  }
  async rotateRefreshToken(input: {
    currentHash: string;
    replacement: RefreshTokenRecord;
    now: Date;
    maxRotations: number;
  }): Promise<RotateResult> {
    const current = this.refresh.get(input.currentHash);
    if (!current) return { status: "missing" };
    if (current.status === "used") {
      for (const item of this.refresh.values())
        if (item.familyId === current.familyId) item.status = "revoked";
      return { status: "reused", record: { ...current } };
    }
    if (current.status === "revoked") return { status: "revoked", record: { ...current } };
    if (current.expiresAt.getTime() <= input.now.getTime())
      return { status: "expired", record: { ...current } };
    if (current.familyExpiresAt.getTime() <= input.now.getTime())
      return { status: "expired", record: { ...current } };
    if (current.rotationCount >= Math.min(current.rotationLimit, input.maxRotations))
      return { status: "limit", record: { ...current } };
    current.status = "used";
    current.usedAt = input.now;
    const replacement = {
      ...input.replacement,
      sub: current.sub,
      familyId: current.familyId,
      familyExpiresAt: current.familyExpiresAt,
      rotationCount: current.rotationCount + 1,
      rotationLimit: Math.min(current.rotationLimit, input.maxRotations),
      expiresAt: new Date(
        Math.min(input.replacement.expiresAt.getTime(), current.familyExpiresAt.getTime())
      )
    };
    this.refresh.set(replacement.tokenHash, replacement);
    return { status: "rotated", record: { ...replacement } };
  }
  async revokeByHash(hash: string, _familyId?: string): Promise<void> {
    const item = this.refresh.get(hash);
    if (item) await this.revokeByFamily(item.familyId);
  }
  async revokeByFamily(familyId: string): Promise<void> {
    for (const item of this.refresh.values())
      if (item.familyId === familyId) item.status = "revoked";
  }
  async revokeBySub(sub: string): Promise<void> {
    for (const item of this.refresh.values()) if (item.sub === sub) item.status = "revoked";
  }
  async revokeByJti(jti: string, expiresAt?: Date): Promise<void> {
    this.revokedJtis.set(jti, expiresAt?.getTime() ?? Number.MAX_SAFE_INTEGER);
  }
  async isJtiRevoked(jti: string): Promise<boolean> {
    const expiry = this.revokedJtis.get(jti);
    if (expiry === undefined) return false;
    if (expiry <= Date.now()) {
      this.revokedJtis.delete(jti);
      return false;
    }
    return true;
  }
}

stores/postgres-token-store.ts uses a caller supplied transaction interface.

import type { RefreshTokenRecord, RotateResult, TokenStore } from "../core/types.js";

export interface QueryResult<Row = Record<string, unknown>> {
  rows: Row[];
  rowCount: number | null;
}
export interface PgTransaction {
  query<Row = Record<string, unknown>>(sql: string, values?: unknown[]): Promise<QueryResult<Row>>;
}
export interface PgLike extends PgTransaction {
  transaction<T>(work: (tx: PgTransaction) => Promise<T>): Promise<T>;
}
type RefreshRow = {
  token_hash: string;
  family_id: string;
  subject: string;
  expires_at: Date;
  family_expires_at: Date;
  rotation_count: number;
  rotation_limit: number;
  status: RefreshTokenRecord["status"];
  created_at: Date;
  used_at: Date | null;
  metadata: Record<string, unknown> | null;
};

/** PostgreSQL store using row locks for one-time refresh-token consumption. */
export class PostgresTokenStore implements TokenStore {
  constructor(
    private readonly db: PgLike,
    private readonly table = "blockend_refresh_tokens",
    private readonly denyTable = "blockend_revoked_jtis"
  ) {
    if (!/^[a-z_][a-z0-9_]*$/i.test(table) || !/^[a-z_][a-z0-9_]*$/i.test(denyTable))
      throw new Error("Unsafe SQL identifier");
  }
  async saveRefreshToken(r: RefreshTokenRecord): Promise<void> {
    await this.db.transaction(async (tx) => {
      await this.setReadCommitted(tx);
      await this.lockSubject(tx, r.sub);
      const clock = await tx.query<{ now: Date }>("SELECT clock_timestamp() AS now");
      const now = clock.rows[0]?.now;
      if (!(now instanceof Date)) throw new Error("PostgreSQL did not return its current time");
      const familyLifetime = r.familyExpiresAt.getTime() - r.createdAt.getTime();
      const refreshLifetime = r.expiresAt.getTime() - r.createdAt.getTime();
      if (familyLifetime <= 0 || refreshLifetime <= 0)
        throw new Error("Refresh family lifetime must be positive");
      const record = {
        ...r,
        createdAt: now,
        familyExpiresAt: new Date(now.getTime() + familyLifetime),
        expiresAt: new Date(
          Math.min(now.getTime() + refreshLifetime, now.getTime() + familyLifetime)
        )
      };
      await tx.query(
        `INSERT INTO ${this.table} (token_hash,family_id,subject,expires_at,family_expires_at,rotation_count,rotation_limit,status,created_at,metadata) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10)`,
        [
          record.tokenHash,
          record.familyId,
          record.sub,
          record.expiresAt,
          record.familyExpiresAt,
          record.rotationCount,
          record.rotationLimit,
          record.status,
          record.createdAt,
          record.metadata ?? null
        ]
      );
    });
  }
  async getRefreshToken(
    tokenHash: string,
    _familyId?: string
  ): Promise<RefreshTokenRecord | undefined> {
    const result = await this.db.query<RefreshRow>(
      `SELECT * FROM ${this.table} WHERE token_hash=$1`,
      [tokenHash]
    );
    const row = result.rows[0];
    return row ? this.fromRow(row) : undefined;
  }
  async rotateRefreshToken(input: {
    currentHash: string;
    replacement: RefreshTokenRecord;
    now: Date;
    maxRotations: number;
  }): Promise<RotateResult> {
    return this.db.transaction(async (tx) => {
      await this.setReadCommitted(tx);
      const identity = await tx.query<{ subject: string }>(
        `SELECT subject FROM ${this.table} WHERE token_hash=$1`,
        [input.currentHash]
      );
      const subject = identity.rows[0]?.subject;
      if (!subject) return { status: "missing" };
      await this.lockSubject(tx, subject);
      const found = await tx.query<RefreshRow>(
        `SELECT * FROM ${this.table} WHERE token_hash=$1 FOR UPDATE`,
        [input.currentHash]
      );
      const row = found.rows[0];
      if (!row) return { status: "missing" };
      const current = this.fromRow(row);
      const clock = await tx.query<{ now: Date }>("SELECT clock_timestamp() AS now");
      const databaseNow = clock.rows[0]?.now;
      if (!(databaseNow instanceof Date))
        throw new Error("PostgreSQL did not return its current time");
      if (current.status === "used") {
        await tx.query(`UPDATE ${this.table} SET status='revoked' WHERE family_id=$1`, [
          current.familyId
        ]);
        return { status: "reused", record: current };
      }
      if (current.status === "revoked") return { status: "revoked", record: current };
      if (current.expiresAt.getTime() <= databaseNow.getTime())
        return { status: "expired", record: current };
      if (current.familyExpiresAt.getTime() <= databaseNow.getTime())
        return { status: "expired", record: current };
      if (current.rotationCount >= Math.min(current.rotationLimit, input.maxRotations))
        return { status: "limit", record: current };
      const refreshLifetime =
        input.replacement.expiresAt.getTime() - input.replacement.createdAt.getTime();
      if (refreshLifetime <= 0) return { status: "expired", record: current };
      await tx.query(`UPDATE ${this.table} SET status='used',used_at=$2 WHERE token_hash=$1`, [
        input.currentHash,
        databaseNow
      ]);
      const next = {
        ...input.replacement,
        sub: current.sub,
        familyId: current.familyId,
        familyExpiresAt: current.familyExpiresAt,
        rotationCount: current.rotationCount + 1,
        rotationLimit: Math.min(current.rotationLimit, input.maxRotations),
        createdAt: databaseNow,
        expiresAt: new Date(
          Math.min(databaseNow.getTime() + refreshLifetime, current.familyExpiresAt.getTime())
        )
      };
      await tx.query(
        `INSERT INTO ${this.table} (token_hash,family_id,subject,expires_at,family_expires_at,rotation_count,rotation_limit,status,created_at,metadata) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10)`,
        [
          next.tokenHash,
          next.familyId,
          next.sub,
          next.expiresAt,
          next.familyExpiresAt,
          next.rotationCount,
          next.rotationLimit,
          next.status,
          next.createdAt,
          next.metadata ?? null
        ]
      );
      return { status: "rotated", record: next };
    });
  }
  async revokeByHash(hash: string): Promise<void> {
    await this.db.transaction(async (tx) => {
      await this.setReadCommitted(tx);
      const found = await tx.query<{ subject: string; family_id: string }>(
        `SELECT subject,family_id FROM ${this.table} WHERE token_hash=$1`,
        [hash]
      );
      const identity = found.rows[0];
      if (!identity) return;
      await this.lockSubject(tx, identity.subject);
      await tx.query(
        `UPDATE ${this.table} SET status='revoked' WHERE family_id=$1 AND status='active'`,
        [identity.family_id]
      );
    });
  }
  async revokeByFamily(id: string): Promise<void> {
    await this.db.transaction(async (tx) => {
      await this.setReadCommitted(tx);
      const subjects = await tx.query<{ subject: string }>(
        `SELECT DISTINCT subject FROM ${this.table} WHERE family_id=$1 ORDER BY subject`,
        [id]
      );
      for (const { subject } of subjects.rows) await this.lockSubject(tx, subject);
      await tx.query(
        `UPDATE ${this.table} SET status='revoked' WHERE family_id=$1 AND status='active'`,
        [id]
      );
    });
  }
  async revokeBySub(sub: string): Promise<void> {
    await this.db.transaction(async (tx) => {
      await this.setReadCommitted(tx);
      await this.lockSubject(tx, sub);
      await tx.query(
        `UPDATE ${this.table} SET status='revoked' WHERE subject=$1 AND status='active'`,
        [sub]
      );
    });
  }
  async revokeByJti(jti: string, expiresAt?: Date): Promise<void> {
    await this.db.query(
      `INSERT INTO ${this.denyTable} (jti,expires_at) VALUES ($1,$2) ON CONFLICT (jti) DO UPDATE SET expires_at=EXCLUDED.expires_at`,
      [jti, expiresAt ?? new Date(Date.now() + 3_600_000)]
    );
    await this.db.query(
      `WITH expired AS (SELECT ctid FROM ${this.denyTable} WHERE expires_at <= NOW() LIMIT 1000) DELETE FROM ${this.denyTable} AS revoked USING expired WHERE revoked.ctid=expired.ctid`
    );
  }
  async isJtiRevoked(jti: string): Promise<boolean> {
    const r = await this.db.query(
      `SELECT 1 FROM ${this.denyTable} WHERE jti=$1 AND expires_at>NOW()`,
      [jti]
    );
    return (r.rowCount ?? 0) > 0;
  }
  private fromRow(r: RefreshRow): RefreshTokenRecord {
    return {
      tokenHash: r.token_hash,
      familyId: r.family_id,
      sub: r.subject,
      expiresAt: r.expires_at,
      familyExpiresAt: r.family_expires_at,
      rotationCount: r.rotation_count,
      rotationLimit: r.rotation_limit,
      status: r.status,
      createdAt: r.created_at,
      ...(r.used_at ? { usedAt: r.used_at } : {}),
      ...(r.metadata ? { metadata: r.metadata } : {})
    };
  }

  private async lockSubject(tx: PgTransaction, subject: string): Promise<void> {
    await tx.query(`SELECT pg_advisory_xact_lock(hashtext($1), hashtext($2))`, [
      "blockend_refresh_subject",
      subject
    ]);
  }

  private async setReadCommitted(tx: PgTransaction): Promise<void> {
    await tx.query("SET TRANSACTION ISOLATION LEVEL READ COMMITTED");
  }
}

stores/redis-token-store.ts uses Lua scripts for refresh rotation and revocation.

import { createHash, createHmac, randomBytes, randomUUID } from "node:crypto";
import type { RefreshTokenRecord, RotateResult, TokenStore } from "../core/types.js";

export interface RedisLike {
  get(key: string): Promise<string | null>;
  hget(key: string, field: string): Promise<string | null>;
  set(key: string, value: string, options?: { NX?: boolean; PX?: number }): Promise<unknown>;
  eval(script: string, input: { keys: string[]; arguments: string[] }): Promise<unknown>;
}
export interface RedisTokenStoreOptions {
  /** Stable base64url-encoded secret shared by every service instance; at least 32 decoded bytes. */
  subjectRoutingKey: Uint8Array | string;
  /** Shared namespace without Redis hash-tag braces. Defaults to `blockend:v2:`. */
  prefix?: string;
}

const SAVE_SCRIPT = `
if redis.call('HEXISTS', KEYS[2], ARGV[1]) == 1 then return 0 end
local time = redis.call('TIME')
local nowMs = tonumber(time[1]) * 1000 + math.floor(tonumber(time[2]) / 1000)
local record = cjson.decode(ARGV[2])
local familyLifetime = tonumber(record.familyExpiresAtMs) - tonumber(record.createdAtMs)
local refreshLifetime = tonumber(record.expiresAtMs) - tonumber(record.createdAtMs)
if familyLifetime <= 0 or refreshLifetime <= 0 then return -1 end
record.createdAtMs = nowMs
record.familyExpiresAtMs = nowMs + familyLifetime
record.expiresAtMs = math.min(nowMs + refreshLifetime, record.familyExpiresAtMs)
redis.call('SETNX', KEYS[1], ARGV[3])
record.subjectGeneration = redis.call('GET', KEYS[1])
redis.call('HSET', KEYS[2], ARGV[1], cjson.encode(record))
local ttl = record.familyExpiresAtMs - nowMs
redis.call('PEXPIRE', KEYS[2], ttl)
if redis.call('PTTL', KEYS[1]) < ttl then redis.call('PEXPIRE', KEYS[1], ttl) end
return 1
`;
const ROTATE_SCRIPT = `
local raw = redis.call('HGET', KEYS[2], ARGV[1])
if not raw then return {'missing', ''} end
local current = cjson.decode(raw)
local generation = redis.call('GET', KEYS[1])
if redis.call('HGET', KEYS[2], '_revoked') == '1' or not generation or
   current.subjectGeneration ~= generation or current.status == 'revoked' then
  return {'revoked', raw}
end
if current.status == 'used' then
  redis.call('HSET', KEYS[2], '_revoked', '1')
  return {'reused', raw}
end
local time = redis.call('TIME')
local nowMs = tonumber(time[1]) * 1000 + math.floor(tonumber(time[2]) / 1000)
if tonumber(current.expiresAtMs) <= nowMs or
   tonumber(current.familyExpiresAtMs) <= nowMs then return {'expired', raw} end
local rotationLimit = math.min(tonumber(current.rotationLimit), tonumber(ARGV[3]))
if tonumber(current.rotationCount) >= rotationLimit then return {'limit', raw} end
if redis.call('HEXISTS', KEYS[2], ARGV[2]) == 1 then return {'collision', ''} end
local replacement = cjson.decode(ARGV[4])
local refreshLifetime = tonumber(replacement.expiresAtMs) - tonumber(replacement.createdAtMs)
if refreshLifetime <= 0 or replacement.familyId ~= current.familyId or replacement.sub ~= current.sub then
  return {'expired', raw}
end
current.status = 'used'; current.usedAtMs = nowMs
redis.call('HSET', KEYS[2], ARGV[1], cjson.encode(current))
replacement.sub = current.sub
replacement.familyId = current.familyId
replacement.subjectGeneration = generation
replacement.familyExpiresAtMs = tonumber(current.familyExpiresAtMs)
replacement.createdAtMs = nowMs
replacement.expiresAtMs = math.min(nowMs + refreshLifetime, tonumber(current.familyExpiresAtMs))
replacement.rotationCount = tonumber(current.rotationCount) + 1
replacement.rotationLimit = rotationLimit
redis.call('HSET', KEYS[2], ARGV[2], cjson.encode(replacement))
local ttl = math.max(1, tonumber(current.familyExpiresAtMs) - nowMs)
redis.call('PEXPIRE', KEYS[2], ttl)
if redis.call('PTTL', KEYS[1]) < ttl then redis.call('PEXPIRE', KEYS[1], ttl) end
return {'rotated', cjson.encode(replacement)}
`;
const REVOKE_FAMILY_SCRIPT = `
if ARGV[1] and redis.call('HEXISTS', KEYS[1], ARGV[1]) == 0 then return 0 end
if redis.call('EXISTS', KEYS[1]) == 1 then
  redis.call('HSET', KEYS[1], '_revoked', '1')
  return 1
end
return 0
`;
const REVOKE_SUBJECT_SCRIPT = `
if redis.call('EXISTS', KEYS[1]) == 0 then return 0 end
local ttl = redis.call('PTTL', KEYS[1])
if ttl > 0 then redis.call('SET', KEYS[1], ARGV[1], 'PX', ttl)
elseif ttl == -1 then redis.call('SET', KEYS[1], ARGV[1], 'KEEPTTL')
else return 0 end
return 1
`;

/** Redis-backed store. Each subject's state is co-located in one cluster slot. */
export class RedisTokenStore implements TokenStore {
  private readonly prefix: string;
  private readonly subjectRoutingKey: Uint8Array;

  constructor(
    private readonly redis: RedisLike,
    options: RedisTokenStoreOptions
  ) {
    this.prefix = options.prefix ?? "blockend:v2:";
    this.subjectRoutingKey =
      typeof options.subjectRoutingKey === "string"
        ? Buffer.from(options.subjectRoutingKey, "base64url")
        : Uint8Array.from(options.subjectRoutingKey);
    if (this.subjectRoutingKey.byteLength < 32)
      throw new Error("Redis subjectRoutingKey must contain at least 32 bytes");
    if (!/^[A-Za-z0-9:_-]{1,128}$/.test(this.prefix))
      throw new Error("Redis prefix must be 1-128 safe characters and must not contain braces");
  }

  createFamilyId(sub: string): string {
    return `${this.subjectRoute(sub)}.${randomUUID()}`;
  }

  async saveRefreshToken(record: RefreshTokenRecord): Promise<void> {
    const route = this.routeFromFamilyId(record.familyId);
    if (route !== this.subjectRoute(record.sub))
      throw new Error("Refresh family routing does not match its subject");
    const saved = await this.redis.eval(SAVE_SCRIPT, {
      keys: [this.subjectKey(route), this.familyKey(record.familyId, route)],
      arguments: [record.tokenHash, this.serialize(record), randomBytes(32).toString("base64url")]
    });
    if (Number(saved) !== 1) throw new Error("Duplicate refresh token hash");
  }

  async getRefreshToken(
    tokenHash: string,
    familyId?: string
  ): Promise<RefreshTokenRecord | undefined> {
    if (!familyId) return undefined;
    const route = this.routeFromFamilyId(familyId);
    const raw = await this.redis.hget(this.familyKey(familyId, route), tokenHash);
    if (!raw) return undefined;
    const record = this.deserialize(raw);
    if (record.tokenHash !== tokenHash || record.familyId !== familyId) return undefined;
    const familyKey = this.familyKey(familyId, route);
    const [familyRevoked, generation] = await Promise.all([
      this.redis.hget(familyKey, "_revoked"),
      this.redis.get(this.subjectKey(route))
    ]);
    if (familyRevoked === "1" || generation !== record.subjectGeneration)
      return { ...record, status: "revoked" };
    return record;
  }

  async rotateRefreshToken(input: {
    currentHash: string;
    replacement: RefreshTokenRecord;
    now: Date;
    maxRotations: number;
  }): Promise<RotateResult> {
    const route = this.routeFromFamilyId(input.replacement.familyId);
    if (route !== this.subjectRoute(input.replacement.sub))
      throw new Error("Refresh family routing does not match its subject");
    const raw = await this.redis.eval(ROTATE_SCRIPT, {
      keys: [this.subjectKey(route), this.familyKey(input.replacement.familyId, route)],
      arguments: [
        input.currentHash,
        input.replacement.tokenHash,
        String(input.maxRotations),
        this.serialize(input.replacement)
      ]
    });
    if (!Array.isArray(raw) || typeof raw[0] !== "string")
      throw new Error("Unexpected Redis rotate response");
    const status = raw[0] as RotateResult["status"] | "collision";
    if (status === "missing") return { status };
    if (status === "collision") throw new Error("Duplicate replacement refresh token hash");
    const record = this.deserialize(String(raw[1]));
    return { status, record } as RotateResult;
  }

  async revokeByHash(hash: string, familyId?: string): Promise<void> {
    if (!familyId) return;
    const route = this.routeFromFamilyId(familyId);
    await this.redis.eval(REVOKE_FAMILY_SCRIPT, {
      keys: [this.familyKey(familyId, route)],
      arguments: [hash]
    });
  }

  async revokeByFamily(id: string): Promise<void> {
    const route = this.routeFromFamilyId(id);
    await this.redis.eval(REVOKE_FAMILY_SCRIPT, {
      keys: [this.familyKey(id, route)],
      arguments: []
    });
  }

  async revokeBySub(sub: string): Promise<void> {
    const route = this.subjectRoute(sub);
    await this.redis.eval(REVOKE_SUBJECT_SCRIPT, {
      keys: [this.subjectKey(route)],
      arguments: [randomBytes(32).toString("base64url")]
    });
  }

  async revokeByJti(jti: string, expiresAt?: Date): Promise<void> {
    await this.redis.set(this.jtiKey(jti), "1", {
      PX: Math.max(1, (expiresAt?.getTime() ?? Date.now() + 3_600_000) - Date.now())
    });
  }

  async isJtiRevoked(jti: string): Promise<boolean> {
    return (await this.redis.get(this.jtiKey(jti))) !== null;
  }

  private subjectRoute(sub: string): string {
    return createHmac("sha256", this.subjectRoutingKey).update(sub, "utf8").digest("base64url");
  }

  private routeFromFamilyId(id: string): string {
    const match =
      /^([A-Za-z0-9_-]{43})\.([0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12})$/i.exec(
        id
      );
    if (!match) throw new Error("Invalid refresh family identifier");
    return match[1]!;
  }

  private subjectKey(route: string): string {
    return `${this.prefix}subject:{${route}}:state`;
  }

  private familyKey(id: string, route: string): string {
    return `${this.prefix}family:{${route}}:${id.slice(route.length + 1)}`;
  }

  private jtiKey(jti: string): string {
    const digest = createHash("sha256").update(jti, "utf8").digest("base64url");
    return `${this.prefix}jti:${digest}`;
  }

  private serialize(record: RefreshTokenRecord): string {
    return JSON.stringify({
      ...record,
      createdAtMs: record.createdAt.getTime(),
      expiresAtMs: record.expiresAt.getTime(),
      familyExpiresAtMs: record.familyExpiresAt.getTime(),
      usedAtMs: record.usedAt?.getTime(),
      createdAt: undefined,
      expiresAt: undefined,
      familyExpiresAt: undefined,
      usedAt: undefined
    });
  }

  private deserialize(raw: string): RefreshTokenRecord {
    const value = JSON.parse(raw) as Record<string, unknown>;
    const record: RefreshTokenRecord = {
      tokenHash: String(value.tokenHash),
      familyId: String(value.familyId),
      sub: String(value.sub),
      status: value.status as RefreshTokenRecord["status"],
      createdAt: new Date(Number(value.createdAtMs)),
      expiresAt: new Date(Number(value.expiresAtMs)),
      familyExpiresAt: new Date(Number(value.familyExpiresAtMs)),
      rotationCount: Number(value.rotationCount),
      rotationLimit: Number(value.rotationLimit),
      ...(typeof value.subjectGeneration === "string"
        ? { subjectGeneration: value.subjectGeneration }
        : {}),
      ...(value.usedAtMs ? { usedAt: new Date(Number(value.usedAtMs)) } : {}),
      ...(value.metadata && typeof value.metadata === "object"
        ? { metadata: value.metadata as Record<string, unknown> }
        : {})
    };
    if (
      !/^[A-Za-z0-9_-]{43}$/.test(record.tokenHash) ||
      !Number.isFinite(record.createdAt.getTime()) ||
      !Number.isFinite(record.expiresAt.getTime()) ||
      !Number.isFinite(record.familyExpiresAt.getTime()) ||
      !Number.isSafeInteger(record.rotationCount) ||
      !Number.isSafeInteger(record.rotationLimit) ||
      typeof record.subjectGeneration !== "string" ||
      !["active", "used", "revoked"].includes(record.status)
    )
      throw new Error("Invalid refresh record in Redis");
    return record;
  }
}

stores/stateless-token-store.ts rejects refresh and server side revocation operations.

import { TokenError } from "../core/errors.js";
import type { RefreshTokenRecord, RotateResult, TokenStore } from "../core/types.js";

/** Access-token-only store. Always request `{ refresh: false }` when using it. */
export class StatelessTokenStore implements TokenStore {
  async saveRefreshToken(_record: RefreshTokenRecord): Promise<void> {
    throw this.unsupported();
  }
  async getRefreshToken(_tokenHash: string): Promise<RefreshTokenRecord | undefined> {
    throw this.unsupported();
  }
  async rotateRefreshToken(_input: {
    currentHash: string;
    replacement: RefreshTokenRecord;
    now: Date;
  }): Promise<RotateResult> {
    throw this.unsupported();
  }
  async revokeByHash(_hash: string): Promise<void> {
    throw this.unsupported();
  }
  async revokeByFamily(_familyId: string): Promise<void> {
    throw this.unsupported();
  }
  async revokeBySub(_sub: string): Promise<void> {
    throw this.unsupported();
  }
  private unsupported(): TokenError {
    return new TokenError(
      "UNSUPPORTED_OPERATION",
      "Refresh tokens and server-side revocation are disabled in stateless mode"
    );
  }
}

Framework adapters

adapters/http.ts maps token errors to HTTP responses and defines protected actions.

import { TokenError } from "../core/errors.js";

export type ProtectedAction = "issue" | "revoke";
export interface HttpAdapterOptions<Request> {
  /** Required to expose privileged issue/revoke endpoints. Return false to deny. */
  authorize?: (action: ProtectedAction, request: Request) => boolean | Promise<boolean>;
  exposeIssue?: boolean;
  exposeVerify?: boolean;
  exposeRevoke?: boolean;
}

export function statusFor(error: unknown): number {
  if (!(error instanceof TokenError)) return 500;
  if (
    error.code === "INVALID_INPUT" ||
    error.code === "CLAIMS_TOO_LARGE" ||
    error.code === "FORBIDDEN_CLAIM"
  )
    return 400;
  if (error.code === "STORE_UNAVAILABLE" || error.code === "KEY_UNAVAILABLE") return 503;
  if (error.code === "UNSUPPORTED_OPERATION") return 501;
  return 401;
}

export function requireObjectBody(value: unknown): Record<string, unknown> {
  if (!value || typeof value !== "object" || Array.isArray(value))
    throw new TokenError("INVALID_INPUT", "Expected a JSON object request body");
  return value as Record<string, unknown>;
}

export function requireStringField(body: unknown, field: string): string {
  const value = requireObjectBody(body)[field];
  if (typeof value !== "string" || value.length === 0)
    throw new TokenError("INVALID_INPUT", `Expected ${field} to be a non-empty string`);
  return value;
}

export function errorBody(error: unknown): { error: { code: string; message: string } } {
  if (error instanceof TokenError) return { error: { code: error.code, message: error.message } };
  return { error: { code: "INTERNAL_ERROR", message: "Token operation failed" } };
}

adapters/express.ts provides an Express router.

import { Router, type Request, type Response } from "express";
import type { TokenService } from "../core/types.js";
import {
  errorBody,
  requireStringField,
  statusFor,
  type HttpAdapterOptions,
  type ProtectedAction
} from "./http.js";

export type ExpressTokenAdapterOptions = HttpAdapterOptions<Request>;

/** Creates opt-in JSON routes. Privileged routes require an authorization callback. */
export function createExpressTokenRouter(
  service: TokenService,
  options: ExpressTokenAdapterOptions = {}
): Router {
  const router = Router();
  const run =
    (operation: (body: unknown) => Promise<unknown>) => async (req: Request, res: Response) => {
      try {
        res.json(await operation(req.body));
      } catch (error) {
        res.status(statusFor(error)).json(errorBody(error));
      }
    };
  const protectedRun =
    (action: ProtectedAction, operation: (body: never) => Promise<unknown>) =>
    async (req: Request, res: Response) => {
      try {
        if (!options.authorize || !(await options.authorize(action, req))) {
          res.status(403).json({ error: { code: "FORBIDDEN", message: "Forbidden" } });
          return;
        }
      } catch {
        res.status(500).json(errorBody(undefined));
        return;
      }
      await run((body) => operation(body as never))(req, res);
    };
  if (options.exposeIssue)
    router.post(
      "/issue",
      protectedRun("issue", (body) => service.issue(body))
    );
  if (options.exposeVerify)
    router.post(
      "/verify",
      run((body) => service.verify(requireStringField(body, "accessToken")))
    );
  router.post(
    "/refresh",
    run((body) => service.refresh(requireStringField(body, "refreshToken")))
  );
  if (options.exposeRevoke)
    router.post(
      "/revoke",
      protectedRun("revoke", async (body) => {
        await service.revoke(body);
        return { revoked: true };
      })
    );
  return router;
}

adapters/fastify.ts provides a Fastify plugin.

import type { FastifyInstance, FastifyPluginAsync, FastifyRequest } from "fastify";
import type { TokenService } from "../core/types.js";
import { errorBody, requireStringField, statusFor, type HttpAdapterOptions } from "./http.js";

export type FastifyTokenAdapterOptions = HttpAdapterOptions<FastifyRequest>;

/** Returns a Fastify plugin with the same security defaults as the Express adapter. */
export function createFastifyTokenPlugin(
  service: TokenService,
  options: FastifyTokenAdapterOptions = {}
): FastifyPluginAsync {
  return async (app: FastifyInstance) => {
    const execute = async (
      reply: {
        code(status: number): { send(value: unknown): unknown };
        send(value: unknown): unknown;
      },
      work: () => Promise<unknown>
    ) => {
      try {
        return reply.send(await work());
      } catch (error) {
        return reply.code(statusFor(error)).send(errorBody(error));
      }
    };
    if (options.exposeIssue)
      app.post("/issue", async (request, reply) => {
        try {
          if (!options.authorize || !(await options.authorize("issue", request)))
            return reply.code(403).send({ error: { code: "FORBIDDEN", message: "Forbidden" } });
        } catch {
          return reply.code(500).send(errorBody(undefined));
        }
        return execute(reply, () => service.issue(request.body as never));
      });
    if (options.exposeVerify)
      app.post("/verify", async (request, reply) =>
        execute(reply, () => service.verify(requireStringField(request.body, "accessToken")))
      );
    app.post("/refresh", async (request, reply) =>
      execute(reply, () => service.refresh(requireStringField(request.body, "refreshToken")))
    );
    if (options.exposeRevoke)
      app.post("/revoke", async (request, reply) => {
        try {
          if (!options.authorize || !(await options.authorize("revoke", request)))
            return reply.code(403).send({ error: { code: "FORBIDDEN", message: "Forbidden" } });
        } catch {
          return reply.code(500).send(errorBody(undefined));
        }
        return execute(reply, async () => {
          await service.revoke(request.body as never);
          return { revoked: true };
        });
      });
  };
}

adapters/hono.ts creates a Hono route group.

import { Hono, type Context } from "hono";
import type { TokenService } from "../core/types.js";
import { TokenError } from "../core/errors.js";
import { errorBody, requireStringField, statusFor, type HttpAdapterOptions } from "./http.js";

export type HonoTokenAdapterOptions = HttpAdapterOptions<Context>;

/** Creates a Hono route group. Mount it under an application-owned prefix. */
export function createHonoTokenRoutes(
  service: TokenService,
  options: HonoTokenAdapterOptions = {}
): Hono {
  const app = new Hono();
  const execute = async (c: Context, work: (body: unknown) => Promise<unknown>) => {
    try {
      return c.json(await work(await c.req.json()));
    } catch (error) {
      if (error instanceof SyntaxError)
        return c.json(errorBody(new TokenError("INVALID_INPUT", "Invalid JSON request body")), 400);
      return c.json(errorBody(error), statusFor(error) as 400);
    }
  };
  if (options.exposeIssue)
    app.post("/issue", async (c) => {
      try {
        if (!options.authorize || !(await options.authorize("issue", c)))
          return c.json({ error: { code: "FORBIDDEN", message: "Forbidden" } }, 403);
      } catch {
        return c.json(errorBody(undefined), 500);
      }
      return execute(c, (body) => service.issue(body as never));
    });
  if (options.exposeVerify)
    app.post("/verify", (c) =>
      execute(c, (body) => service.verify(requireStringField(body, "accessToken")))
    );
  app.post("/refresh", (c) =>
    execute(c, (body) => service.refresh(requireStringField(body, "refreshToken")))
  );
  if (options.exposeRevoke)
    app.post("/revoke", async (c) => {
      try {
        if (!options.authorize || !(await options.authorize("revoke", c)))
          return c.json({ error: { code: "FORBIDDEN", message: "Forbidden" } }, 403);
      } catch {
        return c.json(errorBody(undefined), 500);
      }
      return execute(c, async (body) => {
        await service.revoke(body as never);
        return { revoked: true };
      });
    });
  return app;
}

PostgreSQL schema

Apply the complete idempotent schema at sql/postgres.sql before using PostgresTokenStore. It upgrades older tables with absolute family expiry and rotation limits. Existing families receive a 90-day cap from their first recorded issue; families older than that cap must authenticate again.

CREATE TABLE IF NOT EXISTS blockend_refresh_tokens (
  token_hash text PRIMARY KEY,
  family_id uuid NOT NULL,
  subject text NOT NULL,
  expires_at timestamptz NOT NULL,
  family_expires_at timestamptz NOT NULL,
  rotation_count integer NOT NULL DEFAULT 0 CHECK (rotation_count >= 0),
  rotation_limit integer NOT NULL DEFAULT 10000 CHECK (rotation_limit > 0),
  status text NOT NULL CHECK (status IN ('active', 'used', 'revoked')),
  created_at timestamptz NOT NULL DEFAULT now(),
  used_at timestamptz,
  metadata jsonb
);
-- Upgrade pre-v2 installations. Old families receive the default 90-day absolute cap.
ALTER TABLE blockend_refresh_tokens ADD COLUMN IF NOT EXISTS family_expires_at timestamptz;
ALTER TABLE blockend_refresh_tokens ADD COLUMN IF NOT EXISTS rotation_count integer;
ALTER TABLE blockend_refresh_tokens ADD COLUMN IF NOT EXISTS rotation_limit integer;
WITH family_limits AS (
  SELECT family_id, MIN(created_at) + INTERVAL '90 days' AS family_expires_at,
    GREATEST(COUNT(*) - 1, 0)::integer AS rotation_count
  FROM blockend_refresh_tokens GROUP BY family_id
)
UPDATE blockend_refresh_tokens AS tokens
SET family_expires_at = COALESCE(tokens.family_expires_at, family_limits.family_expires_at),
    rotation_count = COALESCE(tokens.rotation_count, family_limits.rotation_count),
    rotation_limit = COALESCE(tokens.rotation_limit, 10000)
FROM family_limits
WHERE tokens.family_id = family_limits.family_id
  AND (tokens.family_expires_at IS NULL OR tokens.rotation_count IS NULL OR tokens.rotation_limit IS NULL);
ALTER TABLE blockend_refresh_tokens ALTER COLUMN family_expires_at SET NOT NULL;
ALTER TABLE blockend_refresh_tokens ALTER COLUMN rotation_count SET DEFAULT 0;
ALTER TABLE blockend_refresh_tokens ALTER COLUMN rotation_count SET NOT NULL;
ALTER TABLE blockend_refresh_tokens ALTER COLUMN rotation_limit SET DEFAULT 10000;
ALTER TABLE blockend_refresh_tokens ALTER COLUMN rotation_limit SET NOT NULL;
CREATE INDEX IF NOT EXISTS blockend_refresh_family_idx ON blockend_refresh_tokens (family_id);
CREATE INDEX IF NOT EXISTS blockend_refresh_subject_idx ON blockend_refresh_tokens (subject);
CREATE INDEX IF NOT EXISTS blockend_refresh_expiry_idx ON blockend_refresh_tokens (expires_at);
CREATE INDEX IF NOT EXISTS blockend_refresh_family_expiry_idx ON blockend_refresh_tokens (family_expires_at);

CREATE TABLE IF NOT EXISTS blockend_revoked_jtis (
  jti text PRIMARY KEY,
  expires_at timestamptz NOT NULL
);
CREATE INDEX IF NOT EXISTS blockend_revoked_jti_expiry_idx ON blockend_revoked_jtis (expires_at);

Configuration

TokenServiceOptions

OptionTypeDefaultDescription
deploymentMode"development" | "production""production"Development permits localhost HTTP and an omitted audience.
issuerstringRequiredIssuer claim. Production requires an HTTPS URL or URN.
audiencestring | string[]Required in productionAllowed audience values. A token requested with another audience is rejected.
accessTokenTtlSecondsnumber900Clamped to 60 through 3600 seconds.
refreshTokenTtlSecondsnumber604800Clamped to 3600 through 2592000 seconds.
maxRefreshFamilyLifetimeSecondsnumber7776000Absolute family lifetime, from 3600 through 31536000 seconds.
maxRefreshRotationsnumber10000Maximum successful rotations per family, from 1 through 1000000.
clockSkewSecondsnumber30JWT clock tolerance, capped at 300 seconds.
algorithmsSigningAlgorithm[]["RS256", "ES256"]Accepted signing algorithms. Production rejects all HS* algorithms.
keyProviderKeyProviderRequiredSupplies signing and verification keys.
tokenStoreTokenStoreRequiredStores refresh token state and may store revoked JTIs.
maxClaimsBytesnumber4096Maximum serialized custom claims size, clamped to 256 through 16384 bytes.
forbiddenClaimKeysstring[]Built in secret and reserved namesAdditional names rejected recursively in custom claims. Supplying this replaces the built in list.
onEvent(event: TokenEvent) => void—Receives token lifecycle events. Synchronous telemetry errors are ignored.
resolveRefreshContext(sub: string) => Promise<{ claims?: Record<string, unknown>; audience?: string | string[] }>—Reloads current authorization claims during refresh.
now() => Date() => new Date()Clock injection for tests and controlled environments.

issue() input

FieldTypeDefaultDescription
substringRequiredSubject, trimmed, 1 to 512 characters.
claimsRecord<string, unknown>{}Custom JSON claims. Reserved and configured forbidden names are rejected.
tokens.accessbooleantrueWhether to issue an access token.
tokens.refreshbooleantrueWhether to issue a refresh token. At least one token must be requested.
accessTokenTtlSecondsnumberService settingPer-token access lifetime, clamped to 60 through 3600 seconds.
refreshTokenTtlSecondsnumberService settingPer-token refresh lifetime, clamped to 3600 through 2592000 seconds.
audiencestring | string[]Service settingMust be contained in the configured allowed audience.

Adapter options

Each adapter accepts the following options. refresh is always mounted. Other endpoints default to disabled.

OptionTypeDefaultDescription
authorize(action: "issue" | "revoke", request) => boolean | Promise<boolean>—Required to allow exposed issue or revoke routes. Return false to deny.
exposeIssuebooleanfalseMount POST /issue.
exposeVerifybooleanfalseMount POST /verify.
exposeRevokebooleanfalseMount POST /revoke.

Architecture

createTokenService(options)
  └─ validate configuration once

issue(input)
  ├─ validate subject, audience, and custom claims
  ├─ sign a short lived JWT using the active kid
  └─ store only the hash of a random refresh token

verify(accessToken)
  ├─ pin alg and kid to configured verification keys
  ├─ validate iss, aud, exp, iat, jti, and type
  └─ check the optional JTI denylist

refresh(refreshToken)
  ├─ resolve the versioned token's family identifier
  ├─ reload claims and sign before the atomic store operation
  ├─ consume the old hash and store a replacement hash within lifetime and rotation caps
  ├─ revoke the family if an already used token appears again
  ├─ reload current claims when resolveRefreshContext is configured
  └─ sign and return the new access token and opaque refresh token

createTokenService() validates configuration at startup. Every operation validates its input and emits a sanitized TokenEvent. PostgreSQL explicitly uses READ COMMITTED, locks the subject, and locks the current refresh row with SELECT ... FOR UPDATE. Redis places a subject's generation counter and family hashes under one HMAC-derived cluster hash tag. Its Lua scripts receive every key explicitly and do O(1) rotation, reuse detection, and family invalidation. Keep the Redis routing key stable and identical across issuer instances. A Redis subject maps to one slot, while separate subjects distribute across the cluster. Do not read refresh or revocation state from replicas.

Production setup

Use shared state for refresh tokens

Apply sql/postgres.sql, then provide a PgLike wrapper whose transaction() method runs all callback queries in the same database transaction. PostgresTokenStore does not open or manage a pool itself. A normal pg pool needs a transaction wrapper that checks out one client, runs BEGIN, commits on success, rolls back on error, and releases the client.

For Redis, use a client adapter that implements get, hget, set(key, value, { PX }), and eval(script, { keys, arguments }). The default prefix is blockend:v2:. Configure RedisTokenStore with a stable base64url-encoded subject routing key containing at least 32 decoded bytes and an optional prefix without braces. The key must match on every instance and must not be rotated until all families created with it expire or are revoked; otherwise revokeBySub() cannot reach older sessions.

const redisStore = new RedisTokenStore(redisAdapter, {
  prefix: "company-auth:v2:",
  subjectRoutingKey: Buffer.from(process.env.TOKEN_REDIS_SUBJECT_KEY!, "base64url")
});

PostgreSQL transactions set READ COMMITTED before taking advisory locks. Its transaction() callback must use one connection and run no statements before the store's callback statements. Schedule bounded cleanup of rows where family_expires_at <= NOW() and expired JTI rows.

Rotate signing keys

Keep old public keys in the provider while they can still verify unexpired access tokens. Switch signingKid to the new key, wait at least the maximum access lifetime plus clock tolerance, then remove the old key. For a remote JWKS provider, publish the new public key before signing with it.

Protect browser refresh flows

Store browser refresh tokens in Secure, HttpOnly, SameSite cookies. Add CSRF protection to cookie authenticated refresh and revoke requests, and restrict CORS. Do not put refresh tokens in URLs or local storage.

Keep authorization current

The service stores no roles or permissions in refresh records. Configure resolveRefreshContext to reload current authorization data on refresh. If an account has been disabled or deleted, throw from that callback and handle the failed refresh in your application.

Retain and monitor state

Delete expired and used refresh rows after your audit window, and remove expired JTI records on a schedule. Alert on token.reuse_detected, store and key failures, and unusual refresh volume. Events include subject, JTI, family ID, and error code when available, never raw tokens.

Usage

Service setup with local keys

Use a secret manager to supply the PEM values. Local PEM keys are useful when your runtime can protect the private key; use a KMS or HSM provider when private key material must remain outside the process.

import {
  createTokenService,
  LocalKeyProvider,
  PostgresTokenStore
} from "@/blocks/token-service/src";
import { database } from "./database";

export const tokens = createTokenService({
  issuer: "https://auth.example.com",
  audience: ["https://api.example.com"],
  algorithms: ["RS256"],
  accessTokenTtlSeconds: 600,
  keyProvider: new LocalKeyProvider({
    signingKid: "rsa-2026-09",
    keys: [
      {
        kid: "rsa-2026-09",
        alg: "RS256",
        privateKeyPem: process.env.JWT_PRIVATE_KEY_PEM!,
        publicKeyPem: process.env.JWT_PUBLIC_KEY_PEM!
      }
    ]
  }),
  tokenStore: new PostgresTokenStore(database),
  resolveRefreshContext: async (sub) => {
    const user = await database.users.findAuthorizationContext(sub);
    if (!user || user.disabled) throw new Error("User is unavailable");
    return { claims: { role: user.role, permissions: user.permissions } };
  }
});

The database value must implement PgLike, including a transaction wrapper around the same connection for every rotation query.

Express

import express from "express";
import { createExpressTokenRouter } from "@/blocks/token-service/src/adapters/express";
import { tokens } from "./tokens";

const app = express();
app.use(express.json({ limit: "8kb" }));
app.use(
  "/auth/tokens",
  createExpressTokenRouter(tokens, {
    exposeIssue: true,
    exposeRevoke: true,
    authorize: (_action, request) => request.auth?.service === "identity-api"
  })
);

Install your authentication middleware before the router when exposing issue or revoke.

Fastify

import Fastify from "fastify";
import { createFastifyTokenPlugin } from "@/blocks/token-service/src/adapters/fastify";
import { tokens } from "./tokens";

const app = Fastify({ bodyLimit: 8 * 1024 });
await app.register(
  createFastifyTokenPlugin(tokens, {
    exposeIssue: true,
    authorize: (_action, request) => request.user?.service === "identity-api"
  }),
  { prefix: "/auth/tokens" }
);

Fastify's body limit applies before these routes run.

Hono

import { Hono } from "hono";
import { createHonoTokenRoutes } from "@/blocks/token-service/src/adapters/hono";
import { tokens } from "./tokens";

const app = new Hono();
app.route(
  "/auth/tokens",
  createHonoTokenRoutes(tokens, {
    exposeRevoke: true,
    authorize: (_action, context) => context.get("service") === "identity-api"
  })
);

Mount the route group behind your application's authentication and request size middleware.

Access-token-only service

import {
  createTokenService,
  LocalKeyProvider,
  StatelessTokenStore
} from "@/blocks/token-service/src";

const tokens = createTokenService({
  issuer: "https://internal.example.com",
  audience: "https://worker.example.com",
  keyProvider,
  tokenStore: new StatelessTokenStore()
});

const access = await tokens.issue({ sub: "worker-17", tokens: { refresh: false } });

Do not call refresh() or server side revoke operations with StatelessTokenStore; they return UNSUPPORTED_OPERATION.

API reference

createTokenService

function createTokenService(options: TokenServiceOptions): TokenService;

Creates an isolated service and validates its configuration immediately. Invalid configuration throws TokenError with code INVALID_INPUT.

TokenService

MethodInputResult
issue(input)Subject, optional claims, token selection, lifetime, and audiencePromise<TokenPair>
verify(accessToken)Encoded access tokenPromise<VerifyResult>
refresh(refreshToken)Versioned opaque refresh tokenPromise<TokenPair>
revoke(input)Exactly one refresh token, JTI, subject, or family IDPromise<void>

verify() accepts access tokens only. Refresh tokens are versioned opaque bearer credentials and cannot be used as access tokens. Their embedded family-routing hint is not secret; the 256-bit random component is the credential.

Token result types

interface TokenPair {
  accessToken?: string;
  refreshToken?: string;
  tokenType: "Bearer";
  expiresIn?: number;
  refreshExpiresIn?: number;
  issuedAt: number;
}

interface AccessTokenClaims {
  sub: string;
  iss: string;
  aud?: string | string[];
  exp: number;
  iat: number;
  jti: string;
  type: "access";
  [key: string]: unknown;
}

interface VerifyResult {
  valid: true;
  claims: AccessTokenClaims;
}

type RevokeInput =
  | { refreshToken: string }
  | { jti: string }
  | { sub: string }
  | { familyId: string };

Key provider types

SigningAlgorithm is one of RS256, RS384, RS512, PS256, PS384, PS512, ES256, ES384, ES512, EdDSA, HS256, HS384, or HS512. Production configuration rejects the symmetric HS* algorithms.

KeyProvider can expose getSigningKey(), signJwt(), and getVerificationKey({ kid, alg }). A verification only provider may expose just getVerificationKey(). getPublicJwks() is optional.

LocalKey contains kid, alg, privateKeyPem, and publicKeyPem. KmsRsaSigner implements sign(data: Uint8Array): Promise<Uint8Array> and returns an RSASSA-PKCS1-v1_5 SHA-256 signature.

Store types

TokenStore provides createFamilyId?(sub), getRefreshToken(tokenHash, familyId?), and atomic rotateRefreshToken(). Rotation checks the family deadline and rotation ceiling, consumes the old hash, and persists its replacement. Reuse revokes the complete family. Revoking by refresh token also revokes its complete family. The store can implement revokeByJti() and isJtiRevoked() for access-token revocation. MemoryTokenStore, RedisTokenStore, PostgresTokenStore, and StatelessTokenStore are included.

RedisTokenStoreOptions requires subjectRoutingKey as a Uint8Array or base64url string with at least 32 decoded bytes. All issuer instances must share the same key and prefix. RedisLike requires get(key), hget(key, field), set(key, value, options), and eval(script, { keys, arguments }). PgLike requires query(sql, values) and transaction(work). The transaction callback receives a PgTransaction with query(); the store sets READ COMMITTED before issuing its locking statements.

Events and errors

Events use the names token.issued, token.verified, token.refreshed, token.revoked, token.reuse_detected, and token.failed. They carry an ISO timestamp and optional identifiers or error code. A synchronous exception from onEvent does not change token behavior.

Error codeMeaning
INVALID_INPUTConfiguration or operation input failed validation.
INVALID_TOKENToken is malformed, invalid, or does not match the configured issuer, audience, or algorithm.
TOKEN_EXPIREDAccess or refresh token has expired.
TOKEN_REVOKEDToken state is revoked.
REFRESH_REUSEDA used refresh token appeared again and its family was revoked.
REFRESH_LIMITThe refresh family reached its configured rotation ceiling.
CLAIMS_TOO_LARGECustom claims exceed the configured byte limit.
FORBIDDEN_CLAIMCustom claims contain reserved or forbidden data.
KEY_UNAVAILABLEA signing or verification key could not be used.
STORE_UNAVAILABLEA token store operation failed.
UNSUPPORTED_OPERATIONThe selected provider or store does not support the operation.

HTTP routes

RouteDefaultAuthorizationSuccess
POST /refreshEnabledApplication policy and rate limitingRotated token pair
POST /issueDisabledRequired callback when enabledToken pair
POST /verifyDisabledApplication policy when enabled{ valid: true, claims }
POST /revokeDisabledRequired callback when enabled{ revoked: true }

Errors use { error: { code, message } }. Invalid input maps to 400, store/key failures to 503, unsupported operations to 501, and other token errors to 401. A failed protected action returns 403.

The adapter does not add authentication, CSRF, CORS, or rate limiting. Set body size limits in your framework before mounting these routes. Exposing /issue directly to the public internet lets callers mint tokens, so authorize it with your own trusted identity or service credential.

When to use

  • You need one signing and verification policy shared by several services.
  • You need refresh rotation with reuse detection and family revocation.
  • You want to keep key management and refresh token persistence behind replaceable interfaces.
  • You need Express, Fastify, or Hono routes for token operations.

When not to use

  • You need a complete identity provider, login flow, OAuth/OIDC server, account recovery, or authorization policy. This block does not implement those systems.
  • You need refresh state to survive restarts or coordinate across instances but can only use MemoryTokenStore.
  • You need multi region refresh consistency without a shared strongly consistent store. Eventual replicas can accept replayed refresh state.

Testing

The block tests signing and verification, audience checks, access and family expiry, forbidden claims, rotation limits, replay revocation, concurrency, and each framework adapter. Store tests verify transaction query construction and that Redis Lua scripts use bounded state and explicit keys. Live standalone Redis, six-node Redis Cluster, and PostgreSQL integration tests cover concurrent reuse, subject/family revocation, family limits, and JTI revocation.

Set REDIS_URL for standalone Redis, REDIS_CLUSTER_NODES to comma-separated host:port cluster nodes, and TOKEN_SERVICE_POSTGRES_URL to a dedicated test database. The PostgreSQL suite creates and drops a unique schema, so its database role needs schema creation permission. Tests skip a backend when its environment variable is absent.

Install pg and @types/pg as development dependencies to run this integration test outside the Blockend workspace. They are test-only; the store itself accepts the PgLike interface and does not depend on a PostgreSQL driver.

pnpm exec vitest run blocks/token-service
pnpm typecheck:blocks

CI starts PostgreSQL 17, standalone Redis 7, and a six-node Redis 7 Cluster for the full test suite.

  • Password Hashing — Hash user credentials before an application authenticates a user and asks this service to issue tokens.
  • Rate Limiter — Limit refresh and privileged issue requests before they reach token storage or signing infrastructure.
  • Environment Configuration — Validate issuer, audience, and key reference settings before the service starts.

FAQ

Does this block authenticate a user?

No. Call issue() only after your application has authenticated the subject. The block does not look up users or decide permissions.

Does it store refresh tokens in plain text?

No. It returns a versioned token with a 256-bit random secret and stores only its SHA-256 hash. The public routing hint contains no subject value. Treat the whole returned value as a credential and send it only over TLS.

What happens if a refresh token is reused?

The store marks the token family revoked and the service throws REFRESH_REUSED. The client must authenticate again through your application. Already issued access tokens remain valid until expiry unless their JTIs are also revoked.

Can I use MemoryTokenStore in production?

Only for a single process when losing all refresh and revocation state on restart is acceptable. It does not coordinate across processes or hosts.

Can I expose /issue or /revoke?

Yes, but the adapter requires an authorize callback to return true. Keep the route behind your application's identity checks and rate limits.

How do I revoke access tokens immediately?

Use a store with revokeByJti() and isJtiRevoked(), then call revoke({ jti }). This makes each verification consult shared state. Otherwise use short access lifetimes and revoke refresh tokens to stop future access token issuance.

Does a JWKS provider sign tokens?

RemoteJwksKeyProvider is verification only. Configure a local or KMS signer in the issuing service and use the remote JWKS provider in resource services.

On this page