Compare commits
2
Commits
4a276b0af0
...
e6bfb27fa9
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e6bfb27fa9 | ||
|
|
19f506f8bc |
+139
-17
@@ -22,7 +22,7 @@ use async_trait::async_trait;
|
|||||||
use manifest::{Permission, Scope, ScopeConfig, ScopeRule, SharedScope};
|
use manifest::{Permission, Scope, ScopeConfig, ScopeRule, SharedScope};
|
||||||
use sha2::{Digest, Sha256};
|
use sha2::{Digest, Sha256};
|
||||||
use tokio::process::Command;
|
use tokio::process::Command;
|
||||||
use tokio::sync::{Mutex, Notify, broadcast, watch};
|
use tokio::sync::{Mutex, broadcast, watch};
|
||||||
use tokio::task::JoinHandle;
|
use tokio::task::JoinHandle;
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
@@ -54,7 +54,7 @@ fn command_observed_at_ms() -> u64 {
|
|||||||
enum LocalCommand {
|
enum LocalCommand {
|
||||||
Running {
|
Running {
|
||||||
task: JoinHandle<Result<CommandOutput, WorkdirError>>,
|
task: JoinHandle<Result<CommandOutput, WorkdirError>>,
|
||||||
completion: Arc<Notify>,
|
completion: watch::Receiver<bool>,
|
||||||
cancel: watch::Sender<bool>,
|
cancel: watch::Sender<bool>,
|
||||||
},
|
},
|
||||||
Completed(CommandOutput),
|
Completed(CommandOutput),
|
||||||
@@ -666,21 +666,19 @@ impl WorkdirSession for LocalWorkdirSession {
|
|||||||
let id = self.inner.next_command_id.fetch_add(1, Ordering::Relaxed);
|
let id = self.inner.next_command_id.fetch_add(1, Ordering::Relaxed);
|
||||||
let handle = CommandHandle(format!("command-{id}"));
|
let handle = CommandHandle(format!("command-{id}"));
|
||||||
let cwd = self.inner.cwd.clone();
|
let cwd = self.inner.cwd.clone();
|
||||||
let completion = Arc::new(Notify::new());
|
let (completion_tx, completion) = watch::channel(false);
|
||||||
let task_completion = Arc::clone(&completion);
|
|
||||||
let command_id = handle.0.clone();
|
let command_id = handle.0.clone();
|
||||||
let telemetry = self.inner.command_telemetry.clone();
|
let telemetry = self.inner.command_telemetry.clone();
|
||||||
let (cancel, cancel_rx) = watch::channel(false);
|
let (cancel, cancel_rx) = watch::channel(false);
|
||||||
let task = tokio::spawn(async move {
|
let task = tokio::spawn(async move {
|
||||||
let output = run_command(cwd, request, command_id, telemetry, cancel_rx).await;
|
let output = run_command(cwd, request, command_id, telemetry, cancel_rx).await;
|
||||||
task_completion.notify_one();
|
let _ = completion_tx.send(true);
|
||||||
output
|
output
|
||||||
});
|
});
|
||||||
let mut commands = self.inner.commands.lock().await;
|
let mut commands = self.inner.commands.lock().await;
|
||||||
if let Err(error) = self.ensure_open() {
|
if let Err(error) = self.ensure_open() {
|
||||||
let _ = cancel.send(true);
|
let _ = cancel.send(true);
|
||||||
task.abort();
|
task.abort();
|
||||||
completion.notify_one();
|
|
||||||
return Err(error);
|
return Err(error);
|
||||||
}
|
}
|
||||||
commands.insert(
|
commands.insert(
|
||||||
@@ -725,12 +723,14 @@ impl WorkdirSession for LocalWorkdirSession {
|
|||||||
return Err(WorkdirError::UnknownCommand(request.handle.0.clone()));
|
return Err(WorkdirError::UnknownCommand(request.handle.0.clone()));
|
||||||
};
|
};
|
||||||
let completion = match command {
|
let completion = match command {
|
||||||
LocalCommand::Running {
|
LocalCommand::Running { completion, .. }
|
||||||
task, completion, ..
|
if !*completion.borrow() && completion.has_changed().is_ok() =>
|
||||||
} if !task.is_finished() => Some(Arc::clone(completion)),
|
{
|
||||||
|
Some(completion.clone())
|
||||||
|
}
|
||||||
_ => None,
|
_ => None,
|
||||||
};
|
};
|
||||||
if let Some(completion) = completion {
|
if let Some(mut completion) = completion {
|
||||||
if !request.wait {
|
if !request.wait {
|
||||||
return Ok(CommandOutput {
|
return Ok(CommandOutput {
|
||||||
status: CommandStatus::Running,
|
status: CommandStatus::Running,
|
||||||
@@ -742,9 +742,21 @@ impl WorkdirSession for LocalWorkdirSession {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
drop(commands);
|
drop(commands);
|
||||||
completion.notified().await;
|
let _ = completion.changed().await;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
if !request.wait
|
||||||
|
&& matches!(command, LocalCommand::Running { task, .. } if !task.is_finished())
|
||||||
|
{
|
||||||
|
return Ok(CommandOutput {
|
||||||
|
status: CommandStatus::Running,
|
||||||
|
exit_code: None,
|
||||||
|
timed_out: false,
|
||||||
|
content: String::new(),
|
||||||
|
next_cursor: None,
|
||||||
|
truncated: false,
|
||||||
|
});
|
||||||
|
}
|
||||||
break commands
|
break commands
|
||||||
.remove(&request.handle.0)
|
.remove(&request.handle.0)
|
||||||
.expect("command checked above");
|
.expect("command checked above");
|
||||||
@@ -821,14 +833,9 @@ impl WorkdirSession for LocalWorkdirSession {
|
|||||||
};
|
};
|
||||||
for command in commands {
|
for command in commands {
|
||||||
match command {
|
match command {
|
||||||
LocalCommand::Running {
|
LocalCommand::Running { task, cancel, .. } => {
|
||||||
task,
|
|
||||||
completion,
|
|
||||||
cancel,
|
|
||||||
} => {
|
|
||||||
let _ = cancel.send(true);
|
let _ = cancel.send(true);
|
||||||
let _ = task.await;
|
let _ = task.await;
|
||||||
completion.notify_one();
|
|
||||||
}
|
}
|
||||||
LocalCommand::Completed(_) => {}
|
LocalCommand::Completed(_) => {}
|
||||||
}
|
}
|
||||||
@@ -2048,6 +2055,121 @@ mod tests {
|
|||||||
assert!(decoder.pending.is_empty());
|
assert!(decoder.pending.is_empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn command_output_does_not_rewait_before_join_handle_finishes() {
|
||||||
|
let dir = TempDir::new().unwrap();
|
||||||
|
let workdir = make_fs(&dir);
|
||||||
|
let handle = CommandHandle("command-completion-race".into());
|
||||||
|
let (completion_tx, completion) = watch::channel(false);
|
||||||
|
let completion_observer = completion_tx.clone();
|
||||||
|
let (cancel, _cancel_rx) = watch::channel(false);
|
||||||
|
let (start_tx, start_rx) = tokio::sync::oneshot::channel::<()>();
|
||||||
|
let (completion_sent_tx, completion_sent_rx) = tokio::sync::oneshot::channel::<()>();
|
||||||
|
let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
|
||||||
|
let task = tokio::spawn(async move {
|
||||||
|
start_rx.await.unwrap();
|
||||||
|
completion_tx.send(true).unwrap();
|
||||||
|
completion_sent_tx.send(()).unwrap();
|
||||||
|
release_rx.await.unwrap();
|
||||||
|
Ok(CommandOutput {
|
||||||
|
status: CommandStatus::Completed,
|
||||||
|
exit_code: Some(0),
|
||||||
|
timed_out: false,
|
||||||
|
content: "done".into(),
|
||||||
|
next_cursor: None,
|
||||||
|
truncated: false,
|
||||||
|
})
|
||||||
|
});
|
||||||
|
workdir.inner.commands.lock().await.insert(
|
||||||
|
handle.0.clone(),
|
||||||
|
LocalCommand::Running {
|
||||||
|
task,
|
||||||
|
completion,
|
||||||
|
cancel,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
let waiting_workdir = workdir.clone();
|
||||||
|
let waiting_handle = handle.clone();
|
||||||
|
let waiter = tokio::spawn(async move {
|
||||||
|
WorkdirSession::command_output(
|
||||||
|
&waiting_workdir,
|
||||||
|
CommandOutputRequest {
|
||||||
|
handle: waiting_handle,
|
||||||
|
cursor: 0,
|
||||||
|
limit: 1024,
|
||||||
|
wait: true,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
for _ in 0..100 {
|
||||||
|
if completion_observer.receiver_count() >= 2 {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
tokio::task::yield_now().await;
|
||||||
|
}
|
||||||
|
assert_eq!(
|
||||||
|
completion_observer.receiver_count(),
|
||||||
|
2,
|
||||||
|
"waiter must subscribe before completion is published"
|
||||||
|
);
|
||||||
|
|
||||||
|
start_tx.send(()).unwrap();
|
||||||
|
completion_sent_rx.await.unwrap();
|
||||||
|
tokio::task::yield_now().await;
|
||||||
|
assert!(
|
||||||
|
!waiter.is_finished(),
|
||||||
|
"command_output should await the task after observing completion state"
|
||||||
|
);
|
||||||
|
release_tx.send(()).unwrap();
|
||||||
|
|
||||||
|
let output = tokio::time::timeout(Duration::from_secs(1), waiter)
|
||||||
|
.await
|
||||||
|
.expect("command output must not wait for a second completion notification")
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(output.status, CommandStatus::Completed);
|
||||||
|
assert_eq!(output.content, "done");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn command_output_does_not_wait_forever_when_completion_sender_drops() {
|
||||||
|
let dir = TempDir::new().unwrap();
|
||||||
|
let workdir = make_fs(&dir);
|
||||||
|
let handle = CommandHandle("command-completion-drop".into());
|
||||||
|
let (completion_tx, completion) = watch::channel(false);
|
||||||
|
let (cancel, _cancel_rx) = watch::channel(false);
|
||||||
|
let task: JoinHandle<Result<CommandOutput, WorkdirError>> = tokio::spawn(async move {
|
||||||
|
drop(completion_tx);
|
||||||
|
panic!("simulated command task panic");
|
||||||
|
});
|
||||||
|
workdir.inner.commands.lock().await.insert(
|
||||||
|
handle.0.clone(),
|
||||||
|
LocalCommand::Running {
|
||||||
|
task,
|
||||||
|
completion,
|
||||||
|
cancel,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
let result = tokio::time::timeout(
|
||||||
|
Duration::from_secs(1),
|
||||||
|
WorkdirSession::command_output(
|
||||||
|
&workdir,
|
||||||
|
CommandOutputRequest {
|
||||||
|
handle,
|
||||||
|
cursor: 0,
|
||||||
|
limit: 1024,
|
||||||
|
wait: true,
|
||||||
|
},
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("closed completion channel must wake the waiter");
|
||||||
|
assert!(matches!(result, Err(WorkdirError::Unavailable(_))));
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn provider_streams_bounded_command_lifecycle_and_distinct_output() {
|
async fn provider_streams_bounded_command_lifecycle_and_distinct_output() {
|
||||||
let dir = TempDir::new().unwrap();
|
let dir = TempDir::new().unwrap();
|
||||||
|
|||||||
@@ -110,6 +110,16 @@ struct ReviewFindingInput {
|
|||||||
line: Option<u32>,
|
line: Option<u32>,
|
||||||
body: String,
|
body: String,
|
||||||
}
|
}
|
||||||
|
#[derive(Debug, Deserialize)]
|
||||||
|
struct TicketMergeRequestProjection {
|
||||||
|
merge_request: Option<TicketMergeRequestReference>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Deserialize)]
|
||||||
|
struct TicketMergeRequestReference {
|
||||||
|
merge_request_id: String,
|
||||||
|
}
|
||||||
|
|
||||||
impl Kind {
|
impl Kind {
|
||||||
fn enabled(self, config: MergeRequestFeatureConfig) -> bool {
|
fn enabled(self, config: MergeRequestFeatureConfig) -> bool {
|
||||||
match self {
|
match self {
|
||||||
@@ -145,24 +155,22 @@ impl Tool for MergeRequestTool {
|
|||||||
let ws = self.client.workspace_id().ok_or_else(|| {
|
let ws = self.client.workspace_id().ok_or_else(|| {
|
||||||
ToolError::ExecutionFailed("Merge Request tools require Workspace identity".into())
|
ToolError::ExecutionFailed("Merge Request tools require Workspace identity".into())
|
||||||
})?;
|
})?;
|
||||||
|
if matches!(self.kind, Kind::Show) {
|
||||||
|
let value: ShowInput = parse(input)?;
|
||||||
|
nonempty(&value.ticket)?;
|
||||||
|
return self.show_current_merge_request(ws, &value.ticket);
|
||||||
|
}
|
||||||
let (method, path, body) = match self.kind {
|
let (method, path, body) = match self.kind {
|
||||||
Kind::Show | Kind::Readiness => {
|
Kind::Readiness => {
|
||||||
let v: ShowInput = parse(input)?;
|
let v: ShowInput = parse(input)?;
|
||||||
nonempty(&v.ticket)?;
|
nonempty(&v.ticket)?;
|
||||||
(
|
(
|
||||||
WorkspaceRequestMethod::Get,
|
WorkspaceRequestMethod::Get,
|
||||||
format!(
|
format!("/api/w/{ws}/tickets/{}/merge-request/readiness", v.ticket),
|
||||||
"/api/w/{ws}/tickets/{}/merge-request{}",
|
|
||||||
v.ticket,
|
|
||||||
if matches!(self.kind, Kind::Readiness) {
|
|
||||||
"/readiness"
|
|
||||||
} else {
|
|
||||||
""
|
|
||||||
}
|
|
||||||
),
|
|
||||||
None,
|
None,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
Kind::Show => unreachable!("MergeRequestShow is handled above"),
|
||||||
Kind::Open => {
|
Kind::Open => {
|
||||||
let v: OpenInput = parse(input)?;
|
let v: OpenInput = parse(input)?;
|
||||||
nonempty(&v.ticket)?;
|
nonempty(&v.ticket)?;
|
||||||
@@ -225,6 +233,96 @@ impl Tool for MergeRequestTool {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl MergeRequestTool {
|
||||||
|
fn show_current_merge_request(
|
||||||
|
&self,
|
||||||
|
workspace_id: &str,
|
||||||
|
ticket: &str,
|
||||||
|
) -> Result<ToolOutput, ToolError> {
|
||||||
|
let ticket_path = encode_path_segment(ticket);
|
||||||
|
let show_response = self
|
||||||
|
.client
|
||||||
|
.execute(WorkspaceRequest::json(
|
||||||
|
WorkspaceRequestMethod::Post,
|
||||||
|
format!("/api/w/{workspace_id}/tickets/{ticket_path}/show"),
|
||||||
|
json!({"event_limit": 1}).to_string(),
|
||||||
|
))
|
||||||
|
.map_err(|error| ToolError::ExecutionFailed(error.to_string()))?;
|
||||||
|
if !show_response.is_success() {
|
||||||
|
return Err(api_error("Ticket Show API", &show_response));
|
||||||
|
}
|
||||||
|
let projection: TicketMergeRequestProjection = serde_json::from_str(&show_response.body)
|
||||||
|
.map_err(|error| {
|
||||||
|
ToolError::ExecutionFailed(format!(
|
||||||
|
"Ticket Show API returned a malformed Merge Request projection: {error}"
|
||||||
|
))
|
||||||
|
})?;
|
||||||
|
let merge_request = projection.merge_request.ok_or_else(|| {
|
||||||
|
ToolError::ExecutionFailed(format!("Ticket `{ticket}` has no current Merge Request"))
|
||||||
|
})?;
|
||||||
|
nonempty_id("merge_request_id", &merge_request.merge_request_id)?;
|
||||||
|
|
||||||
|
let merge_request_id = encode_path_segment(&merge_request.merge_request_id);
|
||||||
|
let response = self
|
||||||
|
.client
|
||||||
|
.execute(WorkspaceRequest::get(format!(
|
||||||
|
"/api/w/{workspace_id}/merge-requests/{merge_request_id}"
|
||||||
|
)))
|
||||||
|
.map_err(|error| ToolError::ExecutionFailed(error.to_string()))?;
|
||||||
|
if !response.is_success() {
|
||||||
|
return Err(api_error("Merge Request API", &response));
|
||||||
|
}
|
||||||
|
Ok(ToolOutput {
|
||||||
|
summary: self.kind.name().into(),
|
||||||
|
content: Some(response.body),
|
||||||
|
attachments: vec![],
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn api_error(operation: &str, response: &crate::worker::WorkspaceResponse) -> ToolError {
|
||||||
|
ToolError::ExecutionFailed(format!(
|
||||||
|
"{operation} returned HTTP {}: {}",
|
||||||
|
response.status,
|
||||||
|
bounded_body(&response.body)
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn bounded_body(body: &str) -> String {
|
||||||
|
const MAX_CHARS: usize = 4096;
|
||||||
|
let mut chars = body.chars();
|
||||||
|
let bounded: String = chars.by_ref().take(MAX_CHARS).collect();
|
||||||
|
if chars.next().is_some() {
|
||||||
|
format!("{bounded}…")
|
||||||
|
} else {
|
||||||
|
bounded
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn encode_path_segment(value: &str) -> String {
|
||||||
|
let mut encoded = String::with_capacity(value.len());
|
||||||
|
for byte in value.bytes() {
|
||||||
|
if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~') {
|
||||||
|
encoded.push(byte as char);
|
||||||
|
} else {
|
||||||
|
use std::fmt::Write as _;
|
||||||
|
let _ = write!(encoded, "%{byte:02X}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
encoded
|
||||||
|
}
|
||||||
|
|
||||||
|
fn nonempty_id(name: &str, value: &str) -> Result<(), ToolError> {
|
||||||
|
if value.trim().is_empty() || value.chars().any(char::is_control) {
|
||||||
|
Err(ToolError::ExecutionFailed(format!(
|
||||||
|
"Ticket Show API returned an invalid {name}"
|
||||||
|
)))
|
||||||
|
} else {
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn parse<T: serde::de::DeserializeOwned>(v: &str) -> Result<T, ToolError> {
|
fn parse<T: serde::de::DeserializeOwned>(v: &str) -> Result<T, ToolError> {
|
||||||
serde_json::from_str(v).map_err(|e| ToolError::InvalidArgument(e.to_string()))
|
serde_json::from_str(v).map_err(|e| ToolError::InvalidArgument(e.to_string()))
|
||||||
}
|
}
|
||||||
@@ -320,7 +418,103 @@ mod tests {
|
|||||||
use super::*;
|
use super::*;
|
||||||
use crate::feature::FeatureRegistryBuilder;
|
use crate::feature::FeatureRegistryBuilder;
|
||||||
use crate::hook::HookRegistryBuilder;
|
use crate::hook::HookRegistryBuilder;
|
||||||
use crate::worker::TestWorkspaceHttpClient;
|
use crate::worker::{TestWorkspaceHttpClient, WorkspaceClientError, WorkspaceResponse};
|
||||||
|
use std::{collections::VecDeque, sync::Mutex};
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct RecordingWorkspaceClient {
|
||||||
|
responses: Mutex<VecDeque<WorkspaceResponse>>,
|
||||||
|
requests: Mutex<Vec<WorkspaceRequest>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl RecordingWorkspaceClient {
|
||||||
|
fn new(responses: Vec<WorkspaceResponse>) -> Self {
|
||||||
|
Self {
|
||||||
|
responses: Mutex::new(responses.into()),
|
||||||
|
requests: Mutex::new(Vec::new()),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl WorkspaceClient for RecordingWorkspaceClient {
|
||||||
|
fn workspace_id(&self) -> Option<&str> {
|
||||||
|
Some("ws")
|
||||||
|
}
|
||||||
|
|
||||||
|
fn kind(&self) -> &str {
|
||||||
|
"recording"
|
||||||
|
}
|
||||||
|
|
||||||
|
fn is_available(&self) -> bool {
|
||||||
|
true
|
||||||
|
}
|
||||||
|
|
||||||
|
fn execute(
|
||||||
|
&self,
|
||||||
|
request: WorkspaceRequest,
|
||||||
|
) -> Result<WorkspaceResponse, WorkspaceClientError> {
|
||||||
|
self.requests.lock().expect("request lock").push(request);
|
||||||
|
self.responses
|
||||||
|
.lock()
|
||||||
|
.expect("response lock")
|
||||||
|
.pop_front()
|
||||||
|
.ok_or_else(|| WorkspaceClientError::Request("missing test response".into()))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn response(body: serde_json::Value) -> WorkspaceResponse {
|
||||||
|
WorkspaceResponse {
|
||||||
|
status: 200,
|
||||||
|
body: body.to_string(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn show_resolves_ticket_projection_then_reads_canonical_resource() {
|
||||||
|
let client = Arc::new(RecordingWorkspaceClient::new(vec![
|
||||||
|
response(json!({"merge_request":{"merge_request_id":"MR/1"}})),
|
||||||
|
response(json!({"merge_request_id":"MR/1","state":"open"})),
|
||||||
|
]));
|
||||||
|
let tool = MergeRequestTool {
|
||||||
|
kind: Kind::Show,
|
||||||
|
client: client.clone(),
|
||||||
|
};
|
||||||
|
|
||||||
|
let output = tool
|
||||||
|
.execute(r#"{"ticket":"T/1"}"#, ToolExecutionContext::default())
|
||||||
|
.await
|
||||||
|
.expect("show should succeed");
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
output.content.as_deref(),
|
||||||
|
Some(r#"{"merge_request_id":"MR/1","state":"open"}"#)
|
||||||
|
);
|
||||||
|
let requests = client.requests.lock().expect("request lock");
|
||||||
|
assert_eq!(requests.len(), 2);
|
||||||
|
assert_eq!(requests[0].method, WorkspaceRequestMethod::Post);
|
||||||
|
assert_eq!(requests[0].path, "/api/w/ws/tickets/T%2F1/show");
|
||||||
|
assert_eq!(requests[1].method, WorkspaceRequestMethod::Get);
|
||||||
|
assert_eq!(requests[1].path, "/api/w/ws/merge-requests/MR%2F1");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn show_fails_closed_when_ticket_has_no_current_merge_request() {
|
||||||
|
let client = Arc::new(RecordingWorkspaceClient::new(vec![response(
|
||||||
|
json!({"merge_request":null}),
|
||||||
|
)]));
|
||||||
|
let tool = MergeRequestTool {
|
||||||
|
kind: Kind::Show,
|
||||||
|
client: client.clone(),
|
||||||
|
};
|
||||||
|
|
||||||
|
let error = tool
|
||||||
|
.execute(r#"{"ticket":"T1"}"#, ToolExecutionContext::default())
|
||||||
|
.await
|
||||||
|
.expect_err("missing Merge Request must fail");
|
||||||
|
|
||||||
|
assert!(error.to_string().contains("no current Merge Request"));
|
||||||
|
assert_eq!(client.requests.lock().expect("request lock").len(), 1);
|
||||||
|
}
|
||||||
|
|
||||||
fn install(config: MergeRequestFeatureConfig) -> (Vec<String>, Vec<String>) {
|
fn install(config: MergeRequestFeatureConfig) -> (Vec<String>, Vec<String>) {
|
||||||
let client: Arc<dyn WorkspaceClient> =
|
let client: Arc<dyn WorkspaceClient> =
|
||||||
|
|||||||
+114
-24
@@ -337,12 +337,16 @@ impl WorkspaceClient for ReviewerChildWorkspaceClient {
|
|||||||
&self,
|
&self,
|
||||||
mut request: WorkspaceRequest,
|
mut request: WorkspaceRequest,
|
||||||
) -> Result<WorkspaceResponse, WorkspaceClientError> {
|
) -> Result<WorkspaceResponse, WorkspaceClientError> {
|
||||||
let expected_path = format!(
|
let workspace_id = self.workspace_id().unwrap_or_default();
|
||||||
"/api/w/{}/tickets/{}/merge-request/reviews",
|
let ticket_id = &self.context.ticket_id;
|
||||||
self.workspace_id().unwrap_or_default(),
|
let ticket_query_path = format!("/api/w/{workspace_id}/tickets/query");
|
||||||
self.context.ticket_id
|
let ticket_show_path = format!("/api/w/{workspace_id}/tickets/{ticket_id}/show");
|
||||||
);
|
let review_path =
|
||||||
if request.method == WorkspaceRequestMethod::Post && request.path == expected_path {
|
format!("/api/w/{workspace_id}/tickets/{ticket_id}/merge-request/reviews");
|
||||||
|
let read_allowed = request.method == WorkspaceRequestMethod::Get
|
||||||
|
|| (request.method == WorkspaceRequestMethod::Post
|
||||||
|
&& (request.path == ticket_query_path || request.path == ticket_show_path));
|
||||||
|
if request.method == WorkspaceRequestMethod::Post && request.path == review_path {
|
||||||
let body = request.body.take().ok_or_else(|| {
|
let body = request.body.take().ok_or_else(|| {
|
||||||
WorkspaceClientError::Request("review submission requires a JSON body".to_string())
|
WorkspaceClientError::Request("review submission requires a JSON body".to_string())
|
||||||
})?;
|
})?;
|
||||||
@@ -361,7 +365,7 @@ impl WorkspaceClient for ReviewerChildWorkspaceClient {
|
|||||||
serde_json::to_string(&value)
|
serde_json::to_string(&value)
|
||||||
.map_err(|error| WorkspaceClientError::Request(error.to_string()))?,
|
.map_err(|error| WorkspaceClientError::Request(error.to_string()))?,
|
||||||
);
|
);
|
||||||
} else if request.method != WorkspaceRequestMethod::Get {
|
} else if !read_allowed {
|
||||||
return Err(WorkspaceClientError::Unavailable(
|
return Err(WorkspaceClientError::Unavailable(
|
||||||
"Reviewer child Workspace authority is read-only except for its one attested Merge Request review submission".to_string(),
|
"Reviewer child Workspace authority is read-only except for its one attested Merge Request review submission".to_string(),
|
||||||
));
|
));
|
||||||
@@ -483,28 +487,114 @@ impl WorkspaceClient for MarkerWorkspaceClient {
|
|||||||
mod reviewer_client_tests {
|
mod reviewer_client_tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
#[test]
|
#[derive(Debug, Default)]
|
||||||
fn reviewer_child_client_denies_non_review_workspace_mutations() {
|
struct RecordingWorkspaceClient {
|
||||||
let inner: Arc<dyn WorkspaceClient> = Arc::new(MarkerWorkspaceClient {
|
requests: Mutex<Vec<WorkspaceRequest>>,
|
||||||
workspace_id: Some("ws".to_string()),
|
}
|
||||||
kind: "marker".to_string(),
|
|
||||||
available: true,
|
impl WorkspaceClient for RecordingWorkspaceClient {
|
||||||
reason: "forwarded".to_string(),
|
fn workspace_id(&self) -> Option<&str> {
|
||||||
});
|
Some("ws")
|
||||||
let client = ReviewerChildWorkspaceClient::new(
|
}
|
||||||
|
|
||||||
|
fn kind(&self) -> &str {
|
||||||
|
"recording"
|
||||||
|
}
|
||||||
|
|
||||||
|
fn is_available(&self) -> bool {
|
||||||
|
true
|
||||||
|
}
|
||||||
|
|
||||||
|
fn execute(
|
||||||
|
&self,
|
||||||
|
request: WorkspaceRequest,
|
||||||
|
) -> Result<WorkspaceResponse, WorkspaceClientError> {
|
||||||
|
self.requests.lock().expect("recording lock").push(request);
|
||||||
|
Ok(WorkspaceResponse {
|
||||||
|
status: 200,
|
||||||
|
body: "{}".into(),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn reviewer_client(inner: Arc<dyn WorkspaceClient>) -> ReviewerChildWorkspaceClient {
|
||||||
|
ReviewerChildWorkspaceClient::new(
|
||||||
inner,
|
inner,
|
||||||
ReviewerContext {
|
ReviewerContext {
|
||||||
ticket_id: "T1".into(),
|
ticket_id: "T1".into(),
|
||||||
},
|
},
|
||||||
"secret".into(),
|
"secret".into(),
|
||||||
);
|
)
|
||||||
let request = WorkspaceRequest::json(
|
}
|
||||||
WorkspaceRequestMethod::Post,
|
|
||||||
"/api/w/ws/tickets/T1/comments",
|
#[test]
|
||||||
"{}".to_string(),
|
fn reviewer_child_client_allows_typed_ticket_reads() {
|
||||||
);
|
let inner = Arc::new(RecordingWorkspaceClient::default());
|
||||||
let error = client.execute(request).unwrap_err();
|
let client = reviewer_client(inner.clone());
|
||||||
assert!(error.to_string().contains("read-only"));
|
|
||||||
|
for request in [
|
||||||
|
WorkspaceRequest::json(
|
||||||
|
WorkspaceRequestMethod::Post,
|
||||||
|
"/api/w/ws/tickets/query",
|
||||||
|
"{}".to_string(),
|
||||||
|
),
|
||||||
|
WorkspaceRequest::json(
|
||||||
|
WorkspaceRequestMethod::Post,
|
||||||
|
"/api/w/ws/tickets/T1/show",
|
||||||
|
"{}".to_string(),
|
||||||
|
),
|
||||||
|
WorkspaceRequest::get("/api/w/ws/merge-requests/MR1"),
|
||||||
|
] {
|
||||||
|
client.execute(request).expect("read should be forwarded");
|
||||||
|
}
|
||||||
|
|
||||||
|
assert_eq!(inner.requests.lock().expect("recording lock").len(), 3);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn reviewer_child_client_rejects_other_ticket_and_mutation_posts() {
|
||||||
|
let inner: Arc<dyn WorkspaceClient> = Arc::new(RecordingWorkspaceClient::default());
|
||||||
|
let client = reviewer_client(inner);
|
||||||
|
|
||||||
|
for request in [
|
||||||
|
WorkspaceRequest::json(
|
||||||
|
WorkspaceRequestMethod::Post,
|
||||||
|
"/api/w/ws/tickets/T2/show",
|
||||||
|
"{}".to_string(),
|
||||||
|
),
|
||||||
|
WorkspaceRequest::json(
|
||||||
|
WorkspaceRequestMethod::Post,
|
||||||
|
"/api/w/ws/tickets/T1/comments",
|
||||||
|
"{}".to_string(),
|
||||||
|
),
|
||||||
|
WorkspaceRequest::json(
|
||||||
|
WorkspaceRequestMethod::Post,
|
||||||
|
"/api/w/ws/tickets/T1/merge-request",
|
||||||
|
"{}".to_string(),
|
||||||
|
),
|
||||||
|
] {
|
||||||
|
let error = client.execute(request).unwrap_err();
|
||||||
|
assert!(error.to_string().contains("read-only"));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn reviewer_child_client_injects_capability_only_for_attested_review() {
|
||||||
|
let inner = Arc::new(RecordingWorkspaceClient::default());
|
||||||
|
let client = reviewer_client(inner.clone());
|
||||||
|
client
|
||||||
|
.execute(WorkspaceRequest::json(
|
||||||
|
WorkspaceRequestMethod::Post,
|
||||||
|
"/api/w/ws/tickets/T1/merge-request/reviews",
|
||||||
|
r#"{"decision":"approve"}"#.to_string(),
|
||||||
|
))
|
||||||
|
.expect("attested review should be forwarded");
|
||||||
|
|
||||||
|
let requests = inner.requests.lock().expect("recording lock");
|
||||||
|
let body: serde_json::Value =
|
||||||
|
serde_json::from_str(requests[0].body.as_deref().expect("review body"))
|
||||||
|
.expect("review JSON");
|
||||||
|
assert_eq!(body["capability_token"], "secret");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user