use std::path::{Path, PathBuf};
use crate::stored_state::EntryScope;
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Default, Hash)]
pub struct Clock(pub u64);
impl Clock {
pub const ZERO: Clock = Clock(0);
#[inline]
#[must_use]
pub const fn checked_next(self) -> Option<Clock> {
match self.0.checked_add(1) {
Some(value) => Some(Clock(value)),
None => None,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
#[repr(u8)]
pub enum EntryKind {
File = 0,
Dir = 1,
Symlink = 2,
Other = 3,
}
impl EntryKind {
pub const fn from_u8(raw: u8) -> Option<Self> {
match raw {
0 => Some(Self::File),
1 => Some(Self::Dir),
2 => Some(Self::Symlink),
3 => Some(Self::Other),
_ => None,
}
}
#[inline]
pub const fn is_dir(self) -> bool {
matches!(self, Self::Dir)
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default, Hash)]
pub struct Attrs {
pub size: u64,
pub allocated: u64,
pub mtime_ns: i64,
pub ctime_ns: i64,
pub inode: u64,
pub dev: u64,
}
impl Attrs {
#[inline]
pub const fn fingerprint(&self) -> Fingerprint {
Fingerprint {
size: self.size,
mtime_ns: self.mtime_ns,
ctime_ns: self.ctime_ns,
inode: self.inode,
dev: self.dev,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default, Hash)]
pub struct Fingerprint {
pub size: u64,
pub mtime_ns: i64,
pub ctime_ns: i64,
pub inode: u64,
pub dev: u64,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub struct ScanScope {
pub max_depth: Option<usize>,
pub follow_symlinks: bool,
pub one_filesystem: bool,
pub hidden_fingerprint: u64,
pub exclude_special: bool,
pub population: crate::query::IgnoredEntries,
pub ignore_rules_fingerprint: u64,
pub type_rules_fingerprint: u64,
pub reducers_fingerprint: u64,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub struct SemanticIdentity {
pub ignore_rules_fingerprint: u64,
pub type_rules_fingerprint: u64,
pub reducers_fingerprint: u64,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub struct SessionId(pub(crate) u64);
impl SessionId {
pub const fn opaque(self) -> u64 {
self.0
}
pub const fn from_opaque(value: u64) -> Option<Self> {
if value == 0 { None } else { Some(Self(value)) }
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub struct EngineVersion {
pub session: SessionId,
pub sequence: Clock,
pub scope: EntryScope,
pub semantics: SemanticIdentity,
}
impl ScanScope {
pub const fn observes_controls(self) -> bool {
self.ignore_rules_fingerprint != 0
}
pub const fn entry_scope(self) -> EntryScope {
EntryScope {
max_depth: self.max_depth,
follow_symlinks: self.follow_symlinks,
one_filesystem: self.one_filesystem,
hidden_fingerprint: self.hidden_fingerprint,
exclude_special: self.exclude_special,
population: self.population,
control_fingerprint: if matches!(self.population, crate::query::IgnoredEntries::Include)
{
0
} else {
self.ignore_rules_fingerprint
},
}
}
pub const fn semantic_identity(self) -> SemanticIdentity {
SemanticIdentity {
ignore_rules_fingerprint: self.ignore_rules_fingerprint,
type_rules_fingerprint: self.type_rules_fingerprint,
reducers_fingerprint: self.reducers_fingerprint,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Hash, Default)]
pub enum Source {
#[default]
Scanned,
Revalidated,
JournalScoped,
Cached,
}
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Hash, Default)]
#[non_exhaustive]
pub enum Status {
#[default]
Complete,
Partial,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct Provenance {
pub source: Source,
pub observed_at_ns: i64,
pub status: Status,
}
impl Source {
pub const fn is_verified(self) -> bool {
matches!(self, Self::Scanned | Self::Revalidated)
}
}
impl Provenance {
pub const fn scanned(observed_at_ns: i64) -> Self {
Self { source: Source::Scanned, observed_at_ns, status: Status::Complete }
}
#[must_use]
pub fn combine(self, other: Self) -> Self {
Self {
source: self.source.max(other.source),
observed_at_ns: match (self.observed_at_ns, other.observed_at_ns) {
(0, _) | (_, 0) => 0,
(mine, other) => mine.min(other),
},
status: self.status.max(other.status),
}
}
pub const fn is_verified(self) -> bool {
self.source.is_verified()
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum Freshness {
Fresh,
Reconciling,
Stale,
Partial,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum LifecyclePhase {
Discovering,
Reconciling,
Ready,
Watching,
Stopped,
Failed,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum CoverageReason {
Building,
Budget,
Cancelled,
Inaccessible,
Failed,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum Coverage {
Complete,
Partial(CoverageReason),
}
pub const MAX_RETAINED_ISSUES: usize = 64;
pub const MAX_ISSUE_MESSAGE_BYTES: usize = 512;
pub const MAX_ISSUE_PATH_BYTES: usize = 4_096;
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum IssueKind {
Permission,
Disappeared,
InvalidMetadata,
ResourceBudget,
ObservationGap,
ProviderFailure,
}
#[derive(Clone, PartialEq, Eq, Debug, Hash)]
pub struct Issue {
pub kind: IssueKind,
pub path: Option<PathBuf>,
pub message: String,
pub os_error: Option<i32>,
}
impl Issue {
pub fn from_error(error: &Error) -> Self {
match error {
Error::Io { path, source } => Self::from_io(path, source),
other => Self {
kind: IssueKind::ProviderFailure,
path: None,
message: bounded_issue_message(other.to_string()),
os_error: None,
},
}
}
pub(crate) fn from_io(path: &Path, source: &std::io::Error) -> Self {
Self {
kind: match source.kind() {
std::io::ErrorKind::PermissionDenied => IssueKind::Permission,
std::io::ErrorKind::NotFound => IssueKind::Disappeared,
std::io::ErrorKind::InvalidData | std::io::ErrorKind::InvalidInput => {
IssueKind::InvalidMetadata
}
_ => IssueKind::ProviderFailure,
},
path: bounded_issue_path(path),
message: bounded_issue_message(format!("I/O error at {}: {source}", path.display())),
os_error: source.raw_os_error(),
}
}
pub(crate) fn from_error_under(root: &Path, error: &Error) -> Self {
let mut issue = Self::from_error(error);
if let Error::Io { path, .. } = error {
issue.relativize(root, path);
}
issue
}
pub(crate) fn from_io_under(root: &Path, path: &Path, source: &std::io::Error) -> Self {
let mut issue = Self::from_io(path, source);
issue.relativize(root, path);
issue
}
fn relativize(&mut self, root: &Path, path: &Path) {
if let Ok(relative) = path.strip_prefix(root) {
self.path = bounded_issue_path(relative);
}
}
pub(crate) fn resource_budget(max_files: u64) -> Self {
Self {
kind: IssueKind::ResourceBudget,
path: None,
message: format!(
"verified work refused an admissible file after retaining {max_files}"
),
os_error: None,
}
}
pub(crate) fn observation_gap(path: &Path, reason: InvalidateReason) -> Self {
Self {
kind: IssueKind::ObservationGap,
path: bounded_issue_path(path),
message: bounded_issue_message(format!(
"filesystem observation lost precision at {}: {reason:?}",
path.display()
)),
os_error: None,
}
}
pub(crate) fn provider_failure(path: Option<&Path>, message: String) -> Self {
Self {
kind: IssueKind::ProviderFailure,
path: path.and_then(bounded_issue_path),
message: bounded_issue_message(message),
os_error: None,
}
}
}
fn bounded_issue_path(path: &Path) -> Option<PathBuf> {
(path.as_os_str().as_encoded_bytes().len() <= MAX_ISSUE_PATH_BYTES).then(|| path.to_path_buf())
}
fn bounded_issue_message(mut message: String) -> String {
if message.len() <= MAX_ISSUE_MESSAGE_BYTES {
return message;
}
let mut end = MAX_ISSUE_MESSAGE_BYTES;
while !message.is_char_boundary(end) {
end -= 1;
}
message.truncate(end);
message
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default, Hash)]
pub struct IssueSummary {
pub retained: u64,
pub omitted: u64,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default, Hash)]
pub struct DiscoveryProgress {
pub files_retained: u64,
pub directories_complete: u64,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub struct IndexState {
pub phase: LifecyclePhase,
pub coverage: Coverage,
pub freshness: Freshness,
pub source: Source,
pub progress: DiscoveryProgress,
pub issues: IssueSummary,
}
pub const MAX_READ_PROJECTIONS: usize = 16;
pub const MAX_PAGE_ROWS: usize = 4_096;
pub const MAX_PAGE_WORK: u64 = 1_000_000;
pub const MAX_REPORT_VIEWS: usize = 16;
pub const DEFAULT_COUNT_CAP: u64 = 10_000;
pub const MAX_COUNT_CAP: u64 = 1_000_000;
pub const MAX_CONTINUATION_RECORD_BYTES: usize = 64 * 1_024;
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub struct ContinuationId {
pub(crate) session: SessionId,
pub(crate) ordinal: u64,
}
impl ContinuationId {
pub const fn opaque_parts(self) -> (u64, u64) {
(self.session.opaque(), self.ordinal)
}
pub const fn from_opaque_parts(session: u64, ordinal: u64) -> Option<Self> {
match SessionId::from_opaque(session) {
Some(session) if ordinal != 0 => Some(Self { session, ordinal }),
_ => None,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct PageRequest {
pub limit: usize,
pub max_work: u64,
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct PortablePath(String);
impl PortablePath {
pub(crate) fn new(path: String) -> Self {
Self(path)
}
pub fn as_str(&self) -> &str {
&self.0
}
pub fn into_string(self) -> String {
self.0
}
pub(crate) fn retained_heap_bytes(&self) -> usize {
self.0.capacity()
}
}
impl std::borrow::Borrow<str> for PortablePath {
fn borrow(&self) -> &str {
&self.0
}
}
impl std::fmt::Display for PortablePath {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(&self.0)
}
}
impl std::fmt::Debug for PortablePath {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
std::fmt::Debug::fmt(&self.0, formatter)
}
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct EntryValue {
pub path: PathBuf,
pub portable_path: PortablePath,
pub kind: EntryKind,
pub attrs: Attrs,
pub ignored: bool,
pub classification: Option<crate::classify::NameClassification>,
pub rollup: Option<crate::index::PartitionRollUpSummary>,
pub children_complete: Option<bool>,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum Knowledge<T> {
Present(T),
Absent,
Unknown {
reason: CoverageReason,
},
}
#[derive(Clone, Debug)]
pub enum ReadProjection {
Lookup {
path: PathBuf,
},
RollUp {
path: PathBuf,
},
Tree {
path: PathBuf,
depth: crate::query::Bound,
include_ignored: bool,
page: PageRequest,
},
Flat {
selection: crate::query::EntrySelection,
shape: RowShape,
page: PageRequest,
},
Aggregate {
selection: crate::query::EntrySelection,
count_cap: u64,
max_work: u64,
},
Report(ReportRequest),
Continue {
continuation: ContinuationId,
page: PageRequest,
},
Diagnostics,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct ReadDiagnostics {
pub root: PathBuf,
pub scope: ScanScope,
pub entries: u64,
pub issues: Vec<Issue>,
pub controls: crate::control::ControlObservation,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct TreePage {
pub directory: EntryValue,
pub rows: Vec<EntryValue>,
pub next: Option<ContinuationId>,
pub complete: bool,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default, Hash)]
pub enum RowShape {
#[default]
Compact,
Full,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct FlatPage {
pub rows: Vec<EntryValue>,
pub next: Option<ContinuationId>,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum CountResult {
Exact(u64),
AtLeast(u64),
}
#[derive(Clone, Debug)]
pub struct ReportRequest {
pub query: crate::query::Query,
pub now: std::time::SystemTime,
pub max_work: u64,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum LimitedProjection {
Tree,
Flat,
Report,
Aggregate,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct QueryLimit {
pub projection: LimitedProjection,
pub max_work: u64,
pub rows_visited: u64,
}
#[derive(Clone, Debug, Default)]
pub struct ReadRequest {
pub projections: Vec<ReadProjection>,
pub expected: Option<EngineVersion>,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum ProjectionRefusal {
NotADirectory {
path: PathBuf,
},
ContinuationRecordLimit {
attempted: usize,
limit: usize,
},
ContinuationUnavailable,
}
#[derive(Clone, Debug)]
pub enum ProjectionResult {
Lookup(Knowledge<EntryValue>),
RollUp(Knowledge<crate::index::PartitionRollUpSummary>),
Tree(Knowledge<TreePage>),
Flat(FlatPage),
Aggregate(CountResult),
Report(crate::query::Report),
Diagnostics(ReadDiagnostics),
Limit(QueryLimit),
Refused(ProjectionRefusal),
}
#[derive(Clone, Debug)]
pub struct ReadResponse {
pub version: EngineVersion,
pub state: IndexState,
pub results: Vec<ProjectionResult>,
pub work: Work,
pub change_cursor: EngineVersion,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct ChangeRequest {
pub after: EngineVersion,
pub timeout: std::time::Duration,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum ChangeOutcome {
Changes {
commits: Vec<Commit>,
impact: Impact,
},
Idle,
Reset {
impact: Impact,
},
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct ChangePoll {
pub cursor: EngineVersion,
pub version: EngineVersion,
pub state: IndexState,
pub outcome: ChangeOutcome,
pub work: Work,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
#[non_exhaustive]
pub enum RefreshRejection {
OutsideRoot,
BeyondDepth,
NotAdmitted,
UnsafeAncestry,
ResourceBudget,
}
impl RefreshRejection {
pub const fn as_str(self) -> &'static str {
match self {
Self::OutsideRoot => "outside_root",
Self::BeyondDepth => "beyond_depth",
Self::NotAdmitted => "not_admitted",
Self::UnsafeAncestry => "unsafe_ancestry",
Self::ResourceBudget => "resource_budget",
}
}
}
#[derive(Clone, PartialEq, Eq, Debug, Hash)]
pub struct RejectedRefreshPath {
pub path: PathBuf,
pub reason: RefreshRejection,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct RefreshResult {
pub after: EngineVersion,
pub version: EngineVersion,
pub state: IndexState,
pub accepted: Vec<PathBuf>,
pub rejected: Vec<RejectedRefreshPath>,
pub impact: Impact,
pub work: Work,
pub issues: Vec<Issue>,
pub omitted_issues: u64,
}
impl Default for IndexState {
fn default() -> Self {
Self {
phase: LifecyclePhase::Ready,
coverage: Coverage::Complete,
freshness: Freshness::Fresh,
source: Source::Scanned,
progress: DiscoveryProgress::default(),
issues: IssueSummary::default(),
}
}
}
impl Freshness {
pub(crate) const fn rank(self) -> u8 {
match self {
Self::Fresh => 0,
Self::Reconciling => 1,
Self::Stale => 2,
Self::Partial => 3,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum InvalidateReason {
WatchOverflow,
UnpairedRename,
WatchSetupRace,
PeriodicSweep,
VerificationFailed,
UnknownAncestry,
WatchContention,
ControlPopulationChanged,
Requested,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum Op {
Upsert {
path: PathBuf,
kind: EntryKind,
attrs: Attrs,
},
Remove {
path: PathBuf,
},
ControlUpsert {
path: PathBuf,
source: Vec<u8>,
},
ControlRemove {
path: PathBuf,
},
InvalidateSubtree {
path: PathBuf,
reason: InvalidateReason,
},
}
impl Op {
pub fn path(&self) -> &Path {
match self {
Self::Upsert { path, .. }
| Self::Remove { path }
| Self::ControlUpsert { path, .. }
| Self::ControlRemove { path }
| Self::InvalidateSubtree { path, .. } => path,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum PathState {
Absent,
Present {
kind: EntryKind,
attrs: Attrs,
},
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub(crate) struct EntryIdentity {
slot: u32,
generation: u64,
revision: u64,
children_revision: u64,
directory: bool,
}
impl EntryIdentity {
pub(crate) const fn new(
slot: u32,
generation: u64,
revision: u64,
children_revision: u64,
directory: bool,
) -> Self {
Self { slot, generation, revision, children_revision, directory }
}
pub(crate) const fn same_target(self, other: Self, require_structure: bool) -> bool {
self.slot == other.slot
&& self.generation == other.generation
&& self.revision == other.revision
&& (!require_structure || self.children_revision == other.children_revision)
}
pub(crate) const fn same_absence_guard(self, other: Self) -> bool {
self.slot == other.slot
&& self.generation == other.generation
&& self.children_revision == other.children_revision
&& (other.directory || self.revision == other.revision)
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub struct PathExpectation {
pub state: PathState,
entry: Option<EntryIdentity>,
absence_guard: Option<EntryIdentity>,
}
impl PathExpectation {
pub(crate) const fn new(
state: PathState,
entry: Option<EntryIdentity>,
absence_guard: Option<EntryIdentity>,
) -> Self {
Self { state, entry, absence_guard }
}
pub(crate) const fn entry(self) -> Option<EntryIdentity> {
self.entry
}
pub(crate) const fn absence_guard(self) -> Option<EntryIdentity> {
self.absence_guard
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash)]
pub enum Expectation {
Any,
State(PathExpectation),
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct ObservationOp {
pub op: Op,
pub expectation: Expectation,
}
impl ObservationOp {
pub const fn unconditional(op: Op) -> Self {
Self { op, expectation: Expectation::Any }
}
pub const fn if_state(op: Op, expected: PathExpectation) -> Self {
Self { op, expectation: Expectation::State(expected) }
}
}
#[derive(Clone, PartialEq, Eq, Debug, Default)]
pub struct Observation {
pub ops: Vec<ObservationOp>,
}
impl Observation {
pub fn new(ops: Vec<Op>) -> Self {
Self { ops: ops.into_iter().map(ObservationOp::unconditional).collect() }
}
pub const fn from_ops(ops: Vec<ObservationOp>) -> Self {
Self { ops }
}
#[inline]
pub fn is_empty(&self) -> bool {
self.ops.is_empty()
}
#[inline]
pub fn len(&self) -> usize {
self.ops.len()
}
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum EffectiveChange {
Inserted {
path: PathBuf,
kind: EntryKind,
attrs: Attrs,
},
Updated {
path: PathBuf,
kind: EntryKind,
previous: Attrs,
current: Attrs,
},
Removed {
path: PathBuf,
kind: EntryKind,
attrs: Attrs,
},
ControlUpdated {
path: PathBuf,
previous: Option<crate::control::ControlIdentity>,
current: Option<crate::control::ControlIdentity>,
},
ControlRefusalUpdated {
path: PathBuf,
previous: Option<crate::control::ControlRefusalReason>,
current: Option<crate::control::ControlRefusalReason>,
},
Reclassified {
path: PathBuf,
previous_ignored: bool,
current_ignored: bool,
},
Invalidated {
path: PathBuf,
reason: InvalidateReason,
},
}
impl EffectiveChange {
pub fn path(&self) -> &Path {
match self {
Self::Inserted { path, .. }
| Self::Updated { path, .. }
| Self::Removed { path, .. }
| Self::ControlUpdated { path, .. }
| Self::ControlRefusalUpdated { path, .. }
| Self::Reclassified { path, .. }
| Self::Invalidated { path, .. } => path,
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Hash)]
pub enum ImpactDomain {
Topology,
Metadata,
Classification,
Aggregates,
Content,
State,
}
pub const MAX_DIRTY_PATHS: usize = 256;
#[derive(Clone, PartialEq, Eq, Debug, Default)]
pub struct Impact {
pub domains: Vec<ImpactDomain>,
pub dirty_paths: Vec<PathBuf>,
pub all_dirty: bool,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum StateTransition {
Freshness {
path: PathBuf,
previous: Freshness,
current: Freshness,
},
Verified {
path: PathBuf,
},
DirectoryComplete {
path: PathBuf,
},
DirectoryIncomplete {
path: PathBuf,
},
IndexState {
previous: IndexState,
current: IndexState,
},
}
impl StateTransition {
pub fn path(&self) -> &Path {
match self {
Self::Freshness { path, .. }
| Self::Verified { path }
| Self::DirectoryComplete { path }
| Self::DirectoryIncomplete { path } => path,
Self::IndexState { .. } => Path::new(""),
}
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default)]
pub struct Work {
pub observations: u64,
pub unchanged: u64,
pub stale: u64,
pub resource_refused: u64,
pub rows_visited: u64,
pub rows_returned: u64,
pub maintained_index_work: u64,
pub commits_visited: u64,
pub commits_returned: u64,
pub directories_read: u64,
pub entries_visited: u64,
pub files_visited: u64,
pub bytes_visited: u64,
}
const RETAINED_COMMIT_BYTES: usize = 256;
const RETAINED_ITEM_BYTES: usize = 128;
pub const MIN_JOURNAL_CAPACITY_BYTES: usize = 64 * 1024;
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct Commit {
pub clock: Clock,
pub changes: Vec<EffectiveChange>,
pub impact: Impact,
pub state: Vec<StateTransition>,
pub work: Work,
}
impl Commit {
pub fn is_empty(&self) -> bool {
self.changes.is_empty() && self.state.is_empty()
}
pub fn retained_cost(&self) -> usize {
let paths = self
.changes
.iter()
.map(|change| change.path().as_os_str().len())
.chain(self.state.iter().map(|transition| transition.path().as_os_str().len()))
.chain(self.impact.dirty_paths.iter().map(|path| path.as_os_str().len()))
.sum::<usize>();
let items = self.changes.len() + self.state.len() + self.impact.dirty_paths.len();
RETAINED_COMMIT_BYTES + items * RETAINED_ITEM_BYTES + paths
}
}
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error("I/O error at {path}: {source}")]
Io {
path: PathBuf,
#[source]
source: std::io::Error,
},
#[error("path escapes the index root: {0}")]
PathEscapesRoot(PathBuf),
#[error("upsert {path:?} has unknown ancestry; reconcile from {reconcile_from:?}")]
UnknownAncestry {
path: PathBuf,
reconcile_from: PathBuf,
},
#[error("invalid control-file path: {0:?}")]
InvalidControlPath(PathBuf),
#[error("snapshot is not usable: {0}")]
Snapshot(String),
#[error("unsupported scan configuration: {0}")]
UnsupportedScanConfig(&'static str),
#[error(
"this scan did not observe .gitignore control state, so it neither says nor selects \
what is ignored and accepts no control input; scan with read_controls to observe it"
)]
ControlStateNotObserved,
#[error(
"this index's .gitignore limits ({limits}) are not the ones its scan scope was taken \
under; build it with Index::new_with_config from the ScanConfig that made its scope"
)]
ControlLimitsOutsideScope {
limits: crate::control::ControlLimits,
},
#[error(
"journal_capacity_bytes is {requested} bytes, below the {minimum}-byte minimum; it is \
a size in bytes, not a count, so set it to at least {minimum} bytes, or leave it \
unset for the default"
)]
JournalCapacityTooSmall {
requested: usize,
minimum: usize,
},
#[error("scan scope mismatch: index has {indexed:?}, requested {requested:?}")]
ScanScopeMismatch {
indexed: ScanScope,
requested: ScanScope,
},
#[error("subtree {path:?} lies outside scan scope {scope:?}")]
SubtreeOutsideScanScope {
path: PathBuf,
scope: ScanScope,
},
#[error("index lock was poisoned by a panicking writer")]
IndexLockPoisoned,
#[error("the process-local index clock is exhausted")]
ClockExhausted,
#[error(
"counting {path:?} would carry the tree's {counter} total past what a u64 can hold, \
so nothing that includes it was applied or reported"
)]
UnrepresentableTotal {
path: PathBuf,
counter: &'static str,
},
#[error("the process-local opened-index identity space is exhausted")]
OpenedIdentityExhausted,
#[error("the opened index is closed")]
OpenedIndexClosed,
#[error("the opened index is stopped and cannot expand its retained set")]
OpenedIndexStopped,
#[error("priority request contains {attempted} paths; limit is {limit}")]
PriorityPathLimit {
attempted: usize,
limit: usize,
},
#[error("refresh request contains {attempted} paths; limit is {limit}")]
RefreshPathLimit {
attempted: usize,
limit: usize,
},
#[error("read request contains {attempted} projections; limit is {limit}")]
ReadProjectionLimit {
attempted: usize,
limit: usize,
},
#[error("requested opened-index version {requested:?} is unavailable; current is {current:?}")]
VersionUnavailable {
requested: Box<EngineVersion>,
current: Box<EngineVersion>,
},
#[error("page row limit {attempted} is outside 1..={limit}")]
PageRowLimit {
attempted: usize,
limit: usize,
},
#[error("page work limit {attempted} is outside 1..={limit}")]
PageWorkLimit {
attempted: u64,
limit: u64,
},
#[error("tree page depth must be at least one level")]
TreeDepthZero,
#[error(
"flat opened-index pages use fixed portable path order; selection cannot set depth, limit, sort, or reverse"
)]
UnsupportedFlatSelection,
#[error("aggregate count cap {attempted} is outside 1..={limit}")]
CountCapLimit {
attempted: u64,
limit: u64,
},
#[error(transparent)]
InvalidRequest(crate::query::RequestError),
#[error("report request contains {attempted} views or omissions; limit is {limit}")]
ReportViewLimit {
attempted: usize,
limit: usize,
},
#[error(
"the page continuation was not issued by this opened index; continue from a token a \
page of this root returned"
)]
ContinuationUnavailable,
#[error("the opened index continuation identity space is exhausted")]
ContinuationIdentityExhausted,
#[error("the page continuation version {requested:?} is stale; current is {current:?}")]
ContinuationStale {
requested: Box<EngineVersion>,
current: Box<EngineVersion>,
},
#[error("change cursor {requested:?} is unavailable; current is {current:?}")]
ChangeCursorUnavailable {
requested: Box<EngineVersion>,
current: Box<EngineVersion>,
},
#[error("opened-index journal wait state was poisoned by a panic")]
OpenedJournalPoisoned,
#[error("directory completion named an unknown or non-directory path: {0:?}")]
InvalidDirectoryCompletion(PathBuf),
#[error("opened-index lifecycle state was poisoned by a panic")]
OpenedLifecyclePoisoned,
#[error("opened-index worker {worker} panicked")]
OpenedWorkerPanicked {
worker: &'static str,
},
#[error("opened-index worker {worker} failed: {source}")]
OpenedWorkerFailed {
worker: &'static str,
#[source]
source: std::sync::Arc<Error>,
},
#[error("could not start opened-index worker {worker}: {source}")]
OpenedWorkerSpawn {
worker: &'static str,
#[source]
source: std::io::Error,
},
#[error("watch worker stopped before another observation was available")]
WatchStopped,
#[error("watch worker panicked and stopped")]
WatchWorkerPanicked,
#[cfg(all(feature = "watch", test))]
#[error("invalid watch script: {0}")]
WatchScript(String),
#[error("filesystem observation handoff could not establish a complete verified baseline")]
ObservationHandoffIncomplete,
#[cfg(test)]
#[error("prepared commit rejected by {0}")]
CommitRejected(&'static str),
#[error("invalid {kind} {value:?}: {hint}")]
InvalidValue {
kind: &'static str,
value: String,
hint: String,
},
#[error("watch root {watched:?} does not match index root {indexed:?}")]
WatchRootMismatch {
watched: PathBuf,
indexed: PathBuf,
},
}
impl Error {
pub(crate) fn io(path: impl Into<PathBuf>, source: std::io::Error) -> Self {
Self::Io { path: path.into(), source }
}
}
pub type Result<T> = std::result::Result<T, Error>;
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn clock_advances_monotonically() {
let c = Clock::ZERO;
assert_eq!(c.checked_next(), Some(Clock(1)));
assert!(c.checked_next().is_some_and(|next| next > c));
assert_eq!(Clock(u64::MAX).checked_next(), None);
}
#[test]
fn entry_kind_roundtrips_through_its_pinned_value() {
for kind in [EntryKind::File, EntryKind::Dir, EntryKind::Symlink, EntryKind::Other] {
assert_eq!(EntryKind::from_u8(kind as u8), Some(kind));
}
assert_eq!(EntryKind::from_u8(99), None);
}
#[test]
fn fingerprint_ignores_allocated_size_but_tracks_ctime() {
let base = Attrs { size: 10, allocated: 4096, mtime_ns: 5, ctime_ns: 7, inode: 42, dev: 1 };
let repacked = Attrs { allocated: 8192, ..base };
assert_eq!(base.fingerprint(), repacked.fingerprint());
let touched = Attrs { ctime_ns: 9, ..base };
assert_ne!(base.fingerprint(), touched.fingerprint());
}
#[test]
fn scan_scope_separates_admission_from_answer_semantics() {
let base = ScanScope::default();
let changed_admission = ScanScope { exclude_special: !base.exclude_special, ..base };
let changed_semantics =
ScanScope { reducers_fingerprint: base.reducers_fingerprint.wrapping_add(1), ..base };
assert_ne!(base.entry_scope(), changed_admission.entry_scope());
assert_eq!(base.semantic_identity(), changed_admission.semantic_identity());
assert_eq!(base.entry_scope(), changed_semantics.entry_scope());
assert_ne!(base.semantic_identity(), changed_semantics.semantic_identity());
}
}
#[cfg(test)]
mod provenance_tests {
use super::*;
#[test]
fn sources_order_from_most_to_least_trustworthy() {
assert!(Source::Scanned < Source::Revalidated);
assert!(Source::Revalidated < Source::JournalScoped);
assert!(Source::JournalScoped < Source::Cached);
}
#[test]
fn combining_takes_the_least_trustworthy_of_each_fact() {
let verified =
Provenance { source: Source::Scanned, observed_at_ns: 900, status: Status::Complete };
let stale =
Provenance { source: Source::Cached, observed_at_ns: 100, status: Status::Partial };
let combined = verified.combine(stale);
assert_eq!(combined.source, Source::Cached, "weakest source wins");
assert_eq!(combined.observed_at_ns, 100, "oldest observation wins");
assert_eq!(combined.status, Status::Partial, "worst status wins");
assert_eq!(combined, stale.combine(verified), "combination is commutative");
}
#[test]
fn an_unknown_timestamp_makes_the_combination_unknown() {
let known = Provenance::scanned(500);
let unknown = Provenance { observed_at_ns: 0, ..Provenance::scanned(0) };
assert_eq!(known.combine(unknown).observed_at_ns, 0);
assert_eq!(unknown.combine(known).observed_at_ns, 0);
assert_eq!(known.combine(unknown), unknown.combine(known), "combination stays commutative");
}
#[test]
fn only_this_session_counts_as_verified() {
assert!(Provenance::scanned(1).is_verified());
assert!(Provenance { source: Source::Revalidated, ..Provenance::scanned(1) }.is_verified());
assert!(
!Provenance { source: Source::JournalScoped, ..Provenance::scanned(1) }.is_verified()
);
assert!(!Provenance { source: Source::Cached, ..Provenance::scanned(1) }.is_verified());
}
}