返回 CodeWhale
host.rs
根目录 / crates / tui / src / tools / github / host.rs
1 //! Admitted GitHub operations. Core keeps processes, policy, private drafts and receipts.
2 use super::{cli, report, shape, validate_evidence};
3 use crate::tools::spec::{
4 ToolContext, ToolError, ToolResult, optional_bool, optional_str, required_str, required_u64,
5 };
6 use serde::Serialize;
7 use serde_json::{Value, json};
8 use std::sync::{Arc, Mutex};
9 use tokio_util::sync::CancellationToken;
10
11 #[derive(Clone, Serialize)]
12 pub(crate) enum Request {
13 Issue {
14 number: u64,
15 comments: bool,
16 },
17 Pr {
18 number: u64,
19 diff: bool,
20 },
21 Comment {
22 target: String,
23 number: u64,
24 body: String,
25 dry: bool,
26 },
27 Close {
28 target: String,
29 number: u64,
30 comment: Option<String>,
31 dry: bool,
32 allow_dirty: bool,
33 },
34 Draft {
35 input: Value,
36 },
37 Read {
38 session: String,
39 id: String,
40 operator: bool,
41 },
42 }
43 impl Request {
44 pub(super) fn capture(
45 action: &str,
46 input: &Value,
47 context: &ToolContext,
48 ) -> Result<Self, ToolError> {
49 if serde_json::to_vec(input)
50 .map_err(|_| ToolError::invalid_input("Invalid GitHub input"))?
51 .len()
52 > 1024 * 1024
53 {
54 return Err(ToolError::invalid_input("GitHub input exceeds 1 MiB"));
55 }
56 Ok(match action {
57 "issue_context" => Self::Issue {
58 number: required_u64(input, "number")?,
59 comments: optional_bool(input, "include_comments", true)?,
60 },
61 "pr_context" => Self::Pr {
62 number: required_u64(input, "number")?,
63 diff: optional_bool(input, "include_diff", false)?,
64 },
65 "comment" => {
66 validate_evidence(input, false)?;
67 let target = required_str(input, "target")?;
68 if !matches!(target, "issue" | "pr") {
69 return Err(ToolError::invalid_input(
70 "github comment: target must be issue or pr; nothing was posted",
71 ));
72 }
73 Self::Comment {
74 target: target.into(),
75 number: required_u64(input, "number")?,
76 body: required_str(input, "body")?.into(),
77 dry: optional_bool(input, "dry_run", false)?,
78 }
79 }
80 "close_issue" | "close_pr" => {
81 validate_evidence(input, true)?;
82 Self::Close {
83 target: if action == "close_issue" {
84 "issue"
85 } else {
86 "pr"
87 }
88 .into(),
89 number: required_u64(input, "number")?,
90 comment: optional_str(input, "comment")?.map(str::to_string),
91 dry: optional_bool(input, "dry_run", false)?,
92 allow_dirty: optional_bool(input, "allow_dirty", false)?,
93 }
94 }
95 "report_draft" => {
96 report::validate_host_draft(input, context)?;
97 Self::Draft {
98 input: input.clone(),
99 }
100 }
101 "report_read" => {
102 report::validate_host_read(input)?;
103 if context
104 .session_objects
105 .as_ref()
106 .is_none_or(|snapshot| snapshot.session_id != context.state_namespace)
107 {
108 return Err(ToolError::not_available(
109 "Issue reading requires the active Engine session context.",
110 ));
111 }
112 Self::Read {
113 session: context.state_namespace.clone(),
114 id: required_str(input, "report_id")?.into(),
115 operator: false,
116 }
117 }
118 _ => return Err(ToolError::invalid_input("Unknown GitHub action")),
119 })
120 }
121 }
122
123 /// Private completion facts remain in the admitted Core job; no protocol projection.
124 #[derive(Clone, Copy, Default)]
125 enum Mutation {
126 #[default]
127 None,
128 CommentPending,
129 CommentCompleted,
130 ClosePending {
131 comment_completed: bool,
132 },
133 CloseCompleted {
134 comment_completed: bool,
135 },
136 }
137 #[derive(Default)]
138 pub(crate) struct RunState {
139 pub(crate) result: Option<Result<Outcome, ToolError>>,
140 mutation: Mutation,
141 }
142 impl RunState {
143 fn mark(state: &Mutex<Self>, mutation: Mutation) {
144 state.lock().expect("GitHub result lock").mutation = mutation;
145 }
146 fn refused_before_write(state: &Mutex<Self>, error: &ToolError) {
147 // These failures occur before the contained process can be spawned. A
148 // previous successful comment remains a completed effect in a Close job.
149 if matches!(
150 error,
151 ToolError::PermissionDenied { .. }
152 | ToolError::InvalidInput { .. }
153 | ToolError::MissingField { .. }
154 | ToolError::NotAvailable { .. }
155 | ToolError::PathEscape { .. }
156 ) {
157 let mut state = state.lock().expect("GitHub result lock");
158 state.mutation = match state.mutation {
159 Mutation::ClosePending {
160 comment_completed: true,
161 } => Mutation::CommentCompleted,
162 Mutation::CommentPending
163 | Mutation::ClosePending {
164 comment_completed: false,
165 } => Mutation::None,
166 other => other,
167 };
168 }
169 }
170 }
171 impl Mutation {
172 fn failure(self, error: ToolError) -> ToolError {
173 let fact = match self {
174 Self::None => return error,
175 Self::CommentPending => {
176 "GitHub comment outcome is uncertain; inspect the thread before retrying; do not replay automatically"
177 }
178 Self::CommentCompleted => {
179 "GitHub comment completed; inspect the thread and do not replay the comment"
180 }
181 Self::ClosePending {
182 comment_completed: false,
183 } => {
184 "GitHub close outcome is uncertain; inspect the thread before retrying; do not replay automatically"
185 }
186 Self::ClosePending {
187 comment_completed: true,
188 } => {
189 "GitHub comment completed; close outcome is uncertain; inspect the thread and do not replay the comment or close automatically"
190 }
191 Self::CloseCompleted {
192 comment_completed: false,
193 } => "GitHub close completed; inspect the thread and do not replay the close",
194 Self::CloseCompleted {
195 comment_completed: true,
196 } => {
197 "GitHub comment and close completed; inspect the thread and do not replay either write"
198 }
199 };
200 // Only Core-derived effect and failure class survive; never interpolate
201 // arbitrary Host presentation errors or private process diagnostics.
202 match error {
203 ToolError::PermissionDenied { .. } => ToolError::permission_denied(format!(
204 "{fact}; current authority refused completion"
205 )),
206 ToolError::Cancelled { .. } => {
207 ToolError::cancelled(format!("{fact}; completion was cancelled"))
208 }
209 ToolError::NotAvailable { .. } => {
210 ToolError::not_available(format!("{fact}; completion is unavailable"))
211 }
212 ToolError::Timeout { seconds } => ToolError::execution_failed_with_metadata(
213 format!("{fact}; completion timed out"),
214 json!({"github_effect": fact,"status":"timeout","timeout_seconds":seconds}),
215 ),
216 _ => ToolError::execution_failed_with_metadata(
217 format!("{fact}; completion or presentation failed"),
218 json!({"github_effect":fact}),
219 ),
220 }
221 }
222 }
223
224 /// The private Core fault wins over a generic peer failure. A Host refusal,
225 /// forged metadata, timeout or cancellation cannot erase a confirmed write.
226 pub(crate) fn finish_run(
227 state: &Mutex<RunState>,
228 host: Result<ToolResult, ToolError>,
229 ) -> Result<ToolResult, ToolError> {
230 let (owned, mutation) = {
231 let mut state = state.lock().expect("GitHub result lock");
232 (state.result.take(), state.mutation)
233 };
234 match owned {
235 Some(Err(error)) => Err(error),
236 Some(Ok(outcome)) => host
237 .and_then(|host| outcome.finish(host))
238 .map_err(|error| mutation.failure(error)),
239 None => Err(mutation.failure(host.err().unwrap_or_else(|| {
240 ToolError::execution_failed("GitHub driver produced no owned receipt")
241 }))),
242 }
243 }
244
245 pub(crate) struct Outcome {
246 pub projection: Value,
247 pub metadata: Option<Value>,
248 }
249 impl Outcome {
250 pub(crate) fn finish(self, mut result: ToolResult) -> Result<ToolResult, ToolError> {
251 let expected_success = self.projection.get("dirty").and_then(Value::as_bool) != Some(true);
252 if result.success != expected_success {
253 return Err(ToolError::execution_failed(
254 "GitHub presenter status does not match the owned outcome",
255 ));
256 }
257 // The Host owns presentation, never task/report paths, IDs or timestamps.
258 if result
259 .metadata
260 .as_ref()
261 .is_some_and(|value| !value.is_null())
262 {
263 return Err(ToolError::execution_failed(
264 "GitHub presenter returned unowned metadata",
265 ));
266 }
267 result.metadata = self.metadata;
268 Ok(result)
269 }
270 }
271 fn prefix(text: &str, count: usize) -> String {
272 text.chars().take(count).collect()
273 }
274 fn context_outcome(
275 context: &ToolContext,
276 action: &str,
277 number: u64,
278 mut raw: Value,
279 diff: Option<String>,
280 ) -> Result<Outcome, ToolError> {
281 let body = raw.get("body").and_then(Value::as_str).map(str::to_string);
282 let large = body
283 .as_ref()
284 .is_some_and(|body| body.len() > shape::BODY_ARTIFACT_THRESHOLD);
285 let body_artifact = if large {
286 shape::write_artifact_if_needed(
287 context,
288 if action == "issue_context" {
289 "issue_body"
290 } else {
291 "pr_body"
292 },
293 body.as_deref().unwrap_or(""),
294 shape::BODY_ARTIFACT_THRESHOLD,
295 )
296 .map_err(|_| ToolError::execution_failed("GitHub artifact storage is unavailable"))?
297 } else {
298 None
299 };
300 if large {
301 raw["body"] = json!(prefix(body.as_deref().unwrap_or(""), 1200));
302 }
303 let diff_artifact = diff
304 .as_deref()
305 .map(|diff| {
306 shape::write_artifact_if_needed(
307 context,
308 "pr_diff",
309 diff,
310 shape::DIFF_ARTIFACT_THRESHOLD,
311 )
312 })
313 .transpose()
314 .map_err(|_| ToolError::execution_failed("GitHub artifact storage is unavailable"))?
315 .flatten();
316 let mut artifacts = Vec::new();
317 for (path, label, text) in [
318 (
319 body_artifact.as_ref(),
320 if action == "issue_context" {
321 "github_issue_body"
322 } else {
323 "github_pr_body"
324 },
325 body.as_deref(),
326 ),
327 (diff_artifact.as_ref(), "github_pr_diff", diff.as_deref()),
328 ] {
329 if let Some(path) = path {
330 artifacts.push(crate::task_manager::TaskArtifactRef {
331 label: label.into(),
332 path: path.clone(),
333 summary: shape::summarize(text.unwrap_or(""), 900),
334 created_at: chrono::Utc::now(),
335 });
336 }
337 }
338 Ok(Outcome {
339 projection: json!({"action":action,"number":number.to_string(),"raw":raw,"large_body":large,"body_artifact":body_artifact,"diff":diff.map(|diff|prefix(&diff,900)),"diff_artifact":diff_artifact}),
340 metadata: (!artifacts.is_empty()).then(|| json!({"task_updates":{"artifacts":artifacts}})),
341 })
342 }
343 fn report_outcome(report: report::Report, operator: bool) -> Result<Outcome, ToolError> {
344 Ok(Outcome {
345 projection: json!({"action":if operator{"report_review"}else{"report_read"},"report":report.host_snapshot()}),
346 metadata: if operator {
347 None
348 } else {
349 Some(report.host_metadata()?)
350 },
351 })
352 }
353
354 pub(crate) async fn run<F, Fut>(
355 request: Request,
356 context: Option<ToolContext>,
357 cancel: CancellationToken,
358 check: F,
359 state: Arc<Mutex<RunState>>,
360 ) -> Result<Outcome, ToolError>
361 where
362 F: Fn() -> Fut + Clone + Send + 'static,
363 Fut: std::future::Future<Output = Result<(), ToolError>> + Send,
364 {
365 if cancel.is_cancelled() {
366 return Err(ToolError::not_available("GitHub operation cancelled"));
367 }
368 check().await?;
369 match request {
370 Request::Issue { number, comments } => {
371 let context = context
372 .as_ref()
373 .ok_or_else(|| ToolError::not_available("GitHub caller missing"))?;
374 cli::host_git(context, &["rev-parse", "--is-inside-work-tree"], &cancel).await?;
375 let fields = if comments {
376 "number,title,state,author,labels,assignees,milestone,body,comments,url,createdAt,updatedAt"
377 } else {
378 "number,title,state,author,labels,assignees,milestone,body,url,createdAt,updatedAt"
379 };
380 let target = cli::host_target(context, &cancel).await?;
381 check().await?;
382 let raw = cli::host_gh(
383 context,
384 &target,
385 &["issue", "view", &number.to_string(), "--json", fields],
386 None,
387 &cancel,
388 )
389 .await?;
390 check().await?;
391 context_outcome(
392 context,
393 "issue_context",
394 number,
395 serde_json::from_str(&raw)
396 .map_err(|_| ToolError::execution_failed("Invalid GitHub issue JSON"))?,
397 None,
398 )
399 }
400 Request::Pr { number, diff } => {
401 let context = context
402 .as_ref()
403 .ok_or_else(|| ToolError::not_available("GitHub caller missing"))?;
404 cli::host_git(context, &["rev-parse", "--is-inside-work-tree"], &cancel).await?;
405 let target = cli::host_target(context, &cancel).await?;
406 check().await?;
407 let raw=cli::host_gh(context,&target,&["pr","view",&number.to_string(),"--json","number,title,state,author,body,comments,reviews,reviewDecision,statusCheckRollup,baseRefName,headRefName,headRefOid,baseRefOid,files,url,createdAt,updatedAt"],None,&cancel).await?;
408 let diff = if diff {
409 Some(
410 cli::host_gh(
411 context,
412 &target,
413 &["pr", "diff", &number.to_string(), "--patch"],
414 None,
415 &cancel,
416 )
417 .await?,
418 )
419 } else {
420 None
421 };
422 check().await?;
423 context_outcome(
424 context,
425 "pr_context",
426 number,
427 serde_json::from_str(&raw)
428 .map_err(|_| ToolError::execution_failed("Invalid GitHub PR JSON"))?,
429 diff,
430 )
431 }
432 Request::Comment {
433 target,
434 number,
435 body,
436 dry,
437 } => {
438 let context = context
439 .as_ref()
440 .ok_or_else(|| ToolError::not_available("GitHub caller missing"))?;
441 if !dry {
442 let selection = cli::host_target(context, &cancel).await?;
443 check().await?;
444 RunState::mark(&state, Mutation::CommentPending);
445 cli::host_gh(
446 context,
447 &selection,
448 &[&target, "comment", &number.to_string(), "--body-file", "-"],
449 Some(body.as_bytes()),
450 &cancel,
451 )
452 .await
453 .map_err(|error| {
454 RunState::refused_before_write(&state, &error);
455 write_error(
456 error,
457 "GitHub comment outcome is uncertain; inspect the thread before retrying",
458 )
459 })?;
460 RunState::mark(&state, Mutation::CommentCompleted);
461 }
462 check().await.map_err(|error| {
463 if dry {
464 error
465 } else {
466 Mutation::CommentCompleted.failure(error)
467 }
468 })?;
469 let metadata = if dry {
470 None
471 } else {
472 Some(shape::github_event_metadata(
473 "comment",
474 &target,
475 number,
476 shape::summarize(&body, 240),
477 None,
478 shape::write_artifact_if_needed(
479 context,
480 "github_comment",
481 &body,
482 shape::BODY_ARTIFACT_THRESHOLD,
483 ).map_err(|_|ToolError::execution_failed("GitHub comment completed, but its Core artifact receipt could not be saved; do not replay the comment"))?,
484 ))
485 };
486 Ok(Outcome {
487 projection: json!({"action":"comment","number":number.to_string(),"target":target,"dry_run":dry}),
488 metadata,
489 })
490 }
491 Request::Close {
492 target,
493 number,
494 comment,
495 dry,
496 allow_dirty,
497 } => {
498 let context = context
499 .as_ref()
500 .ok_or_else(|| ToolError::not_available("GitHub caller missing"))?;
501 if !allow_dirty {
502 let status = cli::host_git(context, &["status", "--porcelain"], &cancel).await?;
503 if !status.trim().is_empty() {
504 return Ok(Outcome {
505 projection: json!({"action":if target=="issue"{"close_issue"}else{"close_pr"},"number":number.to_string(),"target":target,"dirty":true}),
506 metadata: Some(json!({"dirty_status":status})),
507 });
508 }
509 }
510 if !dry {
511 let selection = cli::host_target(context, &cancel).await?;
512 check().await?;
513 if let Some(comment) = comment.as_ref() {
514 RunState::mark(&state, Mutation::CommentPending);
515 cli::host_gh(context,&selection,&[&target,"comment",&number.to_string(),"--body-file","-"],Some(comment.as_bytes()),&cancel).await.map_err(|error| {
516 RunState::refused_before_write(&state, &error);
517 write_error(error,"GitHub comment outcome is uncertain; inspect the thread before retrying; close was not attempted")
518 })?;
519 RunState::mark(&state, Mutation::CommentCompleted);
520 check().await.map_err(|error| match error {
521 ToolError::PermissionDenied { .. } => ToolError::permission_denied("GitHub comment completed, but caller authority changed; close was not attempted; do not replay the comment"),
522 ToolError::Cancelled { .. } => ToolError::cancelled("GitHub comment completed, but completion was cancelled; close was not attempted; do not replay the comment"),
523 _ => ToolError::not_available("GitHub comment completed, but caller authority changed; close was not attempted; do not replay the comment"),
524 })?;
525 }
526 let number_string = number.to_string();
527 let args = if target == "issue" {
528 vec![
529 "issue",
530 "close",
531 number_string.as_str(),
532 "--reason",
533 "completed",
534 ]
535 } else {
536 vec!["pr", "close", number_string.as_str()]
537 };
538 RunState::mark(
539 &state,
540 Mutation::ClosePending {
541 comment_completed: comment.is_some(),
542 },
543 );
544 if let Err(error) = cli::host_gh(context, &selection, &args, None, &cancel).await {
545 RunState::refused_before_write(&state, &error);
546 return Err(if comment.is_some() {
547 match error {
548 ToolError::PermissionDenied { .. } => ToolError::permission_denied(
549 "GitHub comment completed; close was blocked by current network authority; do not replay the comment",
550 ),
551 ToolError::NotAvailable { .. } => ToolError::not_available(
552 "GitHub comment completed; close was unavailable; do not replay the comment",
553 ),
554 ToolError::Cancelled { .. } => ToolError::cancelled(
555 "GitHub comment completed; close was cancelled and its outcome is uncertain; inspect before retrying and do not replay the comment",
556 ),
557 _ => ToolError::execution_failed(
558 "GitHub comment completed; close outcome is uncertain; inspect before retrying and do not replay the comment",
559 ),
560 }
561 } else {
562 write_error(
563 error,
564 "GitHub close outcome is uncertain; inspect the thread before retrying",
565 )
566 });
567 }
568 RunState::mark(
569 &state,
570 Mutation::CloseCompleted {
571 comment_completed: comment.is_some(),
572 },
573 );
574 }
575 check().await.map_err(|error| {
576 if dry {
577 error
578 } else {
579 Mutation::CloseCompleted {
580 comment_completed: comment.is_some(),
581 }
582 .failure(error)
583 }
584 })?;
585 let metadata = if dry {
586 None
587 } else {
588 Some(shape::github_event_metadata(
589 "close",
590 &target,
591 number,
592 format!(
593 "{} closed as completed with structured evidence",
594 if target == "issue" { "Issue" } else { "PR" }
595 ),
596 None,
597 comment.as_deref().map(|comment|shape::write_artifact_if_needed(context,"github_close_comment",comment,shape::BODY_ARTIFACT_THRESHOLD)).transpose().map_err(|_|ToolError::execution_failed("GitHub close completed, but its Core artifact receipt could not be saved; inspect the thread and do not replay the write"))?.flatten(),
598 ))
599 };
600 Ok(Outcome {
601 projection: json!({"action":if target=="issue"{"close_issue"}else{"close_pr"},"number":number.to_string(),"target":target,"dry_run":dry}),
602 metadata,
603 })
604 }
605 Request::Draft { input } => {
606 let context =
607 context.ok_or_else(|| ToolError::not_available("Report caller missing"))?;
608 let report = report_worker(check.clone(), move || {
609 report::create_host_draft(input, &context)
610 })
611 .await?;
612 check().await?;
613 report_outcome(report, false)
614 }
615 Request::Read {
616 session,
617 id,
618 operator,
619 } => {
620 let report = report_worker(check.clone(), move || report::load(&session, &id)).await?;
621 check().await?;
622 report_outcome(report, operator)
623 }
624 }
625 }
626
627 /// The admitted job/permit remains owned by the actual blocking worker after its
628 /// async waiter is cancelled. Aborting a waiter cannot release quota early.
629 pub(crate) async fn report_worker<A, T>(
630 admission: A,
631 work: impl FnOnce() -> Result<T, ToolError> + Send + 'static,
632 ) -> Result<T, ToolError>
633 where
634 A: Send + 'static,
635 T: Send + 'static,
636 {
637 #[cfg(test)]
638 let scope = crate::test_support::env_scope_ticket();
639 tokio::task::spawn_blocking(move || {
640 let _admission = admission;
641 #[cfg(test)]
642 let _scope = crate::test_support::join_env_scope(scope);
643 work()
644 })
645 .await
646 .map_err(|_| ToolError::execution_failed("Report worker lost"))?
647 }
648
649 fn write_error(error: ToolError, uncertain: &'static str) -> ToolError {
650 match error {
651 ToolError::PermissionDenied { .. }
652 | ToolError::InvalidInput { .. }
653 | ToolError::MissingField { .. }
654 | ToolError::NotAvailable { .. }
655 | ToolError::PathEscape { .. } => error,
656 ToolError::Cancelled { .. } => ToolError::cancelled(uncertain),
657 _ => ToolError::execution_failed(uncertain),
658 }
659 }
660
661 #[cfg(test)]
662 mod tests {
663 use super::*;
664 #[cfg(unix)]
665 use std::sync::atomic::{AtomicUsize, Ordering};
666
667 #[test]
668 fn capture_rejects_invalid_effects_before_admission() {
669 let temp = tempfile::tempdir().unwrap();
670 let context = ToolContext::new(temp.path());
671 let evidence =
672 json!({"files_changed":["src/lib.rs"],"tests_run":["fixture"],"final_status":"green"});
673 for input in [
674 json!({"action":"comment","number":1,"target":"issue","body":"text","dry_run":"true","evidence":evidence}),
675 json!({"action":"comment","number":1,"target":"other","body":"text","evidence":evidence}),
676 json!({"action":"close_issue","number":1,"evidence":evidence,"acceptance_criteria":[]}),
677 json!({"action":"report_read","report_id":"../../other"}),
678 json!({"action":"issue_context","number":1,"include_comments":"true"}),
679 json!({"action":"issue_context","number":1,"extra":"x".repeat(1024*1024)}),
680 ] {
681 assert!(Request::capture(input["action"].as_str().unwrap(), &input, &context).is_err());
682 }
683 let Request::Issue { number, comments } =
684 Request::capture("issue_context", &json!({"number":u64::MAX}), &context).unwrap()
685 else {
686 panic!("issue request")
687 };
688 assert_eq!(number, u64::MAX);
689 assert!(comments);
690 }
691
692 #[test]
693 fn report_read_capture_requires_exact_active_session() {
694 let temp = tempfile::tempdir().unwrap();
695 let context = ToolContext::new(temp.path()).with_state_namespace("session-a");
696 let input =
697 json!({"action":"report_read","report_id":format!("cwreport_{}","a".repeat(64))});
698 assert!(Request::capture("report_read", &input, &context).is_err());
699 let mismatch =
700 context
701 .clone()
702 .with_session_objects(crate::rlm::session::SessionObjectSnapshot::new(
703 "session-b".into(),
704 "model".into(),
705 temp.path().into(),
706 None,
707 vec![],
708 ));
709 assert!(Request::capture("report_read", &input, &mismatch).is_err());
710 let current =
711 context.with_session_objects(crate::rlm::session::SessionObjectSnapshot::new(
712 "session-a".into(),
713 "model".into(),
714 temp.path().into(),
715 None,
716 vec![],
717 ));
718 let Request::Read {
719 session, operator, ..
720 } = Request::capture("report_read", &input, &current).unwrap()
721 else {
722 panic!("report read")
723 };
724 assert_eq!(session, "session-a");
725 assert!(!operator);
726 }
727
728 #[test]
729 fn core_receipts_reject_host_metadata_and_source_forged_artifact_refs() {
730 let temp = tempfile::tempdir().unwrap();
731 let context = ToolContext::new(temp.path());
732 let raw = json!({"body":"small","body_artifact":"/private/forged","comments":[{"body_artifact":"/private/nested"}]});
733 let outcome = context_outcome(&context, "issue_context", 1, raw.clone(), None).unwrap();
734 assert_eq!(outcome.projection["raw"], raw);
735 assert!(outcome.metadata.is_none());
736 assert!(
737 outcome
738 .finish(
739 ToolResult::success("presented")
740 .with_metadata(json!({"task_updates":{"artifacts":[]}}))
741 )
742 .is_err()
743 );
744 let body = format!("{}\u{0085}end", "漢".repeat(1400));
745 let outcome = context_outcome(
746 &context,
747 "pr_context",
748 2,
749 json!({"body":body}),
750 Some("😀".repeat(1000)),
751 )
752 .unwrap();
753 assert_eq!(
754 shape::summarize(outcome.projection["raw"]["body"].as_str().unwrap(), 1200),
755 shape::summarize(&body, 1200)
756 );
757 assert!(outcome.metadata.is_none());
758 }
759
760 #[tokio::test(flavor = "current_thread")]
761 async fn cancelled_report_waiter_retains_production_worker_permit() {
762 let slots = Arc::new(tokio::sync::Semaphore::new(1));
763 let permit = Arc::clone(&slots).acquire_owned().await.unwrap();
764 let (started_tx, started_rx) = tokio::sync::oneshot::channel();
765 let (finish_tx, finish_rx) = std::sync::mpsc::channel();
766 let waiter = tokio::spawn(report_worker(permit, move || {
767 started_tx.send(()).unwrap();
768 finish_rx.recv().unwrap();
769 Ok(())
770 }));
771 started_rx.await.unwrap();
772 waiter.abort();
773 assert!(waiter.await.unwrap_err().is_cancelled());
774 assert_eq!(
775 slots.available_permits(),
776 0,
777 "aborting an async waiter must not free the still running job"
778 );
779 finish_tx.send(()).unwrap();
780 tokio::time::timeout(std::time::Duration::from_secs(2), async {
781 let _permit = slots.acquire().await.unwrap();
782 })
783 .await
784 .unwrap();
785 assert_eq!(slots.available_permits(), 1);
786 }
787
788 #[cfg(unix)]
789 fn recorder(temp: &std::path::Path) -> std::path::PathBuf {
790 use std::os::unix::fs::PermissionsExt;
791 let path = temp.join("gh-record.sh");
792 std::fs::write(
793 &path,
794 r##"#!/bin/sh
795 printf '%s\n' "$@" >> "$CW_GH_TEST_LOG"
796 cat >> "$CW_GH_TEST_STDIN"
797 printf '%s\n' fixture-private-diagnostic >&2
798 if [ -n "$CW_GH_TEST_STARTED" ]; then : > "$CW_GH_TEST_STARTED"; sleep 30; fi
799 if [ "$2" = close ]; then exit "${CW_GH_TEST_CLOSE_EXIT:-${CW_GH_TEST_EXIT:-0}}"; fi
800 exit "${CW_GH_TEST_EXIT:-0}"
801 "##,
802 )
803 .unwrap();
804 std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o700)).unwrap();
805 path
806 }
807
808 #[cfg(unix)]
809 #[tokio::test(flavor = "current_thread")]
810 async fn admitted_write_uses_stdin_exact_repo_and_withholds_private_diagnostics() {
811 let _home = crate::test_support::SealedHome::new();
812 let temp = tempfile::tempdir().unwrap();
813 let log = temp.path().join("args");
814 let input = temp.path().join("stdin");
815 let binary = recorder(temp.path());
816 let _bin = crate::test_support::EnvVarGuard::set("CODEWHALE_GH_BIN", &binary);
817 let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "approved.example/owner/name");
818 let _log = crate::test_support::EnvVarGuard::set("CW_GH_TEST_LOG", &log);
819 let _input = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STDIN", &input);
820 let context = ToolContext::new(temp.path());
821 let body = "private body 漢字";
822 let outcome = run(
823 Request::Comment {
824 target: "pr".into(),
825 number: u64::MAX,
826 body: body.into(),
827 dry: false,
828 },
829 Some(context.clone()),
830 CancellationToken::new(),
831 || async { Ok(()) },
832 Arc::default(),
833 )
834 .await
835 .unwrap();
836 let args = std::fs::read_to_string(&log).unwrap();
837 assert!(args.contains("--body-file\n-\n"));
838 assert!(args.contains("--repo\napproved.example/owner/name\n"));
839 assert!(!args.contains(body));
840 assert_eq!(std::fs::read_to_string(&input).unwrap(), body);
841 assert!(outcome.metadata.is_some());
842 let mut broken = context.clone();
843 let blocked = temp.path().join("private-receipt-location");
844 std::fs::write(&blocked, "not a directory").unwrap();
845 broken.runtime.active_task_id = Some("task".into());
846 broken.runtime.task_data_dir = Some(blocked);
847 let error = run(
848 Request::Comment {
849 target: "pr".into(),
850 number: 1,
851 body: "x".repeat(4001),
852 dry: false,
853 },
854 Some(broken),
855 CancellationToken::new(),
856 || async { Ok(()) },
857 Arc::default(),
858 )
859 .await
860 .err()
861 .unwrap()
862 .to_string();
863 assert!(error.contains("comment completed"));
864 assert!(error.contains("receipt could not be saved"));
865 assert!(!error.contains("private-receipt-location"));
866 let _exit = crate::test_support::EnvVarGuard::set("CW_GH_TEST_EXIT", "7");
867 let error = run(
868 Request::Comment {
869 target: "issue".into(),
870 number: 1,
871 body: body.into(),
872 dry: false,
873 },
874 Some(context),
875 CancellationToken::new(),
876 || async { Ok(()) },
877 Arc::default(),
878 )
879 .await
880 .err()
881 .unwrap()
882 .to_string();
883 assert!(error.contains("uncertain"));
884 assert!(!error.contains(body));
885 assert!(!error.contains("fixture-private-diagnostic"));
886 }
887
888 #[cfg(unix)]
889 #[tokio::test(flavor = "current_thread")]
890 async fn changed_authority_after_comment_refuses_close_without_replaying_comment() {
891 let _home = crate::test_support::SealedHome::new();
892 let temp = tempfile::tempdir().unwrap();
893 let log = temp.path().join("args");
894 let input = temp.path().join("stdin");
895 let _bin = crate::test_support::EnvVarGuard::set("CODEWHALE_GH_BIN", recorder(temp.path()));
896 let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "owner/name");
897 let _log = crate::test_support::EnvVarGuard::set("CW_GH_TEST_LOG", &log);
898 let _input = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STDIN", &input);
899 let checks = Arc::new(AtomicUsize::new(0));
900 let counter = Arc::clone(&checks);
901 let error = run(
902 Request::Close {
903 target: "issue".into(),
904 number: 1,
905 comment: Some("confirmed evidence".into()),
906 dry: false,
907 allow_dirty: true,
908 },
909 Some(ToolContext::new(temp.path())),
910 CancellationToken::new(),
911 move || {
912 let revoked = counter.fetch_add(1, Ordering::SeqCst) >= 2;
913 async move {
914 if revoked {
915 Err(ToolError::not_available("revoked"))
916 } else {
917 Ok(())
918 }
919 }
920 },
921 Arc::default(),
922 )
923 .await
924 .err()
925 .unwrap()
926 .to_string();
927 assert!(error.contains("comment completed"));
928 assert!(error.contains("close was not attempted"));
929 let args = std::fs::read_to_string(&log).unwrap();
930 assert!(args.contains("comment"));
931 assert!(!args.contains("close"));
932 assert_eq!(
933 std::fs::read_to_string(input).unwrap(),
934 "confirmed evidence"
935 );
936 }
937
938 #[cfg(unix)]
939 #[tokio::test(flavor = "current_thread")]
940 async fn deny_or_prompt_policy_and_dry_run_never_invoke_gh() {
941 use crate::network_policy::{DecisionToml, NetworkPolicy, NetworkPolicyDecider};
942 let _home = crate::test_support::SealedHome::new();
943 let temp = tempfile::tempdir().unwrap();
944 let log = temp.path().join("args");
945 let input = temp.path().join("stdin");
946 let _bin = crate::test_support::EnvVarGuard::set("CODEWHALE_GH_BIN", recorder(temp.path()));
947 let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "owner/name");
948 let _log = crate::test_support::EnvVarGuard::set("CW_GH_TEST_LOG", &log);
949 let _input = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STDIN", &input);
950 for default in [DecisionToml::Deny, DecisionToml::Prompt] {
951 let context =
952 ToolContext::new(temp.path()).with_network_policy(NetworkPolicyDecider::new(
953 NetworkPolicy {
954 default,
955 ..Default::default()
956 },
957 None,
958 ));
959 let result = run(
960 Request::Comment {
961 target: "issue".into(),
962 number: 1,
963 body: "no".into(),
964 dry: false,
965 },
966 Some(context),
967 CancellationToken::new(),
968 || async { Ok(()) },
969 Arc::default(),
970 )
971 .await;
972 assert!(matches!(result, Err(ToolError::PermissionDenied { .. })));
973 }
974 assert!(!log.exists());
975 run(
976 Request::Close {
977 target: "pr".into(),
978 number: 1,
979 comment: None,
980 dry: true,
981 allow_dirty: true,
982 },
983 Some(ToolContext::new(temp.path())),
984 CancellationToken::new(),
985 || async { Ok(()) },
986 Arc::default(),
987 )
988 .await
989 .unwrap();
990 assert!(!log.exists());
991 let cancel = CancellationToken::new();
992 cancel.cancel();
993 assert!(
994 run(
995 Request::Comment {
996 target: "issue".into(),
997 number: 1,
998 body: "no".into(),
999 dry: false
1000 },
1001 Some(ToolContext::new(temp.path())),
1002 cancel,
1003 || async { Ok(()) },
1004 Arc::default(),
1005 )
1006 .await
1007 .is_err()
1008 );
1009 assert!(!log.exists());
1010 }
1011 #[test]
1012 fn private_core_fault_wins_over_generic_peer_failure() {
1013 for error in [
1014 ToolError::permission_denied("Core network authority refused before spawn"),
1015 ToolError::cancelled("Core comment cancelled; inspect before retrying"),
1016 ToolError::execution_failed_with_metadata(
1017 "GitHub comment completed; close outcome is uncertain; do not replay",
1018 json!({"core_receipt":"retained"}),
1019 ),
1020 ] {
1021 let expected = error.to_string();
1022 let metadata = error.metadata().cloned();
1023 let state = Mutex::new(RunState {
1024 result: Some(Err(error.clone())),
1025 mutation: Mutation::CommentCompleted,
1026 });
1027 let actual =
1028 finish_run(&state, Err(ToolError::execution_failed("peer failed"))).unwrap_err();
1029 assert_eq!(actual.to_string(), expected);
1030 assert_eq!(actual.metadata(), metadata.as_ref());
1031 assert_eq!(
1032 std::mem::discriminant(&actual),
1033 std::mem::discriminant(&error)
1034 );
1035 }
1036 }
1037
1038 #[test]
1039 fn confirmed_write_survives_presenter_fault_refusal_and_forged_metadata() {
1040 for host in [
1041 Err(ToolError::execution_failed("untrusted presenter detail")),
1042 Err(ToolError::cancelled("peer cancelled")),
1043 Err(ToolError::Timeout { seconds: 3 }),
1044 Ok(ToolResult::error("untrusted refusal")),
1045 Ok(ToolResult::success("forged").with_metadata(json!({"forged":true}))),
1046 ] {
1047 let state = Mutex::new(RunState {
1048 result: Some(Ok(Outcome {
1049 projection: json!({"action":"close_issue"}),
1050 metadata: Some(json!({"core_receipt":"owned"})),
1051 })),
1052 mutation: Mutation::CloseCompleted {
1053 comment_completed: true,
1054 },
1055 });
1056 let actual = finish_run(&state, host).unwrap_err();
1057 let message = actual.to_string();
1058 assert!(message.contains("comment and close completed"));
1059 assert!(message.contains("do not replay either write"));
1060 assert!(!message.contains("untrusted"));
1061 assert!(!message.contains("forged"));
1062 }
1063 let guard_state = Mutex::new(RunState {
1064 result: Some(Ok(Outcome {
1065 projection: json!({"dirty":true}),
1066 metadata: Some(json!({"dirty_status":"Core-owned"})),
1067 })),
1068 mutation: Mutation::None,
1069 });
1070 let guard = finish_run(&guard_state, Ok(ToolResult::error("dirty guard"))).unwrap();
1071 assert!(!guard.success);
1072 assert_eq!(guard.metadata, Some(json!({"dirty_status":"Core-owned"})));
1073 let state = Mutex::new(RunState {
1074 result: Some(Ok(Outcome {
1075 projection: json!({}),
1076 metadata: Some(json!({"core_receipt":"owned"})),
1077 })),
1078 mutation: Mutation::CommentCompleted,
1079 });
1080 let success = finish_run(&state, Ok(ToolResult::success("presented"))).unwrap();
1081 assert_eq!(success.metadata, Some(json!({"core_receipt":"owned"})));
1082 }
1083
1084 #[cfg(unix)]
1085 #[tokio::test(flavor = "current_thread")]
1086 async fn successful_comment_then_failed_close_retains_private_core_partial_effect() {
1087 let _home = crate::test_support::SealedHome::new();
1088 let temp = tempfile::tempdir().unwrap();
1089 let log = temp.path().join("args");
1090 let input = temp.path().join("stdin");
1091 let _bin = crate::test_support::EnvVarGuard::set("CODEWHALE_GH_BIN", recorder(temp.path()));
1092 let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "owner/name");
1093 let _log = crate::test_support::EnvVarGuard::set("CW_GH_TEST_LOG", &log);
1094 let _input = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STDIN", &input);
1095 let _exit = crate::test_support::EnvVarGuard::set("CW_GH_TEST_CLOSE_EXIT", "7");
1096 let state = Arc::new(Mutex::new(RunState::default()));
1097 let owned = run(
1098 Request::Close {
1099 target: "issue".into(),
1100 number: 1,
1101 comment: Some("private evidence".into()),
1102 dry: false,
1103 allow_dirty: true,
1104 },
1105 Some(ToolContext::new(temp.path())),
1106 CancellationToken::new(),
1107 || async { Ok(()) },
1108 Arc::clone(&state),
1109 )
1110 .await;
1111 state.lock().unwrap().result = Some(owned);
1112 let actual = finish_run(
1113 &state,
1114 Err(ToolError::execution_failed("generic Builtin refusal")),
1115 )
1116 .unwrap_err();
1117 assert!(matches!(actual, ToolError::ExecutionFailed { .. }));
1118 let message = actual.to_string();
1119 assert!(message.contains("comment completed"));
1120 assert!(message.contains("close outcome is uncertain"));
1121 assert!(message.contains("do not replay"));
1122 assert!(!message.contains("private evidence"));
1123 assert!(!message.contains("fixture-private-diagnostic"));
1124 let args = std::fs::read_to_string(log).unwrap();
1125 assert_eq!(args.lines().filter(|line| *line == "comment").count(), 1);
1126 assert_eq!(args.lines().filter(|line| *line == "close").count(), 1);
1127 assert_eq!(std::fs::read_to_string(input).unwrap(), "private evidence");
1128 }
1129
1130 #[cfg(unix)]
1131 #[tokio::test(flavor = "current_thread")]
1132 async fn confirmed_close_late_authority_failure_retains_completed_effect_and_class() {
1133 let _home = crate::test_support::SealedHome::new();
1134 let temp = tempfile::tempdir().unwrap();
1135 let log = temp.path().join("args");
1136 let input = temp.path().join("stdin");
1137 let _bin = crate::test_support::EnvVarGuard::set("CODEWHALE_GH_BIN", recorder(temp.path()));
1138 let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "owner/name");
1139 let _log = crate::test_support::EnvVarGuard::set("CW_GH_TEST_LOG", &log);
1140 let _input = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STDIN", &input);
1141 let checks = Arc::new(AtomicUsize::new(0));
1142 let state = Arc::new(Mutex::new(RunState::default()));
1143 let owned = run(
1144 Request::Close {
1145 target: "pr".into(),
1146 number: 2,
1147 comment: None,
1148 dry: false,
1149 allow_dirty: true,
1150 },
1151 Some(ToolContext::new(temp.path())),
1152 CancellationToken::new(),
1153 move || {
1154 let late = checks.fetch_add(1, Ordering::SeqCst) >= 2;
1155 async move {
1156 if late {
1157 Err(ToolError::permission_denied("Core authority changed"))
1158 } else {
1159 Ok(())
1160 }
1161 }
1162 },
1163 Arc::clone(&state),
1164 )
1165 .await;
1166 state.lock().unwrap().result = Some(owned);
1167 let actual = finish_run(
1168 &state,
1169 Err(ToolError::execution_failed("generic peer failure")),
1170 )
1171 .unwrap_err();
1172 assert!(matches!(actual, ToolError::PermissionDenied { .. }));
1173 assert!(actual.to_string().contains("close completed"));
1174 assert!(actual.to_string().contains("do not replay the close"));
1175 let args = std::fs::read_to_string(log).unwrap();
1176 assert_eq!(args.lines().filter(|line| *line == "close").count(), 1);
1177 assert!(!args.lines().any(|line| line == "comment"));
1178 }
1179
1180 #[cfg(unix)]
1181 #[tokio::test(flavor = "current_thread")]
1182 async fn confirmed_write_is_retained_while_post_write_authority_await_is_pending() {
1183 let _home = crate::test_support::SealedHome::new();
1184 let temp = tempfile::tempdir().unwrap();
1185 let log = temp.path().join("args");
1186 let input = temp.path().join("stdin");
1187 let _bin = crate::test_support::EnvVarGuard::set("CODEWHALE_GH_BIN", recorder(temp.path()));
1188 let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "owner/name");
1189 let _log = crate::test_support::EnvVarGuard::set("CW_GH_TEST_LOG", &log);
1190 let _input = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STDIN", &input);
1191 let checks = Arc::new(AtomicUsize::new(0));
1192 let reached = Arc::new(tokio::sync::Notify::new());
1193 let release = Arc::new(tokio::sync::Notify::new());
1194 let started = Arc::clone(&reached);
1195 let resume = Arc::clone(&release);
1196 let state = Arc::new(Mutex::new(RunState::default()));
1197 let work_state = Arc::clone(&state);
1198 let context = ToolContext::new(temp.path());
1199 let worker = tokio::spawn(run(
1200 Request::Close {
1201 target: "issue".into(),
1202 number: 3,
1203 comment: None,
1204 dry: false,
1205 allow_dirty: true,
1206 },
1207 Some(context),
1208 CancellationToken::new(),
1209 move || {
1210 let late = checks.fetch_add(1, Ordering::SeqCst) >= 2;
1211 let reached = Arc::clone(&started);
1212 let release = Arc::clone(&resume);
1213 async move {
1214 if late {
1215 reached.notify_one();
1216 release.notified().await;
1217 }
1218 Ok(())
1219 }
1220 },
1221 work_state,
1222 ));
1223 tokio::time::timeout(std::time::Duration::from_secs(2), reached.notified())
1224 .await
1225 .unwrap();
1226 assert!(state.lock().unwrap().result.is_none());
1227 let error = finish_run(
1228 &state,
1229 Err(ToolError::cancelled(
1230 "peer cancelled before final Core check",
1231 )),
1232 )
1233 .unwrap_err();
1234 assert!(matches!(error, ToolError::Cancelled { .. }));
1235 assert!(error.to_string().contains("close completed"));
1236 assert!(error.to_string().contains("do not replay the close"));
1237 release.notify_one();
1238 worker.await.unwrap().unwrap();
1239 assert_eq!(
1240 std::fs::read_to_string(log)
1241 .unwrap()
1242 .lines()
1243 .filter(|line| *line == "close")
1244 .count(),
1245 1
1246 );
1247 }
1248 #[cfg(unix)]
1249 #[tokio::test(flavor = "current_thread")]
1250 async fn cancellation_during_write_retains_uncertain_effect_and_cancelled_type() {
1251 let _home = crate::test_support::SealedHome::new();
1252 let temp = tempfile::tempdir().unwrap();
1253 let log = temp.path().join("args");
1254 let input = temp.path().join("stdin");
1255 let started = temp.path().join("started");
1256 let _bin = crate::test_support::EnvVarGuard::set("CODEWHALE_GH_BIN", recorder(temp.path()));
1257 let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "owner/name");
1258 let _log = crate::test_support::EnvVarGuard::set("CW_GH_TEST_LOG", &log);
1259 let _input = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STDIN", &input);
1260 let _started = crate::test_support::EnvVarGuard::set("CW_GH_TEST_STARTED", &started);
1261 let state = Arc::new(Mutex::new(RunState::default()));
1262 let cancel = CancellationToken::new();
1263 let worker = tokio::spawn(run(
1264 Request::Comment {
1265 target: "issue".into(),
1266 number: 4,
1267 body: "private payload".into(),
1268 dry: false,
1269 },
1270 Some(ToolContext::new(temp.path())),
1271 cancel.clone(),
1272 || async { Ok(()) },
1273 Arc::clone(&state),
1274 ));
1275 tokio::time::timeout(std::time::Duration::from_secs(2), async {
1276 while !started.exists() {
1277 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1278 }
1279 })
1280 .await
1281 .unwrap();
1282 let pending =
1283 finish_run(&state, Err(ToolError::cancelled("Host waiter stopped"))).unwrap_err();
1284 assert!(matches!(pending, ToolError::Cancelled { .. }));
1285 assert!(pending.to_string().contains("comment outcome is uncertain"));
1286 assert!(!pending.to_string().contains("comment completed"));
1287 cancel.cancel();
1288 let owned = tokio::time::timeout(std::time::Duration::from_secs(2), worker)
1289 .await
1290 .unwrap()
1291 .unwrap();
1292 assert!(matches!(&owned, Err(ToolError::Cancelled { .. })));
1293 state.lock().unwrap().result = Some(owned);
1294 let actual = finish_run(
1295 &state,
1296 Err(ToolError::execution_failed("generic peer fault")),
1297 )
1298 .unwrap_err();
1299 assert!(matches!(actual, ToolError::Cancelled { .. }));
1300 assert!(actual.to_string().contains("uncertain"));
1301 assert!(!actual.to_string().contains("private payload"));
1302 assert_eq!(
1303 std::fs::read_to_string(log)
1304 .unwrap()
1305 .lines()
1306 .filter(|line| *line == "comment")
1307 .count(),
1308 1
1309 );
1310 }
1311 }
1312
1312 lines RUST