use std::{borrow::Cow, mem};
use monty::{Dump, MontyRepl, ReplProgress, ReplStartError, Session, SessionRef, dump};
use monty_type_checking::{SourceFile, TypeChecker};
use monty_types::{
AssertMessageAnnotations, CompileOptions, ExcType, ExtFunctionResult, MontyException, MontyObject, OsFunctionCall,
PrintWriter, PrintWriterCallback, ResourceTracker, TypeCheckState, TypeCheckingConfig,
};
use super::{
FrameError, FrameReader, MAX_FRAME_LEN, WireFunctionCall, check_protocol_version, exceeds_max_frame_len,
exceeds_max_value_depth, future_results_from_proto, pb, write_frame,
};
pub trait EventSink {
fn send(&mut self, event: &pb::ChildEvent) -> Result<(), FrameError>;
}
#[derive(Default)]
pub struct VecEventSink {
frames: Vec<u8>,
}
impl VecEventSink {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn take(&mut self) -> Vec<u8> {
mem::take(&mut self.frames)
}
}
impl EventSink for VecEventSink {
fn send(&mut self, event: &pb::ChildEvent) -> Result<(), FrameError> {
write_frame(&mut self.frames, event)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HandleOutcome {
Continue,
Shutdown,
Fatal,
}
pub fn dispatch_frame(child: &mut Child, request_frame: &[u8]) -> (Vec<u8>, HandleOutcome) {
let mut sink = VecEventSink::new();
let outcome = dispatch_into(child, request_frame, &mut sink);
(sink.take(), outcome)
}
fn dispatch_into(child: &mut Child, request_frame: &[u8], sink: &mut VecEventSink) -> HandleOutcome {
let mut reader = FrameReader::new(request_frame);
match reader.read::<pb::ParentRequest>() {
Ok(Some(request)) => match child.handle(request, sink) {
Ok(outcome) => outcome,
Err(FrameError::FrameTooLarge { len, max }) => {
let _ = sink
.send(&child.fatal_event(&format!("response frame of {len} bytes exceeds maximum of {max} bytes")));
HandleOutcome::Shutdown
}
Err(_) => HandleOutcome::Shutdown,
},
Ok(None) => HandleOutcome::Continue,
Err(FrameError::Decode(err)) => {
let _ = sink.send(&protocol_violation(&format!("malformed request: {err}")));
HandleOutcome::Continue
}
Err(err) => {
let _ = sink.send(&child.fatal_event(&format!("malformed request frame: {err}")));
HandleOutcome::Shutdown
}
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct SessionBudget {
pub max_memory: Option<usize>,
pub type_check: bool,
}
enum SessionState {
Configured(Option<Box<pb::Configure>>),
Ready(Box<MontyRepl>),
Suspended(Box<ReplProgress>),
}
pub struct Child {
state: SessionState,
script_name: String,
type_checker: TypeChecker,
type_check: Option<TypeCheckState>,
}
impl Default for Child {
fn default() -> Self {
Self {
state: SessionState::Configured(None),
script_name: String::new(),
type_checker: TypeChecker::default(),
type_check: None,
}
}
}
impl Child {
pub fn handle(
&mut self,
request: pb::ParentRequest,
sink: &mut dyn EventSink,
) -> Result<HandleOutcome, FrameError> {
let Some(kind) = request.kind else {
sink.send(&protocol_violation("request has no kind"))?;
return Ok(HandleOutcome::Continue);
};
let mut event = match kind {
pb::parent_request::Kind::Configure(configure) => {
if let Err(refusal) = check_protocol_version(configure.protocol_version) {
sink.send(&self.fatal_event(&refusal))?;
return Ok(HandleOutcome::Fatal);
}
self.handle_configure(configure)
}
pb::parent_request::Kind::Feed(feed) => self.handle_repl_feed(feed, sink),
pb::parent_request::Kind::InstallDependencies(_) => error_event(
ExcType::RuntimeError,
"dependency installation is only supported by the CPython worker",
),
pb::parent_request::Kind::ResumeCall(resume) => self.handle_resume_call(resume, sink),
pb::parent_request::Kind::ResumeNameLookup(resume) => self.handle_resume_name_lookup(resume, sink),
pb::parent_request::Kind::ResumeFutures(resume) => self.handle_resume_futures(resume, sink),
pb::parent_request::Kind::Dump(_) => self.handle_dump(),
pb::parent_request::Kind::Load(load) => self.handle_load(&load),
pb::parent_request::Kind::Reset(_) => match self.reset() {
Ok(()) => ok_event(),
Err(err) => {
sink.send(&self.fatal_event(&format!("type-check cleanup failed: {err}")))?;
return Ok(HandleOutcome::Fatal);
}
},
pb::parent_request::Kind::Shutdown(_) => {
sink.send(&ok_event())?;
return Ok(HandleOutcome::Shutdown);
}
};
self.stamp_execution_time(&mut event);
if let Err(err) = sink.send(&event) {
self.recover_send_error(&event, err, sink)?;
}
Ok(HandleOutcome::Continue)
}
#[must_use]
pub fn session_budget(&self) -> SessionBudget {
match &self.state {
SessionState::Configured(Some(config)) => SessionBudget {
max_memory: config
.limits
.as_ref()
.and_then(|limits| limits.max_memory_bytes)
.map(|v| usize::try_from(v).unwrap_or(usize::MAX)),
type_check: config.type_check,
},
SessionState::Configured(None) => SessionBudget::default(),
SessionState::Ready(repl) => self.tracker_budget(repl.tracker()),
SessionState::Suspended(progress) => self.tracker_budget(progress.tracker()),
}
}
fn tracker_budget(&self, tracker: &ResourceTracker) -> SessionBudget {
SessionBudget {
max_memory: tracker.max_memory(),
type_check: self.type_check.is_some(),
}
}
#[must_use]
pub fn fatal_event(&self, message: &str) -> pb::ChildEvent {
let mut event = fatal_error_event(message);
self.stamp_execution_time(&mut event);
event
}
fn recover_send_error(
&mut self,
failed: &pb::ChildEvent,
err: FrameError,
sink: &mut dyn EventSink,
) -> Result<(), FrameError> {
let announces_suspension = matches!(
failed.kind,
Some(
pb::child_event::Kind::FunctionCall(_)
| pb::child_event::Kind::OsCall(_)
| pb::child_event::Kind::NameLookup(_)
| pb::child_event::Kind::ResolveFutures(_)
)
);
match err {
FrameError::FrameTooLarge { len, max } if !announces_suspension => {
let mut event = error_event(
ExcType::RuntimeError,
&format!("result frame of {len} bytes exceeds the maximum of {max} bytes"),
);
self.stamp_execution_time(&mut event);
sink.send(&event)
}
other => Err(other),
}
}
fn stamp_execution_time(&self, event: &mut pb::ChildEvent) {
let tracker = match &self.state {
SessionState::Ready(repl) => repl.tracker(),
SessionState::Suspended(progress) => progress.tracker(),
SessionState::Configured(_) => return,
};
event.total_execution_micros = u64::try_from(tracker.elapsed().as_micros()).unwrap_or(u64::MAX);
event.max_duration_micros = tracker
.max_duration()
.map(|max| u64::try_from(max.as_micros()).unwrap_or(u64::MAX));
}
fn handle_configure(&mut self, configure: pb::Configure) -> pb::ChildEvent {
if matches!(self.state, SessionState::Configured(None)) {
self.state = SessionState::Configured(Some(Box::new(configure)));
ok_event()
} else {
protocol_violation("Configure while a session already exists")
}
}
fn ensure_repl(&mut self) -> Result<(), Box<pb::ChildEvent>> {
let config = match &mut self.state {
SessionState::Configured(config) => config.take(),
SessionState::Ready(_) | SessionState::Suspended(_) => return Ok(()),
};
let Some(config) = config else {
return Err(Box::new(protocol_violation("session has not been configured")));
};
let type_check_config = TypeCheckingConfig::from(config.as_ref());
let pb::Configure {
script_name,
limits,
type_check,
type_check_stubs,
assert_message_annotations,
type_check_format: _,
type_check_color: _,
protocol_version: _,
monty_version: _,
} = *config;
let limits = limits.unwrap_or_default().into();
self.script_name = script_name;
self.type_check = type_check.then(|| TypeCheckState {
committed_stubs: type_check_stubs.unwrap_or_default(),
pending_snippet: None,
config: type_check_config,
});
let options = CompileOptions {
assert_message_annotations: assert_message_annotations.map_or_else(
AssertMessageAnnotations::default,
AssertMessageAnnotations::from_max_bytes,
),
};
self.state = SessionState::Ready(Box::new(MontyRepl::new(
&self.script_name,
ResourceTracker::new(limits),
options,
)));
Ok(())
}
fn handle_repl_feed(&mut self, feed: pb::Feed, sink: &mut dyn EventSink) -> pb::ChildEvent {
if let Err(event) = self.ensure_repl() {
return *event;
}
if !matches!(self.state, SessionState::Ready(_)) {
return protocol_violation("Feed without a session ready for input");
}
if !feed.skip_type_check
&& let Some(event) = self.type_check_feed(&feed.code)
{
return event;
}
let inputs = match named_inputs(feed.inputs) {
Ok(inputs) => inputs,
Err(event) => return *event,
};
let SessionState::Ready(repl) = mem::replace(&mut self.state, SessionState::Configured(None)) else {
unreachable!("checked Ready above");
};
if !feed.skip_type_check
&& let Some(state) = &mut self.type_check
{
state.pending_snippet = Some(feed.code.clone());
}
let mut print = ProtoPrint::new(sink);
let result = repl.feed_start(&feed.code, inputs, PrintWriter::Callback(&mut print));
let event = self.drive(result, &mut print);
print.drain();
event
}
fn handle_resume_call(&mut self, resume: pb::ResumeCall, sink: &mut dyn EventSink) -> pb::ChildEvent {
let expected_call_id = match &self.state {
SessionState::Suspended(progress) => match progress.as_ref() {
ReplProgress::FunctionCall(call) => Some(call.call_id),
ReplProgress::OsCall(call) => Some(call.call_id),
_ => None,
},
_ => None,
};
let Some(call_id) = expected_call_id else {
return protocol_violation("ResumeCall without a suspended function/OS call");
};
if resume.call_id != call_id {
return protocol_violation(&format!(
"ResumeCall call_id {} does not match {call_id}",
resume.call_id
));
}
let Some(wire_result) = resume.result else {
return protocol_violation("ResumeCall has no result");
};
let result: ExtFunctionResult =
if matches!(wire_result.kind, Some(pb::ext_function_result::Kind::NotHandled(_))) {
let SessionState::Suspended(progress) = &self.state else {
unreachable!("checked above");
};
let ReplProgress::OsCall(call) = progress.as_ref() else {
return protocol_violation("NotHandled is only valid answering a suspended OS call");
};
ExtFunctionResult::Error(call.function_call.on_no_handler())
} else {
match wire_result.try_into() {
Ok(result) => result,
Err(err) => return protocol_violation(&format!("invalid result: {err}")),
}
};
let SessionState::Suspended(progress) = mem::replace(&mut self.state, SessionState::Configured(None)) else {
unreachable!("checked above");
};
let mut print = ProtoPrint::new(sink);
let outcome = match *progress {
ReplProgress::FunctionCall(call) => call.resume(result, PrintWriter::Callback(&mut print)),
ReplProgress::OsCall(call) => call.resume(result, PrintWriter::Callback(&mut print)),
_ => unreachable!("checked above"),
};
let event = self.drive(outcome, &mut print);
print.drain();
event
}
fn handle_resume_name_lookup(&mut self, resume: pb::ResumeNameLookup, sink: &mut dyn EventSink) -> pb::ChildEvent {
let SessionState::Suspended(progress) = &self.state else {
return protocol_violation("ResumeNameLookup without a suspended name lookup");
};
if !matches!(progress.as_ref(), ReplProgress::NameLookup(_)) {
return protocol_violation("ResumeNameLookup without a suspended name lookup");
}
let result = match resume.try_into() {
Ok(result) => result,
Err(err) => return protocol_violation(&format!("invalid result: {err}")),
};
let SessionState::Suspended(progress) = mem::replace(&mut self.state, SessionState::Configured(None)) else {
unreachable!("checked above");
};
let ReplProgress::NameLookup(lookup) = *progress else {
unreachable!("checked above");
};
let mut print = ProtoPrint::new(sink);
let outcome = lookup.resume(result, PrintWriter::Callback(&mut print));
let event = self.drive(outcome, &mut print);
print.drain();
event
}
fn handle_resume_futures(&mut self, resume: pb::ResumeFutures, sink: &mut dyn EventSink) -> pb::ChildEvent {
let SessionState::Suspended(progress) = &self.state else {
return protocol_violation("ResumeFutures without suspended futures");
};
if !matches!(progress.as_ref(), ReplProgress::ResolveFutures(_)) {
return protocol_violation("ResumeFutures without suspended futures");
}
let results = match future_results_from_proto(resume.results) {
Ok(results) => results,
Err(err) => return protocol_violation(&format!("invalid results: {err}")),
};
let SessionState::Suspended(progress) = mem::replace(&mut self.state, SessionState::Configured(None)) else {
unreachable!("checked above");
};
let ReplProgress::ResolveFutures(state) = *progress else {
unreachable!("checked above");
};
let mut print = ProtoPrint::new(sink);
let outcome = state.resume(results, PrintWriter::Callback(&mut print));
let event = self.drive(outcome, &mut print);
print.drain();
event
}
fn handle_dump(&mut self) -> pb::ChildEvent {
if let Err(event) = self.ensure_repl() {
return *event;
}
let session = match &self.state {
SessionState::Ready(repl) => SessionRef::Idle(repl),
SessionState::Suspended(progress) => SessionRef::Suspended(progress),
SessionState::Configured(_) => unreachable!("ensure_repl materialized the repl or errored"),
};
match dump(&self.script_name, self.type_check.as_ref(), session) {
Ok(state) => event(pb::child_event::Kind::DumpResult(pb::DumpResult { state })),
Err(err) => protocol_violation(&format!("dump failed: {err}")),
}
}
fn handle_load(&mut self, load: &pb::Load) -> pb::ChildEvent {
if !matches!(self.state, SessionState::Configured(_)) {
return protocol_violation("Load requires a session that has not started (a feed has already run)");
}
let restored = match Dump::load(&load.state) {
Ok(restored) => restored,
Err(err) => return protocol_violation(&format!("failed to load session: {err}")),
};
let Dump {
script_name,
type_check,
state,
} = restored;
let mut event = match state {
Session::Idle(repl) => {
self.state = SessionState::Ready(repl);
ok_event()
}
Session::Running(_) => protocol_violation("dump holds a one-shot run, not a repl session"),
Session::Suspended(progress) => match *progress {
ReplProgress::Complete { repl, value } => {
if exceeds_max_value_depth(&value) {
protocol_violation("dump value exceeds the maximum wire depth")
} else {
self.state = SessionState::Ready(Box::new(repl));
complete_event(value)
}
}
progress => {
if suspension_args_too_deep(&progress) {
protocol_violation("dump suspension arguments exceed the maximum wire depth")
} else {
let event = suspension_event(&progress);
if let Some(message) = oversize_suspension_error_message(&event) {
protocol_violation(&message)
} else {
self.state = SessionState::Suspended(Box::new(progress));
event
}
}
}
},
};
if matches!(self.state, SessionState::Ready(_) | SessionState::Suspended(_)) {
self.script_name = script_name;
self.type_check = type_check;
event.restored_script_name = Some(self.script_name.clone());
}
event
}
fn drive(
&mut self,
mut result: Result<ReplProgress, Box<ReplStartError>>,
print: &mut ProtoPrint,
) -> pb::ChildEvent {
loop {
match result {
Ok(ReplProgress::Complete { repl, value }) => {
self.state = SessionState::Ready(Box::new(repl));
if let Some(state) = &mut self.type_check
&& let Some(snippet) = state.pending_snippet.take()
{
state.committed_stubs.push('\n');
state.committed_stubs.push_str(&snippet);
}
if exceeds_max_value_depth(&value) {
return error_event(ExcType::RuntimeError, "Max output depth exceeded");
}
return complete_event(value);
}
Ok(ReplProgress::OsCall(call)) => {
if os_call_args_too_deep(&call) {
let err =
MontyException::new(ExcType::RuntimeError, Some("Max argument depth exceeded".to_owned()));
result = call.resume(ExtFunctionResult::Error(err), PrintWriter::Callback(print));
continue;
}
let event = suspension_event_os_call(&call);
if let Some(message) = oversize_suspension_error_message(&event) {
return self.abort_feed_with_runtime_error(call.into_repl(), &message);
}
self.state = SessionState::Suspended(Box::new(ReplProgress::OsCall(call)));
return event;
}
Ok(ReplProgress::FunctionCall(call)) => {
if function_call_args_too_deep(&call) {
let err =
MontyException::new(ExcType::RuntimeError, Some("Max argument depth exceeded".to_owned()));
result = call.resume(ExtFunctionResult::Error(err), PrintWriter::Callback(print));
continue;
}
let event = suspension_event_function_call(&call);
if let Some(message) = oversize_suspension_error_message(&event) {
return self.abort_feed_with_runtime_error(call.into_repl(), &message);
}
self.state = SessionState::Suspended(Box::new(ReplProgress::FunctionCall(call)));
return event;
}
Ok(progress) => {
let event = suspension_event(&progress);
self.state = SessionState::Suspended(Box::new(progress));
return event;
}
Err(err) => {
self.state = SessionState::Ready(Box::new(err.repl));
if let Some(state) = &mut self.type_check {
state.pending_snippet = None;
}
return event(pb::child_event::Kind::Error(pb::Error {
exception: Some((&err.error).into()),
}));
}
}
}
}
fn abort_feed_with_runtime_error(&mut self, repl: MontyRepl, message: &str) -> pb::ChildEvent {
self.state = SessionState::Ready(Box::new(repl));
if let Some(state) = &mut self.type_check {
state.pending_snippet = None;
}
error_event(ExcType::RuntimeError, message)
}
fn type_check_feed(&mut self, code: &str) -> Option<pb::ChildEvent> {
let state = self.type_check.as_ref()?;
let stubs =
(!state.committed_stubs.is_empty()).then(|| SourceFile::new(&state.committed_stubs, "repl_type_stubs.pyi"));
match self
.type_checker
.run(&SourceFile::new(code, &self.script_name), stubs.as_ref(), state.config)
{
Ok(None) => None,
Ok(Some(diagnostics)) => Some(event(pb::child_event::Kind::TypingError(pb::TypingError {
diagnostics: diagnostics.to_string(),
}))),
Err(err) => Some(protocol_violation(&format!("type checker failed: {err}"))),
}
}
fn reset(&mut self) -> Result<(), String> {
self.state = SessionState::Configured(None);
self.type_check = None;
self.script_name = String::new();
self.type_checker.reset()
}
}
fn event(kind: pb::child_event::Kind) -> pb::ChildEvent {
pb::ChildEvent {
kind: Some(kind),
..Default::default()
}
}
#[must_use]
pub fn protocol_violation(message: &str) -> pb::ChildEvent {
event(pb::child_event::Kind::Error(pb::Error {
exception: Some(pb::RaisedException {
exc_type: ExcType::RuntimeError.to_string(),
message: Some(format!("protocol violation: {message}")),
traceback: vec![],
data: None,
}),
}))
}
#[must_use]
pub fn fatal_error_event(message: &str) -> pb::ChildEvent {
event(pb::child_event::Kind::FatalError(pb::FatalError {
message: message.to_owned(),
}))
}
fn ok_event() -> pb::ChildEvent {
event(pb::child_event::Kind::Ok(pb::Ok {}))
}
fn error_event(exc_type: ExcType, message: &str) -> pb::ChildEvent {
event(pb::child_event::Kind::Error(pb::Error {
exception: Some(pb::RaisedException {
exc_type: exc_type.to_string(),
message: Some(message.to_owned()),
traceback: vec![],
data: None,
}),
}))
}
fn oversize_suspension_error_message(event: &pb::ChildEvent) -> Option<String> {
exceeds_max_frame_len(event)
.map(|len| format!("argument frame of {len} bytes exceeds the maximum of {MAX_FRAME_LEN} bytes"))
}
fn suspension_event_function_call(call: &monty::ReplFunctionCall) -> pb::ChildEvent {
event(pb::child_event::Kind::FunctionCall(WireFunctionCall {
function_name: call.function_name.clone(),
args: call.args.clone(),
kwargs: call.kwargs.clone(),
call_id: call.call_id,
method_call: call.method_call,
}))
}
fn suspension_event_os_call(call: &monty::ReplOsCall) -> pb::ChildEvent {
event(pb::child_event::Kind::OsCall(pb::OsCall {
call_id: call.call_id,
call: Some(call.function_call.clone().into()),
}))
}
fn complete_event(value: MontyObject) -> pb::ChildEvent {
event(pb::child_event::Kind::Complete(pb::Complete {
value: Some(value.into()),
}))
}
fn suspension_args_too_deep(progress: &ReplProgress) -> bool {
match progress {
ReplProgress::FunctionCall(call) => function_call_args_too_deep(call),
ReplProgress::OsCall(call) => os_call_args_too_deep(call),
_ => false,
}
}
fn function_call_args_too_deep(call: &monty::ReplFunctionCall) -> bool {
call.args.iter().any(exceeds_max_value_depth)
|| call
.kwargs
.iter()
.any(|(k, v)| exceeds_max_value_depth(k) || exceeds_max_value_depth(v))
}
fn os_call_args_too_deep(call: &monty::ReplOsCall) -> bool {
match &call.function_call {
OsFunctionCall::Getenv(args) => exceeds_max_value_depth(&args.default),
_ => false,
}
}
fn suspension_event(progress: &ReplProgress) -> pb::ChildEvent {
match progress {
ReplProgress::FunctionCall(call) => suspension_event_function_call(call),
ReplProgress::OsCall(call) => suspension_event_os_call(call),
ReplProgress::NameLookup(lookup) => event(pb::child_event::Kind::NameLookup(pb::NameLookup {
name: lookup.name.clone(),
})),
ReplProgress::ResolveFutures(state) => event(pb::child_event::Kind::ResolveFutures(pb::ResolveFutures {
pending_call_ids: state.pending_call_ids().to_vec(),
})),
ReplProgress::Complete { .. } => unreachable!("Complete is handled before suspension_event"),
}
}
fn named_inputs(inputs: Vec<pb::NamedValue>) -> Result<Vec<(String, MontyObject)>, Box<pb::ChildEvent>> {
inputs
.into_iter()
.map(|input| {
let value = input
.value
.ok_or_else(|| Box::new(protocol_violation(&format!("input {:?} has no value", input.name))))?;
let value = value
.into_object()
.map_err(|err| Box::new(protocol_violation(&format!("invalid input {:?}: {err}", input.name))))?;
Ok((input.name, value))
})
.collect()
}
struct ProtoPrint<'a> {
buf: String,
sink: &'a mut dyn EventSink,
}
impl<'a> ProtoPrint<'a> {
const FLUSH_BYTES: usize = 8 * 1024;
fn new(sink: &'a mut dyn EventSink) -> Self {
Self {
buf: String::new(),
sink,
}
}
fn flush(&mut self) -> Result<(), MontyException> {
if self.buf.is_empty() {
return Ok(());
}
let event = event(pb::child_event::Kind::Print(pb::Print {
stream: pb::PrintStream::Stdout.into(),
text: mem::take(&mut self.buf),
}));
self.sink.send(&event).map_err(|err| {
MontyException::new(
ExcType::RuntimeError,
Some(format!("failed to stream print output: {err}")),
)
})
}
fn maybe_flush(&mut self) -> Result<(), MontyException> {
if self.buf.ends_with('\n') || self.buf.len() >= Self::FLUSH_BYTES {
self.flush()
} else {
Ok(())
}
}
fn drain(&mut self) {
let _ = self.flush();
}
}
impl PrintWriterCallback for ProtoPrint<'_> {
fn stdout_write(&mut self, output: Cow<'_, str>) -> Result<(), MontyException> {
let mut rest = output.as_ref();
while !rest.is_empty() {
let take = floor_char_boundary(rest, Self::FLUSH_BYTES - self.buf.len());
if take == 0 {
self.flush()?;
continue;
}
self.buf.push_str(&rest[..take]);
rest = &rest[take..];
self.maybe_flush()?;
}
Ok(())
}
fn stdout_push(&mut self, end: char) -> Result<(), MontyException> {
self.buf.push(end);
self.maybe_flush()
}
}
fn floor_char_boundary(s: &str, max: usize) -> usize {
if max >= s.len() {
s.len()
} else {
let mut idx = max;
while !s.is_char_boundary(idx) {
idx -= 1;
}
idx
}
}