From d96b47bd2b520f0944c79713b9add8947e9342a0 Mon Sep 17 00:00:00 2001 From: Anton Borisov Date: Tue, 6 Oct 2026 17:25:57 +0100 Subject: [PATCH] [cpp] Allow starting lookups without blocking --- fluss-rust/bindings/cpp/include/fluss.hpp | 76 +++++ fluss-rust/bindings/cpp/src/lib.rs | 163 ++-------- fluss-rust/bindings/cpp/src/lookup.rs | 292 ++++++++++++++++++ fluss-rust/bindings/cpp/src/table.cpp | 162 +++++++++- .../bindings/cpp/test/test_kv_table.cpp | 143 +++++++++ .../docs/user-guide/cpp/api-reference.md | 31 +- website/docs/apis/cpp/api-reference.md | 31 +- 7 files changed, 750 insertions(+), 148 deletions(-) create mode 100644 fluss-rust/bindings/cpp/src/lookup.rs diff --git a/fluss-rust/bindings/cpp/include/fluss.hpp b/fluss-rust/bindings/cpp/include/fluss.hpp index bc655375c2..d8cfa0182a 100644 --- a/fluss-rust/bindings/cpp/include/fluss.hpp +++ b/fluss-rust/bindings/cpp/include/fluss.hpp @@ -56,6 +56,8 @@ struct BatchScanner; struct UpsertWriter; struct Lookuper; struct PrefixLookuper; +struct PendingLookup; +struct PendingPrefixLookup; struct ScanResultInner; struct GenericRowInner; struct LookupResultInner; @@ -1250,6 +1252,9 @@ class PrefixLookupResult { private: friend class PrefixLookuper; + friend class PendingPrefixLookup; + /// Takes ownership of `inner`, or returns its error. + Result Adopt(ffi::PrefixLookupResultInner* inner); std::shared_ptr data_; }; @@ -1576,12 +1581,15 @@ class LookupResult : public detail::NamedGetters { private: friend class Lookuper; + friend class PendingLookup; size_t Resolve(const std::string& name) const { if (!column_map_) { BuildColumnMap(); } return detail::ResolveColumn(*column_map_, name); } + /// Takes ownership of `inner`, or returns its error. + Result Adopt(ffi::LookupResultInner* inner); void Destroy() noexcept; void BuildColumnMap() const; ffi::LookupResultInner* inner_{nullptr}; @@ -2068,6 +2076,67 @@ class UpsertWriter { std::shared_ptr callback_capacity_; }; +/// A lookup started by Lookuper::Lookup(pk_row, PendingLookup&). Move-only. Destroying it +/// before Wait() returns abandons the lookup. +class PendingLookup { + public: + PendingLookup() noexcept; + ~PendingLookup() noexcept; + + PendingLookup(const PendingLookup&) = delete; + PendingLookup& operator=(const PendingLookup&) = delete; + PendingLookup(PendingLookup&& other) noexcept; + PendingLookup& operator=(PendingLookup&& other) noexcept; + + /// False once Wait() has returned anything but a timeout. + bool Available() const; + + /// Blocks until the row arrives. + Result Wait(LookupResult& out); + + /// Blocks for at most `timeout_ms`. On REQUEST_TIME_OUT the lookup keeps running, and + /// Wait() may be called again. + Result Wait(LookupResult& out, int64_t timeout_ms); + + private: + friend class Lookuper; + explicit PendingLookup(ffi::PendingLookup* pending) noexcept; + // A negative timeout_ms waits without limit. + Result DoWait(LookupResult& out, int64_t timeout_ms); + void Destroy() noexcept; + ffi::PendingLookup* pending_{nullptr}; +}; + +/// A prefix lookup started by PrefixLookuper::PrefixLookup(prefix_row, PendingPrefixLookup&). +/// Behaves like PendingLookup. +class PendingPrefixLookup { + public: + PendingPrefixLookup() noexcept; + ~PendingPrefixLookup() noexcept; + + PendingPrefixLookup(const PendingPrefixLookup&) = delete; + PendingPrefixLookup& operator=(const PendingPrefixLookup&) = delete; + PendingPrefixLookup(PendingPrefixLookup&& other) noexcept; + PendingPrefixLookup& operator=(PendingPrefixLookup&& other) noexcept; + + /// False once Wait() has returned anything but a timeout. + bool Available() const; + + /// Blocks until the rows arrive. + Result Wait(PrefixLookupResult& out); + + /// Blocks for at most `timeout_ms`, like PendingLookup::Wait. + Result Wait(PrefixLookupResult& out, int64_t timeout_ms); + + private: + friend class PrefixLookuper; + explicit PendingPrefixLookup(ffi::PendingPrefixLookup* pending) noexcept; + // A negative timeout_ms waits without limit. + Result DoWait(PrefixLookupResult& out, int64_t timeout_ms); + void Destroy() noexcept; + ffi::PendingPrefixLookup* pending_{nullptr}; +}; + class Lookuper { public: Lookuper() noexcept; @@ -2085,6 +2154,10 @@ class Lookuper { /// go out in the same batches. Result Lookup(const GenericRow& pk_row, LookupResult& out) const; + /// Starts the lookup without blocking; `out` waits for its row. Lookups started this way + /// share batches too, so one thread can keep many in flight. + Result Lookup(const GenericRow& pk_row, PendingLookup& out) const; + private: friend class Table; friend class TableLookup; @@ -2110,6 +2183,9 @@ class PrefixLookuper { /// result arrives. Thread-safe, like Lookuper::Lookup. Result PrefixLookup(const GenericRow& prefix_row, PrefixLookupResult& out) const; + /// Starts the lookup without blocking, like Lookuper::Lookup(pk_row, PendingLookup&). + Result PrefixLookup(const GenericRow& prefix_row, PendingPrefixLookup& out) const; + private: friend class Table; PrefixLookuper(ffi::PrefixLookuper* lookuper) noexcept; diff --git a/fluss-rust/bindings/cpp/src/lib.rs b/fluss-rust/bindings/cpp/src/lib.rs index 69b8f15e92..77f61753b9 100644 --- a/fluss-rust/bindings/cpp/src/lib.rs +++ b/fluss-rust/bindings/cpp/src/lib.rs @@ -15,8 +15,12 @@ // specific language governing permissions and limitations // under the License. +mod lookup; mod types; mod write_callback; +use lookup::{ + PendingLookup, PendingPrefixLookup, delete_pending_lookup, delete_pending_prefix_lookup, +}; use write_callback::ensure_callback_executor; use std::collections::HashMap; @@ -423,6 +427,8 @@ mod ffi { type UpsertWriter; type Lookuper; type PrefixLookuper; + type PendingLookup; + type PendingPrefixLookup; // Opaque types for optimized FFI type ScanResultInner; @@ -703,6 +709,12 @@ mod ffi { // Lookuper unsafe fn delete_lookuper(lookuper: *mut Lookuper); fn lookup(self: &Lookuper, pk_row: &GenericRowInner) -> Box; + fn lookup_async(self: &Lookuper, pk_row: &GenericRowInner) -> FfiPtrResult; + + // PendingLookup. A negative timeout_ms waits without limit. + unsafe fn delete_pending_lookup(pending: *mut PendingLookup); + fn lookup_wait(self: &mut PendingLookup, timeout_ms: i64) -> Box; + fn lookup_is_pending(self: &PendingLookup) -> bool; // LookupResultInner accessors fn lv_has_error(self: &LookupResultInner) -> bool; @@ -763,6 +775,16 @@ mod ffi { self: &PrefixLookuper, prefix_row: &GenericRowInner, ) -> Box; + fn prefix_lookup_async(self: &PrefixLookuper, prefix_row: &GenericRowInner) + -> FfiPtrResult; + + // PendingPrefixLookup. A negative timeout_ms waits without limit. + unsafe fn delete_pending_prefix_lookup(pending: *mut PendingPrefixLookup); + fn prefix_lookup_wait( + self: &mut PendingPrefixLookup, + timeout_ms: i64, + ) -> Box; + fn prefix_lookup_is_pending(self: &PendingPrefixLookup) -> bool; // PrefixLookupResultInner accessors — like LookupResultInner but indexed // by record, since a prefix lookup returns zero-or-more rows. @@ -1026,13 +1048,13 @@ pub struct UpsertWriter { } pub struct Lookuper { - inner: fcore::client::Lookuper, - table_info: fcore::metadata::TableInfo, + inner: Arc, + table_info: Arc, } pub struct PrefixLookuper { - inner: PrefixKeyLookuper, - table_info: TableInfo, + inner: Arc, + table_info: Arc, /// Full-schema indices of the lookup columns, used to compact the input row. lookup_column_indices: Vec, } @@ -2145,8 +2167,8 @@ impl Table { }; let ptr = Box::into_raw(Box::new(Lookuper { - inner: lookuper, - table_info: self.table_info.clone(), + inner: Arc::new(lookuper), + table_info: Arc::new(self.table_info.clone()), })); ok_ptr(ptr as usize) } @@ -2176,8 +2198,8 @@ impl Table { }; let ptr = Box::into_raw(Box::new(PrefixLookuper { - inner: lookuper, - table_info: self.table_info.clone(), + inner: Arc::new(lookuper), + table_info: Arc::new(self.table_info.clone()), lookup_column_indices, })); ok_ptr(ptr as usize) @@ -2373,66 +2395,6 @@ unsafe fn delete_lookuper(lookuper: *mut Lookuper) { } } -impl Lookuper { - fn lookup(&self, pk_row: &GenericRowInner) -> Box { - let schema = self.table_info.get_schema(); - // Compact PK values (set at their full schema positions, e.g. [0, 2]) - // into the dense PK-only row the core KeyEncoder expects. Skips the - // rebuild when the row is already dense and needs no conversion. - let pk_indices = schema.primary_key_indexes(); - let generic_row = - match types::resolve_dense_row_types(&pk_row.row, Some(schema), &pk_indices) { - Ok(r) => r, - Err(e) => { - return Box::new(LookupResultInner::from_error( - CLIENT_ERROR_CODE, - e.to_string(), - )); - } - }; - - let lookup_result = match RUNTIME.block_on(self.inner.lookup(generic_row.as_ref())) { - Ok(r) => r, - Err(e) => { - let ffi_err = err_from_core_error(&e); - return Box::new(LookupResultInner::from_error( - ffi_err.error_code, - ffi_err.error_message, - )); - } - }; - - let columns = self.table_info.get_schema().columns().to_vec(); - match lookup_result.get_single_row() { - Ok(Some(row)) => match types::compacted_row_to_owned(&row, &self.table_info) { - Ok(owned_row) => Box::new(LookupResultInner { - error: None, - found: true, - row: Some(owned_row), - columns, - }), - Err(e) => Box::new(LookupResultInner::from_error( - CLIENT_ERROR_CODE, - e.to_string(), - )), - }, - Ok(None) => Box::new(LookupResultInner { - error: None, - found: false, - row: None, - columns, - }), - Err(e) => { - let ffi_err = err_from_core_error(&e); - Box::new(LookupResultInner::from_error( - ffi_err.error_code, - ffi_err.error_message, - )) - } - } - } -} - // PrefixLookuper implementation unsafe fn delete_prefix_lookuper(lookuper: *mut PrefixLookuper) { if !lookuper.is_null() { @@ -2442,71 +2404,6 @@ unsafe fn delete_prefix_lookuper(lookuper: *mut PrefixLookuper) { } } -impl PrefixLookuper { - fn prefix_lookup(&self, prefix_row: &GenericRowInner) -> Box { - let schema = self.table_info.get_schema(); - // Compact prefix values (set at their full schema positions) into the - // dense, lookup-column-ordered row the core prefix encoder expects. - // Skips the rebuild when the row is already dense and needs no - // conversion. - let generic_row = match types::resolve_dense_row_types( - &prefix_row.row, - Some(schema), - &self.lookup_column_indices, - ) { - Ok(r) => r, - Err(e) => { - return Box::new(PrefixLookupResultInner::from_error( - CLIENT_ERROR_CODE, - e.to_string(), - )); - } - }; - - let lookup_result = match RUNTIME.block_on(self.inner.lookup(generic_row.as_ref())) { - Ok(r) => r, - Err(e) => { - let ffi_err = err_from_core_error(&e); - return Box::new(PrefixLookupResultInner::from_error( - ffi_err.error_code, - ffi_err.error_message, - )); - } - }; - - let lookup_rows = match lookup_result.get_rows() { - Ok(rows) => rows, - Err(e) => { - let ffi_err = err_from_core_error(&e); - return Box::new(PrefixLookupResultInner::from_error( - ffi_err.error_code, - ffi_err.error_message, - )); - } - }; - - let mut rows = Vec::with_capacity(lookup_rows.len()); - for row in &lookup_rows { - match types::compacted_row_to_owned(row, &self.table_info) { - Ok(owned_row) => rows.push(owned_row), - Err(e) => { - return Box::new(PrefixLookupResultInner::from_error( - CLIENT_ERROR_CODE, - e.to_string(), - )); - } - } - } - - let columns = self.table_info.get_schema().columns().to_vec(); - Box::new(PrefixLookupResultInner { - error: None, - rows, - columns, - }) - } -} - // LogScanner implementation unsafe fn delete_log_scanner(scanner: *mut LogScanner) { if !scanner.is_null() { diff --git a/fluss-rust/bindings/cpp/src/lookup.rs b/fluss-rust/bindings/cpp/src/lookup.rs new file mode 100644 index 0000000000..d0d289e4c6 --- /dev/null +++ b/fluss-rust/bindings/cpp/src/lookup.rs @@ -0,0 +1,292 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Lookups for C++: blocking, or started in the background and waited for +//! through a handle. + +use std::future::Future; +use std::sync::Arc; +use std::time::Duration; + +use fluss as fcore; +use fluss::metadata::TableInfo; +use fluss::row::GenericRow; +use fluss::rpc::FlussError as CoreFlussError; +use tokio::task::JoinHandle; + +use crate::{ + CLIENT_ERROR_CODE, GenericRowInner, LookupResultInner, Lookuper, PrefixLookupResultInner, + PrefixLookuper, RUNTIME, client_err_ptr, err_from_core_error, ffi, ok_ptr, types, +}; + +/// A lookup result for C++, or the error in its place. +trait Outcome: Send + 'static { + fn from_error(code: i32, message: String) -> Self; +} + +impl Outcome for Box { + fn from_error(code: i32, message: String) -> Self { + Box::new(LookupResultInner::from_error(code, message)) + } +} + +impl Outcome for Box { + fn from_error(code: i32, message: String) -> Self { + Box::new(PrefixLookupResultInner::from_error(code, message)) + } +} + +/// A lookup running in the background. Dropping it abandons the lookup. +struct Pending { + task: Option>, +} + +impl Pending { + fn start(lookup: impl Future + Send + 'static) -> Self { + Self { + task: Some(RUNTIME.spawn(lookup)), + } + } + + /// Waits up to `timeout_ms` for the outcome, or without limit when it is + /// negative. A timeout leaves the lookup running, to be waited for again. + fn wait(&mut self, timeout_ms: i64) -> T { + let Some(task) = self.task.as_mut() else { + return T::from_error( + CLIENT_ERROR_CODE, + "Lookup was already waited for".to_string(), + ); + }; + let joined = match u64::try_from(timeout_ms) { + Ok(ms) => { + let timeout = Duration::from_millis(ms); + match RUNTIME.block_on(async { tokio::time::timeout(timeout, task).await }) { + Ok(joined) => joined, + Err(_) => { + return T::from_error( + CoreFlussError::RequestTimeOut.code(), + "Lookup did not complete within the wait timeout".to_string(), + ); + } + } + } + Err(_) => RUNTIME.block_on(task), + }; + self.task = None; + joined.unwrap_or_else(|e| T::from_error(CLIENT_ERROR_CODE, format!("Lookup failed: {e}"))) + } + + fn is_pending(&self) -> bool { + self.task.is_some() + } +} + +impl Drop for Pending { + fn drop(&mut self) { + if let Some(task) = self.task.take() { + task.abort(); + } + } +} + +pub struct PendingLookup(Pending>); + +pub struct PendingPrefixLookup(Pending>); + +impl PendingLookup { + pub(crate) fn lookup_wait(&mut self, timeout_ms: i64) -> Box { + self.0.wait(timeout_ms) + } + + pub(crate) fn lookup_is_pending(&self) -> bool { + self.0.is_pending() + } +} + +impl PendingPrefixLookup { + pub(crate) fn prefix_lookup_wait(&mut self, timeout_ms: i64) -> Box { + self.0.wait(timeout_ms) + } + + pub(crate) fn prefix_lookup_is_pending(&self) -> bool { + self.0.is_pending() + } +} + +pub(crate) unsafe fn delete_pending_lookup(pending: *mut PendingLookup) { + if !pending.is_null() { + unsafe { + drop(Box::from_raw(pending)); + } + } +} + +pub(crate) unsafe fn delete_pending_prefix_lookup(pending: *mut PendingPrefixLookup) { + if !pending.is_null() { + unsafe { + drop(Box::from_raw(pending)); + } + } +} + +impl Lookuper { + /// The lookup key: the primary-key values, dense as the core key encoder + /// expects them, and owned, so a lookup can outlive the C++ row. + fn key(&self, pk_row: &GenericRowInner) -> Result, String> { + let schema = self.table_info.get_schema(); + let pk_indices = schema.primary_key_indexes(); + types::resolve_dense_row_types(&pk_row.row, Some(schema), &pk_indices) + .map(|row| row.into_owned().into_owned()) + .map_err(|e| e.to_string()) + } + + fn run( + &self, + key: GenericRow<'static>, + ) -> impl Future> + Send + 'static { + let inner = Arc::clone(&self.inner); + let table_info = Arc::clone(&self.table_info); + async move { to_lookup_result(inner.lookup(&key).await, &table_info) } + } + + pub(crate) fn lookup(&self, pk_row: &GenericRowInner) -> Box { + match self.key(pk_row) { + Ok(key) => RUNTIME.block_on(self.run(key)), + Err(e) => Box::new(LookupResultInner::from_error(CLIENT_ERROR_CODE, e)), + } + } + + pub(crate) fn lookup_async(&self, pk_row: &GenericRowInner) -> ffi::FfiPtrResult { + match self.key(pk_row) { + Ok(key) => { + let pending = PendingLookup(Pending::start(self.run(key))); + ok_ptr(Box::into_raw(Box::new(pending)) as usize) + } + Err(e) => client_err_ptr(e), + } + } +} + +impl PrefixLookuper { + /// The lookup key: the prefix values, dense and in lookup-column order as the + /// core prefix encoder expects them, and owned, so a lookup can outlive the + /// C++ row. + fn key(&self, prefix_row: &GenericRowInner) -> Result, String> { + let schema = self.table_info.get_schema(); + types::resolve_dense_row_types(&prefix_row.row, Some(schema), &self.lookup_column_indices) + .map(|row| row.into_owned().into_owned()) + .map_err(|e| e.to_string()) + } + + fn run( + &self, + key: GenericRow<'static>, + ) -> impl Future> + Send + 'static { + let inner = Arc::clone(&self.inner); + let table_info = Arc::clone(&self.table_info); + async move { to_prefix_lookup_result(inner.lookup(&key).await, &table_info) } + } + + pub(crate) fn prefix_lookup( + &self, + prefix_row: &GenericRowInner, + ) -> Box { + match self.key(prefix_row) { + Ok(key) => RUNTIME.block_on(self.run(key)), + Err(e) => Box::new(PrefixLookupResultInner::from_error(CLIENT_ERROR_CODE, e)), + } + } + + pub(crate) fn prefix_lookup_async(&self, prefix_row: &GenericRowInner) -> ffi::FfiPtrResult { + match self.key(prefix_row) { + Ok(key) => { + let pending = PendingPrefixLookup(Pending::start(self.run(key))); + ok_ptr(Box::into_raw(Box::new(pending)) as usize) + } + Err(e) => client_err_ptr(e), + } + } +} + +fn to_lookup_result( + result: fcore::error::Result, + table_info: &TableInfo, +) -> Box { + let lookup_result = match result { + Ok(r) => r, + Err(e) => return core_error(&e), + }; + let columns = table_info.get_schema().columns().to_vec(); + match lookup_result.get_single_row() { + Ok(Some(row)) => match types::compacted_row_to_owned(&row, table_info) { + Ok(owned_row) => Box::new(LookupResultInner { + error: None, + found: true, + row: Some(owned_row), + columns, + }), + Err(e) => Box::new(LookupResultInner::from_error( + CLIENT_ERROR_CODE, + e.to_string(), + )), + }, + Ok(None) => Box::new(LookupResultInner { + error: None, + found: false, + row: None, + columns, + }), + Err(e) => core_error(&e), + } +} + +fn to_prefix_lookup_result( + result: fcore::error::Result, + table_info: &TableInfo, +) -> Box { + let lookup_result = match result { + Ok(r) => r, + Err(e) => return core_error(&e), + }; + let lookup_rows = match lookup_result.get_rows() { + Ok(rows) => rows, + Err(e) => return core_error(&e), + }; + let mut rows = Vec::with_capacity(lookup_rows.len()); + for row in &lookup_rows { + match types::compacted_row_to_owned(row, table_info) { + Ok(owned_row) => rows.push(owned_row), + Err(e) => { + return Box::new(PrefixLookupResultInner::from_error( + CLIENT_ERROR_CODE, + e.to_string(), + )); + } + } + } + let columns = table_info.get_schema().columns().to_vec(); + Box::new(PrefixLookupResultInner { + error: None, + rows, + columns, + }) +} + +fn core_error(e: &fcore::error::Error) -> T { + let ffi_err = err_from_core_error(e); + T::from_error(ffi_err.error_code, ffi_err.error_message) +} diff --git a/fluss-rust/bindings/cpp/src/table.cpp b/fluss-rust/bindings/cpp/src/table.cpp index 36b228273e..8de4f08c0c 100644 --- a/fluss-rust/bindings/cpp/src/table.cpp +++ b/fluss-rust/bindings/cpp/src/table.cpp @@ -19,6 +19,7 @@ #include +#include #include #include #include @@ -1045,6 +1046,16 @@ void LookupResult::Destroy() noexcept { } } +Result LookupResult::Adopt(ffi::LookupResultInner* inner) { + auto owned = rust::Box::from_raw(inner); + if (owned->lv_has_error()) { + return utils::make_error(owned->lv_error_code(), std::string(owned->lv_error_message())); + } + Destroy(); + inner_ = owned.into_raw(); + return utils::make_ok(); +} + LookupResult::LookupResult(LookupResult&& other) noexcept : inner_(other.inner_), column_map_(std::move(other.column_map_)) { other.inner_ = nullptr; @@ -1986,16 +1997,72 @@ Result Lookuper::Lookup(const GenericRow& pk_row, LookupResult& out) const { if (!pk_row.Available()) { return utils::make_client_error("GenericRow not available"); } + return out.Adopt(lookuper_->lookup(*pk_row.inner_).into_raw()); +} - auto result_box = lookuper_->lookup(*pk_row.inner_); - if (result_box->lv_has_error()) { - return utils::make_error(result_box->lv_error_code(), - std::string(result_box->lv_error_message())); +Result Lookuper::Lookup(const GenericRow& pk_row, PendingLookup& out) const { + if (!Available()) { + return utils::make_client_error("Lookuper not available"); } + if (!pk_row.Available()) { + return utils::make_client_error("GenericRow not available"); + } + auto ffi_result = lookuper_->lookup_async(*pk_row.inner_); + auto result = utils::from_ffi_result(ffi_result.result); + if (result.Ok()) { + out = PendingLookup(utils::ptr_from_ffi(ffi_result)); + } + return result; +} - out.Destroy(); - out.inner_ = result_box.into_raw(); - return utils::make_ok(); +// ============================================================================ +// PendingLookup +// ============================================================================ + +PendingLookup::PendingLookup() noexcept = default; + +PendingLookup::PendingLookup(ffi::PendingLookup* pending) noexcept : pending_(pending) {} + +PendingLookup::~PendingLookup() noexcept { Destroy(); } + +void PendingLookup::Destroy() noexcept { + if (pending_) { + ffi::delete_pending_lookup(pending_); + pending_ = nullptr; + } +} + +PendingLookup::PendingLookup(PendingLookup&& other) noexcept : pending_(other.pending_) { + other.pending_ = nullptr; +} + +PendingLookup& PendingLookup::operator=(PendingLookup&& other) noexcept { + if (this != &other) { + Destroy(); + pending_ = other.pending_; + other.pending_ = nullptr; + } + return *this; +} + +bool PendingLookup::Available() const { return pending_ != nullptr; } + +Result PendingLookup::Wait(LookupResult& out) { return DoWait(out, -1); } + +Result PendingLookup::Wait(LookupResult& out, int64_t timeout_ms) { + return DoWait(out, std::max(timeout_ms, 0)); +} + +Result PendingLookup::DoWait(LookupResult& out, int64_t timeout_ms) { + if (!Available()) { + return utils::make_client_error("PendingLookup not available"); + } + auto result = out.Adopt(pending_->lookup_wait(timeout_ms).into_raw()); + // A timeout leaves the lookup running; any other outcome ends it. + if (!pending_->lookup_is_pending()) { + Destroy(); + } + return result; } // ============================================================================ @@ -2037,27 +2104,96 @@ Result PrefixLookuper::PrefixLookup(const GenericRow& prefix_row, PrefixLookupRe if (!prefix_row.Available()) { return utils::make_client_error("GenericRow not available"); } + return out.Adopt(lookuper_->prefix_lookup(*prefix_row.inner_).into_raw()); +} - auto result_box = lookuper_->prefix_lookup(*prefix_row.inner_); - if (result_box->plv_has_error()) { - return utils::make_error(result_box->plv_error_code(), - std::string(result_box->plv_error_message())); +Result PrefixLookuper::PrefixLookup(const GenericRow& prefix_row, PendingPrefixLookup& out) const { + if (!Available()) { + return utils::make_client_error("PrefixLookuper not available"); + } + if (!prefix_row.Available()) { + return utils::make_client_error("GenericRow not available"); + } + auto ffi_result = lookuper_->prefix_lookup_async(*prefix_row.inner_); + auto result = utils::from_ffi_result(ffi_result.result); + if (result.Ok()) { + out = PendingPrefixLookup(utils::ptr_from_ffi(ffi_result)); + } + return result; +} + +Result PrefixLookupResult::Adopt(ffi::PrefixLookupResultInner* inner) { + auto owned = rust::Box::from_raw(inner); + if (owned->plv_has_error()) { + return utils::make_error(owned->plv_error_code(), std::string(owned->plv_error_message())); } // Take ownership of the FFI box first (~PrefixData calls from_raw), so the // column-map loop below can't leak it if a string/map allocation throws. // The map is built eagerly and shared by all PrefixRowViews. - auto data = std::make_shared(result_box.into_raw(), detail::ColumnMap{}); + auto data = std::make_shared(owned.into_raw(), detail::ColumnMap{}); auto col_count = data->raw->plv_field_count(); for (size_t i = 0; i < col_count; ++i) { auto name = data->raw->plv_column_name(i); data->columns[std::string(name.data(), name.size())] = { i, static_cast(data->raw->plv_column_type(i))}; } - out.data_ = std::move(data); + data_ = std::move(data); return utils::make_ok(); } +// ============================================================================ +// PendingPrefixLookup +// ============================================================================ + +PendingPrefixLookup::PendingPrefixLookup() noexcept = default; + +PendingPrefixLookup::PendingPrefixLookup(ffi::PendingPrefixLookup* pending) noexcept + : pending_(pending) {} + +PendingPrefixLookup::~PendingPrefixLookup() noexcept { Destroy(); } + +void PendingPrefixLookup::Destroy() noexcept { + if (pending_) { + ffi::delete_pending_prefix_lookup(pending_); + pending_ = nullptr; + } +} + +PendingPrefixLookup::PendingPrefixLookup(PendingPrefixLookup&& other) noexcept + : pending_(other.pending_) { + other.pending_ = nullptr; +} + +PendingPrefixLookup& PendingPrefixLookup::operator=(PendingPrefixLookup&& other) noexcept { + if (this != &other) { + Destroy(); + pending_ = other.pending_; + other.pending_ = nullptr; + } + return *this; +} + +bool PendingPrefixLookup::Available() const { return pending_ != nullptr; } + +Result PendingPrefixLookup::Wait(PrefixLookupResult& out) { return DoWait(out, -1); } + +Result PendingPrefixLookup::Wait(PrefixLookupResult& out, int64_t timeout_ms) { + return DoWait(out, std::max(timeout_ms, 0)); +} + +Result PendingPrefixLookup::DoWait(PrefixLookupResult& out, int64_t timeout_ms) { + if (!Available()) { + return utils::make_client_error("PendingPrefixLookup not available"); + } + auto result = out.Adopt(pending_->prefix_lookup_wait(timeout_ms).into_raw()); + // A timeout leaves the lookup running; any other outcome ends it. + if (!pending_->prefix_lookup_is_pending()) { + Destroy(); + } + return result; +} + // ============================================================================ // LogScanner // ============================================================================ diff --git a/fluss-rust/bindings/cpp/test/test_kv_table.cpp b/fluss-rust/bindings/cpp/test/test_kv_table.cpp index 2de9402646..0f078e4475 100644 --- a/fluss-rust/bindings/cpp/test/test_kv_table.cpp +++ b/fluss-rust/bindings/cpp/test/test_kv_table.cpp @@ -2065,3 +2065,146 @@ TEST_F(KvTableTest, AllSupportedDatatypes) { ASSERT_OK(adm.DropTable(table_path, false)); } + +namespace { + +// Ids 0..63, each with sequence numbers 0 and 1, bucketed by id: the full key +// (id, seq) serves a Lookuper and the prefix (id) a PrefixLookuper. +void CreateSequenceTable(fluss::Admin& adm, fluss::Connection& conn, const fluss::TablePath& path) { + auto schema = fluss::Schema::NewBuilder() + .AddColumn("id", DataType::Int()) + .AddColumn("seq", DataType::BigInt()) + .AddColumn("name", DataType::String()) + .SetPrimaryKeys({"id", "seq"}) + .Build(); + auto table_descriptor = fluss::TableDescriptor::NewBuilder() + .SetSchema(schema) + .SetBucketCount(3) + .SetBucketKeys({"id"}) + .SetProperty("table.replication.factor", "1") + .Build(); + fluss_test::CreateTable(adm, path, table_descriptor); + + fluss::Table table; + ASSERT_OK(conn.GetTable(path, table)); + fluss::UpsertWriter upsert_writer; + ASSERT_OK(table.NewUpsert().CreateWriter(upsert_writer)); + for (int32_t id = 0; id < 64; ++id) { + for (int64_t seq = 0; seq < 2; ++seq) { + fluss::GenericRow row(3); + row.SetInt32(0, id); + row.SetInt64(1, seq); + row.SetString(2, std::to_string(id) + "-" + std::to_string(seq)); + ASSERT_OK(upsert_writer.Upsert(row)); + } + } + ASSERT_OK(upsert_writer.Flush()); +} + +fluss::GenericRow Key(int32_t id, int64_t seq) { + fluss::GenericRow row(3); + row.SetInt32(0, id); + row.SetInt64(1, seq); + return row; +} + +fluss::GenericRow Prefix(int32_t id) { + fluss::GenericRow row(3); + row.SetInt32(0, id); + return row; +} + +// A connection whose lookups each wait `batch_timeout_ms` for their batch to fill, +// so a lookup is still running when a shorter timeout passes. +void ConnectSlowLookups(fluss::Connection& out, uint64_t batch_timeout_ms) { + fluss::Configuration config; + config.bootstrap_servers = fluss_test::FlussTestEnvironment::Instance()->GetBootstrapServers(); + config.lookup_batch_timeout_ms = batch_timeout_ms; + ASSERT_OK(fluss::Connection::Create(config, out)); +} + +} // namespace + +TEST_F(KvTableTest, OneThreadKeepsManyLookupsInFlight) { + fluss::TablePath path("fluss", "test_pending_lookups_cpp"); + CreateSequenceTable(admin(), connection(), path); + fluss::Table table; + ASSERT_OK(connection().GetTable(path, table)); + fluss::Lookuper lookuper; + ASSERT_OK(table.NewLookup().CreateLookuper(lookuper)); + fluss::PrefixLookuper prefix_lookuper; + ASSERT_OK(table.NewPrefixLookup({"id"}, prefix_lookuper)); + + // The first lookups also fetch metadata. + fluss::LookupResult first; + ASSERT_OK(lookuper.Lookup(Key(0, 0), first)); + fluss::PrefixLookupResult first_rows; + ASSERT_OK(prefix_lookuper.PrefixLookup(Prefix(0), first_rows)); + + // One at a time, these 192 lookups take ~19 s, since each waits for a batch. + constexpr int32_t kKeys = 96; + std::vector pending(kKeys); + std::vector pending_rows(kKeys); + auto started = std::chrono::steady_clock::now(); + for (int32_t id = 0; id < kKeys; ++id) { + ASSERT_OK(lookuper.Lookup(Key(id, 1), pending[id])); + ASSERT_OK(prefix_lookuper.PrefixLookup(Prefix(id), pending_rows[id])); + } + for (int32_t id = 0; id < kKeys; ++id) { + fluss::LookupResult result; + ASSERT_OK(pending[id].Wait(result)); + EXPECT_FALSE(pending[id].Available()); + if (id < 64) { + ASSERT_TRUE(result.Found()) << "id=" << id; + EXPECT_EQ(result.GetString("name"), std::to_string(id) + "-1"); + } else { + EXPECT_FALSE(result.Found()) << "id=" << id; + } + fluss::PrefixLookupResult rows; + ASSERT_OK(pending_rows[id].Wait(rows)); + EXPECT_EQ(rows.Size(), id < 64 ? 2u : 0u) << "id=" << id; + } + EXPECT_LT(std::chrono::steady_clock::now() - started, std::chrono::seconds(5)); + + // A lookup that was waited for is used up. + fluss::LookupResult again; + EXPECT_FALSE(pending[0].Wait(again).Ok()); + + ASSERT_OK(admin().DropTable(path, false)); +} + +TEST_F(KvTableTest, WaitingForALookupWithATimeout) { + fluss::TablePath path("fluss", "test_lookup_wait_timeout_cpp"); + CreateSequenceTable(admin(), connection(), path); + + // Each lookup through this connection takes about 2 s. + fluss::Connection slow; + ConnectSlowLookups(slow, 2000); + fluss::Table table; + ASSERT_OK(slow.GetTable(path, table)); + fluss::Lookuper lookuper; + ASSERT_OK(table.NewLookup().CreateLookuper(lookuper)); + fluss::PrefixLookuper prefix_lookuper; + ASSERT_OK(table.NewPrefixLookup({"id"}, prefix_lookuper)); + + // A wait that times out leaves the lookup running, to be waited for again. + fluss::PendingLookup pending; + ASSERT_OK(lookuper.Lookup(Key(1, 0), pending)); + fluss::LookupResult result; + auto timed_out = pending.Wait(result, 100); + EXPECT_EQ(timed_out.error_code, fluss::ErrorCode::REQUEST_TIME_OUT); + EXPECT_TRUE(timed_out.IsRetriable()); + ASSERT_TRUE(pending.Available()); + ASSERT_OK(pending.Wait(result)); + EXPECT_EQ(result.GetString("name"), "1-0"); + + fluss::PendingPrefixLookup pending_rows; + ASSERT_OK(prefix_lookuper.PrefixLookup(Prefix(1), pending_rows)); + fluss::PrefixLookupResult rows; + EXPECT_EQ(pending_rows.Wait(rows, 100).error_code, fluss::ErrorCode::REQUEST_TIME_OUT); + ASSERT_TRUE(pending_rows.Available()); + ASSERT_OK(pending_rows.Wait(rows)); + EXPECT_EQ(rows.Size(), 2u); + + ASSERT_OK(admin().DropTable(path, false)); +} diff --git a/fluss-rust/website/docs/user-guide/cpp/api-reference.md b/fluss-rust/website/docs/user-guide/cpp/api-reference.md index 9fb6f6fd0e..a9ea9cc5b2 100644 --- a/fluss-rust/website/docs/user-guide/cpp/api-reference.md +++ b/fluss-rust/website/docs/user-guide/cpp/api-reference.md @@ -535,14 +535,43 @@ Performs point lookups by primary key. Obtained from `TableLookup::CreateLookupe | Method | Description | |---------------------------------------------------------------------|-----------------------------| | `Lookup(const GenericRow& pk_row, LookupResult& out) const -> Result` | Lookup a row by primary key | +| `Lookup(const GenericRow& pk_row, PendingLookup& out) const -> Result` | Start a lookup without blocking; `out` waits for its row | + +Each blocking lookup waits for its batch, up to `lookup_batch_timeout_ms`. To look up many keys from one thread, start them all, then wait for each: + +```cpp +std::vector pending(keys.size()); +for (size_t i = 0; i < keys.size(); ++i) { + if (!lookuper.Lookup(keys[i], pending[i]).Ok()) { + // This lookup did not start. + } +} +for (auto& lookup : pending) { + fluss::LookupResult result; + if (lookup.Available() && lookup.Wait(result).Ok() && result.Found()) { + // Use result. + } +} +``` + +## `PendingLookup` + +A lookup started by `Lookuper::Lookup(pk_row, PendingLookup&)`. Move-only. Destroying it before `Wait` returns abandons the lookup. + +| Method | Description | +|---|---| +| `Wait(LookupResult& out) -> Result` | Block until the row arrives | +| `Wait(LookupResult& out, int64_t timeout_ms) -> Result` | Block for at most `timeout_ms`. After the retriable `REQUEST_TIME_OUT` the lookup keeps running, and `Wait` can be called again | +| `Available() const -> bool` | Whether `Wait` can still be called | ## `PrefixLookuper` -Performs prefix (bucket-key) lookups, returning all rows whose primary key starts with the given prefix. Obtained from `Table::NewPrefixLookup()`. Like a `Lookuper`, it can be shared by threads. See the [Prefix Lookup example](./example/prefix-lookup.md). +Performs prefix (bucket-key) lookups, returning all rows whose primary key starts with the given prefix. Obtained from `Table::NewPrefixLookup()`. Like a `Lookuper`, it can be shared by threads, and it can start a lookup without blocking. `PendingPrefixLookup` has the same methods as `PendingLookup`. See the [Prefix Lookup example](./example/prefix-lookup.md). | Method | Description | |---------------------------------------------------------------------------------------|-----------------------------------------------| | `PrefixLookup(const GenericRow& prefix_row, PrefixLookupResult& out) const -> Result` | Look up all rows matching the prefix columns | +| `PrefixLookup(const GenericRow& prefix_row, PendingPrefixLookup& out) const -> Result` | Start a lookup without blocking; `out` waits for its rows | ## `LogScanner` diff --git a/website/docs/apis/cpp/api-reference.md b/website/docs/apis/cpp/api-reference.md index 0a254ae977..b61809554a 100644 --- a/website/docs/apis/cpp/api-reference.md +++ b/website/docs/apis/cpp/api-reference.md @@ -187,14 +187,43 @@ Performs point lookups by primary key. Obtained from `TableLookup::CreateLookupe | Method | Description | |---------------------------------------------------------------------|-----------------------------| | `Lookup(const GenericRow& pk_row, LookupResult& out) const -> Result` | Lookup a row by primary key | +| `Lookup(const GenericRow& pk_row, PendingLookup& out) const -> Result` | Start a lookup without blocking; `out` waits for its row | + +Each blocking lookup waits for its batch, up to `lookup_batch_timeout_ms`. To look up many keys from one thread, start them all, then wait for each: + +```cpp +std::vector pending(keys.size()); +for (size_t i = 0; i < keys.size(); ++i) { + if (!lookuper.Lookup(keys[i], pending[i]).Ok()) { + // This lookup did not start. + } +} +for (auto& lookup : pending) { + fluss::LookupResult result; + if (lookup.Available() && lookup.Wait(result).Ok() && result.Found()) { + // Use result. + } +} +``` + +## `PendingLookup` + +A lookup started by `Lookuper::Lookup(pk_row, PendingLookup&)`. Move-only. Destroying it before `Wait` returns abandons the lookup. + +| Method | Description | +|---|---| +| `Wait(LookupResult& out) -> Result` | Block until the row arrives | +| `Wait(LookupResult& out, int64_t timeout_ms) -> Result` | Block for at most `timeout_ms`. After the retriable `REQUEST_TIME_OUT` the lookup keeps running, and `Wait` can be called again | +| `Available() const -> bool` | Whether `Wait` can still be called | ## `PrefixLookuper` -Performs prefix (bucket-key) lookups, returning all rows whose primary key starts with the given prefix. Obtained from `Table::NewPrefixLookup()`. Like a `Lookuper`, it can be shared by threads. See the [Prefix Lookup example](./example/prefix-lookup.md). +Performs prefix (bucket-key) lookups, returning all rows whose primary key starts with the given prefix. Obtained from `Table::NewPrefixLookup()`. Like a `Lookuper`, it can be shared by threads, and it can start a lookup without blocking. `PendingPrefixLookup` has the same methods as `PendingLookup`. See the [Prefix Lookup example](./example/prefix-lookup.md). | Method | Description | |---------------------------------------------------------------------------------------|-----------------------------------------------| | `PrefixLookup(const GenericRow& prefix_row, PrefixLookupResult& out) const -> Result` | Look up all rows matching the prefix columns | +| `PrefixLookup(const GenericRow& prefix_row, PendingPrefixLookup& out) const -> Result` | Start a lookup without blocking; `out` waits for its rows | ## `LogScanner`