File size: 4,136 Bytes
afa0cbf
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
use std::future::pending;
use std::sync::Arc;
use std::time::Duration;

use codex_code_mode_protocol::DEFAULT_EXEC_YIELD_TIME_MS;
use pretty_assertions::assert_eq;

use super::MAX_ERROR_BYTES;
use super::RequestError;
use super::enforce;
use super::startup;

#[tokio::test(start_paused = true)]
async fn stalled_transport_fails_after_its_deadline() {
    let task = tokio::spawn(enforce(
        "termination",
        Duration::ZERO,
        pending::<Result<(), tonic::Status>>(),
    ));
    tokio::task::yield_now().await;
    tokio::time::advance(Duration::from_secs(61)).await;

    assert!(matches!(
        task.await.expect("deadline task"),
        Err(RequestError::TimedOut(message))
            if message == "gRPC code-mode host timed out waiting for termination response"
    ));
}

#[tokio::test(start_paused = true)]
async fn requested_runtime_duration_is_added_to_the_transport_deadline() {
    let task = tokio::spawn(enforce(
        "wait",
        Duration::from_secs(120),
        pending::<Result<(), tonic::Status>>(),
    ));
    tokio::task::yield_now().await;
    tokio::time::advance(Duration::from_secs(61)).await;
    tokio::task::yield_now().await;
    assert!(!task.is_finished());

    tokio::time::advance(Duration::from_secs(120)).await;
    assert!(matches!(
        task.await.expect("deadline task"),
        Err(RequestError::TimedOut(_))
    ));
}

#[tokio::test(start_paused = true)]
async fn default_execution_yield_and_grace_extend_the_outcome_deadline() {
    let runtime_timeout =
        Duration::from_millis(DEFAULT_EXEC_YIELD_TIME_MS).saturating_add(Duration::from_secs(1));
    let task = tokio::spawn(enforce(
        "execution outcome",
        runtime_timeout,
        pending::<Result<(), tonic::Status>>(),
    ));
    tokio::task::yield_now().await;

    tokio::time::advance(Duration::from_secs(70)).await;
    tokio::task::yield_now().await;
    assert!(!task.is_finished());

    tokio::time::advance(Duration::from_secs(2)).await;
    assert!(matches!(
        task.await.expect("execution outcome deadline task"),
        Err(RequestError::TimedOut(message))
            if message == "gRPC code-mode host timed out waiting for execution outcome response"
    ));
}

#[tokio::test]
async fn transport_status_is_preserved() {
    let result = enforce("wait", Duration::ZERO, async {
        Err::<(), _>(tonic::Status::not_found("missing"))
    })
    .await;

    match result {
        Err(RequestError::Failed(error)) => {
            assert_eq!(error.code(), tonic::Code::NotFound);
            assert_eq!(error.message(), "missing");
        }
        _ => panic!("expected the original gRPC status"),
    }
}

#[tokio::test]
async fn transport_status_messages_are_bounded_at_utf8_boundaries() {
    let error = startup("session opening", async {
        Err::<(), _>(tonic::Status::internal("🦀".repeat(MAX_ERROR_BYTES)))
    })
    .await
    .expect_err("oversized gRPC status must fail");

    assert!(error.len() <= MAX_ERROR_BYTES);
    assert!(error.starts_with("gRPC code-mode session opening failed:"));
    assert!(error.ends_with("..."));
}

#[tokio::test(start_paused = true)]
async fn stalled_channel_acquisition_times_out_and_remains_retryable() {
    let channel = Arc::new(tokio::sync::OnceCell::new());
    let stalled_channel = Arc::clone(&channel);
    let stalled = tokio::spawn(async move {
        startup("transport connection", async {
            stalled_channel
                .get_or_try_init(pending::<Result<usize, String>>)
                .await
                .copied()
        })
        .await
    });
    tokio::task::yield_now().await;
    tokio::time::advance(Duration::from_secs(61)).await;

    assert_eq!(
        stalled.await.expect("channel connection task"),
        Err("gRPC code-mode host timed out waiting for transport connection response".to_string())
    );
    assert_eq!(
        startup("transport connection", async {
            channel
                .get_or_try_init(|| async { Ok::<_, String>(42usize) })
                .await
                .copied()
        })
        .await,
        Ok(42)
    );
}