Skip to content
Merged
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
274 changes: 193 additions & 81 deletions rust/log-service/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -481,6 +481,10 @@ pub struct Metrics {
conditional_write_admission_rejections: opentelemetry::metrics::Counter<u64>,
/// The number of conditional push retries after a stale required fragment start.
conditional_write_required_start_retries: opentelemetry::metrics::Counter<u64>,
/// Time spent waiting to pin the next fragment generation.
conditional_write_fragment_pin_wait_us: opentelemetry::metrics::Histogram<u64>,
/// Time a fragment pin is held open across validation and the pinned append.
conditional_write_fragment_pin_hold_us: opentelemetry::metrics::Histogram<u64>,
/// The number of successful conditional writes with an inserted offset.
conditional_write_success_with_offset: opentelemetry::metrics::Counter<u64>,
/// The number of successful conditional writes with durable contention and no offset.
Expand Down Expand Up @@ -528,6 +532,12 @@ impl Metrics {
conditional_write_required_start_retries: meter
.u64_counter("conditional_write_required_start_retries")
.build(),
conditional_write_fragment_pin_wait_us: meter
.u64_histogram("conditional_write_fragment_pin_wait_us")
.build(),
conditional_write_fragment_pin_hold_us: meter
.u64_histogram("conditional_write_fragment_pin_hold_us")
.build(),
conditional_write_success_with_offset: meter
.u64_counter("conditional_write_success_with_offset")
.build(),
Expand Down Expand Up @@ -2553,100 +2563,140 @@ impl LogServer {
messages.push(buf);
}
let record_count = messages.len() as i32;
let first_inserted_record_offset =
if let Some(conditional_write) = conditional_write.as_ref() {
let admission_predicate = conditional_admission_predicate(conditional_write);
let mut append_result = None;
let mut validation_start =
LogPosition::from_offset(conditional_write.observed_log_offset);
let conditional_push_max_retries = self.config.conditional_push_retry_attempts();
for attempt in 0..conditional_push_max_retries {
let validated_tail = self
.validate_committed_log_for_conditional_write(
topology_name.as_ref(),
collection_id,
conditional_write,
validation_start,
)
.await?;
let append_options = AppendOptions::new(write_id_metadata.clone())
.with_admission_predicate(Arc::clone(&admission_predicate))
.with_required_fragment_start(validated_tail);
match log
.append_many_with_options(messages.clone(), Some(append_options))
.await
{
Ok(offset) => {
self.metrics
.conditional_write_success_with_offset
.add(1, &[]);
tracing::debug!(
%collection_id,
first_inserted_record_offset = offset.offset(),
"conditional write append succeeded"
);
append_result = Some(Some(offset));
break;
}
Err(wal3::Error::LogContentionDurable) => {
self.metrics
.conditional_write_success_without_offset
.add(1, &[]);
tracing::debug!(
%collection_id,
"conditional write append durably contended"
);
append_result = Some(None);
break;
}
Err(wal3::Error::AdmissionRejected) => {
self.metrics
.conditional_write_admission_rejections
.add(1, &[]);
self.metrics
.conditional_write_in_flight_conflicts
.add(1, &[]);
tracing::info!(
%collection_id,
"conditional write rejected by in-flight admission"
);
return Err(Status::aborted(CONDITIONAL_WRITE_CONFLICT_MESSAGE));
}
let message_bytes = messages.iter().map(Vec::len).sum::<usize>();
let first_inserted_record_offset = if let Some(conditional_write) =
conditional_write.as_ref()
{
let admission_predicate = conditional_admission_predicate(conditional_write);
let mut append_result = None;
let mut validation_start =
LogPosition::from_offset(conditional_write.observed_log_offset);
let conditional_push_max_retries = self.config.conditional_push_retry_attempts();
for attempt in 0..conditional_push_max_retries {
// Replicated logs serialize conditional writes in Spanner. S3-backed logs pin
// their in-memory fragment before validation so concurrent transactions can join
// the same publish.
let (fragment_pin, pin_acquired_at) = if topology_name.is_none() {
let pin_started = Instant::now();
let pin_result = log
.acquire_fragment_pin(message_bytes)
.instrument(tracing::info_span!("acquire_conditional_fragment_pin"))
.await;
self.metrics.conditional_write_fragment_pin_wait_us.record(
u64::try_from(pin_started.elapsed().as_micros()).unwrap_or(u64::MAX),
&[],
);
match pin_result {
Ok(pin) => (Some(pin), Some(Instant::now())),
Err(wal3::Error::LogContentionRetry)
if attempt + 1 < conditional_push_max_retries =>
{
validation_start = validated_tail;
self.metrics
.conditional_write_required_start_retries
.add(1, &[]);
tracing::info!(
%collection_id,
attempt = attempt + 1,
"conditional write missed required start; retrying validation"
"conditional write fragment pin contended; retrying"
);
continue;
}
Err(err) => return Err(push_append_error_to_status(err)),
}
} else {
(None, None)
};
let validated_tail = self
.validate_committed_log_for_conditional_write(
topology_name.as_ref(),
collection_id,
conditional_write,
validation_start,
)
.await?;
let append_options = AppendOptions::new(write_id_metadata.clone())
.with_admission_predicate(Arc::clone(&admission_predicate))
.with_required_fragment_start(validated_tail);
let append_outcome = log
.append_many_with_options(messages.clone(), Some(append_options), fragment_pin)
.await;
// The pin is consumed by the append above; the hold spans validation plus the
// pinned append, which is how long publishing was held open for this log.
if let Some(acquired_at) = pin_acquired_at {
self.metrics.conditional_write_fragment_pin_hold_us.record(
u64::try_from(acquired_at.elapsed().as_micros()).unwrap_or(u64::MAX),
&[],
);
}
match append_result {
Some(append_result) => append_result,
None => {
return Err(Status::internal(
"conditional push retry loop exited without an append result",
));
match append_outcome {
Ok(offset) => {
self.metrics
.conditional_write_success_with_offset
.add(1, &[]);
tracing::debug!(
%collection_id,
first_inserted_record_offset = offset.offset(),
"conditional write append succeeded"
);
append_result = Some(Some(offset));
break;
}
Err(wal3::Error::LogContentionDurable) => {
self.metrics
.conditional_write_success_without_offset
.add(1, &[]);
tracing::debug!(
%collection_id,
"conditional write append durably contended"
);
append_result = Some(None);
break;
}
Err(wal3::Error::AdmissionRejected) => {
self.metrics
.conditional_write_admission_rejections
.add(1, &[]);
self.metrics
.conditional_write_in_flight_conflicts
.add(1, &[]);
tracing::info!(
%collection_id,
"conditional write rejected by in-flight admission"
);
return Err(Status::aborted(CONDITIONAL_WRITE_CONFLICT_MESSAGE));
}
Err(wal3::Error::LogContentionRetry)
if attempt + 1 < conditional_push_max_retries =>
{
validation_start = validated_tail;
self.metrics
.conditional_write_required_start_retries
.add(1, &[]);
tracing::info!(
%collection_id,
attempt = attempt + 1,
"conditional write missed required start; retrying validation"
);
}
}
} else {
let append_options = AppendOptions::new(write_id_metadata);
match log
.append_many_with_options(messages, Some(append_options))
.await
{
Ok(offset) => Some(offset),
Err(wal3::Error::LogContentionDurable) => None,
Err(err) => return Err(push_append_error_to_status(err)),
}
};
}
match append_result {
Some(append_result) => append_result,
None => {
return Err(Status::internal(
"conditional push retry loop exited without an append result",
));
}
}
} else {
let append_options = AppendOptions::new(write_id_metadata);
match log
.append_many_with_options(messages, Some(append_options), None)
.await
{
Ok(offset) => Some(offset),
Err(wal3::Error::LogContentionDurable) => None,
Err(err) => return Err(push_append_error_to_status(err)),
}
};
let first_inserted_record_offset = first_inserted_record_offset
.map(|offset| {
offset.offset().try_into().map_err(|_| {
Expand Down Expand Up @@ -6739,32 +6789,59 @@ mod tests {
#[derive(Debug)]
struct FakeLogWriter {
append_results: Mutex<Vec<Result<LogPosition, wal3::Error>>>,
fragment_pin_errors: Mutex<Vec<wal3::Error>>,
fragment_pin_attempts: AtomicU64,
observed_options: Mutex<Vec<Option<AppendOptions>>>,
}

impl FakeLogWriter {
fn new(append_results: Vec<Result<LogPosition, wal3::Error>>) -> Arc<Self> {
Arc::new(Self {
append_results: Mutex::new(append_results),
fragment_pin_errors: Mutex::new(Vec::new()),
fragment_pin_attempts: AtomicU64::new(0),
observed_options: Mutex::new(Vec::new()),
})
}

fn with_fragment_pin_errors(fragment_pin_errors: Vec<wal3::Error>) -> Arc<Self> {
Arc::new(Self {
append_results: Mutex::new(Vec::new()),
fragment_pin_errors: Mutex::new(fragment_pin_errors),
fragment_pin_attempts: AtomicU64::new(0),
observed_options: Mutex::new(Vec::new()),
})
}
}

#[async_trait::async_trait]
impl LogWriterTrait for FakeLogWriter {
async fn acquire_fragment_pin(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Head of line blocking. If one transaction is very slow, and holds the generation open, it'll block all future generations right? Or can multiple generations be open at once?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If one transaction is very slow you will get HOL on the generation. This is not a problem as there are a finite number of requests allowed in at once, and wal3 alredy does HOL blocking on the batch.

&self,
_reserved_bytes: usize,
) -> Result<wal3::FragmentPin, wal3::Error> {
self.fragment_pin_attempts.fetch_add(1, Ordering::Relaxed);
let mut fragment_pin_errors = self.fragment_pin_errors.lock();
if fragment_pin_errors.is_empty() {
unreachable!("fake log writer does not pin fragments")
}
Err(fragment_pin_errors.remove(0))
}

async fn append_with_options(
&self,
message: Vec<u8>,
options: Option<AppendOptions>,
) -> Result<LogPosition, wal3::Error> {
self.append_many_with_options(vec![message], options).await
self.append_many_with_options(vec![message], options, None)
.await
}

async fn append_many_with_options(
&self,
_messages: Vec<Vec<u8>>,
options: Option<AppendOptions>,
_pin: Option<wal3::FragmentPin>,
) -> Result<LogPosition, wal3::Error> {
self.observed_options.lock().push(options);
let mut append_results = self.append_results.lock();
Expand Down Expand Up @@ -7357,6 +7434,41 @@ mod tests {
);
}

#[tokio::test]
async fn conditional_push_retries_fragment_pin_contention() {
let mut log_server = log_server_for_installed_writer_tests();
log_server.config.conditional_push_max_retries = 2;
let collection_id = CollectionUuid::new();
let fake_log = FakeLogWriter::with_fragment_pin_errors(vec![
wal3::Error::LogContentionRetry,
wal3::Error::Backoff,
]);
let fake_log_trait: Arc<dyn LogWriterTrait> = fake_log.clone();
let _handle = install_active_log(&log_server, collection_id, fake_log_trait).await;

let err = log_server
.push_logs(Request::new(conditional_push_logs_request(
"dbname",
collection_id,
&[test_operation_record("doc-1")],
PushLogsCondition {
observed_log_offset: 0,
read_ids: vec![],
},
)))
.await
.expect_err("conditional push should return the second fragment pin error");

assert_eq!(2, fake_log.fragment_pin_attempts.load(Ordering::Relaxed));
assert_eq!(Code::ResourceExhausted, err.code());
assert_eq!(
Some("batching"),
err.metadata()
.get(BACKOFF_REASON_MD_KEY)
.and_then(|value| value.to_str().ok())
);
}

fn run_k8s_async_test(future: impl Future<Output = ()> + Send + 'static) {
let runtime = Runtime::new().unwrap();
std::thread::Builder::new()
Expand Down
Loading
Loading