Skip to content
Open
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
5 changes: 2 additions & 3 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions src/cdp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1741,6 +1741,7 @@ fn endpoint_owner_process_ids(port: u16, _process_id: u32) -> Result<BTreeSet<u3
if written >= descriptors.len() || written % 8 != 0 {
return Err(CdpError::StaleTarget);
}
#[allow(clippy::chunks_exact_to_as_chunks)]
for descriptor in descriptors[..written].chunks_exact(8) {
let file_descriptor =
i32::from_ne_bytes(descriptor[..4].try_into().map_err(|_| CdpError::Protocol)?);
Expand Down Expand Up @@ -1851,6 +1852,7 @@ fn endpoint_owner_process_ids(port: u16, _process_id: u32) -> Result<BTreeSet<u3
return Err(CdpError::Protocol);
}
let mut owners = BTreeSet::new();
#[allow(clippy::chunks_exact_to_as_chunks)]
for row in table[4..].chunks_exact(TCP_ROW_BYTES).take(row_count) {
if u32::from_ne_bytes(row[..4].try_into().map_err(|_| CdpError::Protocol)?) == TCP_LISTEN
&& row[4..8] == Ipv4Addr::LOCALHOST.octets()
Expand Down
93 changes: 58 additions & 35 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6583,41 +6583,16 @@ impl OperationLedger {
return Ok(None);
};
let finished_at_ms = now_ms();
let delivery_route = delivery_route.unwrap_or(DeliveryRoute::Unknown);
let session_isolation = session_isolation.unwrap_or(SessionIsolation::Unknown);
let interaction_mode = interaction_mode.unwrap_or(InteractionMode::Unknown);
let acknowledgement = ActionAck {
protocol_version: PROTOCOL_VERSION,
operation_id: operation_id.to_string(),
sequence: 2,
action_hash: action_hash.clone(),
replayed: false,
state: AckState::Terminal {
terminal: Box::new(Terminal::OutcomeUnknown {
receipt: Receipt {
protocol_version: PROTOCOL_VERSION,
action_name: action_name.unwrap_or_else(|| "unknown".to_string()),
action_hash,
started_at_ms: claimed_at_ms,
finished_at_ms,
backend: "unknown".to_string(),
fallback_chain: Vec::new(),
delivery_route,
session_isolation,
interaction_mode,
context_preservation: recovered_context_preservation(
interaction_mode,
session_isolation,
),
effect: Effect::Unknown,
before: None,
after: None,
warnings: Vec::new(),
},
message: interrupted_outcome_message(),
}),
},
};
let acknowledgement = generate_interrupted_ack(
operation_id,
action_hash,
claimed_at_ms,
finished_at_ms,
action_name,
delivery_route,
session_isolation,
interaction_mode,
);
self.finish(&acknowledgement)?;
Ok(Some(acknowledgement))
}
Expand Down Expand Up @@ -6781,6 +6756,54 @@ fn repair_jsonl_tail<T: DeserializeOwned>(
Ok(())
}

#[allow(clippy::too_many_arguments)]
fn generate_interrupted_ack(
operation_id: &str,
action_hash: String,
claimed_at_ms: i64,
finished_at_ms: i64,
action_name: Option<String>,
delivery_route: Option<DeliveryRoute>,
session_isolation: Option<SessionIsolation>,
interaction_mode: Option<InteractionMode>,
) -> ActionAck {
let delivery_route = delivery_route.unwrap_or(DeliveryRoute::Unknown);
let session_isolation = session_isolation.unwrap_or(SessionIsolation::Unknown);
let interaction_mode = interaction_mode.unwrap_or(InteractionMode::Unknown);
ActionAck {
protocol_version: PROTOCOL_VERSION,
operation_id: operation_id.to_string(),
sequence: 2,
action_hash: action_hash.clone(),
replayed: false,
state: AckState::Terminal {
terminal: Box::new(Terminal::OutcomeUnknown {
receipt: Receipt {
protocol_version: PROTOCOL_VERSION,
action_name: action_name.unwrap_or_else(|| "unknown".to_string()),
action_hash,
started_at_ms: claimed_at_ms,
finished_at_ms,
backend: "unknown".to_string(),
fallback_chain: Vec::new(),
delivery_route,
session_isolation,
interaction_mode,
context_preservation: recovered_context_preservation(
interaction_mode,
session_isolation,
),
effect: Effect::Unknown,
before: None,
after: None,
warnings: Vec::new(),
},
message: interrupted_outcome_message(),
}),
},
}
}

fn persistence_unknown_ack(mut acknowledgement: ActionAck) -> ActionAck {
if let AckState::Terminal { terminal } = &mut acknowledgement.state
&& matches!(&**terminal, Terminal::Succeeded { .. })
Expand Down
2 changes: 1 addition & 1 deletion tests/cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ fn run(arguments: &[&str], stdin: &str) -> std::process::Output {
.stderr(Stdio::piped())
.spawn()
.expect("spawn CLI");
if let Some(mut child_stdin) = child.stdin.take() {
if let Some(mut child_stdin) = child.stdin.take() {
let _ = child_stdin.write_all(stdin.as_bytes());
}
child.wait_with_output().expect("CLI output")
Expand Down
Loading