diff --git a/moli-renderer-v8/src/context_bootstrap/stream_adapter/utils.rs b/moli-renderer-v8/src/context_bootstrap/stream_adapter/utils.rs index 47dda2f1bd..bd2634eac6 100644 --- a/moli-renderer-v8/src/context_bootstrap/stream_adapter/utils.rs +++ b/moli-renderer-v8/src/context_bootstrap/stream_adapter/utils.rs @@ -121,7 +121,7 @@ pub(in crate::context_bootstrap) fn publish_required_stream_promise_reactions<'s publish_required_stream_value(scope, None, "promise reaction attachment", role) } -fn publish_required_stream_value( +pub(in crate::context_bootstrap::stream_adapter) fn publish_required_stream_value( scope: &mut v8::PinScope<'_, '_>, value: Option, operation: &'static str, diff --git a/moli-renderer-v8/src/context_bootstrap/stream_adapter/writable.rs b/moli-renderer-v8/src/context_bootstrap/stream_adapter/writable.rs index cb172ff12f..ee2ae57711 100644 --- a/moli-renderer-v8/src/context_bootstrap/stream_adapter/writable.rs +++ b/moli-renderer-v8/src/context_bootstrap/stream_adapter/writable.rs @@ -1,6 +1,8 @@ use super::*; +mod finish_context; use crate::text_codec::{TextCodecStore, TextDecodeError}; use crate::util::get_private_value; +use finish_context::{TransformFinishContext, TransformFinishWithReason}; use moli_streams::queue::{QueueBounds, QueueRemainderPlan}; use moli_streams::strategy::StrategySnapshot; use moli_streams::transform::{ @@ -760,7 +762,14 @@ pub(in crate::context_bootstrap::stream_adapter) fn transform_stream_readable_ca rv.set(finish_promise.into()); return; }; - attach_transform_source_cancel_reactions(scope, cancel_promise, writable, residence, reason); + attach_transform_source_cancel_reactions( + scope, + cancel_promise, + writable, + readable, + residence, + reason, + ); rv.set(finish_promise.into()); } @@ -902,13 +911,21 @@ fn attach_transform_source_cancel_reactions<'s>( scope: &mut v8::PinScope<'s, '_>, cancel_promise: v8::Local<'s, v8::Promise>, writable: v8::Local<'s, v8::Object>, + readable: v8::Local<'s, v8::Object>, residence: v8::Local<'s, v8::Object>, reason: v8::Local<'s, v8::Value>, ) { - let data = v8::Array::new(scope, 3); - let _ = data.set_index(scope, 0, writable.into()); - let _ = data.set_index(scope, 1, residence.into()); - let _ = data.set_index(scope, 2, reason); + let StreamOwnerPublication::Published(data) = (TransformFinishWithReason { + finish: TransformFinishContext { + writable, + readable, + residence, + }, + reason, + }) + .into_callback_data(scope) else { + return; + }; publish_required_stream_promise_reactions( scope, cancel_promise, @@ -926,20 +943,19 @@ fn transform_source_cancel_fulfilled_callback<'s>( args: v8::FunctionCallbackArguments<'s>, mut rv: v8::ReturnValue<'_, v8::Value>, ) { - let Some((writable, residence, reason)) = - transform_source_cancel_reaction_values(scope, args.data()) + let StreamOwnerPublication::Published(TransformFinishWithReason { finish, reason }) = + TransformFinishWithReason::from_callback_data(scope, args.data()) else { rv.set_undefined(); return; }; - let Some(readable) = stream_slot_object(scope, writable, WRITABLE_STREAM_TARGET_READABLE_SLOT) - .filter(|readable| !readable.is_null_or_undefined()) - else { - reject_pending_read(scope, residence, reason); - rv.set_undefined(); - return; - }; - apply_transform_source_cancel_fulfillment(scope, writable, readable, residence, reason); + apply_transform_source_cancel_fulfillment( + scope, + finish.writable, + finish.readable, + finish.residence, + reason, + ); rv.set_undefined(); } @@ -970,20 +986,18 @@ fn transform_source_cancel_rejected_callback<'s>( args: v8::FunctionCallbackArguments<'s>, mut rv: v8::ReturnValue<'_, v8::Value>, ) { - let Some((writable, residence, _)) = - transform_source_cancel_reaction_values(scope, args.data()) + let StreamOwnerPublication::Published(context) = + TransformFinishWithReason::from_callback_data(scope, args.data()) else { rv.set_undefined(); return; }; + let TransformFinishContext { + writable, + readable, + residence, + } = context.finish; let error = args.get(0); - let Some(readable) = stream_slot_object(scope, writable, WRITABLE_STREAM_TARGET_READABLE_SLOT) - .filter(|readable| !readable.is_null_or_undefined()) - else { - reject_pending_read(scope, residence, error); - rv.set_undefined(); - return; - }; match transform_stream_snapshot(scope, writable, readable) .plan_finish_settlement(FinishOperation::ReadableCancel, AlgorithmOutcome::Rejected) { @@ -996,27 +1010,6 @@ fn transform_source_cancel_rejected_callback<'s>( rv.set_undefined(); } -fn transform_source_cancel_reaction_values<'s>( - scope: &mut v8::PinScope<'s, '_>, - data: v8::Local<'s, v8::Value>, -) -> Option<( - v8::Local<'s, v8::Object>, - v8::Local<'s, v8::Object>, - v8::Local<'s, v8::Value>, -)> { - let data = v8::Local::::try_from(data).ok()?; - let writable = data - .get_index(scope, 0) - .and_then(|value| v8::Local::::try_from(value).ok())?; - let residence = data - .get_index(scope, 1) - .and_then(|value| v8::Local::::try_from(value).ok())?; - let reason = data - .get_index(scope, 2) - .unwrap_or_else(|| v8::undefined(scope).into()); - Some((writable, residence, reason)) -} - fn transform_stream_sink_abort_algorithm<'s>( scope: &mut v8::PinScope<'s, '_>, writable: v8::Local<'s, v8::Object>, @@ -1069,11 +1062,17 @@ fn transform_stream_sink_abort_algorithm_in_relevant_realm<'s>( ); return Some(finish_promise.into()); }; - let data = v8::Array::new(scope, 4); - let _ = data.set_index(scope, 0, writable.into()); - let _ = data.set_index(scope, 1, readable.into()); - let _ = data.set_index(scope, 2, residence.into()); - let _ = data.set_index(scope, 3, reason); + let StreamOwnerPublication::Published(data) = (TransformFinishWithReason { + finish: TransformFinishContext { + writable, + readable, + residence, + }, + reason, + }) + .into_callback_data(scope) else { + return Some(finish_promise.into()); + }; if matches!( publish_required_stream_promise_reactions( scope, @@ -1096,18 +1095,19 @@ fn transform_sink_abort_fulfilled_callback<'s>( args: v8::FunctionCallbackArguments<'s>, mut rv: v8::ReturnValue<'_, v8::Value>, ) { - let Some((readable, residence, reason)) = - transform_sink_abort_reaction_values(scope, args.data()) + let StreamOwnerPublication::Published(TransformFinishWithReason { finish, reason }) = + TransformFinishWithReason::from_callback_data(scope, args.data()) else { rv.set_undefined(); return; }; - let Some(writable) = transform_sink_abort_reaction_writable(scope, args.data()) else { - reject_pending_read(scope, residence, reason); - rv.set_undefined(); - return; - }; - apply_transform_sink_abort_fulfillment(scope, writable, readable, residence, reason); + apply_transform_sink_abort_fulfillment( + scope, + finish.writable, + finish.readable, + finish.residence, + reason, + ); rv.set_undefined(); } @@ -1138,17 +1138,18 @@ fn transform_sink_abort_rejected_callback<'s>( args: v8::FunctionCallbackArguments<'s>, mut rv: v8::ReturnValue<'_, v8::Value>, ) { - let Some((readable, residence, _)) = transform_sink_abort_reaction_values(scope, args.data()) + let StreamOwnerPublication::Published(context) = + TransformFinishWithReason::from_callback_data(scope, args.data()) else { rv.set_undefined(); return; }; + let TransformFinishContext { + writable, + readable, + residence, + } = context.finish; let error = args.get(0); - let Some(writable) = transform_sink_abort_reaction_writable(scope, args.data()) else { - reject_pending_read(scope, residence, error); - rv.set_undefined(); - return; - }; match transform_stream_snapshot(scope, writable, readable) .plan_finish_settlement(FinishOperation::WritableAbort, AlgorithmOutcome::Rejected) { @@ -1161,36 +1162,6 @@ fn transform_sink_abort_rejected_callback<'s>( rv.set_undefined(); } -fn transform_sink_abort_reaction_values<'s>( - scope: &mut v8::PinScope<'s, '_>, - data: v8::Local<'s, v8::Value>, -) -> Option<( - v8::Local<'s, v8::Object>, - v8::Local<'s, v8::Object>, - v8::Local<'s, v8::Value>, -)> { - let data = v8::Local::::try_from(data).ok()?; - let readable = data - .get_index(scope, 1) - .and_then(|value| v8::Local::::try_from(value).ok())?; - let residence = data - .get_index(scope, 2) - .and_then(|value| v8::Local::::try_from(value).ok())?; - let reason = data - .get_index(scope, 3) - .unwrap_or_else(|| v8::undefined(scope).into()); - Some((readable, residence, reason)) -} - -fn transform_sink_abort_reaction_writable<'s>( - scope: &mut v8::PinScope<'s, '_>, - data: v8::Local<'s, v8::Value>, -) -> Option> { - let data = v8::Local::::try_from(data).ok()?; - data.get_index(scope, 0) - .and_then(|value| v8::Local::::try_from(value).ok()) -} - fn process_transform_writable_queue<'s>( scope: &mut v8::PinScope<'s, '_>, writable: v8::Local<'s, v8::Object>, @@ -2690,10 +2661,14 @@ fn attach_transform_sink_close_reactions<'s>( readable: v8::Local<'s, v8::Object>, residence: v8::Local<'s, v8::Object>, ) { - let data = v8::Array::new(scope, 3); - let _ = data.set_index(scope, 0, writable.into()); - let _ = data.set_index(scope, 1, readable.into()); - let _ = data.set_index(scope, 2, residence.into()); + let StreamOwnerPublication::Published(data) = (TransformFinishContext { + writable, + readable, + residence, + }) + .into_callback_data(scope) else { + return; + }; publish_required_stream_promise_reactions( scope, flush_promise, @@ -2711,8 +2686,11 @@ fn transform_sink_close_fulfilled_callback<'s>( args: v8::FunctionCallbackArguments<'s>, mut rv: v8::ReturnValue<'_, v8::Value>, ) { - let Some((writable, readable, residence)) = - transform_sink_close_reaction_values(scope, args.data()) + let StreamOwnerPublication::Published(TransformFinishContext { + writable, + readable, + residence, + }) = TransformFinishContext::from_callback_data(scope, args.data()) else { rv.set_undefined(); return; @@ -2742,8 +2720,11 @@ fn transform_sink_close_rejected_callback<'s>( args: v8::FunctionCallbackArguments<'s>, mut rv: v8::ReturnValue<'_, v8::Value>, ) { - let Some((writable, readable, residence)) = - transform_sink_close_reaction_values(scope, args.data()) + let StreamOwnerPublication::Published(TransformFinishContext { + writable, + readable, + residence, + }) = TransformFinishContext::from_callback_data(scope, args.data()) else { rv.set_undefined(); return; @@ -2761,27 +2742,6 @@ fn transform_sink_close_rejected_callback<'s>( rv.set_undefined(); } -fn transform_sink_close_reaction_values<'s>( - scope: &mut v8::PinScope<'s, '_>, - data: v8::Local<'s, v8::Value>, -) -> Option<( - v8::Local<'s, v8::Object>, - v8::Local<'s, v8::Object>, - v8::Local<'s, v8::Object>, -)> { - let data = v8::Local::::try_from(data).ok()?; - let writable = data - .get_index(scope, 0) - .and_then(|value| v8::Local::::try_from(value).ok())?; - let readable = data - .get_index(scope, 1) - .and_then(|value| v8::Local::::try_from(value).ok())?; - let residence = data - .get_index(scope, 2) - .and_then(|value| v8::Local::::try_from(value).ok())?; - Some((writable, readable, residence)) -} - fn attach_transform_writable_close_settlement<'s>( scope: &mut v8::PinScope<'s, '_>, writable: v8::Local<'s, v8::Object>, diff --git a/moli-renderer-v8/src/context_bootstrap/stream_adapter/writable/finish_context.rs b/moli-renderer-v8/src/context_bootstrap/stream_adapter/writable/finish_context.rs new file mode 100644 index 0000000000..bf30a0577b --- /dev/null +++ b/moli-renderer-v8/src/context_bootstrap/stream_adapter/writable/finish_context.rs @@ -0,0 +1,169 @@ +//! GC-visible data shared by the two reactions of a transform finish operation. +//! These Rust views never outlive their HandleScope; the V8 callback owns the +//! private-slot carrier across promise jobs. + +use moli_webapi_declare::WebApiObject; + +use super::super::utils::{StreamOwnerPublication, publish_required_stream_value}; +use crate::util::{get_private_object, private_key}; + +const WRITABLE_SLOT: &str = "__moliTransformFinishWritable"; +const READABLE_SLOT: &str = "__moliTransformFinishReadable"; +const RESIDENCE_SLOT: &str = "__moliTransformFinishResidence"; +const REASON_SLOT: &str = "__moliTransformFinishReason"; + +#[derive(WebApiObject)] +#[webapi(plain)] +struct FinishDeclaration<'s> { + #[webapi(slot = WRITABLE_SLOT)] + writable: v8::Local<'s, v8::Object>, + #[webapi(slot = READABLE_SLOT)] + readable: v8::Local<'s, v8::Object>, + #[webapi(slot = RESIDENCE_SLOT)] + residence: v8::Local<'s, v8::Object>, +} + +#[derive(WebApiObject)] +#[webapi(fragment)] +struct ReasonDeclaration<'s> { + #[webapi(slot = REASON_SLOT)] + reason: v8::Local<'s, v8::Value>, +} + +pub(super) struct TransformFinishContext<'s> { + pub(super) writable: v8::Local<'s, v8::Object>, + pub(super) readable: v8::Local<'s, v8::Object>, + pub(super) residence: v8::Local<'s, v8::Object>, +} + +impl<'s> TransformFinishContext<'s> { + pub(super) fn into_callback_data( + self, + scope: &mut v8::PinScope<'s, '_>, + ) -> StreamOwnerPublication> { + let data = FinishDeclaration::new(self.writable, self.readable, self.residence) + .bind(scope) + .ok(); + publish_required_stream_value(scope, data, "callback data creation", "transform finish") + } + + fn decode(scope: &mut v8::PinScope<'s, '_>, data: v8::Local<'s, v8::Value>) -> Option { + let data = v8::Local::::try_from(data).ok()?; + Some(Self { + writable: get_private_object(scope, data, WRITABLE_SLOT)?, + readable: get_private_object(scope, data, READABLE_SLOT)?, + residence: get_private_object(scope, data, RESIDENCE_SLOT)?, + }) + } + + pub(super) fn from_callback_data( + scope: &mut v8::PinScope<'s, '_>, + data: v8::Local<'s, v8::Value>, + ) -> StreamOwnerPublication { + let context = Self::decode(scope, data); + publish_required_stream_value(scope, context, "callback data decoding", "transform finish") + } +} + +pub(super) struct TransformFinishWithReason<'s> { + pub(super) finish: TransformFinishContext<'s>, + pub(super) reason: v8::Local<'s, v8::Value>, +} + +impl<'s> TransformFinishWithReason<'s> { + pub(super) fn into_callback_data( + self, + scope: &mut v8::PinScope<'s, '_>, + ) -> StreamOwnerPublication> { + let StreamOwnerPublication::Published(data) = self.finish.into_callback_data(scope) else { + return StreamOwnerPublication::OwnerTerminating; + }; + let data = ReasonDeclaration::new(self.reason) + .initialize(scope, data) + .ok() + .map(|_| data); + publish_required_stream_value( + scope, + data, + "callback data creation", + "transform finish reason", + ) + } + + pub(super) fn from_callback_data( + scope: &mut v8::PinScope<'s, '_>, + data: v8::Local<'s, v8::Value>, + ) -> StreamOwnerPublication { + let context = (|| { + let finish = TransformFinishContext::decode(scope, data)?; + let data = v8::Local::::try_from(data).ok()?; + // get_private_value treats undefined as absent. A cancel/abort + // reason may be undefined, but the slot itself must be present. + let key = private_key(scope, REASON_SLOT)?; + if data.has_private(scope, key) != Some(true) { + return None; + } + let reason = data.get_private(scope, key)?; + Some(Self { finish, reason }) + })(); + publish_required_stream_value( + scope, + context, + "callback data decoding", + "transform finish with reason", + ) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn finish<'s>(scope: &mut v8::PinScope<'s, '_>) -> TransformFinishContext<'s> { + TransformFinishContext { + writable: v8::Object::new(scope), + readable: v8::Object::new(scope), + residence: v8::Object::new(scope), + } + } + + #[test] + fn undefined_cancel_reason_survives_callback_data() { + moli_v8_test_util::ensure_v8(); + let mut isolate = v8::Isolate::new(Default::default()); + let scope = std::pin::pin!(v8::HandleScope::new(&mut isolate)); + let scope = &mut scope.init(); + let context = v8::Context::new(scope, Default::default()); + let scope = &mut v8::ContextScope::new(scope, context); + let context = TransformFinishWithReason { + finish: finish(scope), + reason: v8::undefined(scope).into(), + }; + let StreamOwnerPublication::Published(data) = context.into_callback_data(scope) else { + panic!("live context must publish callback data") + }; + let StreamOwnerPublication::Published(context) = + TransformFinishWithReason::from_callback_data(scope, data.into()) + else { + panic!("undefined is a valid cancel reason") + }; + assert!(context.reason.is_undefined()); + } + + #[test] + #[should_panic(expected = "callback data decoding for `transform finish with reason`")] + fn close_context_cannot_supply_a_missing_cancel_reason() { + moli_v8_test_util::ensure_v8(); + let mut isolate = v8::Isolate::new(Default::default()); + let scope = std::pin::pin!(v8::HandleScope::new(&mut isolate)); + let scope = &mut scope.init(); + let context = v8::Context::new(scope, Default::default()); + let scope = &mut v8::ContextScope::new(scope, context); + let StreamOwnerPublication::Published(data) = finish(scope).into_callback_data(scope) + else { + panic!("live context must publish callback data") + }; + TransformFinishWithReason::from_callback_data(scope, data.into()) + .finish_at_owner_boundary(); + } +} diff --git a/moli-renderer-v8/src/script_vm/tests/streams/transform_finish.rs b/moli-renderer-v8/src/script_vm/tests/streams/transform_finish.rs index 5a26c7eae1..4df2a609ab 100644 --- a/moli-renderer-v8/src/script_vm/tests/streams/transform_finish.rs +++ b/moli-renderer-v8/src/script_vm/tests/streams/transform_finish.rs @@ -1,5 +1,66 @@ use super::*; +#[test] +fn transform_finish_reactions_keep_their_context_alive_across_gc() { + for operation in ["cancel", "abort", "close"] { + let mut vm = stream_test_vm(); + vm.eval(&format!("globalThis.__finishOperation = {operation:?}")) + .expect("finish operation"); + vm.eval( + r#" +(() => { + const state = globalThis.__finishState = { calls: 0, closed: false, finished: false }; + const gate = new Promise(resolve => { globalThis.__releaseFinish = resolve; }); + const stream = new TransformStream({ + cancel() { state.calls++; return gate; }, + flush() { state.calls++; return gate; } + }); + const reason = { marker: 'original reason' }; + const expectedReason = new WeakRef(reason); + const operation = __finishOperation; + const closed = operation === 'cancel' + ? stream.writable.getWriter().closed + : stream.readable.getReader().closed; + closed.then( + () => { state.closed = operation === 'close'; }, + error => { state.closed = error === expectedReason.deref() && error.marker === 'original reason'; } + ); + const finish = operation === 'cancel' ? stream.readable.cancel(reason) + : operation === 'abort' ? stream.writable.abort(reason) : stream.writable.close(); + finish.then( + () => { state.finished = true; }, + error => { state.finished = String(error); } + ); +})() +"#, + ) + .expect("pending finish setup"); + assert_eq!( + vm.eval("JSON.stringify(__finishState)") + .expect("pending finish state"), + r#"{"calls":1,"closed":false,"finished":false}"#, + "{operation} must wait for its callback promise" + ); + + vm.renderer_document_isolate + .clone() + .with_entered_renderer_document_isolate(|isolate| { + isolate.low_memory_notification(); + Ok(()) + }) + .expect("collect while only the pending reactions retain their context"); + + vm.eval("__releaseFinish()") + .expect("release finish callback promise"); + assert_eq!( + vm.eval("JSON.stringify(__finishState)") + .expect("finish after GC"), + r#"{"calls":1,"closed":true,"finished":true}"#, + "{operation} must retain its streams, residence and original reason" + ); + } +} + #[test] fn transform_finish_pending_cancel_observes_writable_error_at_fulfillment() { let mut vm = stream_test_vm();