Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 5 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -102,12 +102,16 @@ hmac = "0.12"
hkdf = "0.12"
subtle = "2"
base64 = "0.23"
getrandom = { version = "0.4", features = ["std"] }
getrandom = { version = "0.4", features = ["std", "wasm_js"] }
# wasm builds pin the older generations: `getrandom` 0.3 needs `wasm_js` and 0.2
# needs `js` to reach the browser CSPRNG. Both are aliased so one crate can hold
# all three generations at once.
getrandom_03 = { package = "getrandom", version = "0.3", features = ["wasm_js"] }
getrandom_02 = { package = "getrandom", version = "0.2", features = ["js"] }
# The browser clock and timer bindings, reached only from
# `wasm32-unknown-unknown`; declared here so a crate depends on it the same way
# it depends on every other shared dependency.
js-sys = "0.3"
aes-gcm = "0.10"
argon2 = "0.5"
rsa = "0.9"
Expand Down
8 changes: 3 additions & 5 deletions nodedb-array/src/sync/hlc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@
//! spec's u48 intent while keeping byte alignment simple.

use std::sync::Mutex;
use std::time::{SystemTime, UNIX_EPOCH};

use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -165,11 +164,10 @@ impl HlcGenerator {

/// Return the current wall-clock milliseconds since Unix epoch.
fn now_ms() -> ArrayResult<u64> {
SystemTime::now()
.duration_since(UNIX_EPOCH)
nodedb_types::clock::since_epoch()
.map(|d| d.as_millis() as u64)
.map_err(|e| ArrayError::InvalidHlc {
detail: format!("system clock before Unix epoch: {e}"),
.ok_or_else(|| ArrayError::InvalidHlc {
detail: "system clock before Unix epoch".to_owned(),
})
}

Expand Down
9 changes: 4 additions & 5 deletions nodedb-crdt/src/dead_letter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@
//! 3. A machine-readable `CompensationHint` suggesting how to fix it.

use std::collections::VecDeque;
use std::time::{SystemTime, UNIX_EPOCH};

use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -91,8 +90,9 @@ pub struct DeadLetter {

/// Timestamp when the rejection occurred (unix millis).
///
/// Node-local and NON-DETERMINISTIC: sourced from `SystemTime::now()` at
/// enqueue time, so two replicas that reject the same delta will record
/// Node-local and NON-DETERMINISTIC: sourced from
/// `nodedb_types::clock::since_epoch()` at enqueue time, so two replicas
/// that reject the same delta will record
/// different values. Safe today because the DLQ is a per-replica in-memory
/// structure that is never replicated or compared across nodes. If the DLQ
/// is ever replicated, snapshotted into shared state, or diffed between
Expand Down Expand Up @@ -158,8 +158,7 @@ impl DeadLetterQueue {
let id = self.next_id;
self.next_id += 1;

let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
let now = nodedb_types::clock::since_epoch()
.unwrap_or_default()
.as_millis() as u64;

Expand Down
23 changes: 5 additions & 18 deletions nodedb-crdt/src/validator/policy_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -121,8 +121,7 @@ impl Validator {
max_retries,
ttl_secs,
} => {
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
let now_ms = nodedb_types::clock::since_epoch()
.unwrap_or_default()
.as_millis() as u64;

Expand Down Expand Up @@ -286,10 +285,7 @@ mod tests {
fields: vec![("email".into(), LoroValue::String("a@b.com".into()))],
};

let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
let now = nodedb_types::clock::since_epoch().unwrap().as_millis() as u64;

let resolution = validator
.validate_with_policy(
Expand Down Expand Up @@ -342,10 +338,7 @@ mod tests {
],
};

let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
let now = nodedb_types::clock::since_epoch().unwrap().as_millis() as u64;

let resolution = validator
.validate_with_policy(
Expand Down Expand Up @@ -387,10 +380,7 @@ mod tests {
fields: vec![("author_id".into(), LoroValue::String("u1".into()))],
};

let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
let now = nodedb_types::clock::since_epoch().unwrap().as_millis() as u64;

let resolution = validator
.validate_with_policy(
Expand Down Expand Up @@ -435,10 +425,7 @@ mod tests {
fields: vec![("email".into(), LoroValue::String("a@b.com".into()))],
};

let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64;
let now = nodedb_types::clock::since_epoch().unwrap().as_millis() as u64;

let resolution = validator
.validate_with_policy(
Expand Down
6 changes: 2 additions & 4 deletions nodedb-crdt/src/validator/validate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -133,8 +133,7 @@ impl Validator {
// Check auth expiry: agents that accumulated deltas offline must
// re-authenticate before syncing.
if auth.auth_expires_at > 0 {
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
let now_ms = nodedb_types::clock::since_epoch()
.unwrap_or_default()
.as_millis() as u64;
if now_ms > auth.auth_expires_at {
Expand All @@ -147,8 +146,7 @@ impl Validator {

self.verify_delta_auth(&change.collection, &auth, &delta_bytes)?;

let hlc_timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
let hlc_timestamp = nodedb_types::clock::since_epoch()
.unwrap_or_default()
.as_millis() as u64;

Expand Down
13 changes: 3 additions & 10 deletions nodedb-sql/src/parser/preprocess/temporal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -251,8 +251,7 @@ fn parse_temporal_expr(token: &str) -> Result<i64, TemporalParseError> {

// NOW() — case-insensitive
if t.to_uppercase() == "NOW()" {
let ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
let ms = nodedb_types::clock::since_epoch()
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
return Ok(ms);
Expand Down Expand Up @@ -579,17 +578,11 @@ mod tests {

#[test]
fn as_of_system_time_now() {
let before = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
let before = nodedb_types::clock::since_epoch().unwrap().as_millis() as i64;
let ex = extract("SELECT * FROM array_slice('g', '{}') AS OF SYSTEM TIME NOW()")
.unwrap()
.unwrap();
let after = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
let after = nodedb_types::clock::since_epoch().unwrap().as_millis() as i64;
let ts = ex
.temporal
.system_time
Expand Down
4 changes: 1 addition & 3 deletions nodedb-sql/src/planner/defaults/kind.rs
Original file line number Diff line number Diff line change
Expand Up @@ -64,9 +64,7 @@ pub(super) fn cuid2_with_length(len: usize) -> crate::Result<String> {

/// Render the current wall-clock instant the way `DEFAULT NOW()` stores it.
fn now_rfc3339() -> String {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default();
let now = nodedb_types::clock::since_epoch().unwrap_or_default();
chrono::DateTime::from_timestamp_millis(now.as_millis() as i64)
.map(|dt| dt.to_rfc3339())
.unwrap_or_else(|| now.as_millis().to_string())
Expand Down
11 changes: 11 additions & 0 deletions nodedb-types/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -41,3 +41,14 @@ sha2 = { workspace = true }

[dev-dependencies]
toml = { workspace = true }

# The browser clock, reached only by `clock::since_epoch`'s
# `wasm32-unknown-unknown` arm. Every other target, `wasm32-wasip1` included,
# reads the clock through std.
[target.'cfg(all(target_arch = "wasm32", target_os = "unknown"))'.dependencies]
js-sys = { workspace = true }
# `getrandom` reaches this crate transitively through `aes-gcm`, and 0.3 refuses
# `wasm32-unknown-unknown` unless `wasm_js` is on and the matching cfg is set.
# Naming it here is what puts the feature in this crate's resolution.
getrandom_02 = { workspace = true }
getrandom_03 = { workspace = true }
106 changes: 106 additions & 0 deletions nodedb-types/src/clock.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,106 @@
// SPDX-License-Identifier: BUSL-1.1

//! Wall-clock access, one target split.
//!
//! `wasm32-unknown-unknown` has no std clock: `SystemTime::now()` panics with
//! "time not implemented on this platform". Builds for that target read the
//! host clock through `js_sys` instead; every other target, `wasm32-wasip1`
//! included, uses std.
//!
//! Callers keep their own conversion, saturation, and pre-epoch fallback —
//! this module only answers "how long since the Unix epoch".

use std::time::Duration;

/// A reading of `Date.now()`, in milliseconds, as an elapsed `Duration`.
///
/// `None` when the reading is negative, i.e. the clock stands before the Unix
/// epoch. Both arms agree on that contract, which is the whole point: callers
/// map the absence to an error, so a target that answered `Some(0)` instead
/// would store a plausible zero timestamp and lose the failure.
///
/// `Date.now()` returns an `f64`, and `as u64` **saturates** a negative value to
/// `0` rather than wrapping, so the guard has to come before the conversion —
/// converting first is exactly the bug this exists to prevent. Declared for
/// every target so the guard itself is testable on the host, where the
/// `js_sys` arm never compiles and could otherwise only be reviewed by eye.
#[cfg_attr(not(test), allow(dead_code))]
fn duration_from_epoch_millis(millis: f64) -> Option<Duration> {
if millis < 0.0 {
return None;
}
Some(Duration::from_millis(millis as u64))
}

/// Time elapsed since the Unix epoch.
///
/// `None` when the system clock reads earlier than the epoch. On
/// `wasm32-unknown-unknown` the value comes from `Date.now()` and therefore
/// has millisecond resolution.
#[cfg(all(target_arch = "wasm32", target_os = "unknown"))]
pub fn since_epoch() -> Option<Duration> {
duration_from_epoch_millis(js_sys::Date::now())
}

#[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
pub fn since_epoch() -> Option<Duration> {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.ok()
}

#[cfg(test)]
mod tests {
use super::{duration_from_epoch_millis, since_epoch};
use std::time::Duration;

/// The helper exists so callers never reach `SystemTime::now()` on a
/// target that has no clock. It must answer on every target it builds for.
#[test]
fn since_epoch_is_available() {
let elapsed = since_epoch().expect("clock reads after the Unix epoch");
assert!(
elapsed.as_secs() > 1_600_000_000,
"clock reads a date after 2020: {elapsed:?}"
);
}

/// Callers assume the clock does not go backwards between two reads
/// (HLC monotonicity, retry stamps, auth expiry).
#[test]
fn since_epoch_does_not_go_backwards() {
let first = since_epoch().expect("clock");
let second = since_epoch().expect("clock");
assert!(second >= first, "{second:?} is before {first:?}");
}

/// A clock standing before the Unix epoch must read as absent, on every
/// target.
///
/// This is the arm `wasm32-unknown-unknown` reaches, and it cannot be
/// exercised by running that target here — hence the guard living in a
/// function the host can call. Without it `f64 as u64` saturates a negative
/// reading to `0`, so the browser arm would answer `Some(0)` where the std
/// arm answers `None`, and a caller mapping absence to an error would store
/// a zero timestamp instead of reporting the fault.
#[test]
fn a_pre_epoch_reading_is_absent_not_zero() {
assert_eq!(duration_from_epoch_millis(-1.0), None);
assert_eq!(duration_from_epoch_millis(-1_700_000_000_000.0), None);
assert_eq!(
duration_from_epoch_millis(-0.5),
None,
"a sub-millisecond pre-epoch reading is still before the epoch"
);
}

/// The epoch itself and anything after it are elapsed time, unchanged.
#[test]
fn a_post_epoch_reading_is_the_elapsed_time() {
assert_eq!(duration_from_epoch_millis(0.0), Some(Duration::ZERO));
assert_eq!(
duration_from_epoch_millis(1_700_000_000_000.0),
Some(Duration::from_millis(1_700_000_000_000))
);
}
}
24 changes: 11 additions & 13 deletions nodedb-types/src/datetime/timestamp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -74,20 +74,18 @@ impl NdbDateTime {
/// `i64::MAX` (year ~292,277 CE) rather than wrapping — clocks that far
/// in the future simply report the maximum representable timestamp.
pub fn now() -> Self {
let dur = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_else(|_| {
use std::sync::atomic::{AtomicBool, Ordering};
static LOGGED: AtomicBool = AtomicBool::new(false);
if !LOGGED.swap(true, Ordering::Relaxed) {
tracing::error!(
module = module_path!(),
"system clock is before UNIX_EPOCH; using 0 (epoch) \
let dur = crate::clock::since_epoch().unwrap_or_else(|| {
use std::sync::atomic::{AtomicBool, Ordering};
static LOGGED: AtomicBool = AtomicBool::new(false);
if !LOGGED.swap(true, Ordering::Relaxed) {
tracing::error!(
module = module_path!(),
"system clock is before UNIX_EPOCH; using 0 (epoch) \
— check NTP/RTC configuration"
);
}
std::time::Duration::ZERO
});
);
}
std::time::Duration::ZERO
});
Self {
micros: i64::try_from(dur.as_micros()).unwrap_or(i64::MAX),
}
Expand Down
6 changes: 2 additions & 4 deletions nodedb-types/src/hlc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,6 @@

use std::cmp::Ordering;
use std::sync::Mutex;
use std::time::{SystemTime, UNIX_EPOCH};

use serde::{Deserialize, Serialize};

Expand Down Expand Up @@ -189,10 +188,9 @@ impl HlcClock {
}

fn wall_now_ns() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
crate::clock::since_epoch()
.map(|d| d.as_nanos() as u64)
.unwrap_or_else(|_| {
.unwrap_or_else(|| {
use std::sync::atomic::{AtomicBool, Ordering};
static LOGGED: AtomicBool = AtomicBool::new(false);
if !LOGGED.swap(true, Ordering::Relaxed) {
Expand Down
1 change: 1 addition & 0 deletions nodedb-types/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
//! crate.

pub mod audit_dml;
pub mod clock;
pub mod config;
pub mod quota;

Expand Down
Loading
Loading