mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-02 10:05:41 +00:00
chore(toolchain): switch to stable Rust 1.96.1 and remove all nightly feature gates (#9303)
* chore(toolchain): switch to stable Rust 1.96.1 and remove all nightly feature gates Move the workspace from the pinned nightly-2026-03-21 to stable 1.96.1 and drop all 23 '#![feature]' gates across 13 crates, rewriting the still-unstable API usages with stable equivalents: - try_blocks: closures / an async block (table, query, common-function, servers) - duration_constructors: Duration::from_secs(n * 86400) / (n * 60) - iterator_try_collect: collect::<Result<Vec<_>, _>>() - box_patterns: as_deref() + matches! chains (sql) - error_iter: error_chain_root() source-chain walker (common-error); sources() includes the error itself, so the walker never panics - int_roundings: div_floor -> div_euclid (equal for positive divisors) - iter_partition_in_place: stable sort_by_key partition helper (index) - hash_set_entry: HashSet::insert bool / contains+insert - trait_alias: lifetime-parameterized dyn FnOnce type aliases (puffin) - string_from_utf8_lossy_owned: from_utf8_lossy(&v).into_owned() - never_type: Infallible (common-recordbatch) - debug_closure_helpers: closure-backed DebugFmt newtype (mito2) - binary_heap_pop_if: peek().is_some_and() + pop() - exclusive_wrapper: drop Exclusive; C: Send + Unpin already in bounds - stmt_expr_attributes: stale gate, no usages Also fix release-dev-builder-images.yaml, which parsed rust-toolchain.toml with a date-only regex and would produce empty image versions with a stable channel; it now extracts the full channel token. Dev-builder images verified against stable 1.96.1 (image build, default-toolchain behavior, binstall/nextest, riscv64 and android targets, in-image cargo check). Validated on 1.96.1: cargo check --workspace --all-targets, clippy --workspace --all-targets --all-features -D warnings, cargo fmt --check, and nextest on all 13 affected crates (4586 passed). Part of #9289. Depends on #9298 (fuzz nightly quarantine) merging first. Signed-off-by: Ning Sun <sunning@greptime.com> * chore: update flake checksum * chore: use wild for linker in flake --------- Signed-off-by: Ning Sun <sunning@greptime.com>
This commit is contained in:
@@ -59,7 +59,10 @@ jobs:
|
||||
commitShortSHA=`echo ${{ github.sha }} | cut -c1-8`
|
||||
buildTime=`date +%Y%m%d%H%M%S`
|
||||
BUILD_VERSION="$commitShortSHA-$buildTime"
|
||||
RUST_TOOLCHAIN_VERSION=$(cat rust-toolchain.toml | grep -Eo '[0-9]{4}-[0-9]{2}-[0-9]{2}')
|
||||
# Extract the full channel string (e.g. "1.96.1" or
|
||||
# "nightly-2026-03-21") so the image version always carries the
|
||||
# toolchain it was built with.
|
||||
RUST_TOOLCHAIN_VERSION=$(grep -E '^channel' rust-toolchain.toml | cut -d'"' -f2)
|
||||
IMAGE_VERSION="${RUST_TOOLCHAIN_VERSION}-${BUILD_VERSION}"
|
||||
echo "VERSION=${IMAGE_VERSION}" >> $GITHUB_ENV
|
||||
echo "version=$IMAGE_VERSION" >> $GITHUB_OUTPUT
|
||||
|
||||
@@ -22,7 +22,7 @@
|
||||
lib = nixpkgs.lib;
|
||||
rustToolchain = fenix.packages.${system}.fromToolchainName {
|
||||
name = (lib.importTOML ./rust-toolchain.toml).toolchain.channel;
|
||||
sha256 = "sha256-rboGKQLH4eDuiY01SINOqmXUFUNr9F4awoFZGzib17o=";
|
||||
sha256 = "sha256-h+t2xTBz5yt2YIO+1VMIIGlCU7gyp2LYOFvaV1nwOXU=";
|
||||
};
|
||||
in
|
||||
{
|
||||
@@ -34,7 +34,7 @@
|
||||
gcc
|
||||
protobuf
|
||||
gnumake
|
||||
mold
|
||||
wild
|
||||
(rustToolchain.withComponents [
|
||||
"cargo"
|
||||
"clippy"
|
||||
|
||||
+1
-1
@@ -1,2 +1,2 @@
|
||||
[toolchain]
|
||||
channel = "nightly-2026-03-21"
|
||||
channel = "1.96.1"
|
||||
|
||||
@@ -21,6 +21,20 @@ use std::sync::Arc;
|
||||
use serde::{Deserialize, Deserializer, Serializer};
|
||||
use snafu::{FromString, Snafu};
|
||||
|
||||
/// Returns the root cause of an error's source chain, i.e. the last error
|
||||
/// reachable via [`std::error::Error::source`]. For an error without any
|
||||
/// source the error itself is returned.
|
||||
///
|
||||
/// Mirrors `err.sources().last().unwrap()` (unstable `error_iter`, which
|
||||
/// yields the error itself followed by its sources).
|
||||
fn error_chain_root(err: &dyn std::error::Error) -> &dyn std::error::Error {
|
||||
let mut root = err;
|
||||
while let Some(source) = root.source() {
|
||||
root = source;
|
||||
}
|
||||
root
|
||||
}
|
||||
|
||||
use crate::status_code::StatusCode;
|
||||
|
||||
/// Describes whether an error instance is safe and useful to retry.
|
||||
@@ -148,7 +162,7 @@ pub trait ErrorExt: StackError {
|
||||
_ => {
|
||||
let error = self.last();
|
||||
if let Some(external_error) = error.source() {
|
||||
let external_root = external_error.sources().last().unwrap();
|
||||
let external_root = error_chain_root(external_error);
|
||||
|
||||
if error.transparent() {
|
||||
format!("{external_root}")
|
||||
@@ -169,7 +183,7 @@ pub trait ErrorExt: StackError {
|
||||
{
|
||||
let error = self.last();
|
||||
if let Some(external_error) = error.source() {
|
||||
let external_root = external_error.sources().last().unwrap();
|
||||
let external_root = error_chain_root(external_error);
|
||||
Some(external_root)
|
||||
} else {
|
||||
None
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(error_iter)]
|
||||
|
||||
pub mod ext;
|
||||
pub mod mock;
|
||||
pub mod status_code;
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(duration_constructors)]
|
||||
|
||||
pub mod context;
|
||||
pub mod error;
|
||||
pub mod event_table;
|
||||
|
||||
@@ -108,9 +108,9 @@ fn event_type_filter_is_all(event_types: &EventTypeFilterRef) -> bool {
|
||||
/// The time interval for flushing batched events to the event handler.
|
||||
pub const DEFAULT_FLUSH_INTERVAL_SECONDS: Duration = Duration::from_secs(5);
|
||||
/// The default TTL(90 days) for the events table.
|
||||
const DEFAULT_EVENTS_TABLE_TTL: Duration = Duration::from_days(90);
|
||||
const DEFAULT_EVENTS_TABLE_TTL: Duration = Duration::from_secs(90 * 86400);
|
||||
/// The default compaction time window for the events table.
|
||||
pub const DEFAULT_COMPACTION_TIME_WINDOW: Duration = Duration::from_days(1);
|
||||
pub const DEFAULT_COMPACTION_TIME_WINDOW: Duration = Duration::from_secs(86400);
|
||||
// The capacity of the tokio channel for transmitting events to background processor.
|
||||
const DEFAULT_CHANNEL_SIZE: usize = 2048;
|
||||
// The size of the buffer for batching events before flushing to event handler.
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(try_blocks)]
|
||||
|
||||
mod admin;
|
||||
mod flush_flow;
|
||||
mod macros;
|
||||
|
||||
@@ -1119,11 +1119,12 @@ mod test {
|
||||
];
|
||||
|
||||
for (query, expected) in cases {
|
||||
let result: Result<()> = try {
|
||||
let result: Result<()> = (|| {
|
||||
let parser = ParserContext { stack: vec![] };
|
||||
let ast = parser.parse_pattern(query)?;
|
||||
let _ast = ast.transform_ast()?;
|
||||
};
|
||||
Ok(())
|
||||
})();
|
||||
|
||||
assert!(result.is_err(), "{query}");
|
||||
let actual_error = result.unwrap_err().to_string();
|
||||
|
||||
@@ -570,7 +570,7 @@ impl MetricCollector {
|
||||
}
|
||||
|
||||
impl ExecutionPlanVisitor for MetricCollector {
|
||||
type Error = !;
|
||||
type Error = std::convert::Infallible;
|
||||
|
||||
fn pre_visit(&mut self, plan: &dyn ExecutionPlan) -> std::result::Result<bool, Self::Error> {
|
||||
// skip if no metric available
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(never_type)]
|
||||
|
||||
pub mod adapter;
|
||||
pub mod cursor;
|
||||
pub mod error;
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(duration_constructors)]
|
||||
|
||||
pub mod logging;
|
||||
mod macros;
|
||||
pub mod metric;
|
||||
|
||||
@@ -348,7 +348,7 @@ impl Default for SlowQueryOptions {
|
||||
record_type: SlowQueriesRecordType::SystemTable,
|
||||
threshold: Duration::from_secs(30),
|
||||
sample_ratio: 1.0,
|
||||
ttl: Duration::from_days(90),
|
||||
ttl: Duration::from_secs(90 * 86400),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -234,7 +234,7 @@ mod tests {
|
||||
create_topic_timeout: Duration::from_secs(30),
|
||||
},
|
||||
auto_create_topics: true,
|
||||
auto_prune_interval: Duration::from_mins(30),
|
||||
auto_prune_interval: Duration::from_secs(30 * 60),
|
||||
auto_prune_logical_delete: false,
|
||||
auto_prune_parallelism: 10,
|
||||
flush_trigger_size: ReadableSize::mb(512),
|
||||
|
||||
@@ -38,7 +38,7 @@ pub const DEFAULT_BACKOFF_CONFIG: BackoffConfig = BackoffConfig {
|
||||
};
|
||||
|
||||
/// Default interval for auto WAL pruning.
|
||||
pub const DEFAULT_AUTO_PRUNE_INTERVAL: Duration = Duration::from_mins(30);
|
||||
pub const DEFAULT_AUTO_PRUNE_INTERVAL: Duration = Duration::from_secs(30 * 60);
|
||||
/// Default mode for auto WAL pruning.
|
||||
pub const DEFAULT_AUTO_PRUNE_LOGICAL_DELETE: bool = false;
|
||||
/// Default limit for concurrent auto pruning tasks.
|
||||
|
||||
@@ -161,10 +161,10 @@ impl BloomFilterCreator {
|
||||
|
||||
if let Some(elem) = elem {
|
||||
let old_len = self.cur_seg_distinct_elems.len();
|
||||
// A borrowed entry lookup avoids hashing new values twice and only
|
||||
// allocates when the value is absent from the current segment.
|
||||
self.cur_seg_distinct_elems
|
||||
.get_or_insert_with(elem, <[u8]>::to_vec);
|
||||
// Only allocate when the value is absent from the current segment.
|
||||
if !self.cur_seg_distinct_elems.contains(elem) {
|
||||
self.cur_seg_distinct_elems.insert(elem.to_vec());
|
||||
}
|
||||
if self.cur_seg_distinct_elems.len() != old_len {
|
||||
self.cur_seg_distinct_elems_mem_usage += elem.len();
|
||||
self.global_memory_usage
|
||||
|
||||
@@ -16,3 +16,13 @@ pub mod fst_apply;
|
||||
pub mod fst_values_mapper;
|
||||
pub mod index_apply;
|
||||
pub mod predicate;
|
||||
|
||||
/// Partitions `slice` in place so that elements matching `pred` come first,
|
||||
/// preserving the original relative order within each group. Returns the
|
||||
/// number of matching elements.
|
||||
///
|
||||
/// Stable replacement for the unstable `Iterator::partition_in_place`.
|
||||
pub(crate) fn partition_in_place<T>(slice: &mut [T], mut pred: impl FnMut(&T) -> bool) -> usize {
|
||||
slice.sort_by_key(|x| !pred(x));
|
||||
slice.iter().filter(|x| pred(x)).count()
|
||||
}
|
||||
|
||||
@@ -68,16 +68,16 @@ impl KeysFstApplier {
|
||||
}
|
||||
|
||||
fn split_at_in_lists(predicates: &mut [Predicate]) -> (&mut [Predicate], &mut [Predicate]) {
|
||||
let in_list_index = predicates
|
||||
.iter_mut()
|
||||
.partition_in_place(|p| matches!(p, Predicate::InList(_)));
|
||||
let in_list_index = crate::inverted_index::search::partition_in_place(predicates, |p| {
|
||||
matches!(p, Predicate::InList(_))
|
||||
});
|
||||
predicates.split_at_mut(in_list_index)
|
||||
}
|
||||
|
||||
fn split_at_ranges(predicates: &mut [Predicate]) -> (&mut [Predicate], &mut [Predicate]) {
|
||||
let range_index = predicates
|
||||
.iter_mut()
|
||||
.partition_in_place(|p| matches!(p, Predicate::Range(_)));
|
||||
let range_index = crate::inverted_index::search::partition_in_place(predicates, |p| {
|
||||
matches!(p, Predicate::Range(_))
|
||||
});
|
||||
predicates.split_at_mut(range_index)
|
||||
}
|
||||
|
||||
|
||||
@@ -127,9 +127,10 @@ impl PredicatesIndexApplier {
|
||||
let mut fst_appliers = Vec::with_capacity(predicates.len());
|
||||
|
||||
// InList predicates are applied first to benefit from higher selectivity.
|
||||
let in_list_index = predicates
|
||||
.iter_mut()
|
||||
.partition_in_place(|(_, ps)| ps.iter().any(|p| matches!(p, Predicate::InList(_))));
|
||||
let in_list_index =
|
||||
crate::inverted_index::search::partition_in_place(&mut predicates, |(_, ps)| {
|
||||
ps.iter().any(|p| matches!(p, Predicate::InList(_)))
|
||||
});
|
||||
let mut iter = predicates.into_iter();
|
||||
for _ in 0..in_list_index {
|
||||
let (column_name, predicates) = iter.next().unwrap();
|
||||
|
||||
@@ -12,9 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(iter_partition_in_place)]
|
||||
#![feature(hash_set_entry)]
|
||||
|
||||
pub mod bitmap;
|
||||
pub mod bloom_filter;
|
||||
pub mod error;
|
||||
|
||||
@@ -38,7 +38,7 @@ impl Default for SoftDropGcOptions {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
enable: false,
|
||||
retention: Duration::from_days(7),
|
||||
retention: Duration::from_secs(7 * 86400),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -242,7 +242,7 @@ mod tests {
|
||||
|
||||
assert!(!options.experimental_soft_drop.enable);
|
||||
assert_eq!(
|
||||
Duration::from_days(7),
|
||||
Duration::from_secs(7 * 86400),
|
||||
options.experimental_soft_drop.retention
|
||||
);
|
||||
}
|
||||
@@ -254,7 +254,7 @@ mod tests {
|
||||
enable: true,
|
||||
experimental_soft_drop: SoftDropGcOptions {
|
||||
enable: true,
|
||||
retention: Duration::from_days(1),
|
||||
retention: Duration::from_secs(86400),
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
@@ -269,7 +269,7 @@ mod tests {
|
||||
enable: true,
|
||||
experimental_soft_drop: SoftDropGcOptions {
|
||||
enable: true,
|
||||
retention: Duration::from_days(1),
|
||||
retention: Duration::from_secs(86400),
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
@@ -284,7 +284,7 @@ mod tests {
|
||||
let options = GcSchedulerOptions {
|
||||
experimental_soft_drop: SoftDropGcOptions {
|
||||
enable: true,
|
||||
retention: Duration::from_days(1),
|
||||
retention: Duration::from_secs(86400),
|
||||
},
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
@@ -154,9 +154,9 @@ fn align_ts(ts: i64, interval: Duration) -> i64 {
|
||||
impl PersistStatsHandler {
|
||||
/// Creates a new [`PersistStatsHandler`].
|
||||
pub fn new(inserter: Box<dyn Inserter>, mut persist_interval: Duration) -> Self {
|
||||
if persist_interval < Duration::from_mins(10) {
|
||||
if persist_interval < Duration::from_secs(10 * 60) {
|
||||
warn!("persist_interval is less than 10 minutes, set to 10 minutes");
|
||||
persist_interval = Duration::from_mins(10);
|
||||
persist_interval = Duration::from_secs(10 * 60);
|
||||
}
|
||||
|
||||
Self {
|
||||
|
||||
@@ -12,9 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(hash_set_entry)]
|
||||
#![feature(duration_constructors)]
|
||||
|
||||
pub mod bootstrap;
|
||||
pub mod cache_invalidator;
|
||||
pub mod cluster;
|
||||
|
||||
@@ -122,7 +122,7 @@ impl Default for StatsPersistenceOptions {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
ttl: Duration::ZERO,
|
||||
interval: Duration::from_mins(10),
|
||||
interval: Duration::from_secs(10 * 60),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -92,7 +92,7 @@ use crate::utils::database::DatabaseOperator;
|
||||
use crate::utils::insert_forwarder::InsertForwarder;
|
||||
|
||||
/// The time window for twcs compaction of the region stats table.
|
||||
const REGION_STATS_TABLE_TWCS_COMPACTION_TIME_WINDOW: Duration = Duration::from_days(1);
|
||||
const REGION_STATS_TABLE_TWCS_COMPACTION_TIME_WINDOW: Duration = Duration::from_secs(86400);
|
||||
|
||||
// TODO(fys): try use derive_builder macro
|
||||
pub struct MetasrvBuilder {
|
||||
|
||||
@@ -107,7 +107,7 @@ pub struct PersistentContext {
|
||||
}
|
||||
|
||||
fn default_timeout() -> Duration {
|
||||
Duration::from_mins(2)
|
||||
Duration::from_secs(2 * 60)
|
||||
}
|
||||
|
||||
impl PersistentContext {
|
||||
|
||||
@@ -13,7 +13,6 @@
|
||||
// limitations under the License.
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::collections::hash_set::Entry;
|
||||
use std::fmt::{Debug, Formatter};
|
||||
use std::sync::{Arc, RwLock};
|
||||
|
||||
@@ -46,15 +45,13 @@ impl WalPruneProcedureTracker {
|
||||
/// consume acquire a semaphore permit for the given topic name.
|
||||
pub fn insert_running_procedure(&self, topic_name: String) -> Option<WalPruneProcedureGuard> {
|
||||
let mut running_procedures = self.running_procedures.write().unwrap();
|
||||
match running_procedures.entry(topic_name.clone()) {
|
||||
Entry::Occupied(_) => None,
|
||||
Entry::Vacant(entry) => {
|
||||
entry.insert();
|
||||
Some(WalPruneProcedureGuard {
|
||||
topic_name,
|
||||
running_procedures: self.running_procedures.clone(),
|
||||
})
|
||||
}
|
||||
if running_procedures.insert(topic_name.clone()) {
|
||||
Some(WalPruneProcedureGuard {
|
||||
topic_name,
|
||||
running_procedures: self.running_procedures.clone(),
|
||||
})
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -475,7 +475,11 @@ where
|
||||
for item in items {
|
||||
let (start, _) = item.range();
|
||||
|
||||
while let Some(run) = runs_sort_by_end.pop_if(|x| x.run.end.unwrap() <= start) {
|
||||
while runs_sort_by_end
|
||||
.peek()
|
||||
.is_some_and(|x| x.run.end.unwrap() <= start)
|
||||
{
|
||||
let run = runs_sort_by_end.pop().unwrap();
|
||||
runs_sort_by_index.push(Wrapper(run));
|
||||
}
|
||||
|
||||
|
||||
@@ -16,10 +16,6 @@
|
||||
//!
|
||||
//! Mito is the a region engine to store timeseries data.
|
||||
|
||||
#![feature(debug_closure_helpers)]
|
||||
#![feature(duration_constructors)]
|
||||
#![feature(binary_heap_pop_if)]
|
||||
|
||||
#[cfg(any(test, feature = "test"))]
|
||||
#[cfg_attr(feature = "test", allow(unused))]
|
||||
pub mod test_util;
|
||||
|
||||
@@ -43,7 +43,7 @@ use crate::memtable::{KeyValues, MemtableBuilderRef, MemtableId, MemtableRef};
|
||||
use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
|
||||
|
||||
/// Initial time window if not specified.
|
||||
const INITIAL_TIME_WINDOW: Duration = Duration::from_days(1);
|
||||
const INITIAL_TIME_WINDOW: Duration = Duration::from_secs(86400);
|
||||
|
||||
/// A partition holds rows with timestamps between `[min, max)`.
|
||||
#[derive(Debug, Clone)]
|
||||
|
||||
+33
-17
@@ -298,20 +298,33 @@ fn is_false(value: &bool) -> bool {
|
||||
!*value
|
||||
}
|
||||
|
||||
/// Formats a debug field with a custom closure, as a stable replacement for
|
||||
/// the unstable `DebugStruct::field_with`.
|
||||
struct DebugFmt<F: Fn(&mut Formatter<'_>) -> fmt::Result>(F);
|
||||
|
||||
impl<F: Fn(&mut Formatter<'_>) -> fmt::Result> std::fmt::Debug for DebugFmt<F> {
|
||||
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
|
||||
(self.0)(f)
|
||||
}
|
||||
}
|
||||
|
||||
impl Debug for FileMeta {
|
||||
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
|
||||
let mut debug_struct = f.debug_struct("FileMeta");
|
||||
debug_struct
|
||||
.field("region_id", &self.region_id)
|
||||
.field_with("file_id", |f| write!(f, "{} ", self.file_id))
|
||||
.field_with("time_range", |f| {
|
||||
write!(
|
||||
f,
|
||||
"({}, {}) ",
|
||||
self.time_range.0.to_iso8601_string(),
|
||||
self.time_range.1.to_iso8601_string()
|
||||
)
|
||||
})
|
||||
.field("file_id", &DebugFmt(|f| write!(f, "{} ", self.file_id)))
|
||||
.field(
|
||||
"time_range",
|
||||
&DebugFmt(|f| {
|
||||
write!(
|
||||
f,
|
||||
"({}, {}) ",
|
||||
self.time_range.0.to_iso8601_string(),
|
||||
self.time_range.1.to_iso8601_string()
|
||||
)
|
||||
}),
|
||||
)
|
||||
.field("level", &self.level)
|
||||
.field("file_size", &ReadableSize(self.file_size))
|
||||
.field(
|
||||
@@ -327,14 +340,17 @@ impl Debug for FileMeta {
|
||||
debug_struct
|
||||
.field("num_rows", &self.num_rows)
|
||||
.field("num_row_groups", &self.num_row_groups)
|
||||
.field_with("sequence", |f| match self.sequence {
|
||||
None => {
|
||||
write!(f, "None")
|
||||
}
|
||||
Some(seq) => {
|
||||
write!(f, "{}", seq)
|
||||
}
|
||||
})
|
||||
.field(
|
||||
"sequence",
|
||||
&DebugFmt(|f| match self.sequence {
|
||||
None => {
|
||||
write!(f, "None")
|
||||
}
|
||||
Some(seq) => {
|
||||
write!(f, "{}", seq)
|
||||
}
|
||||
}),
|
||||
)
|
||||
.field("partition_expr", &self.partition_expr)
|
||||
.field("num_series", &self.num_series);
|
||||
if self.primary_key_min.is_some() || self.primary_key_max.is_some() {
|
||||
|
||||
@@ -741,9 +741,9 @@ fn resolve_value(
|
||||
&ConcreteDataType::string_datatype(),
|
||||
schema_info,
|
||||
)?;
|
||||
Some(ValueData::StringValue(String::from_utf8_lossy_owned(
|
||||
v.to_vec(),
|
||||
)))
|
||||
Some(ValueData::StringValue(
|
||||
String::from_utf8_lossy(&v).into_owned(),
|
||||
))
|
||||
}
|
||||
|
||||
VrlValue::Regex(v) => {
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(string_from_utf8_lossy_owned)]
|
||||
|
||||
mod dispatcher;
|
||||
pub mod error;
|
||||
mod etl;
|
||||
|
||||
@@ -56,7 +56,7 @@ impl TableSuffixTemplate {
|
||||
let v = val.get(key.as_str())?;
|
||||
match v {
|
||||
VrlValue::Integer(v) => Some(v.to_string()),
|
||||
VrlValue::Bytes(v) => Some(String::from_utf8_lossy_owned(v.to_vec())),
|
||||
VrlValue::Bytes(v) => Some(String::from_utf8_lossy(v).into_owned()),
|
||||
_ => None,
|
||||
}
|
||||
})
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(trait_alias)]
|
||||
|
||||
pub mod blob_metadata;
|
||||
pub mod error;
|
||||
pub mod file_format;
|
||||
|
||||
@@ -43,13 +43,13 @@ pub type DirWriterProviderRef = Box<dyn DirWriterProvider + Send>;
|
||||
///
|
||||
/// `Stager` will provide a `BoxWriter` that the caller of `get_blob`
|
||||
/// can use to write the blob into the staging area.
|
||||
pub trait InitBlobFn = FnOnce(BoxWriter) -> WriteResult;
|
||||
pub type InitBlobFn<'a> = dyn FnOnce(BoxWriter) -> WriteResult + Send + Sync + 'a;
|
||||
|
||||
/// Function that initializes a directory.
|
||||
///
|
||||
/// `Stager` will provide a `DirWriterProvider` that the caller of `get_dir`
|
||||
/// can use to write files inside the directory into the staging area.
|
||||
pub trait InitDirFn = FnOnce(DirWriterProviderRef) -> WriteResult;
|
||||
pub type InitDirFn<'a> = dyn FnOnce(DirWriterProviderRef) -> WriteResult + Send + Sync + 'a;
|
||||
|
||||
/// `Stager` manages the staging area for the puffin files.
|
||||
#[async_trait]
|
||||
@@ -67,7 +67,7 @@ pub trait Stager: Send + Sync {
|
||||
&self,
|
||||
handle: &Self::FileHandle,
|
||||
key: &str,
|
||||
init_factory: Box<dyn InitBlobFn + Send + Sync + 'a>,
|
||||
init_factory: Box<InitBlobFn<'a>>,
|
||||
) -> Result<Self::Blob>;
|
||||
|
||||
/// Retrieves a directory, initializing it if necessary using the provided `init_fn`.
|
||||
@@ -79,7 +79,7 @@ pub trait Stager: Send + Sync {
|
||||
&self,
|
||||
handle: &Self::FileHandle,
|
||||
key: &str,
|
||||
init_fn: Box<dyn InitDirFn + Send + Sync + 'a>,
|
||||
init_fn: Box<InitDirFn<'a>>,
|
||||
) -> Result<(Self::Dir, DirMetrics)>;
|
||||
|
||||
/// Stores a directory in the staging area.
|
||||
|
||||
@@ -146,7 +146,7 @@ impl<H: ToString + Clone + Send + Sync> Stager for BoundedStager<H> {
|
||||
&self,
|
||||
handle: &Self::FileHandle,
|
||||
key: &str,
|
||||
init_fn: Box<dyn InitBlobFn + Send + Sync + 'a>,
|
||||
init_fn: Box<InitBlobFn<'a>>,
|
||||
) -> Result<Self::Blob> {
|
||||
let handle_str = handle.to_string();
|
||||
let cache_key = Self::encode_cache_key(&handle_str, key);
|
||||
@@ -202,7 +202,7 @@ impl<H: ToString + Clone + Send + Sync> Stager for BoundedStager<H> {
|
||||
&self,
|
||||
handle: &Self::FileHandle,
|
||||
key: &str,
|
||||
init_fn: Box<dyn InitDirFn + Send + Sync + 'a>,
|
||||
init_fn: Box<InitDirFn<'a>>,
|
||||
) -> Result<(Self::Dir, DirMetrics)> {
|
||||
let handle_str = handle.to_string();
|
||||
|
||||
@@ -332,10 +332,7 @@ impl<H> BoundedStager<H> {
|
||||
BASE64_URL_SAFE.encode(hash)
|
||||
}
|
||||
|
||||
async fn write_blob(
|
||||
target_path: &PathBuf,
|
||||
init_fn: Box<dyn InitBlobFn + Send + Sync + '_>,
|
||||
) -> Result<u64> {
|
||||
async fn write_blob(target_path: &PathBuf, init_fn: Box<InitBlobFn<'_>>) -> Result<u64> {
|
||||
// To guarantee the atomicity of writing the file, we need to write
|
||||
// the file to a temporary file first...
|
||||
let tmp_path = target_path.with_extension(TMP_EXTENSION);
|
||||
@@ -354,10 +351,7 @@ impl<H> BoundedStager<H> {
|
||||
Ok(size)
|
||||
}
|
||||
|
||||
async fn write_dir(
|
||||
target_path: &PathBuf,
|
||||
init_fn: Box<dyn InitDirFn + Send + Sync + '_>,
|
||||
) -> Result<u64> {
|
||||
async fn write_dir(target_path: &PathBuf, init_fn: Box<InitDirFn<'_>>) -> Result<u64> {
|
||||
// To guarantee the atomicity of writing the directory, we need to write
|
||||
// the directory to a temporary directory first...
|
||||
let tmp_base = target_path.with_extension(TMP_EXTENSION);
|
||||
|
||||
@@ -12,12 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(int_roundings)]
|
||||
#![feature(try_blocks)]
|
||||
#![feature(stmt_expr_attributes)]
|
||||
#![feature(iterator_try_collect)]
|
||||
#![feature(box_patterns)]
|
||||
|
||||
mod analyze;
|
||||
pub mod datafusion;
|
||||
pub mod dist_plan;
|
||||
|
||||
@@ -142,7 +142,7 @@ impl LogQueryPlanner {
|
||||
let exprs = filters
|
||||
.iter()
|
||||
.filter_map(|filter| self.build_filters(filter, schema).transpose())
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
if exprs.is_empty() {
|
||||
Ok(None)
|
||||
} else {
|
||||
@@ -153,7 +153,7 @@ impl LogQueryPlanner {
|
||||
let exprs = filters
|
||||
.iter()
|
||||
.filter_map(|filter| self.build_filters(filter, schema).transpose())
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
if exprs.is_empty() {
|
||||
Ok(None)
|
||||
} else {
|
||||
@@ -191,7 +191,7 @@ impl LogQueryPlanner {
|
||||
self.build_content_filter_with_expr(col_expr.clone(), filter, &df_schema)
|
||||
.transpose()
|
||||
})
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
|
||||
if filter_exprs.is_empty() {
|
||||
return Ok(Some(col_expr.is_true()));
|
||||
@@ -285,7 +285,7 @@ impl LogQueryPlanner {
|
||||
self.build_content_filter_with_expr(col_expr.clone(), filter, schema)
|
||||
.transpose()
|
||||
})
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
|
||||
if exprs.is_empty() {
|
||||
return Ok(None);
|
||||
@@ -326,19 +326,19 @@ impl LogQueryPlanner {
|
||||
let args = args
|
||||
.iter()
|
||||
.map(|expr| self.log_expr_to_df_expr(expr, schema))
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
if let Some(alias) = alias {
|
||||
Ok(aggr_fn.call(args).alias(alias))
|
||||
} else {
|
||||
Ok(aggr_fn.call(args))
|
||||
}
|
||||
})
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
|
||||
let group_exprs = by
|
||||
.iter()
|
||||
.map(|expr| self.log_expr_to_df_expr(expr, schema))
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
|
||||
Ok((aggr_expr, group_exprs))
|
||||
}
|
||||
@@ -380,7 +380,7 @@ impl LogQueryPlanner {
|
||||
let args = args
|
||||
.iter()
|
||||
.map(|expr| self.log_expr_to_df_expr(expr, schema))
|
||||
.try_collect::<Vec<_>>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
let func = self.session_state.scalar_functions().get(name).context(
|
||||
UnknownScalarFunctionSnafu {
|
||||
name: name.to_string(),
|
||||
|
||||
@@ -199,12 +199,12 @@ fn fetch_partition_range(input: Arc<dyn ExecutionPlan>) -> DataFusionResult<Opti
|
||||
Ok(Transformed::no(plan))
|
||||
})?;
|
||||
|
||||
let result = try {
|
||||
ScannerInfo {
|
||||
let result = (|| {
|
||||
Some(ScannerInfo {
|
||||
partition_ranges: partition_ranges?,
|
||||
tag_columns: tag_columns?,
|
||||
}
|
||||
};
|
||||
})
|
||||
})();
|
||||
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
@@ -208,9 +208,11 @@ impl DfLogicalPlanner {
|
||||
);
|
||||
|
||||
// TODO(LFC): Remove this when Datafusion supports **both** the syntax and implementation of "explain with format".
|
||||
if let datafusion::sql::parser::Statement::Statement(
|
||||
box datafusion::sql::sqlparser::ast::Statement::Explain { .. },
|
||||
) = &mut df_stmt
|
||||
if let datafusion::sql::parser::Statement::Statement(stmt) = &mut df_stmt
|
||||
&& matches!(
|
||||
stmt.as_ref(),
|
||||
datafusion::sql::sqlparser::ast::Statement::Explain { .. }
|
||||
)
|
||||
{
|
||||
UnimplementedSnafu {
|
||||
operation: "EXPLAIN with FORMAT using raw datafusion planner",
|
||||
|
||||
@@ -967,7 +967,9 @@ fn produce_align_time(
|
||||
// make modify_map for range_fn[i]
|
||||
for (row, hash) in by_columns_hash.iter().enumerate() {
|
||||
let ts = ts_column.value(row);
|
||||
let ith_slot = (ts - align_to).div_floor(align);
|
||||
let diff = ts - align_to;
|
||||
// `div_euclid` equals `div_floor` for positive divisors (`align`).
|
||||
let ith_slot = diff.div_euclid(align);
|
||||
let mut align_ts = ith_slot * align + align_to;
|
||||
while align_ts <= ts && ts < align_ts + range {
|
||||
modify_map
|
||||
|
||||
@@ -807,7 +807,7 @@ fn find_slice_from_range(
|
||||
.map_err(|e| DataFusionError::External(Box::new(e) as _))
|
||||
})
|
||||
})
|
||||
.try_collect::<_, Vec<_>, _>()?;
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
|
||||
let (min_val, max_val) = (typed_sorted_range[0].clone(), typed_sorted_range[1].clone());
|
||||
|
||||
@@ -2682,7 +2682,10 @@ mod test {
|
||||
let exec_stream = exec.execute(0, Arc::new(TaskContext::default())).unwrap();
|
||||
|
||||
let real_output = exec_stream.collect::<Vec<_>>().await;
|
||||
let real_output: Vec<_> = real_output.into_iter().try_collect().unwrap();
|
||||
let real_output: Vec<_> = real_output
|
||||
.into_iter()
|
||||
.collect::<std::result::Result<Vec<_>, _>>()
|
||||
.unwrap();
|
||||
real_output
|
||||
}
|
||||
}
|
||||
|
||||
@@ -231,7 +231,7 @@ impl PrometheusJsonResponse {
|
||||
) -> Self {
|
||||
// Hold the collector while a streaming result is consumed.
|
||||
let collector = query_id.and_then(get_promql_annotation_collector);
|
||||
let response: Result<Self> = try {
|
||||
let response: Result<Self> = async {
|
||||
let result = result?;
|
||||
let mut resp =
|
||||
match result.data {
|
||||
@@ -259,8 +259,9 @@ impl PrometheusJsonResponse {
|
||||
resp.resp_metrics = re;
|
||||
}
|
||||
|
||||
resp
|
||||
};
|
||||
Ok(resp)
|
||||
}
|
||||
.await;
|
||||
|
||||
let result_type_string = result_type.to_string();
|
||||
|
||||
|
||||
@@ -12,9 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(try_blocks)]
|
||||
#![feature(exclusive_wrapper)]
|
||||
|
||||
use datafusion_expr::LogicalPlan;
|
||||
use sql::statements::statement::Statement;
|
||||
// Re-export for use in add_service! macro
|
||||
|
||||
@@ -13,7 +13,6 @@
|
||||
// limitations under the License.
|
||||
|
||||
use std::fmt::Debug;
|
||||
use std::sync::Exclusive;
|
||||
|
||||
use ::auth::{
|
||||
BEARER_TOKEN_USER, Identity, Password, PgAuthInfo, PgScramSha256Verifier, UserInfoRef,
|
||||
@@ -252,7 +251,7 @@ impl StartupHandler for PostgresServerHandlerInner {
|
||||
auth::save_startup_parameters_to_metadata(client, startup);
|
||||
|
||||
// check if db is valid
|
||||
match resolve_db_info(Exclusive::new(client), self.query_handler.clone()).await? {
|
||||
match resolve_db_info(client, self.query_handler.clone()).await? {
|
||||
DbResolution::Resolved(catalog, schema) => {
|
||||
let metadata = client.metadata_mut();
|
||||
let _ = metadata.insert(super::METADATA_CATALOG.to_owned(), catalog);
|
||||
@@ -595,13 +594,13 @@ enum DbResolution {
|
||||
|
||||
/// A function extracted to resolve lifetime and readability issues:
|
||||
async fn resolve_db_info<C>(
|
||||
client: Exclusive<&mut C>,
|
||||
client: &mut C,
|
||||
query_handler: ServerSqlQueryHandlerRef,
|
||||
) -> PgWireResult<DbResolution>
|
||||
where
|
||||
C: ClientInfo + Unpin + Send,
|
||||
{
|
||||
let db_ref = client.into_inner().metadata().get(super::METADATA_DATABASE);
|
||||
let db_ref = client.metadata().get(super::METADATA_DATABASE);
|
||||
if let Some(db) = db_ref {
|
||||
let (catalog, schema) = parse_catalog_and_schema_from_db_string(db);
|
||||
if query_handler
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(box_patterns)]
|
||||
|
||||
pub mod ast;
|
||||
pub mod dialect;
|
||||
pub mod error;
|
||||
|
||||
@@ -14,7 +14,7 @@
|
||||
|
||||
use serde::Serialize;
|
||||
use sqlparser::ast::{
|
||||
Insert as SpInsert, ObjectName, ObjectNamePart, Parens, Query, SetExpr, Statement, TableObject,
|
||||
Insert as SpInsert, ObjectName, ObjectNamePart, Parens, SetExpr, Statement, TableObject,
|
||||
UnaryOperator, ValueWithSpan, Values,
|
||||
};
|
||||
use sqlparser::parser::ParserError;
|
||||
@@ -70,14 +70,18 @@ impl Insert {
|
||||
/// Extracts the literal insert statement body if possible
|
||||
pub fn values_body(&self) -> Result<Vec<Vec<Value>>> {
|
||||
match &self.inner {
|
||||
Statement::Insert(SpInsert {
|
||||
source:
|
||||
Some(box Query {
|
||||
body: box SetExpr::Values(Values { rows, .. }),
|
||||
..
|
||||
}),
|
||||
..
|
||||
}) => sql_exprs_to_values(rows),
|
||||
Statement::Insert(SpInsert { source, .. }) => {
|
||||
let rows = source
|
||||
.as_deref()
|
||||
.and_then(|query| match query.body.as_ref() {
|
||||
SetExpr::Values(Values { rows, .. }) => Some(rows),
|
||||
_ => None,
|
||||
});
|
||||
match rows {
|
||||
Some(rows) => sql_exprs_to_values(rows),
|
||||
None => unreachable!(),
|
||||
}
|
||||
}
|
||||
_ => unreachable!(),
|
||||
}
|
||||
}
|
||||
@@ -86,36 +90,37 @@ impl Insert {
|
||||
/// The rules is the same as function `values_body()`.
|
||||
pub fn can_extract_values(&self) -> bool {
|
||||
match &self.inner {
|
||||
Statement::Insert(SpInsert {
|
||||
source:
|
||||
Some(box Query {
|
||||
body: box SetExpr::Values(Values { rows, .. }),
|
||||
..
|
||||
}),
|
||||
..
|
||||
}) => rows.iter().all(|es| {
|
||||
es.iter().all(|expr| match expr {
|
||||
Expr::Value(_) => true,
|
||||
Expr::Identifier(ident) => {
|
||||
if ident.quote_style.is_none() {
|
||||
ident.value.to_lowercase() == "default"
|
||||
} else {
|
||||
ident.quote_style == Some('"')
|
||||
}
|
||||
}
|
||||
Expr::UnaryOp { op, expr } => {
|
||||
matches!(op, UnaryOperator::Minus | UnaryOperator::Plus)
|
||||
&& matches!(
|
||||
&**expr,
|
||||
Expr::Value(ValueWithSpan {
|
||||
value: Value::Number(_, _),
|
||||
..
|
||||
})
|
||||
)
|
||||
}
|
||||
_ => false,
|
||||
Statement::Insert(SpInsert { source, .. }) => source
|
||||
.as_deref()
|
||||
.and_then(|query| match query.body.as_ref() {
|
||||
SetExpr::Values(Values { rows, .. }) => Some(rows),
|
||||
_ => None,
|
||||
})
|
||||
}),
|
||||
.is_some_and(|rows| {
|
||||
rows.iter().all(|es| {
|
||||
es.iter().all(|expr| match expr {
|
||||
Expr::Value(_) => true,
|
||||
Expr::Identifier(ident) => {
|
||||
if ident.quote_style.is_none() {
|
||||
ident.value.to_lowercase() == "default"
|
||||
} else {
|
||||
ident.quote_style == Some('"')
|
||||
}
|
||||
}
|
||||
Expr::UnaryOp { op, expr } => {
|
||||
matches!(op, UnaryOperator::Minus | UnaryOperator::Plus)
|
||||
&& matches!(
|
||||
&**expr,
|
||||
Expr::Value(ValueWithSpan {
|
||||
value: Value::Number(_, _),
|
||||
..
|
||||
})
|
||||
)
|
||||
}
|
||||
_ => false,
|
||||
})
|
||||
})
|
||||
}),
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
@@ -124,7 +129,7 @@ impl Insert {
|
||||
pub fn has_non_values_query_source(&self) -> bool {
|
||||
match &self.inner {
|
||||
Statement::Insert(SpInsert {
|
||||
source: Some(box query),
|
||||
source: Some(query),
|
||||
..
|
||||
}) => !matches!(&*query.body, SetExpr::Values(_)),
|
||||
_ => false,
|
||||
@@ -134,9 +139,9 @@ impl Insert {
|
||||
pub fn query_body(&self) -> Result<Option<GtQuery>> {
|
||||
Ok(match &self.inner {
|
||||
Statement::Insert(SpInsert {
|
||||
source: Some(box query),
|
||||
source: Some(query),
|
||||
..
|
||||
}) => Some(query.clone().try_into()?),
|
||||
}) => Some(query.as_ref().clone().try_into()?),
|
||||
_ => None,
|
||||
})
|
||||
}
|
||||
@@ -383,11 +388,9 @@ mod tests {
|
||||
let q = insert.query_body().unwrap().unwrap();
|
||||
assert!(insert.has_non_values_query_source());
|
||||
assert!(matches!(
|
||||
q.inner,
|
||||
Query {
|
||||
body: box SetExpr::Select { .. },
|
||||
..
|
||||
}
|
||||
&q.inner,
|
||||
sqlparser::ast::Query { body, .. }
|
||||
if matches!(body.as_ref(), SetExpr::Select { .. })
|
||||
));
|
||||
}
|
||||
_ => unreachable!(),
|
||||
|
||||
@@ -12,8 +12,6 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
#![feature(try_blocks)]
|
||||
|
||||
pub mod dist_table;
|
||||
pub mod error;
|
||||
pub mod metadata;
|
||||
|
||||
@@ -144,8 +144,8 @@ impl RegionScanExec {
|
||||
))
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let ts_col: Option<PhysicalSortExpr> = try {
|
||||
PhysicalSortExpr::new(
|
||||
let ts_col: Option<PhysicalSortExpr> = (|| {
|
||||
Some(PhysicalSortExpr::new(
|
||||
Arc::new(
|
||||
Column::new_with_schema(
|
||||
&metadata.time_index_column().column_schema.name,
|
||||
@@ -157,8 +157,8 @@ impl RegionScanExec {
|
||||
descending: false,
|
||||
nulls_first: true,
|
||||
},
|
||||
)
|
||||
};
|
||||
))
|
||||
})();
|
||||
|
||||
let eq_props = match request.distribution {
|
||||
Some(TimeSeriesDistribution::PerSeries) => {
|
||||
|
||||
Reference in New Issue
Block a user