refactor(streams): type transform finish callback data

This commit is contained in:
ldm0
2026-09-12 01:53:51 +08:00
parent d05fc1199e
commit a36916fdcb
4 changed files with 314 additions and 124 deletions
@@ -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<T>(
pub(in crate::context_bootstrap::stream_adapter) fn publish_required_stream_value<T>(
scope: &mut v8::PinScope<'_, '_>,
value: Option<T>,
operation: &'static str,
@@ -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::<v8::Array>::try_from(data).ok()?;
let writable = data
.get_index(scope, 0)
.and_then(|value| v8::Local::<v8::Object>::try_from(value).ok())?;
let residence = data
.get_index(scope, 1)
.and_then(|value| v8::Local::<v8::Object>::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::<v8::Array>::try_from(data).ok()?;
let readable = data
.get_index(scope, 1)
.and_then(|value| v8::Local::<v8::Object>::try_from(value).ok())?;
let residence = data
.get_index(scope, 2)
.and_then(|value| v8::Local::<v8::Object>::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<v8::Local<'s, v8::Object>> {
let data = v8::Local::<v8::Array>::try_from(data).ok()?;
data.get_index(scope, 0)
.and_then(|value| v8::Local::<v8::Object>::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::<v8::Array>::try_from(data).ok()?;
let writable = data
.get_index(scope, 0)
.and_then(|value| v8::Local::<v8::Object>::try_from(value).ok())?;
let readable = data
.get_index(scope, 1)
.and_then(|value| v8::Local::<v8::Object>::try_from(value).ok())?;
let residence = data
.get_index(scope, 2)
.and_then(|value| v8::Local::<v8::Object>::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>,
@@ -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<v8::Local<'s, v8::Object>> {
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<Self> {
let data = v8::Local::<v8::Object>::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<Self> {
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<v8::Local<'s, v8::Object>> {
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<Self> {
let context = (|| {
let finish = TransformFinishContext::decode(scope, data)?;
let data = v8::Local::<v8::Object>::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();
}
}
@@ -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();