#![no_std]
#![cfg_attr(test, allow(clippy::indexing_slicing))]
extern crate alloc;
use alloc::collections::BTreeMap;
use alloc::format;
use alloc::string::String;
use alloc::vec::Vec;
use base64::Engine;
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
use crypto_secretbox::aead::{Aead, KeyInit};
use crypto_secretbox::{Key as SecretboxKey, Nonce as SecretboxNonce, XSalsa20Poly1305};
use ed25519_dalek::{Signature, Verifier, VerifyingKey};
use rand_core::TryRng;
use salsa20::cipher::consts::{U10, U16};
use salsa20::hsalsa;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use sha2::{Digest, Sha512_256};
use x25519_dalek::{PublicKey as X25519PublicKey, StaticSecret};
use zeroize::Zeroizing;
const HEADER_JSON: &str = r#"{"typ":"JWT","alg":"ed25519-nkey"}"#;
const XKEY_VERSION_V1: &[u8; 4] = b"xkv1";
const XKEY_NONCE_LEN: usize = 24;
const XKEY_TAG_LEN: usize = 16;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum NkeyType {
Account,
Cluster,
Server,
Operator,
User,
Curve,
}
impl NkeyType {
const ALL: [Self; 6] = [
Self::Account,
Self::Cluster,
Self::Server,
Self::Operator,
Self::User,
Self::Curve,
];
#[must_use]
pub const fn letter(self) -> char {
match self {
Self::Account => 'A',
Self::Cluster => 'C',
Self::Server => 'N',
Self::Operator => 'O',
Self::User => 'U',
Self::Curve => 'X',
}
}
#[must_use]
pub const fn prefix_byte(self) -> u8 {
match self {
Self::Account => 0, Self::Cluster => 2 << 3, Self::Server => 13 << 3, Self::Operator => 14 << 3, Self::User => 20 << 3, Self::Curve => 23 << 3, }
}
#[must_use]
pub fn from_letter(letter: char) -> Option<Self> {
Self::ALL.into_iter().find(|role| role.letter() == letter)
}
#[must_use]
fn from_prefix_byte(byte: u8) -> Option<Self> {
Self::ALL
.into_iter()
.find(|role| role.prefix_byte() == byte)
}
}
impl core::fmt::Display for NkeyType {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(f, "{}", self.letter())
}
}
#[derive(Debug)]
pub enum Error {
BadPublicKeyLen(usize),
UnsupportedPrefix(char),
UnexpectedPrefix {
expected: NkeyType,
actual: NkeyType,
},
Json(serde_json::Error),
InvalidClaims(String),
MalformedJwt(String),
BadSignatureLen(usize),
UnexpectedXKeyPrefix(NkeyType),
BadXKeyVersion,
BadXKeyCiphertextLen(usize),
XKeyRandomnessFailed,
XKeySealFailed,
XKeyOpenFailed,
}
impl core::fmt::Display for Error {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::BadPublicKeyLen(n) => write!(f, "expected a 32-byte Ed25519 public key, got {n}"),
Self::UnsupportedPrefix(c) => write!(f, "unsupported NKey prefix letter: {c}"),
Self::UnexpectedPrefix { expected, actual } => {
write!(f, "expected NKey role {expected}, got {actual}")
}
Self::Json(e) => write!(f, "json error: {e}"),
Self::InvalidClaims(message) | Self::MalformedJwt(message) => f.write_str(message),
Self::BadSignatureLen(n) => write!(f, "expected a 64-byte Ed25519 signature, got {n}"),
Self::UnexpectedXKeyPrefix(actual) => {
write!(f, "expected curve xkey prefix X, got {actual}")
}
Self::BadXKeyVersion => f.write_str("NATS xkey ciphertext has unsupported version"),
Self::BadXKeyCiphertextLen(n) => {
write!(f, "NATS xkey ciphertext is too short: {n} bytes")
}
Self::XKeyRandomnessFailed => f.write_str("NATS xkey randomness failed"),
Self::XKeySealFailed => f.write_str("NATS xkey seal failed"),
Self::XKeyOpenFailed => f.write_str("NATS xkey open failed"),
}
}
}
impl core::error::Error for Error {}
impl From<serde_json::Error> for Error {
fn from(e: serde_json::Error) -> Self {
Self::Json(e)
}
}
fn crc16(data: &[u8]) -> u16 {
let mut crc: u16 = 0;
for &b in data {
crc ^= u16::from(b) << 8;
for _ in 0..8 {
crc = if crc & 0x8000 != 0 {
(crc << 1) ^ 0x1021
} else {
crc << 1
};
}
}
crc
}
fn encode_nkey(prefix: u8, public: &[u8]) -> Result<String, Error> {
if public.len() != 32 {
return Err(Error::BadPublicKeyLen(public.len()));
}
let mut raw = Vec::with_capacity(1 + 32 + 2);
raw.push(prefix);
raw.extend_from_slice(public);
let crc = crc16(&raw);
raw.push((crc & 0xff) as u8); raw.push((crc >> 8) as u8);
Ok(data_encoding::BASE32_NOPAD.encode(&raw))
}
pub fn encode_public(role: NkeyType, public: &[u8]) -> Result<String, Error> {
encode_nkey(role.prefix_byte(), public)
}
pub fn decode_public(nkey: &str) -> Result<(NkeyType, [u8; 32]), Error> {
let raw = data_encoding::BASE32_NOPAD
.decode(nkey.as_bytes())
.map_err(|_| Error::BadPublicKeyLen(0))?;
let Some((body, crc_bytes)) = raw.split_at_checked(33) else {
return Err(Error::BadPublicKeyLen(raw.len()));
};
let (Ok(crc_bytes), Some((&prefix, key))) =
(<[u8; 2]>::try_from(crc_bytes), body.split_first())
else {
return Err(Error::BadPublicKeyLen(raw.len()));
};
let Ok(key) = <[u8; 32]>::try_from(key) else {
return Err(Error::BadPublicKeyLen(raw.len()));
};
let crc = u16::from_le_bytes(crc_bytes);
if crc != crc16(body) {
return Err(Error::BadPublicKeyLen(raw.len()));
}
let role = NkeyType::from_prefix_byte(prefix)
.ok_or_else(|| Error::UnsupportedPrefix(nkey.chars().next().unwrap_or('?')))?;
Ok((role, key))
}
pub fn require_public_prefix(nkey: &str, expected: NkeyType) -> Result<[u8; 32], Error> {
let (actual, key) = decode_public(nkey)?;
if actual != expected {
return Err(Error::UnexpectedPrefix { expected, actual });
}
Ok(key)
}
pub fn verify_public_signature(
nkey: &str,
message: &[u8],
signature: &[u8],
) -> Result<bool, Error> {
let (_, public_key) = decode_public(nkey)?;
let signature =
<[u8; 64]>::try_from(signature).map_err(|_| Error::BadSignatureLen(signature.len()))?;
Ok(verify_ed25519(&public_key, message, &signature))
}
#[must_use]
pub fn xkey_public_from_private(private: &Zeroizing<[u8; 32]>) -> [u8; 32] {
let secret = StaticSecret::from(**private);
X25519PublicKey::from(&secret).to_bytes()
}
pub fn seal_nats_curve(
sender_private: &Zeroizing<[u8; 32]>,
recipient_public_xkey: &str,
plaintext: &[u8],
rng: &mut impl TryRng,
) -> Result<Vec<u8>, Error> {
let recipient_public = decode_xkey_public(recipient_public_xkey)?;
let mut nonce = [0u8; XKEY_NONCE_LEN];
rng.try_fill_bytes(&mut nonce)
.map_err(|_| Error::XKeyRandomnessFailed)?;
let ciphertext = box_crypt(
sender_private,
&recipient_public,
&nonce,
plaintext,
BoxMode::Seal,
)?;
let mut out = Vec::with_capacity(XKEY_VERSION_V1.len() + nonce.len() + ciphertext.len());
out.extend_from_slice(XKEY_VERSION_V1);
out.extend_from_slice(&nonce);
out.extend_from_slice(&ciphertext);
Ok(out)
}
pub fn open_nats_curve(
recipient_private: &Zeroizing<[u8; 32]>,
sender_public_xkey: &str,
ciphertext: &[u8],
) -> Result<Zeroizing<Vec<u8>>, Error> {
let sender_public = decode_xkey_public(sender_public_xkey)?;
let Some(rest) = ciphertext.strip_prefix(XKEY_VERSION_V1) else {
return Err(Error::BadXKeyVersion);
};
if rest.len() <= XKEY_NONCE_LEN + XKEY_TAG_LEN {
return Err(Error::BadXKeyCiphertextLen(ciphertext.len()));
}
let Some((nonce, body)) = rest.split_at_checked(XKEY_NONCE_LEN) else {
return Err(Error::BadXKeyCiphertextLen(ciphertext.len()));
};
let nonce: [u8; XKEY_NONCE_LEN] = nonce
.try_into()
.map_err(|_| Error::BadXKeyCiphertextLen(ciphertext.len()))?;
box_crypt(
recipient_private,
&sender_public,
&nonce,
body,
BoxMode::Open,
)
.map(Zeroizing::new)
}
fn decode_xkey_public(nkey: &str) -> Result<[u8; 32], Error> {
let (actual, key) = decode_public(nkey)?;
if actual != NkeyType::Curve {
return Err(Error::UnexpectedXKeyPrefix(actual));
}
Ok(key)
}
#[derive(Clone, Copy)]
enum BoxMode {
Seal,
Open,
}
fn box_crypt(
private: &Zeroizing<[u8; 32]>,
peer_public: &[u8; 32],
nonce: &[u8; XKEY_NONCE_LEN],
input: &[u8],
mode: BoxMode,
) -> Result<Vec<u8>, Error> {
let private = StaticSecret::from(**private);
let peer = X25519PublicKey::from(*peer_public);
let shared_secret = private.diffie_hellman(&peer);
if !shared_secret.was_contributory() {
return match mode {
BoxMode::Seal => Err(Error::XKeySealFailed),
BoxMode::Open => Err(Error::XKeyOpenFailed),
};
}
let shared = Zeroizing::new(shared_secret.to_bytes());
let key = nats_box_key(&shared);
let cipher = XSalsa20Poly1305::new(&key);
let nonce = SecretboxNonce::from(*nonce);
match mode {
BoxMode::Seal => cipher
.encrypt(&nonce, input)
.map_err(|_| Error::XKeySealFailed),
BoxMode::Open => cipher
.decrypt(&nonce, input)
.map_err(|_| Error::XKeyOpenFailed),
}
}
#[allow(deprecated)]
fn nats_box_key(shared: &Zeroizing<[u8; 32]>) -> Zeroizing<SecretboxKey> {
let input = crypto_secretbox::aead::generic_array::GenericArray::<u8, U16>::default();
let key = SecretboxKey::clone_from_slice(shared.as_slice());
Zeroizing::new(hsalsa::<U10>(&key, &input))
}
#[derive(Deserialize)]
struct CompactHeader {
typ: String,
alg: String,
}
#[derive(Deserialize)]
struct CompactClaims {
iss: String,
sub: String,
#[serde(default)]
iat: Option<u64>,
#[serde(default)]
exp: Option<u64>,
#[serde(default)]
nbf: Option<u64>,
#[serde(default)]
nats: Option<CompactNatsClaims>,
}
#[derive(Deserialize)]
struct CompactNatsClaims {
#[serde(rename = "type")]
kind: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NatsJwtClaims {
pub issuer: String,
pub subject: String,
pub issued_at: Option<u64>,
pub expires_at: Option<u64>,
pub not_before: Option<u64>,
pub nats_type: Option<String>,
pub raw: Value,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DecodedNatsJwt {
signing_input: String,
signature: [u8; 64],
claims: NatsJwtClaims,
issuer_role: NkeyType,
issuer_public_key: [u8; 32],
}
impl DecodedNatsJwt {
#[must_use]
pub fn signing_input(&self) -> &str {
&self.signing_input
}
#[must_use]
pub const fn signature(&self) -> &[u8; 64] {
&self.signature
}
#[must_use]
pub const fn claims(&self) -> &NatsJwtClaims {
&self.claims
}
#[must_use]
pub const fn issuer_role(&self) -> NkeyType {
self.issuer_role
}
#[must_use]
pub const fn issuer_public_key(&self) -> &[u8; 32] {
&self.issuer_public_key
}
#[must_use]
pub fn verify_signature(&self, public_key: &[u8; 32]) -> bool {
verify_ed25519(public_key, self.signing_input.as_bytes(), &self.signature)
}
pub fn verify_with_candidates<'a, I>(
&self,
candidates: I,
now_unix: u64,
) -> Result<NatsJwtValidation, Error>
where
I: IntoIterator<Item = CandidateSigner<'a>>,
{
for candidate in candidates {
let resolved = self.resolve_candidate(candidate)?;
if let Some(matched) = resolved {
return Ok(self.validation_for_match(matched, now_unix));
}
}
Ok(NatsJwtValidation {
reason: NatsJwtValidationReason::UnknownSigner,
matched_signer: None,
})
}
fn resolve_candidate(
&self,
candidate: CandidateSigner<'_>,
) -> Result<Option<MatchedNatsSigner>, Error> {
match candidate {
CandidateSigner::Nkey(nkey) => {
let (role, public_key) = decode_public(nkey)?;
if role == self.issuer_role && public_key == self.issuer_public_key {
Ok(Some(MatchedNatsSigner {
role: Some(role),
public_key,
}))
} else {
Ok(None)
}
}
CandidateSigner::RawPublicKey(public_key) => {
if public_key == &self.issuer_public_key {
Ok(Some(MatchedNatsSigner {
role: None,
public_key: *public_key,
}))
} else {
Ok(None)
}
}
}
}
fn validation_for_match(
&self,
matched_signer: MatchedNatsSigner,
now_unix: u64,
) -> NatsJwtValidation {
if !self.verify_signature(&matched_signer.public_key) {
return NatsJwtValidation {
reason: NatsJwtValidationReason::BadSignature,
matched_signer: None,
};
}
if let Some(exp) = self.claims.expires_at
&& exp <= now_unix
{
return NatsJwtValidation {
reason: NatsJwtValidationReason::Expired,
matched_signer: None,
};
}
if let Some(nbf) = self.claims.not_before
&& nbf > now_unix
{
return NatsJwtValidation {
reason: NatsJwtValidationReason::NotYetValid,
matched_signer: None,
};
}
NatsJwtValidation {
reason: NatsJwtValidationReason::Valid,
matched_signer: Some(matched_signer),
}
}
}
#[derive(Debug, Clone, Copy)]
pub enum CandidateSigner<'a> {
Nkey(&'a str),
RawPublicKey(&'a [u8; 32]),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MatchedNatsSigner {
pub role: Option<NkeyType>,
pub public_key: [u8; 32],
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct NatsJwtValidation {
pub reason: NatsJwtValidationReason,
pub matched_signer: Option<MatchedNatsSigner>,
}
impl NatsJwtValidation {
#[must_use]
pub const fn is_valid(self) -> bool {
matches!(self.reason, NatsJwtValidationReason::Valid)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NatsJwtValidationReason {
Valid,
UnknownSigner,
BadSignature,
Expired,
NotYetValid,
}
pub fn decode_nats_jwt(token: &str) -> Result<DecodedNatsJwt, Error> {
let mut parts = token.split('.');
let Some(header_segment) = parts.next() else {
return Err(Error::MalformedJwt("jwt is empty".into()));
};
let Some(claims_segment) = parts.next() else {
return Err(Error::MalformedJwt(
"jwt must contain header, claims, and signature".into(),
));
};
let Some(signature_segment) = parts.next() else {
return Err(Error::MalformedJwt(
"jwt must contain header, claims, and signature".into(),
));
};
if parts.next().is_some() {
return Err(Error::MalformedJwt(
"jwt must contain exactly three segments".into(),
));
}
let header: CompactHeader = decode_url_json(header_segment, "header")?;
if header.typ != "JWT" || header.alg != "ed25519-nkey" {
return Err(Error::MalformedJwt(
"jwt header must be typ=JWT and alg=ed25519-nkey".into(),
));
}
let claims_raw = decode_url_bytes(claims_segment, "claims")?;
let claims_value: Value = serde_json::from_slice(&claims_raw)
.map_err(|e| Error::MalformedJwt(format!("invalid jwt claims json: {e}")))?;
let claims: CompactClaims = serde_json::from_value(claims_value.clone())
.map_err(|e| Error::MalformedJwt(format!("invalid jwt claims: {e}")))?;
let (issuer_role, issuer_public_key) = decode_public(&claims.iss)
.map_err(|e| Error::MalformedJwt(format!("invalid jwt issuer nkey: {e}")))?;
let signature = decode_url_bytes(signature_segment, "signature")?;
let signature_len = signature.len();
let Ok(signature) = <[u8; 64]>::try_from(signature) else {
return Err(Error::BadSignatureLen(signature_len));
};
Ok(DecodedNatsJwt {
signing_input: format!("{header_segment}.{claims_segment}"),
signature,
claims: NatsJwtClaims {
issuer: claims.iss,
subject: claims.sub,
issued_at: claims.iat,
expires_at: claims.exp,
not_before: claims.nbf,
nats_type: claims.nats.and_then(|nats| nats.kind),
raw: claims_value,
},
issuer_role,
issuer_public_key,
})
}
fn decode_url_json<T>(segment: &str, name: &str) -> Result<T, Error>
where
T: DeserializeOwned,
{
let bytes = decode_url_bytes(segment, name)?;
serde_json::from_slice(&bytes)
.map_err(|e| Error::MalformedJwt(format!("invalid jwt {name} json: {e}")))
}
fn decode_url_bytes(segment: &str, name: &str) -> Result<Vec<u8>, Error> {
URL_SAFE_NO_PAD
.decode(segment)
.map_err(|e| Error::MalformedJwt(format!("invalid jwt {name} base64url: {e}")))
}
fn verify_ed25519(public_key: &[u8; 32], message: &[u8], signature: &[u8; 64]) -> bool {
let signature = Signature::from_bytes(signature);
VerifyingKey::from_bytes(public_key).is_ok_and(|key| key.verify(message, &signature).is_ok())
}
#[derive(Serialize, Default, Clone)]
pub struct Permission {
#[serde(skip_serializing_if = "Vec::is_empty")]
pub allow: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub deny: Vec<String>,
}
#[derive(Serialize)]
struct NatsUser {
#[serde(rename = "pub")]
publish: Permission,
sub: Permission,
subs: i64,
data: i64,
payload: i64,
#[serde(skip_serializing_if = "Option::is_none")]
issuer_account: Option<String>,
#[serde(rename = "type")]
kind: &'static str,
version: i64,
}
#[derive(Serialize)]
struct ClaimsHash<'a> {
#[serde(skip_serializing_if = "Option::is_none")]
aud: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
exp: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
jti: Option<&'a str>,
iat: u64,
iss: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
name: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
nbf: Option<u64>,
sub: &'a str,
}
#[derive(Serialize, Default, Clone)]
pub struct Permissions {
#[serde(rename = "pub")]
pub publish: Permission,
pub sub: Permission,
#[serde(skip_serializing_if = "Option::is_none")]
pub resp: Option<ResponsePermission>,
}
#[derive(Serialize, Clone)]
pub struct ResponsePermission {
pub max: i64,
pub ttl: i64,
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "lowercase")]
pub enum ExportType {
Stream,
Service,
}
#[derive(Serialize, Clone)]
pub struct AccountImport {
#[serde(skip_serializing_if = "str::is_empty")]
pub name: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub subject: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub account: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub token: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub to: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub local_subject: String,
#[serde(rename = "type")]
pub kind: ExportType,
#[serde(skip_serializing_if = "is_false")]
pub share: bool,
#[serde(skip_serializing_if = "is_false")]
pub allow_trace: bool,
}
#[derive(Serialize, Clone)]
pub struct ServiceLatency {
pub sampling: SamplingRate,
pub results: String,
}
#[derive(Clone)]
pub enum SamplingRate {
Headers,
Percent(u8),
}
impl Serialize for SamplingRate {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
match self {
Self::Headers => serializer.serialize_str("headers"),
Self::Percent(percent) => serializer.serialize_u8(*percent),
}
}
}
#[derive(Serialize, Clone)]
pub struct AccountExport {
#[serde(skip_serializing_if = "str::is_empty")]
pub name: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub subject: String,
#[serde(rename = "type")]
pub kind: ExportType,
#[serde(skip_serializing_if = "is_false")]
pub token_req: bool,
#[serde(skip_serializing_if = "BTreeMap::is_empty")]
pub revocations: BTreeMap<String, i64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub response_type: Option<ResponseType>,
#[serde(skip_serializing_if = "is_zero_i64")]
pub response_threshold: i64,
#[serde(skip_serializing_if = "Option::is_none")]
pub service_latency: Option<ServiceLatency>,
#[serde(skip_serializing_if = "is_zero_u32")]
pub account_token_position: u32,
#[serde(skip_serializing_if = "is_false")]
pub advertise: bool,
#[serde(skip_serializing_if = "is_false")]
pub allow_trace: bool,
#[serde(skip_serializing_if = "str::is_empty")]
pub description: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub info_url: String,
}
#[derive(Serialize, Clone)]
pub enum ResponseType {
Singleton,
Stream,
Chunked,
}
#[derive(Serialize, Default, Clone)]
pub struct NatsLimits {
#[serde(skip_serializing_if = "is_zero_i64")]
pub subs: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub data: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub payload: i64,
}
#[derive(Serialize, Default, Clone)]
pub struct AccountLimits {
#[serde(skip_serializing_if = "is_zero_i64")]
pub imports: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub exports: i64,
#[serde(skip_serializing_if = "is_false")]
pub wildcards: bool,
#[serde(skip_serializing_if = "is_false")]
pub disallow_bearer: bool,
#[serde(skip_serializing_if = "is_zero_i64")]
pub conn: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub leaf: i64,
}
#[derive(Serialize, Default, Clone)]
pub struct JetStreamLimits {
#[serde(skip_serializing_if = "is_zero_i64")]
pub mem_storage: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub disk_storage: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub streams: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub consumer: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub max_ack_pending: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub mem_max_stream_bytes: i64,
#[serde(skip_serializing_if = "is_zero_i64")]
pub disk_max_stream_bytes: i64,
#[serde(skip_serializing_if = "is_false")]
pub max_bytes_required: bool,
}
#[derive(Serialize, Default, Clone)]
pub struct OperatorLimits {
#[serde(flatten)]
pub nats: NatsLimits,
#[serde(flatten)]
pub account: AccountLimits,
#[serde(flatten)]
pub jetstream: JetStreamLimits,
#[serde(skip_serializing_if = "BTreeMap::is_empty")]
pub tiered_limits: BTreeMap<String, JetStreamLimits>,
}
impl OperatorLimits {
#[must_use]
pub fn unlimited() -> Self {
Self {
nats: NatsLimits {
subs: -1,
data: -1,
payload: -1,
},
account: AccountLimits {
imports: -1,
exports: -1,
wildcards: true,
conn: -1,
leaf: -1,
..AccountLimits::default()
},
..Self::default()
}
}
}
#[derive(Serialize, Clone)]
pub struct WeightedMapping {
pub subject: String,
#[serde(skip_serializing_if = "is_zero_u8")]
pub weight: u8,
#[serde(skip_serializing_if = "str::is_empty")]
pub cluster: String,
}
#[derive(Serialize, Default, Clone)]
pub struct ExternalAuthorization {
#[serde(skip_serializing_if = "Vec::is_empty")]
pub auth_users: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub allowed_accounts: Vec<String>,
#[serde(skip_serializing_if = "str::is_empty")]
pub xkey: String,
}
#[derive(Serialize, Clone)]
pub struct MsgTrace {
#[serde(skip_serializing_if = "str::is_empty")]
pub dest: String,
#[serde(skip_serializing_if = "is_zero_u8")]
pub sampling: u8,
}
#[derive(Serialize, Clone)]
#[serde(rename_all = "lowercase")]
pub enum ClusterTraffic {
System,
Owner,
}
#[derive(Serialize, Default, Clone)]
pub struct AccountClaims {
#[serde(skip_serializing_if = "Vec::is_empty")]
pub imports: Vec<AccountImport>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub exports: Vec<AccountExport>,
#[serde(skip_serializing_if = "operator_limits_is_empty")]
pub limits: OperatorLimits,
#[serde(rename = "signing_keys", skip_serializing_if = "Vec::is_empty")]
pub signing_keys: Vec<String>,
#[serde(skip_serializing_if = "BTreeMap::is_empty")]
pub revocations: BTreeMap<String, i64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub default_permissions: Option<Permissions>,
#[serde(skip_serializing_if = "BTreeMap::is_empty")]
pub mappings: BTreeMap<String, Vec<WeightedMapping>>,
#[serde(skip_serializing_if = "external_authorization_is_empty")]
pub authorization: ExternalAuthorization,
#[serde(skip_serializing_if = "Option::is_none")]
pub trace: Option<MsgTrace>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cluster_traffic: Option<ClusterTraffic>,
}
#[derive(Serialize)]
struct NatsAccount {
#[serde(flatten)]
claims: AccountClaims,
#[serde(rename = "type")]
kind: &'static str,
version: i64,
}
#[derive(Serialize, Default, Clone)]
pub struct OperatorClaims {
#[serde(rename = "signing_keys", skip_serializing_if = "Vec::is_empty")]
pub signing_keys: Vec<String>,
#[serde(rename = "account_server_url", skip_serializing_if = "str::is_empty")]
pub account_server_url: String,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub operator_service_urls: Vec<String>,
#[serde(rename = "system_account", skip_serializing_if = "str::is_empty")]
pub system_account: String,
#[serde(skip_serializing_if = "str::is_empty")]
pub assert_server_version: String,
#[serde(skip_serializing_if = "is_false")]
pub strict_signing_key_usage: bool,
}
#[derive(Serialize)]
struct NatsOperator {
#[serde(flatten)]
claims: OperatorClaims,
#[serde(rename = "type")]
kind: &'static str,
version: i64,
}
#[allow(clippy::trivially_copy_pass_by_ref)]
const fn is_false(value: &bool) -> bool {
!*value
}
#[allow(clippy::trivially_copy_pass_by_ref)]
const fn is_zero_i64(value: &i64) -> bool {
*value == 0
}
#[allow(clippy::trivially_copy_pass_by_ref)]
const fn is_zero_u32(value: &u32) -> bool {
*value == 0
}
#[allow(clippy::trivially_copy_pass_by_ref)]
const fn is_zero_u8(value: &u8) -> bool {
*value == 0
}
fn external_authorization_is_empty(value: &ExternalAuthorization) -> bool {
value.auth_users.is_empty() && value.allowed_accounts.is_empty() && value.xkey.is_empty()
}
fn operator_limits_is_empty(value: &OperatorLimits) -> bool {
value.nats.subs == 0
&& value.nats.data == 0
&& value.nats.payload == 0
&& value.account.imports == 0
&& value.account.exports == 0
&& !value.account.wildcards
&& !value.account.disallow_bearer
&& value.account.conn == 0
&& value.account.leaf == 0
&& value.jetstream.mem_storage == 0
&& value.jetstream.disk_storage == 0
&& value.jetstream.streams == 0
&& value.jetstream.consumer == 0
&& value.jetstream.max_ack_pending == 0
&& value.jetstream.mem_max_stream_bytes == 0
&& value.jetstream.disk_max_stream_bytes == 0
&& !value.jetstream.max_bytes_required
&& value.tiered_limits.is_empty()
}
#[derive(Serialize)]
struct NatsRole {
#[serde(rename = "type")]
kind: &'static str,
version: i64,
}
#[derive(Serialize)]
struct TokenClaimsGeneric<'a, N>
where
N: Serialize,
{
jti: String,
iat: u64,
iss: &'a str,
#[serde(skip_serializing_if = "str::is_empty")]
name: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
exp: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
nbf: Option<u64>,
nats: N,
sub: &'a str,
}
fn build_signing_input<N: Serialize>(
iss: &str,
sub: &str,
name: &str,
iat: u64,
exp: Option<u64>,
nats: N,
) -> Result<String, Error> {
let jti = jti_for_standard_claims(iss, sub, Some(name), iat, exp, None, None)?;
let claims = TokenClaimsGeneric {
jti,
iat,
iss,
name,
exp,
nbf: None,
nats,
sub,
};
let header = URL_SAFE_NO_PAD.encode(HEADER_JSON.as_bytes());
let payload = URL_SAFE_NO_PAD.encode(serde_json::to_vec(&claims)?);
Ok(format!("{header}.{payload}"))
}
pub fn jti_for_standard_claims(
iss: &str,
sub: &str,
name: Option<&str>,
iat: u64,
exp: Option<u64>,
aud: Option<&str>,
nbf: Option<u64>,
) -> Result<String, Error> {
let hash_src = ClaimsHash {
aud: aud.filter(|value| !value.is_empty()),
exp: exp.filter(|value| *value != 0),
jti: None,
iat,
iss,
name: name.filter(|value| !value.is_empty()),
nbf: nbf.filter(|value| *value != 0),
sub,
};
let mut hasher = Sha512_256::new();
hasher.update(serde_json::to_vec(&hash_src)?);
Ok(data_encoding::BASE32_NOPAD.encode(&hasher.finalize()))
}
pub fn signing_input_from_claims(claims: &Value) -> Result<String, Error> {
if !claims.is_object() {
return Err(Error::InvalidClaims(
"nats jwt claims must be a JSON object".into(),
));
}
let header = URL_SAFE_NO_PAD.encode(HEADER_JSON.as_bytes());
let payload = URL_SAFE_NO_PAD.encode(serde_json::to_vec(claims)?);
Ok(format!("{header}.{payload}"))
}
#[derive(Default, Clone)]
pub struct UserPermissions {
pub pub_allow: Vec<String>,
pub pub_deny: Vec<String>,
pub sub_allow: Vec<String>,
pub sub_deny: Vec<String>,
}
pub struct UserJwt {
pub issuer: String,
pub issuer_account: Option<String>,
pub subject_user: String,
pub name: String,
pub issued_at: u64,
pub expires: Option<u64>,
pub permissions: UserPermissions,
}
impl UserJwt {
pub fn signing_input(&self) -> Result<String, Error> {
let nats = NatsUser {
publish: Permission {
allow: self.permissions.pub_allow.clone(),
deny: self.permissions.pub_deny.clone(),
},
sub: Permission {
allow: self.permissions.sub_allow.clone(),
deny: self.permissions.sub_deny.clone(),
},
subs: -1,
data: -1,
payload: -1,
issuer_account: self.issuer_account.clone(),
kind: "user",
version: 2,
};
build_signing_input(
&self.issuer,
&self.subject_user,
&self.name,
self.issued_at,
self.expires,
nats,
)
}
}
pub struct AccountJwt {
pub issuer: String,
pub subject_account: String,
pub name: String,
pub issued_at: u64,
pub expires: Option<u64>,
pub signing_keys: Vec<String>,
pub claims: AccountClaims,
}
impl AccountJwt {
pub fn signing_input(&self) -> Result<String, Error> {
let nats = NatsAccount {
claims: {
let mut claims = self.claims.clone();
if claims.signing_keys.is_empty() {
claims.signing_keys.clone_from(&self.signing_keys);
}
claims
},
kind: "account",
version: 2,
};
build_signing_input(
&self.issuer,
&self.subject_account,
&self.name,
self.issued_at,
self.expires,
nats,
)
}
}
pub struct OperatorJwt {
pub issuer: String,
pub subject_operator: String,
pub name: String,
pub issued_at: u64,
pub expires: Option<u64>,
pub signing_keys: Vec<String>,
pub account_server_url: String,
pub system_account: String,
pub claims: OperatorClaims,
}
impl OperatorJwt {
pub fn signing_input(&self) -> Result<String, Error> {
let nats = NatsOperator {
claims: {
let mut claims = self.claims.clone();
if claims.signing_keys.is_empty() {
claims.signing_keys.clone_from(&self.signing_keys);
}
if claims.account_server_url.is_empty() {
claims
.account_server_url
.clone_from(&self.account_server_url);
}
if claims.system_account.is_empty() {
claims.system_account.clone_from(&self.system_account);
}
claims
},
kind: "operator",
version: 2,
};
build_signing_input(
&self.issuer,
&self.subject_operator,
&self.name,
self.issued_at,
self.expires,
nats,
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RoleKind {
Server,
Curve,
Signer,
}
impl RoleKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Server => "server",
Self::Curve => "curve",
Self::Signer => "signer",
}
}
}
impl core::fmt::Display for RoleKind {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.write_str(self.as_str())
}
}
pub struct RoleJwt {
pub issuer: String,
pub subject: String,
pub name: String,
pub issued_at: u64,
pub expires: Option<u64>,
pub kind: RoleKind,
}
impl RoleJwt {
pub fn signing_input(&self) -> Result<String, Error> {
let nats = NatsRole {
kind: self.kind.as_str(),
version: 2,
};
build_signing_input(
&self.issuer,
&self.subject,
&self.name,
self.issued_at,
self.expires,
nats,
)
}
}
#[must_use]
pub fn assemble(signing_input: &str, signature: &[u8]) -> String {
format!("{signing_input}.{}", URL_SAFE_NO_PAD.encode(signature))
}
pub fn format_user_creds(jwt: &str, seed: &str) -> Result<String, Error> {
let jwt = single_line_creds_field("NATS user JWT", jwt)?;
let seed = single_line_creds_field("NATS user NKey seed", seed)?;
Ok(format!(
"-----BEGIN NATS USER JWT-----\n{jwt}\n------END NATS USER JWT------\n\n************************* IMPORTANT *************************\nNKEY Seed printed below can be used to sign and prove identity.\nNKEYs are sensitive and should be treated as secrets.\n\n-----BEGIN USER NKEY SEED-----\n{seed}\n------END USER NKEY SEED------\n\n*************************************************************\n"
))
}
fn single_line_creds_field<'a>(name: &str, value: &'a str) -> Result<&'a str, Error> {
let trimmed = value.trim();
if trimmed.is_empty() {
return Err(Error::InvalidClaims(format!("{name} must not be empty")));
}
if trimmed.contains('\n') || trimmed.contains('\r') {
return Err(Error::InvalidClaims(format!(
"{name} must be a single line"
)));
}
Ok(trimmed)
}
#[cfg(test)]
mod tests {
use super::*;
use alloc::string::ToString;
use alloc::vec;
use nkeys::{KeyPair, XKey};
use serde_json::json;
#[derive(Default)]
struct TestRng(u64);
impl rand_core::TryRng for TestRng {
type Error = core::convert::Infallible;
fn try_next_u32(&mut self) -> Result<u32, Self::Error> {
let mut bytes = [0u8; 4];
self.try_fill_bytes(&mut bytes)?;
Ok(u32::from_le_bytes(bytes))
}
fn try_next_u64(&mut self) -> Result<u64, Self::Error> {
let mut bytes = [0u8; 8];
self.try_fill_bytes(&mut bytes)?;
Ok(u64::from_le_bytes(bytes))
}
fn try_fill_bytes(&mut self, dst: &mut [u8]) -> Result<(), Self::Error> {
for byte in dst {
*byte = self.0.to_le_bytes()[0].wrapping_mul(37).wrapping_add(0x5a);
self.0 = self.0.wrapping_add(1);
}
Ok(())
}
}
fn signed_token(signing_input: &str, signer: &KeyPair) -> String {
let signature = signer.sign(signing_input.as_bytes()).expect("sign");
assemble(signing_input, &signature)
}
#[test]
fn format_user_creds_renders_canonical_nsc_document() {
let creds = format_user_creds(" jwt.token.sig ", " SUUSERSEED ")
.expect("credentials document must render");
assert_eq!(
creds,
"-----BEGIN NATS USER JWT-----\njwt.token.sig\n------END NATS USER JWT------\n\n************************* IMPORTANT *************************\nNKEY Seed printed below can be used to sign and prove identity.\nNKEYs are sensitive and should be treated as secrets.\n\n-----BEGIN USER NKEY SEED-----\nSUUSERSEED\n------END USER NKEY SEED------\n\n*************************************************************\n"
);
}
#[test]
fn format_user_creds_rejects_empty_or_multiline_fields() {
assert!(format_user_creds("", "SUUSERSEED").is_err());
assert!(format_user_creds("jwt.token.sig", "").is_err());
assert!(format_user_creds("jwt\nsecond", "SUUSERSEED").is_err());
assert!(format_user_creds("jwt.token.sig", "SU\nseed").is_err());
}
fn assert_valid_token(token: &str, signer: &KeyPair, expected_kind: &str) {
let signer_public = signer.public_key();
let decoded = decode_nats_jwt(token).expect("decode jwt");
assert_eq!(decoded.claims().issuer, signer_public);
assert_eq!(decoded.claims().nats_type.as_deref(), Some(expected_kind));
assert!(decoded.verify_signature(decoded.issuer_public_key()));
let validation = decoded
.verify_with_candidates(
[CandidateSigner::Nkey(&signer_public)],
decoded.claims().issued_at.unwrap_or_default(),
)
.expect("validate");
assert_eq!(validation.reason, NatsJwtValidationReason::Valid);
assert!(validation.is_valid());
assert!(validation.matched_signer.is_some());
let raw_validation = decoded
.verify_with_candidates(
[CandidateSigner::RawPublicKey(decoded.issuer_public_key())],
decoded.claims().issued_at.unwrap_or_default(),
)
.expect("validate raw");
assert_eq!(raw_validation.reason, NatsJwtValidationReason::Valid);
assert!(raw_validation.is_valid());
assert!(raw_validation.matched_signer.is_some());
}
#[test]
fn account_nkey_encoding_matches_nkeys() {
let account = KeyPair::new_account();
let expected = account.public_key(); let (role, raw) = decode_public(&expected).expect("decode");
assert_eq!(role, NkeyType::Account);
assert_eq!(encode_public(NkeyType::Account, &raw).unwrap(), expected);
}
#[test]
fn user_jwt_verifies_under_issuer_account_key() {
let account = KeyPair::new_account();
let user = KeyPair::new_user();
let jwt = UserJwt {
issuer: account.public_key(),
issuer_account: None,
subject_user: user.public_key(),
name: "bob".into(),
issued_at: 1_782_000_000,
expires: None,
permissions: UserPermissions::default(),
};
let signing_input = jwt.signing_input().expect("signing input");
let signature = account.sign(signing_input.as_bytes()).expect("sign");
let token = assemble(&signing_input, &signature);
let parts: Vec<&str> = token.split('.').collect();
assert_eq!(parts.len(), 3);
let header = URL_SAFE_NO_PAD.decode(parts[0]).unwrap();
assert_eq!(header, HEADER_JSON.as_bytes());
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
assert_eq!(claims["iss"], account.public_key());
assert_eq!(claims["sub"], user.public_key());
assert_eq!(claims["nats"]["type"], "user");
assert_eq!(claims["nats"]["version"], 2);
assert!(claims["jti"].as_str().unwrap().len() >= 52);
let issuer = KeyPair::from_public_key(claims["iss"].as_str().unwrap()).unwrap();
issuer
.verify(signing_input.as_bytes(), &signature)
.expect("signature must verify under iss account key");
}
#[test]
fn user_jwt_sets_nats_issuer_account_when_signing_key() {
let signing = KeyPair::new_account();
let account_identity = KeyPair::new_account();
let user = KeyPair::new_user();
let jwt = UserJwt {
issuer: signing.public_key(),
issuer_account: Some(account_identity.public_key()),
subject_user: user.public_key(),
name: "svc".into(),
issued_at: 1_782_000_000,
expires: None,
permissions: UserPermissions::default(),
};
let signing_input = jwt.signing_input().expect("signing input");
let parts: Vec<&str> = signing_input.split('.').collect();
assert_eq!(parts.len(), 2);
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
assert_eq!(claims["iss"], signing.public_key());
assert_eq!(
claims["nats"]["issuer_account"],
account_identity.public_key()
);
let plain = UserJwt {
issuer: signing.public_key(),
issuer_account: None,
subject_user: user.public_key(),
name: "svc".into(),
issued_at: 1_782_000_000,
expires: None,
permissions: UserPermissions::default(),
};
let plain_input = plain.signing_input().expect("signing input");
let plain_parts: Vec<&str> = plain_input.split('.').collect();
let plain_claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(plain_parts[1]).unwrap()).unwrap();
assert!(plain_claims["nats"].get("issuer_account").is_none());
}
#[test]
fn encode_public_round_trips_each_role() {
let account = KeyPair::new_account();
let (_, raw) = decode_public(&account.public_key()).unwrap();
for role in [
NkeyType::Account,
NkeyType::Cluster,
NkeyType::Server,
NkeyType::Operator,
NkeyType::User,
NkeyType::Curve,
] {
let encoded = encode_public(role, &raw).unwrap();
assert!(
encoded.starts_with(role.letter()),
"{role} key should start with {}, got {encoded}",
role.letter()
);
let (decoded_role, decoded_raw) = decode_public(&encoded).unwrap();
assert_eq!(decoded_role, role);
assert_eq!(decoded_raw, raw);
}
}
#[test]
fn nats_curve_box_round_trips_and_validates_wire_shape() {
let sender_private = Zeroizing::new([0x11; 32]);
let receiver_private = Zeroizing::new([0x22; 32]);
let sender_public =
encode_public(NkeyType::Curve, &xkey_public_from_private(&sender_private))
.expect("sender xkey encodes");
let receiver_public = encode_public(
NkeyType::Curve,
&xkey_public_from_private(&receiver_private),
)
.expect("receiver xkey encodes");
let mut rng = TestRng::default();
let boxed = seal_nats_curve(&sender_private, &receiver_public, b"payload", &mut rng)
.expect("seal succeeds");
assert!(boxed.starts_with(XKEY_VERSION_V1));
assert_eq!(
boxed.len(),
XKEY_VERSION_V1.len() + XKEY_NONCE_LEN + XKEY_TAG_LEN + b"payload".len()
);
let opened =
open_nats_curve(&receiver_private, &sender_public, &boxed).expect("open succeeds");
assert_eq!(opened.as_slice(), b"payload");
let mut tampered = boxed;
if let Some(byte) = tampered.last_mut() {
*byte ^= 0x80;
}
assert!(matches!(
open_nats_curve(&receiver_private, &sender_public, &tampered),
Err(Error::XKeyOpenFailed)
));
}
#[test]
fn nats_curve_box_interops_with_nkeys_xkey_both_directions() {
let sender_private = Zeroizing::new([0x33; 32]);
let receiver_private = Zeroizing::new([0x44; 32]);
let sender = XKey::new_from_raw(*sender_private);
let receiver = XKey::new_from_raw(*receiver_private);
let mut rng = TestRng::default();
let basil_box = seal_nats_curve(
&sender_private,
&receiver.public_key(),
b"from basil",
&mut rng,
)
.expect("basil seal succeeds");
let nkeys_opened = receiver
.open(&basil_box, &sender)
.expect("nkeys opens basil box");
assert_eq!(nkeys_opened, b"from basil");
let nkeys_box = sender
.seal(b"from nkeys", &receiver)
.expect("nkeys seal succeeds");
let basil_opened = open_nats_curve(&receiver_private, &sender.public_key(), &nkeys_box)
.expect("basil opens nkeys box");
assert_eq!(basil_opened.as_slice(), b"from nkeys");
}
#[test]
fn nkey_type_letter_round_trips() {
for role in NkeyType::ALL {
assert_eq!(NkeyType::from_letter(role.letter()), Some(role));
}
assert_eq!(NkeyType::from_letter('Z'), None);
assert_eq!(NkeyType::from_letter('S'), None);
}
#[test]
fn require_public_prefix_rejects_wrong_role() {
let user = KeyPair::new_user();
let err =
require_public_prefix(&user.public_key(), NkeyType::Account).expect_err("wrong role");
assert!(matches!(
err,
Error::UnexpectedPrefix {
expected: NkeyType::Account,
actual: NkeyType::User
}
));
require_public_prefix(&user.public_key(), NkeyType::User).expect("user role");
}
#[test]
fn account_jwt_verifies_under_issuer_operator_key() {
let operator = KeyPair::new_operator();
let account = KeyPair::new_account();
let signing = KeyPair::new_account();
let jwt = AccountJwt {
issuer: operator.public_key(),
subject_account: account.public_key(),
name: "acme".into(),
issued_at: 1_782_000_000,
expires: Some(1_782_003_600),
signing_keys: vec![signing.public_key()],
claims: AccountClaims::default(),
};
let signing_input = jwt.signing_input().expect("signing input");
let signature = operator.sign(signing_input.as_bytes()).expect("sign");
let token = assemble(&signing_input, &signature);
let parts: Vec<&str> = token.split('.').collect();
assert_eq!(parts.len(), 3);
assert_eq!(
URL_SAFE_NO_PAD.decode(parts[0]).unwrap(),
HEADER_JSON.as_bytes()
);
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
assert_eq!(claims["iss"], operator.public_key());
assert_eq!(claims["sub"], account.public_key());
assert_eq!(claims["nats"]["type"], "account");
assert_eq!(claims["nats"]["version"], 2);
assert_eq!(claims["nats"]["signing_keys"][0], signing.public_key());
assert_eq!(claims["exp"], 1_782_003_600u64);
let issuer = KeyPair::from_public_key(claims["iss"].as_str().unwrap()).unwrap();
issuer
.verify(signing_input.as_bytes(), &signature)
.expect("must verify under iss operator key");
}
#[test]
fn operator_jwt_self_signed_verifies_and_omits_empty_fields() {
let operator = KeyPair::new_operator();
let sys = KeyPair::new_account();
let jwt = OperatorJwt {
issuer: operator.public_key(),
subject_operator: operator.public_key(),
name: "root-op".into(),
issued_at: 1_782_000_000,
expires: None,
signing_keys: Vec::new(),
account_server_url: "nats://localhost:4222".into(),
system_account: sys.public_key(),
claims: OperatorClaims::default(),
};
let signing_input = jwt.signing_input().expect("signing input");
let signature = operator.sign(signing_input.as_bytes()).expect("sign");
let token = assemble(&signing_input, &signature);
let parts: Vec<&str> = token.split('.').collect();
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
assert_eq!(claims["iss"], operator.public_key());
assert_eq!(claims["sub"], operator.public_key());
assert_eq!(claims["nats"]["type"], "operator");
assert_eq!(
claims["nats"]["account_server_url"],
"nats://localhost:4222"
);
assert_eq!(claims["nats"]["system_account"], sys.public_key());
assert!(claims.get("exp").is_none());
assert!(claims["nats"].get("signing_keys").is_none());
let issuer = KeyPair::from_public_key(claims["iss"].as_str().unwrap()).unwrap();
issuer
.verify(signing_input.as_bytes(), &signature)
.expect("must verify under iss operator key");
}
#[test]
#[allow(clippy::too_many_lines)]
fn account_jwt_serializes_rich_claim_surface() {
let operator = KeyPair::new_operator();
let account = KeyPair::new_account();
let signing = KeyPair::new_account();
let exporting = KeyPair::new_account();
let user = KeyPair::new_user();
let mut revocations = BTreeMap::new();
revocations.insert(user.public_key(), 1_782_000_001);
let mut export_revocations = BTreeMap::new();
export_revocations.insert("*".to_string(), 1_782_000_002);
let mut tiered_limits = BTreeMap::new();
tiered_limits.insert(
"gold".to_string(),
JetStreamLimits {
mem_storage: 64,
disk_storage: 128,
streams: 3,
consumer: 4,
max_ack_pending: 5,
mem_max_stream_bytes: 6,
disk_max_stream_bytes: 7,
max_bytes_required: true,
},
);
let jwt = AccountJwt {
issuer: operator.public_key(),
subject_account: account.public_key(),
name: "tenant".into(),
issued_at: 1_782_000_000,
expires: None,
signing_keys: vec![signing.public_key()],
claims: AccountClaims {
imports: vec![AccountImport {
name: "orders-in".into(),
subject: "orders.*".into(),
account: exporting.public_key(),
token: "activation.jwt".into(),
to: String::new(),
local_subject: "tenant.$1".into(),
kind: ExportType::Stream,
share: false,
allow_trace: true,
}],
exports: vec![AccountExport {
name: "lookup".into(),
subject: "svc.lookup".into(),
kind: ExportType::Service,
token_req: true,
revocations: export_revocations,
response_type: Some(ResponseType::Stream),
response_threshold: 1_000_000,
service_latency: Some(ServiceLatency {
sampling: SamplingRate::Headers,
results: "latency.results".into(),
}),
account_token_position: 0,
advertise: true,
allow_trace: true,
description: "lookup service".into(),
info_url: "https://example.test/lookup".into(),
}],
limits: OperatorLimits {
nats: NatsLimits {
subs: 100,
data: 1_024,
payload: 256,
},
account: AccountLimits {
imports: 10,
exports: 11,
wildcards: true,
disallow_bearer: true,
conn: 12,
leaf: 13,
},
jetstream: JetStreamLimits::default(),
tiered_limits,
},
signing_keys: Vec::new(),
revocations,
default_permissions: Some(Permissions {
publish: Permission {
allow: vec!["pub.>".into()],
deny: vec!["pub.secret".into()],
},
sub: Permission {
allow: vec!["sub.>".into()],
deny: Vec::new(),
},
resp: Some(ResponsePermission {
max: 1,
ttl: 2_000_000_000,
}),
}),
mappings: BTreeMap::from([(
"legacy.>".into(),
vec![WeightedMapping {
subject: "modern.>".into(),
weight: 50,
cluster: "edge".into(),
}],
)]),
authorization: ExternalAuthorization {
auth_users: vec![user.public_key()],
allowed_accounts: vec![account.public_key()],
xkey: encode_public(NkeyType::Curve, &[9; 32]).unwrap(),
},
trace: Some(MsgTrace {
dest: "trace.out".into(),
sampling: 25,
}),
cluster_traffic: Some(ClusterTraffic::System),
},
};
let signing_input = jwt.signing_input().expect("signing input");
let parts: Vec<&str> = signing_input.split('.').collect();
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
let nats = &claims["nats"];
assert_eq!(nats["type"], "account");
assert_eq!(nats["imports"][0]["type"], "stream");
assert_eq!(nats["imports"][0]["local_subject"], "tenant.$1");
assert_eq!(nats["exports"][0]["type"], "service");
assert_eq!(nats["exports"][0]["response_type"], "Stream");
assert_eq!(nats["exports"][0]["service_latency"]["sampling"], "headers");
assert_eq!(nats["limits"]["imports"], 10);
assert_eq!(nats["limits"]["wildcards"], true);
assert_eq!(nats["limits"]["tiered_limits"]["gold"]["mem_storage"], 64);
assert_eq!(nats["revocations"][user.public_key()], 1_782_000_001);
assert_eq!(nats["default_permissions"]["pub"]["allow"][0], "pub.>");
assert_eq!(nats["default_permissions"]["resp"]["ttl"], 2_000_000_000i64);
assert_eq!(nats["mappings"]["legacy.>"][0]["subject"], "modern.>");
assert_eq!(nats["authorization"]["auth_users"][0], user.public_key());
assert_eq!(nats["trace"]["dest"], "trace.out");
assert_eq!(nats["cluster_traffic"], "system");
assert_eq!(nats["signing_keys"][0], signing.public_key());
}
#[test]
fn unlimited_operator_limits_serialize_as_a_real_not_deny_all_limits_block() {
let operator = KeyPair::new_operator();
let account = KeyPair::new_account();
let jwt = AccountJwt {
issuer: operator.public_key(),
subject_account: account.public_key(),
name: "tenant".into(),
issued_at: 1_782_000_000,
expires: None,
signing_keys: Vec::new(),
claims: AccountClaims {
limits: OperatorLimits::unlimited(),
..AccountClaims::default()
},
};
let signing_input = jwt.signing_input().expect("signing input");
let parts: Vec<&str> = signing_input.split('.').collect();
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
let limits = &claims["nats"]["limits"];
assert!(limits.is_object());
assert_eq!(limits["conn"], -1);
assert_eq!(limits["subs"], -1);
assert_eq!(limits["leaf"], -1);
assert_eq!(limits["data"], -1);
assert_eq!(limits["payload"], -1);
assert_eq!(limits["imports"], -1);
assert_eq!(limits["exports"], -1);
assert_eq!(limits["wildcards"], true);
assert!(limits["mem_storage"].is_null());
}
#[test]
fn operator_jwt_serializes_service_urls_and_strict_signing_keys() {
let operator = KeyPair::new_operator();
let signing = KeyPair::new_operator();
let jwt = OperatorJwt {
issuer: operator.public_key(),
subject_operator: operator.public_key(),
name: "root".into(),
issued_at: 1_782_000_000,
expires: None,
signing_keys: vec![signing.public_key()],
account_server_url: "https://accounts.example.test/jwt/v1".into(),
system_account: String::new(),
claims: OperatorClaims {
operator_service_urls: vec![
"nats://nats.example.test:4222".into(),
"tls://nats.example.test:4443".into(),
],
assert_server_version: "2.11.0".into(),
strict_signing_key_usage: true,
..OperatorClaims::default()
},
};
let signing_input = jwt.signing_input().expect("signing input");
let parts: Vec<&str> = signing_input.split('.').collect();
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
let nats = &claims["nats"];
assert_eq!(nats["type"], "operator");
assert_eq!(nats["signing_keys"][0], signing.public_key());
assert_eq!(
nats["account_server_url"],
"https://accounts.example.test/jwt/v1"
);
assert_eq!(
nats["operator_service_urls"][0],
"nats://nats.example.test:4222"
);
assert_eq!(nats["assert_server_version"], "2.11.0");
assert_eq!(nats["strict_signing_key_usage"], true);
}
#[test]
fn role_jwt_verifies_and_carries_kind() {
let operator = KeyPair::new_operator();
let server = KeyPair::new_server();
let jwt = RoleJwt {
issuer: operator.public_key(),
subject: server.public_key(),
name: "nats-server".into(),
issued_at: 1_782_000_000,
expires: Some(1_782_003_600),
kind: RoleKind::Server,
};
let signing_input = jwt.signing_input().expect("signing input");
let signature = operator.sign(signing_input.as_bytes()).expect("sign");
let token = assemble(&signing_input, &signature);
let parts: Vec<&str> = token.split('.').collect();
let claims: serde_json::Value =
serde_json::from_slice(&URL_SAFE_NO_PAD.decode(parts[1]).unwrap()).unwrap();
assert_eq!(claims["iss"], operator.public_key());
assert_eq!(claims["sub"], server.public_key());
assert_eq!(claims["nats"]["type"], "server");
assert_eq!(claims["nats"]["version"], 2);
let issuer = KeyPair::from_public_key(claims["iss"].as_str().unwrap()).unwrap();
issuer
.verify(signing_input.as_bytes(), &signature)
.expect("must verify under iss operator key");
}
#[test]
fn decode_and_validate_known_good_user_account_operator_tokens() {
let account = KeyPair::new_account();
let user = KeyPair::new_user();
let user_jwt = UserJwt {
issuer: account.public_key(),
issuer_account: None,
subject_user: user.public_key(),
name: "bob".into(),
issued_at: 1_782_000_000,
expires: Some(1_782_003_600),
permissions: UserPermissions::default(),
};
let user_input = user_jwt.signing_input().expect("user input");
let user_token = signed_token(&user_input, &account);
assert_valid_token(&user_token, &account, "user");
let operator = KeyPair::new_operator();
let tenant = KeyPair::new_account();
let account_jwt = AccountJwt {
issuer: operator.public_key(),
subject_account: tenant.public_key(),
name: "tenant".into(),
issued_at: 1_782_000_000,
expires: Some(1_782_003_600),
signing_keys: Vec::new(),
claims: AccountClaims::default(),
};
let account_input = account_jwt.signing_input().expect("account input");
let account_token = signed_token(&account_input, &operator);
assert_valid_token(&account_token, &operator, "account");
let operator_jwt = OperatorJwt {
issuer: operator.public_key(),
subject_operator: operator.public_key(),
name: "root".into(),
issued_at: 1_782_000_000,
expires: Some(1_782_003_600),
signing_keys: Vec::new(),
account_server_url: String::new(),
system_account: String::new(),
claims: OperatorClaims::default(),
};
let operator_input = operator_jwt.signing_input().expect("operator input");
let operator_token = signed_token(&operator_input, &operator);
assert_valid_token(&operator_token, &operator, "operator");
}
#[test]
fn validation_rejects_expired_and_not_yet_valid_tokens() {
let account = KeyPair::new_account();
let user = KeyPair::new_user();
let expired = UserJwt {
issuer: account.public_key(),
issuer_account: None,
subject_user: user.public_key(),
name: "old".into(),
issued_at: 100,
expires: Some(200),
permissions: UserPermissions::default(),
};
let expired_input = expired.signing_input().expect("expired input");
let expired_token = signed_token(&expired_input, &account);
let decoded_expired = decode_nats_jwt(&expired_token).expect("decode expired");
let account_public = account.public_key();
let expired_validation = decoded_expired
.verify_with_candidates([CandidateSigner::Nkey(&account_public)], 200)
.expect("validate expired");
assert_eq!(expired_validation.reason, NatsJwtValidationReason::Expired);
assert!(!expired_validation.is_valid());
assert!(expired_validation.matched_signer.is_none());
let claims = json!({
"jti": "manual",
"iat": 100_u64,
"iss": account.public_key(),
"name": "future",
"nbf": 300_u64,
"nats": {"type": "user", "version": 2},
"sub": user.public_key(),
});
let future_input = signing_input_from_claims(&claims).expect("future input");
let future_token = signed_token(&future_input, &account);
let decoded_future = decode_nats_jwt(&future_token).expect("decode future");
let future_validation = decoded_future
.verify_with_candidates([CandidateSigner::Nkey(&account_public)], 299)
.expect("validate future");
assert_eq!(
future_validation.reason,
NatsJwtValidationReason::NotYetValid
);
assert!(!future_validation.is_valid());
assert!(future_validation.matched_signer.is_none());
}
#[test]
fn validation_rejects_tampered_signature_and_unknown_signer() {
let account = KeyPair::new_account();
let other = KeyPair::new_account();
let user = KeyPair::new_user();
let jwt = UserJwt {
issuer: account.public_key(),
issuer_account: None,
subject_user: user.public_key(),
name: "bob".into(),
issued_at: 1_782_000_000,
expires: None,
permissions: UserPermissions::default(),
};
let input = jwt.signing_input().expect("input");
let mut signature = account.sign(input.as_bytes()).expect("sign");
signature[0] ^= 1;
let tampered = assemble(&input, &signature);
let account_public = account.public_key();
let decoded_tampered = decode_nats_jwt(&tampered).expect("decode tampered");
let bad_signature = decoded_tampered
.verify_with_candidates([CandidateSigner::Nkey(&account_public)], 1_782_000_000)
.expect("validate tampered");
assert_eq!(bad_signature.reason, NatsJwtValidationReason::BadSignature);
assert!(!bad_signature.is_valid());
assert!(bad_signature.matched_signer.is_none());
let valid = signed_token(&input, &account);
let decoded_valid = decode_nats_jwt(&valid).expect("decode valid");
let other_public = other.public_key();
let unknown = decoded_valid
.verify_with_candidates([CandidateSigner::Nkey(&other_public)], 1_782_000_000)
.expect("validate unknown");
assert_eq!(unknown.reason, NatsJwtValidationReason::UnknownSigner);
assert!(!unknown.is_valid());
assert!(unknown.matched_signer.is_none());
}
#[test]
fn public_nkey_verifies_raw_signature() {
let user = KeyPair::new_user();
let message = b"basil sealed invocation digest";
let signature = user.sign(message).expect("sign");
assert!(
verify_public_signature(&user.public_key(), message, &signature)
.expect("valid public nkey")
);
let other = KeyPair::new_user();
assert!(
!verify_public_signature(&other.public_key(), message, &signature)
.expect("valid other public nkey")
);
assert!(matches!(
verify_public_signature(&user.public_key(), message, &[0_u8; 63]),
Err(Error::BadSignatureLen(63))
));
assert!(matches!(
verify_public_signature("not-an-nkey", message, &signature),
Err(Error::BadPublicKeyLen(_) | Error::UnsupportedPrefix(_))
));
}
#[test]
fn decode_rejects_malformed_tokens() {
assert!(matches!(
decode_nats_jwt("not-a-jwt"),
Err(Error::MalformedJwt(_))
));
let claims = URL_SAFE_NO_PAD.encode(br#"{"iss":"O","sub":"U"}"#);
let bad_header = format!(
"{}.{}.{}",
URL_SAFE_NO_PAD.encode(br#"{"alg":"none"}"#),
claims,
""
);
assert!(matches!(
decode_nats_jwt(&bad_header),
Err(Error::MalformedJwt(_))
));
let header = URL_SAFE_NO_PAD.encode(HEADER_JSON.as_bytes());
let bad_signature = format!("{header}.{claims}.{}", URL_SAFE_NO_PAD.encode([1_u8, 2]));
assert!(matches!(
decode_nats_jwt(&bad_signature),
Err(Error::MalformedJwt(_) | Error::BadSignatureLen(_))
));
}
}