#![allow(dead_code)]
use std::{
future::Future,
pin::Pin,
sync::{
Mutex,
atomic::{AtomicBool, AtomicU64, Ordering},
},
task::{Context, Poll, Waker},
thread::{self, JoinHandle},
};
use saddle_observability::file::{
CompletionResult, CompletionRole, CompletionSignal, FixedCompletion, FixedFailure,
FixedFileSink, FixedNotification, PreparedFixedFileCore,
};
const WRITER_NAME: &str = "sdl-log-writer";
const WRITER_STACK_BYTES: usize = 1_048_576;
static WRITER_LEASED: AtomicBool = AtomicBool::new(false);
static REGISTRY_UNHEALTHY: AtomicBool = AtomicBool::new(false);
static NOTIFICATIONS: [NotificationSlot; 3] = [
NotificationSlot::new(),
NotificationSlot::new(),
NotificationSlot::new(),
];
struct NotificationSlot {
registered: AtomicBool,
generation: AtomicU64,
signalled: AtomicBool,
waker: Mutex<Option<Waker>>,
}
impl NotificationSlot {
const fn new() -> Self {
Self {
registered: AtomicBool::new(false),
generation: AtomicU64::new(0),
signalled: AtomicBool::new(false),
waker: Mutex::new(None),
}
}
fn register(&self) -> Result<(), ManagedWriterError> {
self.registered
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.map(|_| ())
.map_err(|_| ManagedWriterError::NotificationBusy)
}
fn clear(&self) {
self.registered.store(false, Ordering::Release);
self.signalled.store(false, Ordering::Release);
match self.waker.lock() {
Ok(mut waker) => *waker = None,
Err(_) => REGISTRY_UNHEALTHY.store(true, Ordering::Release),
}
}
fn reset(&self) {
self.clear();
self.generation.store(0, Ordering::Release);
}
}
fn fixed_notify(token: usize, signal: CompletionSignal) {
let Some(slot) = NOTIFICATIONS.get(token) else {
REGISTRY_UNHEALTHY.store(true, Ordering::Release);
return;
};
if token != signal.role as usize || signal.generation == 0 {
REGISTRY_UNHEALTHY.store(true, Ordering::Release);
return;
}
let previous = slot.generation.swap(signal.generation, Ordering::AcqRel);
if signal.generation <= previous {
REGISTRY_UNHEALTHY.store(true, Ordering::Release);
}
slot.signalled.store(true, Ordering::Release);
if let Ok(mut registered) = slot.waker.try_lock() {
if let Some(waker) = registered.take() {
waker.wake();
}
}
}
fn fixed_notifications() -> [FixedNotification; 3] {
std::array::from_fn(|token| FixedNotification::new(token, fixed_notify))
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ManagedWriterError {
AlreadyOwned,
ThreadSpawn,
ThreadPanic,
Completion(FixedFailure),
CompletionStale,
CompletionRecycle,
NotificationBusy,
NotificationUnhealthy,
Submit,
}
struct CompletionWait<const BLOCKS: usize, const BYTES: usize, const COMMANDS: usize> {
completion: Option<FixedCompletion<BLOCKS, BYTES, COMMANDS>>,
role: CompletionRole,
registered: bool,
}
impl<const BLOCKS: usize, const BYTES: usize, const COMMANDS: usize>
CompletionWait<BLOCKS, BYTES, COMMANDS>
{
fn new(completion: FixedCompletion<BLOCKS, BYTES, COMMANDS>, role: CompletionRole) -> Self {
Self {
completion: Some(completion),
role,
registered: false,
}
}
fn slot(&self) -> &'static NotificationSlot {
&NOTIFICATIONS[self.role as usize]
}
fn unregister(&mut self) {
if self.registered {
self.slot().clear();
self.registered = false;
}
}
}
impl<const BLOCKS: usize, const BYTES: usize, const COMMANDS: usize> Future
for CompletionWait<BLOCKS, BYTES, COMMANDS>
{
type Output = Result<(), ManagedWriterError>;
fn poll(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
if REGISTRY_UNHEALTHY.load(Ordering::Acquire) {
self.unregister();
return Poll::Ready(Err(ManagedWriterError::NotificationUnhealthy));
}
let completion = self
.completion
.as_ref()
.expect("pending completion remains owned");
match completion.result() {
CompletionResult::Ready(result) => {
self.unregister();
let completion = self.completion.take().expect("completion remains owned");
if completion.recycle().is_err() {
return Poll::Ready(Err(ManagedWriterError::CompletionRecycle));
}
Poll::Ready(result.map_err(ManagedWriterError::Completion))
}
CompletionResult::Stale => {
self.unregister();
Poll::Ready(Err(ManagedWriterError::CompletionStale))
}
CompletionResult::Pending => {
if !self.registered {
if let Err(error) = self.slot().register() {
return Poll::Ready(Err(error));
}
self.registered = true;
}
let slot = self.slot();
match slot.waker.lock() {
Ok(mut waker) => {
if waker
.as_ref()
.is_none_or(|registered| !registered.will_wake(context.waker()))
{
*waker = Some(context.waker().clone());
}
}
Err(_) => {
REGISTRY_UNHEALTHY.store(true, Ordering::Release);
self.unregister();
return Poll::Ready(Err(ManagedWriterError::NotificationUnhealthy));
}
}
if slot.signalled.swap(false, Ordering::AcqRel) {
context.waker().wake_by_ref();
}
Poll::Pending
}
}
}
}
impl<const BLOCKS: usize, const BYTES: usize, const COMMANDS: usize> Drop
for CompletionWait<BLOCKS, BYTES, COMMANDS>
{
fn drop(&mut self) {
self.unregister();
}
}
struct ManagedFileWriter<
const BLOCKS: usize,
const BYTES: usize,
const COMMANDS: usize,
const PATH: usize,
const DIRENT_BYTES: usize,
const ENTRIES: usize,
> {
sink: FixedFileSink<BLOCKS, BYTES, COMMANDS>,
writer: Option<JoinHandle<Result<(), FixedFailure>>>,
shutdown_complete: bool,
}
impl<
const BLOCKS: usize,
const BYTES: usize,
const COMMANDS: usize,
const PATH: usize,
const DIRENT_BYTES: usize,
const ENTRIES: usize,
> ManagedFileWriter<BLOCKS, BYTES, COMMANDS, PATH, DIRENT_BYTES, ENTRIES>
{
async fn start_until<D>(
prepared: PreparedFixedFileCore<BLOCKS, BYTES, COMMANDS, PATH, DIRENT_BYTES, ENTRIES>,
deadline: D,
) -> Result<Self, ManagedWriterError>
where
D: Future<Output = ()>,
{
if WRITER_LEASED
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return Err(ManagedWriterError::AlreadyOwned);
}
REGISTRY_UNHEALTHY.store(false, Ordering::Release);
for slot in &NOTIFICATIONS {
slot.reset();
}
let writer = match thread::Builder::new()
.name(WRITER_NAME.to_owned())
.stack_size(WRITER_STACK_BYTES)
.spawn(move || prepared.writer.run())
{
Ok(writer) => writer,
Err(_) => {
WRITER_LEASED.store(false, Ordering::Release);
return Err(ManagedWriterError::ThreadSpawn);
}
};
let startup = CompletionWait::new(prepared.startup, CompletionRole::Startup);
tokio::pin!(startup);
tokio::pin!(deadline);
let startup = tokio::select! {
biased;
() = &mut deadline => std::process::abort(),
result = &mut startup => result,
};
if let Err(error) = startup {
let _ = writer.join();
WRITER_LEASED.store(false, Ordering::Release);
return Err(error);
}
Ok(Self {
sink: prepared.sink,
writer: Some(writer),
shutdown_complete: false,
})
}
async fn flush(&self) -> Result<(), ManagedWriterError> {
let completion = self
.sink
.try_flush()
.map_err(|_| ManagedWriterError::Submit)?;
CompletionWait::new(completion, CompletionRole::Flush).await
}
async fn shutdown_until<D>(mut self, deadline: D) -> Result<(), ManagedWriterError>
where
D: Future<Output = ()>,
{
let completion = self
.sink
.try_shutdown()
.map_err(|_| ManagedWriterError::Submit)?;
let shutdown = CompletionWait::new(completion, CompletionRole::Shutdown);
tokio::pin!(shutdown);
tokio::pin!(deadline);
tokio::select! {
biased;
() = &mut deadline => std::process::abort(),
result = &mut shutdown => result?,
}
let writer = self.writer.take().ok_or(ManagedWriterError::ThreadPanic)?;
let result = writer.join().map_err(|_| ManagedWriterError::ThreadPanic)?;
result.map_err(ManagedWriterError::Completion)?;
self.shutdown_complete = true;
WRITER_LEASED.store(false, Ordering::Release);
Ok(())
}
}
impl<
const BLOCKS: usize,
const BYTES: usize,
const COMMANDS: usize,
const PATH: usize,
const DIRENT_BYTES: usize,
const ENTRIES: usize,
> Drop for ManagedFileWriter<BLOCKS, BYTES, COMMANDS, PATH, DIRENT_BYTES, ENTRIES>
{
fn drop(&mut self) {
if !self.shutdown_complete {
std::process::abort();
}
}
}
#[cfg(test)]
mod tests {
use std::{
os::unix::ffi::OsStrExt,
path::{Path, PathBuf},
sync::MutexGuard,
};
use saddle_admission::{ProcessLedger, ResourceConfig};
use saddle_observability::file::{
DeploymentFilesystemServiceAttestation, FilesystemDeploymentIdentity,
FilesystemServiceRates, FilesystemTarget, FilesystemWorkDomain, FixedFileLimits,
PreparedFixedFileCore, fixed_core_layout, prepare_fixed_file_core,
verify_filesystem_service,
};
use super::*;
const BLOCKS: usize = 2;
const BYTES: usize = 256;
const COMMANDS: usize = 4;
const PATH: usize = 256;
const DIRENT: usize = 4096;
const ENTRIES: usize = 32;
const BUILD_IDENTITY: [u8; 32] = [6; 32];
type Prepared = PreparedFixedFileCore<BLOCKS, BYTES, COMMANDS, PATH, DIRENT, ENTRIES>;
type Managed = ManagedFileWriter<BLOCKS, BYTES, COMMANDS, PATH, DIRENT, ENTRIES>;
static TEST_SERIAL: Mutex<()> = Mutex::new(());
struct TestFilesystemService;
impl DeploymentFilesystemServiceAttestation for TestFilesystemService {
fn target(&self) -> FilesystemTarget {
FilesystemTarget::LinuxX86_64
}
fn approved(&self) -> bool {
true
}
fn reproducible_measurement(&self) -> bool {
true
}
fn conservative_upper_bound(&self) -> bool {
true
}
fn work_domain(&self) -> FilesystemWorkDomain {
FilesystemWorkDomain {
max_write_ops: u64::MAX,
max_write_bytes: u64::MAX,
max_file_sync_ops: u64::MAX,
max_file_sync_bytes: u64::MAX,
max_directory_sync_ops: u64::MAX,
max_hard_link_ops: u64::MAX,
max_rename_ops: 0,
max_unlink_ops: u64::MAX,
max_create_ops: u64::MAX,
max_scan_ops: u64::MAX,
max_scan_entries: u64::MAX,
max_parallel_filesystem_ops: 1,
advisory_lock: true,
hard_link_no_replace: true,
directory_fsync: true,
data_fsync: true,
}
}
fn service_rates(&self) -> FilesystemServiceRates {
FilesystemServiceRates {
write_op_nanos: 1,
write_byte_nanos: 1,
file_sync_op_nanos: 1,
file_sync_byte_nanos: 1,
directory_sync_op_nanos: 1,
hard_link_op_nanos: 1,
rename_op_nanos: 0,
unlink_op_nanos: 1,
create_op_nanos: 1,
scan_op_nanos: 1,
scan_entry_nanos: 1,
}
}
fn calibration_identity(&self) -> [u8; 32] {
[1; 32]
}
fn environment_identity(&self) -> [u8; 32] {
[2; 32]
}
fn filesystem_identity(&self) -> [u8; 32] {
[3; 32]
}
fn mount_identity(&self) -> [u8; 32] {
[4; 32]
}
fn service_attestation(&self) -> [u8; 32] {
[5; 32]
}
fn build_identity(&self) -> [u8; 32] {
BUILD_IDENTITY
}
}
fn lock_tests() -> MutexGuard<'static, ()> {
TEST_SERIAL
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn test_directory(name: &str) -> PathBuf {
std::env::temp_dir().join(format!(
"saddle-runtime-observability-{name}-{}",
std::process::id()
))
}
fn ledger_and_prepared(directory: &Path) -> (ProcessLedger, Prepared) {
let layout = fixed_core_layout::<BLOCKS, BYTES, COMMANDS, PATH, DIRENT, ENTRIES>().unwrap();
let ledger_state = ProcessLedger::minimum_process_state_reserve(1).unwrap();
let queue_state = ProcessLedger::observability_queue_state_reserve(BLOCKS).unwrap();
let process_state_reserve = ledger_state + queue_state + layout.total;
let framework_reserve = BLOCKS * BYTES + 64;
let managed_capacity = BYTES;
let fixed = 1 + framework_reserve + 64 + process_state_reserve + 1 + 1;
let ledger = ProcessLedger::new(ResourceConfig {
managed_capacity,
entry_reserve: 1,
framework_reserve,
task_reserve: 64,
process_state_reserve,
system_estimate: 1,
safety_margin: 1,
process_limit: managed_capacity + fixed,
max_active_requests: 1,
})
.unwrap();
let domain = ledger.prepare_observability_queue(BLOCKS, BYTES).unwrap();
let prepared = prepare_fixed_file_core(
domain,
directory.as_os_str().as_bytes(),
layout.total,
FixedFileLimits::candidate(),
verify_filesystem_service(
TestFilesystemService,
FilesystemDeploymentIdentity {
target: FilesystemTarget::LinuxX86_64,
environment: [2; 32],
filesystem: [3; 32],
mount: [4; 32],
build: BUILD_IDENTITY,
},
)
.unwrap(),
fixed_notifications(),
)
.unwrap();
(ledger, prepared)
}
fn runtime() -> tokio::runtime::Runtime {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
}
#[test]
fn real_entry_start_flush_shutdown_and_join_use_fixed_slots() {
let _serial = lock_tests();
let directory = test_directory("normal");
let _ = std::fs::remove_dir_all(&directory);
std::fs::create_dir_all(&directory).unwrap();
let (ledger, prepared) = ledger_and_prepared(&directory);
let runtime = runtime();
runtime.block_on(async {
let managed = Managed::start_until(prepared, std::future::pending())
.await
.unwrap();
assert!(WRITER_LEASED.load(Ordering::Acquire));
managed.flush().await.unwrap();
managed
.shutdown_until(std::future::pending())
.await
.unwrap();
});
assert!(!WRITER_LEASED.load(Ordering::Acquire));
assert!(NOTIFICATIONS.iter().all(|slot| {
!slot.registered.load(Ordering::Acquire) && slot.waker.lock().unwrap().is_none()
}));
assert!(ledger.try_shutdown().unwrap().healthy);
let _ = std::fs::remove_dir_all(directory);
}
#[test]
fn startup_failure_is_not_published_and_releases_ledger_owner() {
let _serial = lock_tests();
let directory = test_directory("missing-parent").join("logs");
let _ = std::fs::remove_dir_all(directory.parent().expect("test directory has a parent"));
let (ledger, prepared) = ledger_and_prepared(&directory);
let runtime = runtime();
assert!(matches!(
runtime.block_on(Managed::start_until(prepared, std::future::pending())),
Err(ManagedWriterError::Completion(FixedFailure::Directory))
));
assert!(!WRITER_LEASED.load(Ordering::Acquire));
assert!(ledger.try_shutdown().unwrap().healthy);
}
#[test]
fn cancelled_fixed_wait_releases_the_only_role_slot() {
let _serial = lock_tests();
let directory = test_directory("cancel");
let _ = std::fs::remove_dir_all(&directory);
std::fs::create_dir_all(&directory).unwrap();
let (ledger, prepared) = ledger_and_prepared(&directory);
let runtime = runtime();
runtime.block_on(async {
let managed = Managed::start_until(prepared, std::future::pending())
.await
.unwrap();
let waker = Waker::noop();
let mut context = Context::from_waker(waker);
let mut cancelled_pending_wait = false;
for _ in 0..256 {
let completion = managed.sink.try_flush().unwrap();
let mut wait = Box::pin(CompletionWait::new(completion, CompletionRole::Flush));
match wait.as_mut().poll(&mut context) {
Poll::Pending => {
drop(wait);
cancelled_pending_wait = true;
break;
}
Poll::Ready(Ok(())) => {}
Poll::Ready(Err(error)) => {
panic!("flush wait must remain healthy before cancellation: {error:?}")
}
}
}
assert!(
cancelled_pending_wait,
"writer completed every flush before cancellation could be observed"
);
assert!(
!NOTIFICATIONS[CompletionRole::Flush as usize]
.registered
.load(Ordering::Acquire)
);
managed
.shutdown_until(std::future::pending())
.await
.unwrap();
});
assert!(ledger.try_shutdown().unwrap().healthy);
let _ = std::fs::remove_dir_all(directory);
}
}