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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 76 additions & 0 deletions fluss-rust/bindings/cpp/include/fluss.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,8 @@ struct BatchScanner;
struct UpsertWriter;
struct Lookuper;
struct PrefixLookuper;
struct PendingLookup;
struct PendingPrefixLookup;
struct ScanResultInner;
struct GenericRowInner;
struct LookupResultInner;
Expand Down Expand Up @@ -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<const detail::PrefixData> data_;
};

Expand Down Expand Up @@ -1576,12 +1581,15 @@ class LookupResult : public detail::NamedGetters<LookupResult> {

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};
Expand Down Expand Up @@ -2068,6 +2076,67 @@ class UpsertWriter {
std::shared_ptr<ffi::WriteCallbackCapacity> 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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down
163 changes: 30 additions & 133 deletions fluss-rust/bindings/cpp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -423,6 +427,8 @@ mod ffi {
type UpsertWriter;
type Lookuper;
type PrefixLookuper;
type PendingLookup;
type PendingPrefixLookup;

// Opaque types for optimized FFI
type ScanResultInner;
Expand Down Expand Up @@ -703,6 +709,12 @@ mod ffi {
// Lookuper
unsafe fn delete_lookuper(lookuper: *mut Lookuper);
fn lookup(self: &Lookuper, pk_row: &GenericRowInner) -> Box<LookupResultInner>;
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<LookupResultInner>;
fn lookup_is_pending(self: &PendingLookup) -> bool;

// LookupResultInner accessors
fn lv_has_error(self: &LookupResultInner) -> bool;
Expand Down Expand Up @@ -763,6 +775,16 @@ mod ffi {
self: &PrefixLookuper,
prefix_row: &GenericRowInner,
) -> Box<PrefixLookupResultInner>;
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<PrefixLookupResultInner>;
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.
Expand Down Expand Up @@ -1026,13 +1048,13 @@ pub struct UpsertWriter {
}

pub struct Lookuper {
inner: fcore::client::Lookuper,
table_info: fcore::metadata::TableInfo,
inner: Arc<fcore::client::Lookuper>,
table_info: Arc<TableInfo>,
}

pub struct PrefixLookuper {
inner: PrefixKeyLookuper,
table_info: TableInfo,
inner: Arc<PrefixKeyLookuper>,
table_info: Arc<TableInfo>,
/// Full-schema indices of the lookup columns, used to compact the input row.
lookup_column_indices: Vec<usize>,
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -2373,66 +2395,6 @@ unsafe fn delete_lookuper(lookuper: *mut Lookuper) {
}
}

impl Lookuper {
fn lookup(&self, pk_row: &GenericRowInner) -> Box<LookupResultInner> {
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() {
Expand All @@ -2442,71 +2404,6 @@ unsafe fn delete_prefix_lookuper(lookuper: *mut PrefixLookuper) {
}
}

impl PrefixLookuper {
fn prefix_lookup(&self, prefix_row: &GenericRowInner) -> Box<PrefixLookupResultInner> {
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() {
Expand Down
Loading
Loading