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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion libsql-server/proto/namespace_fence.proto
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,8 @@ message DrainPolicy {

message FrozenBoundary {
string log_id = 1;
uint64 frame_no = 2;
// Absent when the replication log has no frames.
optional uint64 frame_no = 2;
}

message LegacyBlocks {
Expand Down
3 changes: 3 additions & 0 deletions libsql-server/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,9 @@ pub struct MetaStoreConfig {
/// How long receipts of finished fence operations are kept. `None` is the default of
/// 30 days.
pub namespace_fence_receipt_retention: Option<Duration>,
/// How long `AcquireSourceWriteFence` waits for active writers when the request names no
/// drain policy. `None` is the default of 30 seconds.
pub namespace_fence_default_write_drain: Option<Duration>,
}

#[derive(Debug, Clone)]
Expand Down
70 changes: 62 additions & 8 deletions libsql-server/src/connection/connection_core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ use crate::connection::legacy::open_conn_active_checkpoint;
use crate::error::Error;
use crate::metrics::{PROGRAM_EXEC_COUNT, QUERY_CANCELED, VACUUM_COUNT, WAL_CHECKPOINT_COUNT};
use crate::namespace::broadcasters::BroadcasterHandle;
use crate::namespace::fence::controller::FenceConnState;
use crate::namespace::fence::state::OperationClass;
use crate::namespace::meta_store::MetaStoreHandle;
use crate::namespace::ResolveNamespacePathFn;
use crate::query_analysis::StmtKind;
Expand All @@ -24,6 +26,15 @@ use super::program::{DescribeCol, DescribeParam, DescribeResponse, Program, Vm};

pub type GetCurrentFrameNo = Arc<dyn Fn() -> Option<FrameNo> + Send + Sync + 'static>;

/// What [`CoreConnection::vacuum_if_needed_above`] did.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum VacuumOutcome {
Vacuumed,
NotNeeded,
/// Skipped because the namespace fence denies normal writes.
Fenced,
}

/// The base connection type, shared between legacy and libsql-wal implementations
pub(super) struct CoreConnection<W> {
conn: libsql_sys::Connection<W>,
Expand All @@ -37,6 +48,8 @@ pub(super) struct CoreConnection<W> {
broadcaster: BroadcasterHandle,
hooked: bool,
canceled: Arc<AtomicBool>,
/// Shared with this connection's WAL wrapper (`docs/NAMESPACE_FENCE.md` section 7.4).
fence: Arc<FenceConnState>,
}

fn update_stats(
Expand Down Expand Up @@ -68,6 +81,7 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
get_current_frame_no: GetCurrentFrameNo,
block_writes: Arc<AtomicBool>,
resolve_attach_path: ResolveNamespacePathFn,
fence: Arc<FenceConnState>,
) -> Result<Self> {
let conn = open_conn_active_checkpoint(
path,
Expand Down Expand Up @@ -113,6 +127,7 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
hooked: false,
canceled,
get_current_frame_no,
fence,
};

for ext in extensions.iter() {
Expand Down Expand Up @@ -188,17 +203,21 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
pgm: Program,
mut builder: B,
) -> Result<B> {
let (config, stats, block_writes, resolve_attach_path) = {
let (config, stats, block_writes, resolve_attach_path, fence) = {
let mut lock = this.lock();
let config = lock.config_store.get();
let stats = lock.stats.clone();
let block_writes = lock.block_writes.clone();
let resolve_attach_path = lock.resolve_attach_path.clone();
let fence = lock.fence.clone();

lock.update_hooks();

(config, stats, block_writes, resolve_attach_path)
(config, stats, block_writes, resolve_attach_path, fence)
};
// The program is admitted under the gate's current write generation; a write
// transaction it opens must start under the same one (section 8.1).
fence.begin_program();

builder.init(&this.lock().builder_config)?;
let mut vm = Vm::new(
Expand Down Expand Up @@ -229,7 +248,8 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
update_stats(&stats, sql, rows_read, rows_written, mem_used, elapsed)
},
resolve_attach_path,
);
)
.with_fence(fence);

let mut has_timeout = false;
while !vm.finished() {
Expand Down Expand Up @@ -296,21 +316,45 @@ impl<W: Wal + Send + 'static> CoreConnection<W> {
}

pub(super) fn vacuum_if_needed(&self) -> Result<()> {
// NOTICE: don't bother vacuuming if we don't have at least 256MiB of data
self.vacuum_if_needed_above(65536).map(|_| ())
}

/// `VACUUM` if the database has at least `min_pages` pages and more than half of them are
/// free. `VACUUM` is not maintenance: it takes a write transaction and produces replicated
/// frames, so it is skipped whenever the namespace fence denies normal writes
/// (`docs/NAMESPACE_FENCE.md` section 7.3), including when the fence closes between the
/// check and the `VACUUM` itself.
pub(super) fn vacuum_if_needed_above(&self, min_pages: i64) -> Result<VacuumOutcome> {
self.fence.begin_program();
if let Err(e) = self.fence.controller().permits(OperationClass::Vacuum) {
tracing::debug!("skipping vacuum: {e}");
return Ok(VacuumOutcome::Fenced);
}
let page_count = self
.conn
.query_row("PRAGMA page_count", (), |row| row.get::<_, i64>(0))?;
let freelist_count = self
.conn
.query_row("PRAGMA freelist_count", (), |row| row.get::<_, i64>(0))?;
// NOTICE: don't bother vacuuming if we don't have at least 256MiB of data
if page_count >= 65536 && freelist_count * 2 > page_count {
let outcome = if page_count >= min_pages && freelist_count * 2 > page_count {
tracing::info!("Vacuuming: pages={page_count} freelist={freelist_count}");
self.conn.execute("VACUUM", ())?;
if let Err(e) = self.conn.execute("VACUUM", ()) {
return match self.fence.take_denial() {
Some(denial) => {
tracing::debug!("skipping vacuum: {denial}");
Ok(VacuumOutcome::Fenced)
}
None => Err(e.into()),
};
}
VacuumOutcome::Vacuumed
} else {
tracing::trace!("Not vacuuming: pages={page_count} freelist={freelist_count}");
}
VacuumOutcome::NotNeeded
};
VACUUM_COUNT.increment(1);
Ok(())
Ok(outcome)
}

pub(super) fn describe(&self, sql: &str) -> crate::Result<DescribeResponse> {
Expand Down Expand Up @@ -394,6 +438,7 @@ mod test {
use crate::auth::Authenticated;
use crate::connection::legacy::MakeLegacyConnection;
use crate::connection::{Connection as _, RequestContext, TXN_TIMEOUT};
use crate::namespace::fence::controller::FenceController;
use crate::namespace::meta_store::{metastore_connection_maker, MetaStore};
use crate::namespace::NamespaceName;
use crate::query_result_builder::test::{test_driver, TestBuilder};
Expand All @@ -415,6 +460,10 @@ mod test {
hooked: false,
canceled: Arc::new(false.into()),
get_current_frame_no: Arc::new(|| None),
fence: FenceConnState::new(
FenceController::unfenced(Default::default()),
crate::namespace::fence::state::OperationClass::NormalWrite,
),
};

let conn = Arc::new(Mutex::new(conn));
Expand Down Expand Up @@ -454,6 +503,7 @@ mod test {
Default::default(),
Arc::new(|_| unreachable!()),
Arc::new(|| Sqlite3WalManager::default()),
FenceController::unfenced(Default::default()),
)
.await
.unwrap();
Expand Down Expand Up @@ -500,6 +550,7 @@ mod test {
Default::default(),
Arc::new(|_| unreachable!()),
Arc::new(|| Sqlite3WalManager::default()),
FenceController::unfenced(Default::default()),
)
.await
.unwrap();
Expand Down Expand Up @@ -551,6 +602,7 @@ mod test {
Default::default(),
Arc::new(|_| unreachable!()),
Arc::new(|| Sqlite3WalManager::default()),
FenceController::unfenced(Default::default()),
)
.await
.unwrap();
Expand Down Expand Up @@ -634,6 +686,7 @@ mod test {
Default::default(),
Arc::new(|_| unreachable!()),
Arc::new(|| Sqlite3WalManager::default()),
FenceController::unfenced(Default::default()),
)
.await
.unwrap();
Expand Down Expand Up @@ -727,6 +780,7 @@ mod test {
Default::default(),
Arc::new(|_| unreachable!()),
Arc::new(|| Sqlite3WalManager::default()),
FenceController::unfenced(Default::default()),
)
.await
.unwrap();
Expand Down
Loading
Loading