alien-bindings 3.3.22

Alien direct in-process resource bindings
Documentation
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
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
1353
1354
1355
1356
1357
1358
1359
1360
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
1373
1374
1375
1376
1377
1378
1379
1380
1381
1382
1383
1384
1385
1386
1387
1388
1389
1390
1391
1392
1393
1394
1395
1396
1397
1398
1399
1400
1401
1402
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
1464
1465
1466
1467
1468
1469
1470
1471
1472
1473
1474
1475
1476
1477
1478
1479
1480
1481
1482
1483
1484
1485
1486
1487
1488
1489
1490
1491
1492
1493
1494
1495
1496
1497
1498
1499
1500
1501
1502
1503
1504
1505
1506
1507
1508
1509
1510
1511
1512
1513
1514
1515
1516
1517
1518
1519
1520
1521
1522
1523
1524
1525
1526
1527
1528
1529
1530
1531
1532
1533
1534
1535
1536
1537
1538
1539
1540
1541
1542
1543
1544
1545
1546
1547
1548
1549
1550
1551
1552
1553
1554
1555
1556
1557
1558
1559
1560
1561
1562
1563
1564
1565
1566
//! GCP Agent Platform sandbox provider.
//!
//! Sandboxes are `sandboxEnvironments` created under a durable reasoning engine and reached from
//! outside the guest through the `:execute` proxy, which forwards one request to the agent's
//! `POST /` envelope and returns its body verbatim. So every command, file operation and health
//! check is one envelope over that proxy, and the lifecycle verbs are long-running operations
//! polled to completion.

use std::collections::{BTreeMap, VecDeque};
use std::sync::Arc;
use std::time::Duration;

use async_trait::async_trait;
use base64::engine::general_purpose::STANDARD as BASE64;
use base64::Engine as _;
use futures::stream::{self, BoxStream};
use serde::Deserialize;
use serde_json::json;
use tracing::warn;

use crate::error::{ErrorData, Result};
use crate::traits::{
    Binding, CommandOutput, CreateSandboxRequest, JobError, JobExit, JobPoll, JobStart,
    PreviewCapability, ResolvedSandbox, RunCommandRequest, Sandbox, SandboxInstance, SandboxState,
};
use alien_core::{SandboxCapabilities, SandboxEgress};
use alien_error::{AlienError, Context, ContextError};
use alien_gcp_clients::gcp::agent_platform::{
    AgentPlatformApi, AgentPlatformErrorData, EgressControlConfig, SandboxCreateRequest,
    SandboxEnvironment, SandboxSnapshot,
};
use alien_gcp_clients::gcp::longrunning::{Operation, OperationResult};

/// The envelope protocol version this provider speaks. It matches the agent's `PROTOCOL_VERSION`;
/// a peer that answers a different one is refused rather than guessed at.
const AGENT_PROTOCOL_VERSION: u32 = 1;

/// The proxy holds one `:execute` request open for roughly this long, so a command whose deadline
/// is within it runs synchronously and anything longer is detached as a job and polled. Set below
/// the measured ceiling, because a command that overruns a synchronous execute is lost, where an
/// overrun job is still reachable by a later poll.
const MAX_SYNCHRONOUS_TIMEOUT: Duration = Duration::from_secs(30);

/// Longest sandbox id this provider will place in a proxy URL. A bound on what is handed back to a
/// caller, not on what the API mints — the names seen are far shorter.
const MAX_SANDBOX_ID: usize = 63;

/// How long a created sandbox has to reach `STATE_RUNNING`, and how often that is checked.
const SANDBOX_READY_ATTEMPTS: u32 = 150;
const SANDBOX_READY_INTERVAL: Duration = Duration::from_secs(2);

/// How long a lifecycle operation (`create`, `:pause`, `:resume`, `:snapshot`) is polled before it
/// is reported incomplete rather than waited on forever.
const OPERATION_POLL_ATTEMPTS: u32 = 150;
const OPERATION_POLL_INTERVAL: Duration = Duration::from_secs(2);

/// How long `terminate` polls the sandbox to `not-found`, turning an accepted delete into a
/// confirmed one.
const TERMINATE_POLL_ATTEMPTS: u32 = 30;
const TERMINATE_POLL_INTERVAL: Duration = Duration::from_secs(2);

/// How often a detached job is polled for new output.
const JOB_POLL_INTERVAL: Duration = Duration::from_secs(1);

/// The grace a job's poll loop allows past the command's own deadline before it cancels the job:
/// the agent kills the command at the deadline and the next poll reports it, and this covers the
/// round trips to observe that.
const JOB_POLL_GRACE: Duration = Duration::from_secs(15);

const CREATE: &str = "sandbox.create";
const GET: &str = "sandbox.get";
const GET_OR_CREATE: &str = "sandbox.getOrCreate";
const RUN_COMMAND: &str = "sandbox.runCommand";
const JOB_START: &str = "sandbox.jobStart";
const JOB_POLL: &str = "sandbox.jobPoll";
const JOB_CANCEL: &str = "sandbox.jobCancel";
const TERMINATE: &str = "sandbox.terminate";

/// The generation of a sandbox whose live container identity was not established: a state with no
/// reachable agent, or a bulk `list` that does not probe each sandbox. Never a value
/// `generation_from_boot_id` returns, so a real identity is always distinguishable from an
/// unprobed one.
const NO_GENERATION: u64 = 0;

/// A single health probe is bounded to this, because the client sets no per-request timeout and an
/// agent that accepts the connection but never answers would otherwise hang `get()` and `create()`
/// forever. Set above the proxy's ~30s synchronous window (see `MAX_SYNCHRONOUS_TIMEOUT`) rather
/// than tight to the round trip: too tight reports a healthy sandbox unreachable, and
/// `get_or_create` then provisions a fresh sandbox and loses the caller's filesystem — the failure
/// this task exists to prevent — where too loose only delays an already-broken sandbox.
const AGENT_PROBE_BUDGET: Duration = Duration::from_secs(60);

/// Maps a declared egress mode onto the template's `egressControlConfig`, or refuses one the API
/// cannot express.
///
/// `internetAccess` is a single boolean, so `AllowDomains` has no representation and is refused
/// rather than approximated into `allow` (which would open more than was asked) or `deny` (which
/// would close a caller out of hosts it named). Not called by the runtime verbs — the template is
/// pre-created — but this is the mapping the template controller uses, kept beside the provider so
/// the two agree on what a mode means. `sandbox_label` names the offending sandbox in the refusal.
pub fn egress_control_config(
    sandbox_label: &str,
    egress: &SandboxEgress,
) -> Result<EgressControlConfig> {
    let Some(internet_access) = egress.internet_access_switch() else {
        return Err(AlienError::new(ErrorData::InvalidInput {
            operation_context: "sandbox.template".to_string(),
            details: format!(
                "sandbox '{sandbox_label}' asked for domain-scoped egress, which Agent \
                 Platform cannot express; it offers only 'allow' (open) and 'deny' (closed)"
            ),
            field_name: Some("egress".to_string()),
        }));
    };

    Ok(EgressControlConfig {
        internet_access: Some(internet_access),
        extra: Default::default(),
    })
}

/// A Sandbox backed by the Vertex AI Agent Platform.
#[derive(Debug)]
pub struct GcpAgentPlatformSandbox {
    client: Arc<dyn AgentPlatformApi>,
    /// Bare reasoning-engine id the client interpolates into its paths. The binding may carry a
    /// full resource name, so it is reduced to its last segment once, here.
    engine: String,
    /// Template every sandbox is cut from, as a resource name the create body carries unchanged.
    template: String,
    /// Sandbox lifetime in seconds, from the declaration; absent takes the service default.
    max_lifetime_seconds: Option<u32>,
}

impl GcpAgentPlatformSandbox {
    /// The `ttl` a create is sent with: what the caller asked for, never above what the
    /// declaration allows.
    fn lifetime_seconds(&self, timeout_ms: Option<u64>, operation: &str) -> Result<Option<u32>> {
        match timeout_ms {
            Some(timeout_ms) => {
                super::requested_lifetime_seconds(timeout_ms, self.max_lifetime_seconds, operation)
                    .map(Some)
            }
            None => Ok(self.max_lifetime_seconds),
        }
    }

    /// Builds a provider bound to one engine and template.
    ///
    /// The engine is normalised to its last path segment because the client builds the full
    /// resource path itself; passing the whole name would double it and address nothing.
    pub fn new(
        client: Arc<dyn AgentPlatformApi>,
        engine: String,
        template: String,
        max_lifetime_seconds: Option<u32>,
    ) -> Self {
        let engine = engine.rsplit('/').next().unwrap_or(&engine).to_string();
        Self {
            client,
            engine,
            template,
            max_lifetime_seconds,
        }
    }

    /// The engine id sent to the client. Exists so a test can prove the binding's full resource
    /// name was reduced to a bare segment — a doubled path is invisible against the mock otherwise.
    #[cfg(test)]
    pub(crate) fn engine(&self) -> &str {
        &self.engine
    }

    fn unsupported(&self, capability: &str, reason: &str) -> AlienError<ErrorData> {
        AlienError::new(ErrorData::OperationNotSupported {
            operation: capability.to_string(),
            reason: reason.to_string(),
        })
    }

    /// A sandbox id that stays a single path segment.
    ///
    /// The id is interpolated into the proxy URL, so one carrying `/`, `..`, `?` or `#` would
    /// address a different sandbox — a resource the same engine grant can reach. The API mints
    /// these; this bounds the ones a caller hands back.
    fn checked_sandbox_id(operation: &str, sandbox_id: &str) -> Result<()> {
        if is_addressable_id(sandbox_id) {
            return Ok(());
        }
        Err(AlienError::new(ErrorData::InvalidInput {
            operation_context: operation.to_string(),
            details: format!(
                "sandbox id '{sandbox_id}' must be a single segment of letters, digits, '-' and \
                 '_', at most {MAX_SANDBOX_ID} characters"
            ),
            field_name: Some("sandboxId".to_string()),
        }))
    }

    /// Reads a sandbox, or `None` when it is gone, without judging it.
    async fn read_sandbox(
        &self,
        operation: &str,
        sandbox_id: &str,
    ) -> Result<Option<SandboxEnvironment>> {
        match self.client.get_sandbox(&self.engine, sandbox_id).await {
            Ok(sandbox) => Ok(Some(sandbox)),
            Err(error) if is_not_found(&error) => Ok(None),
            Err(error) => Err(error.context(ErrorData::SandboxUnreachable {
                operation: operation.to_string(),
                reason: "the Agent Platform API did not answer a sandbox read".to_string(),
            })),
        }
    }

    /// Polls a lifecycle operation to completion, returning its response payload.
    ///
    /// Bounded rather than open-ended: a caller waiting forever is its own outage, and the
    /// operation name is carried so an incomplete one can be resumed rather than lost.
    async fn await_operation(
        &self,
        operation: &str,
        started: Operation,
    ) -> Result<serde_json::Value> {
        let Some(name) = started.name.clone() else {
            return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
                provider: "gcp-agent-platform".to_string(),
                binding_name: operation.to_string(),
                field: "name".to_string(),
                response_json: "the operation carried no resource name to poll".to_string(),
            }));
        };

        let mut current = started;
        for _ in 0..OPERATION_POLL_ATTEMPTS {
            if current.done == Some(true) {
                return finish_operation(operation, &name, current);
            }
            tokio::time::sleep(OPERATION_POLL_INTERVAL).await;
            current =
                self.client
                    .get_operation(&name)
                    .await
                    .context(ErrorData::SandboxUnreachable {
                        operation: operation.to_string(),
                        reason: format!("could not read operation '{name}'"),
                    })?;
        }

        if current.done == Some(true) {
            return finish_operation(operation, &name, current);
        }
        Err(AlienError::new(ErrorData::SandboxUnreachable {
            operation: operation.to_string(),
            reason: format!("operation '{name}' did not complete within its polling budget"),
        }))
    }

    /// Sends one envelope through the `:execute` proxy and returns the agent's body verbatim.
    ///
    /// A client error is a transport failure — the proxy could not deliver or the API refused. A
    /// body the op's parser cannot read is the agent's own reason, handled by each verb. A
    /// not-found is reported as a gone sandbox so a caller does not read it as a live one.
    async fn execute_op(
        &self,
        sandbox_id: &str,
        operation: &str,
        envelope: serde_json::Value,
    ) -> Result<Vec<u8>> {
        let body = serde_json::to_vec(&envelope).map_err(|error| {
            AlienError::new(ErrorData::SerializationFailed {
                message: format!("could not encode the {operation} envelope: {error}"),
            })
        })?;

        self.client
            .execute(&self.engine, sandbox_id, &body)
            .await
            .map_err(|error| Self::execute_failed(operation, error))
    }

    fn execute_failed(
        operation: &str,
        error: AlienError<AgentPlatformErrorData>,
    ) -> AlienError<ErrorData> {
        if is_not_found(&error) {
            return error.context(ErrorData::SandboxCommandFailed {
                failure: "sandboxGone".to_string(),
                reason: format!("{operation}: the sandbox does not exist"),
            });
        }
        // The client does not tell a delivered-but-failed call apart from an undelivered one, so
        // a `:execute` carrying a command leaves its outcome unestablished. The cause stays on the
        // chain rather than in `reason`, keeping a redacted request body out of an externally
        // visible message.
        if operation == RUN_COMMAND || operation == JOB_START {
            return error.context(ErrorData::SandboxOutcomeUnknown {
                operation: operation.to_string(),
                reason: "the sandbox did not complete the call".to_string(),
            });
        }
        error.context(ErrorData::SandboxCommandFailed {
            failure: "executeFailed".to_string(),
            reason: format!("{operation} could not be completed against the sandbox"),
        })
    }

    /// Confirms the agent answers and speaks the protocol, and returns the sandbox's generation.
    ///
    /// A sandbox can report `STATE_RUNNING` while every `:execute` fails, so a state read is not a
    /// health check; the agent has to answer for the sandbox to be usable. The reply carries the
    /// container boot id, from which the generation is derived so a caller can detect a container
    /// that was replaced under a stable sandbox name.
    async fn probe_agent(&self, operation: &str, sandbox_id: &str) -> Result<u64> {
        // Mapped to unreachable whatever the failure — a refused delivery, a probe that outran its
        // budget, an unparseable body, a protocol mismatch — because a health probe is idempotent
        // and the caller acts on the same thing each way: the agent cannot be reached, so
        // `get_or_create` provisions a fresh one rather than destroying a sandbox it did not create.
        let unreachable = |reason: String| {
            AlienError::new(ErrorData::SandboxUnreachable {
                operation: operation.to_string(),
                reason,
            })
        };

        let body = tokio::time::timeout(
            AGENT_PROBE_BUDGET,
            self.client.execute(
                &self.engine,
                sandbox_id,
                &serde_json::to_vec(&json!({ "v": AGENT_PROTOCOL_VERSION, "op": "health" }))
                    .unwrap_or_default(),
            ),
        )
        .await
        .map_err(|_| {
            unreachable(format!(
                "the sandbox's agent did not answer a health probe within {}s",
                AGENT_PROBE_BUDGET.as_secs()
            ))
        })?
        .map_err(|error| {
            error.context(ErrorData::SandboxUnreachable {
                operation: operation.to_string(),
                reason: "the sandbox's agent did not answer a health probe".to_string(),
            })
        })?;

        #[derive(Deserialize)]
        #[serde(rename_all = "camelCase")]
        struct Health {
            protocol_version: u32,
            boot_id: String,
        }

        let health: Health = serde_json::from_slice(&body).map_err(|_| {
            unreachable(format!(
                "the sandbox's agent answered a health probe with a body this provider cannot \
                 read: {}",
                truncated(&body)
            ))
        })?;

        if health.protocol_version != AGENT_PROTOCOL_VERSION {
            return Err(unreachable(format!(
                "the sandbox's agent speaks protocol {} where this provider speaks {}",
                health.protocol_version, AGENT_PROTOCOL_VERSION
            )));
        }
        // An agent that answers without a boot id cannot be told apart from a replaced container,
        // so the sandbox is refused rather than reconnected to a possibly-blank one.
        if health.boot_id.is_empty() {
            return Err(unreachable(
                "the sandbox's agent reported no container boot id, so its identity cannot be \
                 established"
                    .to_string(),
            ));
        }
        Ok(generation_from_boot_id(&health.boot_id))
    }

    /// Deletes a sandbox the caller will never receive, keeping the reason it is discarded.
    ///
    /// Every failure after the sandbox exists reaches here, so `create` has one delete rather than
    /// one beside each `?`. The delete's own failure names the leak without replacing the finding
    /// that caused it. A not-found delete is already success in the client.
    async fn discard(
        &self,
        sandbox_id: &str,
        reason: AlienError<ErrorData>,
    ) -> AlienError<ErrorData> {
        let Err(error) = self.client.delete_sandbox(&self.engine, sandbox_id).await else {
            return reason;
        };
        warn!(
            sandbox = %sandbox_id,
            %error,
            "could not delete a sandbox that was never handed to its caller"
        );
        reason.context(ErrorData::SandboxCommandFailed {
            failure: "sandboxLeftBehind".to_string(),
            reason: format!(
                "sandbox '{sandbox_id}' was not handed to its caller and could not be deleted, so \
                 it is still running"
            ),
        })
    }

    /// Waits for a created sandbox to reach `STATE_RUNNING`, confirms its agent answers, and returns
    /// the sandbox's generation.
    ///
    /// The running record is judged, not the create accept: a sandbox still coming up need not be
    /// addressable yet, and reading that as a failure would delete every one that answered early.
    async fn settle(&self, sandbox_id: &str) -> Result<u64> {
        for _ in 0..SANDBOX_READY_ATTEMPTS {
            let Some(sandbox) = self.read_sandbox(CREATE, sandbox_id).await? else {
                return Err(AlienError::new(ErrorData::SandboxCommandFailed {
                    failure: "sandboxGone".to_string(),
                    reason: format!("sandbox '{sandbox_id}' disappeared while it was coming up"),
                }));
            };
            match sandbox_state(CREATE, sandbox.state.as_deref())? {
                SandboxState::Running => {
                    return self.probe_agent(CREATE, sandbox_id).await;
                }
                SandboxState::Terminated => {
                    return Err(AlienError::new(ErrorData::SandboxCommandFailed {
                        failure: "sandboxTerminated".to_string(),
                        reason: format!(
                            "sandbox '{sandbox_id}' reached a terminal state while starting"
                        ),
                    }));
                }
                // Waited on rather than woken: a fresh sandbox has no idle-suspend policy to pause
                // it before its first command — the binding carries no such field — so a suspended
                // reading here is a transient step on the way up, not a resting state to resume.
                SandboxState::Starting | SandboxState::Paused => {}
            }
            tokio::time::sleep(SANDBOX_READY_INTERVAL).await;
        }
        Err(AlienError::new(ErrorData::SandboxUnreachable {
            operation: CREATE.to_string(),
            reason: format!(
                "sandbox '{sandbox_id}' was not running after {}s",
                SANDBOX_READY_ATTEMPTS as u64 * SANDBOX_READY_INTERVAL.as_secs()
            ),
        }))
    }

    /// Runs a command inside the proxy's synchronous window, streaming the buffered NDJSON body.
    async fn run_synchronous(
        &self,
        sandbox_id: &str,
        request: &RunCommandRequest,
    ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
        let envelope = exec_envelope("exec", sandbox_id, request);
        let body = self.execute_op(sandbox_id, RUN_COMMAND, envelope).await?;
        let frames = parse_exec_frames(&body)?;
        Ok(Box::pin(stream::iter(frames)))
    }

    /// Runs a command as a detached job whose output is polled for until it ends.
    async fn run_detached(
        &self,
        sandbox_id: &str,
        request: RunCommandRequest,
    ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
        let timeout = request.timeout;
        let started = self.start_job(sandbox_id, request).await?;

        let state = JobPollState {
            client: self.client.clone(),
            engine: self.engine.clone(),
            sandbox_id: sandbox_id.to_string(),
            job_id: started.job_id,
            since_seq: None,
            pending: VecDeque::new(),
            finished: false,
            stopped: false,
            deadline_at: tokio::time::Instant::now() + timeout + JOB_POLL_GRACE,
        };

        Ok(Box::pin(stream::unfold(state, job_poll_step)))
    }

    /// Refuses a command the agent would refuse anyway, before a call is spent on it.
    fn checked_command(operation: &str, request: &RunCommandRequest) -> Result<()> {
        if request.command.is_empty() {
            return Err(AlienError::new(ErrorData::InvalidInput {
                operation_context: operation.to_string(),
                details: "a command must name a program to run".to_string(),
                field_name: Some("command".to_string()),
            }));
        }
        // Refused rather than defaulted, and refused where it floors to zero milliseconds too: the
        // agent rejects a `timeoutMs` of 0, and a defaulted timeout is a hang waiting for a slow
        // day in a sandbox running code the caller does not control.
        if timeout_millis(request.timeout) == 0 {
            return Err(AlienError::new(ErrorData::SandboxCommandFailed {
                failure: "invalidRequest".to_string(),
                reason: "a command must carry a timeout of at least one millisecond".to_string(),
            }));
        }
        Ok(())
    }
}

impl Binding for GcpAgentPlatformSandbox {}

#[async_trait]
impl Sandbox for GcpAgentPlatformSandbox {
    fn as_any(&self) -> &dyn std::any::Any {
        self
    }

    /// The platform's row, unnarrowed. `sandboxLifetime` stays true even with no declared ttl:
    /// the API always sets `expireTime` on output, so an undeclared sandbox still carries a
    /// deadline the platform enforces.
    fn capabilities(&self) -> SandboxCapabilities {
        SandboxCapabilities::gcp_agent_platform()
    }

    async fn create(&self, request: CreateSandboxRequest) -> Result<SandboxInstance> {
        // A sandbox takes no environment of its own: `SandboxCreateRequest` has no env field,
        // so silently dropping one would run the caller's code without the variables it asked for.
        // They travel per command through `run_command` instead.
        // `OperationNotSupported`, not `InvalidInput`: the value is fine, the backend has nowhere
        // to put it. AWS answers the identical condition the same way, and a portable caller
        // branching on the code must not get two answers for one situation.
        if !request.env.is_empty() {
            return Err(AlienError::new(ErrorData::OperationNotSupported {
                operation: CREATE.to_string(),
                reason: "Agent Platform sandboxes take no sandbox-level env; set env per command \
                         instead"
                    .to_string(),
            }));
        }

        // Same reason as `env` above: nowhere to carry a tenant key, so accepting one would
        // silently merge tenants into one sandbox.
        if request.tenant_key.is_some() {
            return Err(AlienError::new(ErrorData::OperationNotSupported {
                operation: CREATE.to_string(),
                reason: "Agent Platform sandboxes take no tenantKey; create one sandbox per \
                         tenant instead"
                    .to_string(),
            }));
        }

        let ttl = self
            .lifetime_seconds(request.timeout_ms, CREATE)?
            .map(|seconds| format!("{seconds}s"));

        let started = self
            .client
            .create_sandbox(
                &self.engine,
                SandboxCreateRequest {
                    display_name: request.sandbox_id.clone(),
                    sandbox_environment_template: Some(self.template.clone()),
                    sandbox_environment_snapshot: None,
                    ttl,
                },
            )
            .await
            .context(ErrorData::SandboxUnreachable {
                operation: CREATE.to_string(),
                reason: "the Agent Platform API refused a sandbox create".to_string(),
            })?;

        let created: SandboxEnvironment = serde_json::from_value(
            self.await_operation(CREATE, started).await?,
        )
        .map_err(|error| {
            AlienError::new(ErrorData::UnexpectedResponseFormat {
                provider: "gcp-agent-platform".to_string(),
                binding_name: CREATE.to_string(),
                field: "response".to_string(),
                response_json: format!("the create operation resolved to a non-sandbox: {error}"),
            })
        })?;

        // The caller's requested id is not authoritative — the API allocates the name, and the
        // last segment is the id every later verb addresses it by. One this client cannot send is
        // one nothing can reach or reap, so an unreadable name is reported without a delete it
        // cannot target.
        let Some(sandbox_id) = created.name.as_deref().and_then(sandbox_segment) else {
            return Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
                provider: "gcp-agent-platform".to_string(),
                binding_name: CREATE.to_string(),
                field: "name".to_string(),
                response_json: format!("{:?}", created.name),
            }));
        };
        let sandbox_id = sandbox_id.to_string();

        // Past here a sandbox exists the caller has no id for, so every failure deletes it.
        match self.settle(&sandbox_id).await {
            Ok(generation) => Ok(SandboxInstance {
                sandbox_id,
                state: SandboxState::Running,
                generation,
            }),
            Err(error) => Err(self.discard(&sandbox_id, error).await),
        }
    }

    async fn get(&self, sandbox_id: &str) -> Result<Option<SandboxInstance>> {
        Self::checked_sandbox_id(GET, sandbox_id)?;
        let Some(sandbox) = self.read_sandbox(GET, sandbox_id).await? else {
            return Ok(None);
        };

        let state = sandbox_state(GET, sandbox.state.as_deref())?;
        // Only a running sandbox carries a reachable agent, and a state read is not health: a
        // running record whose agent does not answer is not reported as usable. A non-running
        // sandbox has no live container to identify, so it carries no generation.
        let generation = if state == SandboxState::Running {
            self.probe_agent(GET, sandbox_id).await?
        } else {
            NO_GENERATION
        };

        Ok(Some(SandboxInstance {
            sandbox_id: sandbox_id.to_string(),
            state,
            generation,
        }))
    }

    async fn get_or_create(&self, request: CreateSandboxRequest) -> Result<ResolvedSandbox> {
        if let Some(id) = request.sandbox_id.as_deref() {
            // A running, reachable sandbox is handed back; anything else is served by a fresh
            // sandbox rather than by destroying one this call did not create, which may be
            // another revision's.
            match self.get(id).await {
                Ok(Some(sandbox)) if sandbox.state == SandboxState::Running => {
                    return Ok(ResolvedSandbox::found(sandbox))
                }
                // The ordinary resting state for a reconnect: a suspended sandbox is woken and
                // confirmed, and handed back if it comes up healthy. A wake this call made that
                // cannot be confirmed is put back to sleep before a fresh sandbox is provisioned —
                // the paused one may be another revision's, and a second live sandbox beside it is
                // a leak the caller never receives an id for.
                Ok(Some(sandbox)) if sandbox.state == SandboxState::Paused => {
                    if self.resume(id).await.is_ok() {
                        match self.get(id).await {
                            Ok(Some(woken)) if woken.state == SandboxState::Running => {
                                return Ok(ResolvedSandbox::found(woken))
                            }
                            _ => {
                                // The wake could not be undone: leaving it live beside a fresh
                                // sandbox is a leak the caller gets no id for. Fail so the woken
                                // sandbox stays identifiable rather than provisioning a second one.
                                if let Err(error) = self.pause(id).await {
                                    return Err(error.context(ErrorData::SandboxCommandFailed {
                                        failure: "resumeRollbackFailed".to_string(),
                                        reason: format!(
                                            "{GET_OR_CREATE}: woke sandbox '{id}' but could not \
                                             confirm it healthy or put it back to sleep"
                                        ),
                                    }));
                                }
                            }
                        }
                    }
                }
                // Still coming up, or already being woken by someone else. Waited for rather than
                // replaced: the sandbox keeps starting either way, and a second one beside it is
                // the leak the arm above exists to avoid. A slow data plane is answered with the
                // failure, never by provisioning more of it.
                Ok(Some(sandbox)) if sandbox.state == SandboxState::Starting => {
                    let generation = self.settle(id).await?;
                    return Ok(ResolvedSandbox::found(SandboxInstance {
                        sandbox_id: id.to_string(),
                        state: SandboxState::Running,
                        generation,
                    }));
                }
                Ok(_) => {}
                Err(error) if error.code == "SANDBOX_UNREACHABLE" => {}
                Err(error) => {
                    return Err(error.context(ErrorData::SandboxCommandFailed {
                        failure: "getOrCreateFailed".to_string(),
                        reason: format!("{GET_OR_CREATE}: reaching sandbox '{id}' failed"),
                    }))
                }
            }
        }

        self.create(request).await.map(ResolvedSandbox::created)
    }

    async fn list(&self) -> Result<Vec<SandboxInstance>> {
        let sandboxes = self.client.list_sandboxes(&self.engine).await.context(
            ErrorData::SandboxUnreachable {
                operation: "sandbox.list".to_string(),
                reason: "the Agent Platform API did not answer a sandbox list".to_string(),
            },
        )?;

        // A sandbox this provider cannot fully read — an unaddressable name or an unrecognised
        // state — is left out rather than surfaced as a handle to nothing or failing the whole
        // enumeration; one odd sandbox must not hide every other from an orphan sweep. Both halves
        // are skipped for the same reason, so leniency is consistent across the record.
        Ok(sandboxes
            .into_iter()
            .filter_map(|sandbox| {
                let sandbox_id = sandbox.name.as_deref().and_then(sandbox_segment)?;
                let state = sandbox_state("sandbox.list", sandbox.state.as_deref()).ok()?;
                // A bulk list does not probe each agent, so it reports no generation; a caller that
                // needs one reads the single sandbox through `get`.
                Some(SandboxInstance {
                    sandbox_id: sandbox_id.to_string(),
                    state,
                    generation: NO_GENERATION,
                })
            })
            .collect())
    }

    async fn run_command(
        &self,
        sandbox_id: &str,
        request: RunCommandRequest,
    ) -> Result<BoxStream<'static, Result<CommandOutput>>> {
        Self::checked_sandbox_id(RUN_COMMAND, sandbox_id)?;
        Self::checked_command(RUN_COMMAND, &request)?;

        // The synchronous window is the proxy's, not the command's: a command that outlives one
        // `:execute` is detached as a job so a later poll can still reach its output.
        if request.timeout <= MAX_SYNCHRONOUS_TIMEOUT {
            self.run_synchronous(sandbox_id, &request).await
        } else {
            self.run_detached(sandbox_id, request).await
        }
    }

    async fn start_job(&self, sandbox_id: &str, request: RunCommandRequest) -> Result<JobStart> {
        Self::checked_sandbox_id(JOB_START, sandbox_id)?;
        Self::checked_command(JOB_START, &request)?;

        let envelope = exec_envelope("jobStart", sandbox_id, &request);
        let body = self.execute_op(sandbox_id, JOB_START, envelope).await?;

        // The `:execute` succeeded, so the job was accepted and is running; only its id could not
        // be read. Nothing can poll or cancel it after this, and the command's own deadline is
        // what bounds it — so the outcome is unestablished rather than a reply that failed to read.
        let started: JobStartReply = serde_json::from_slice(&body).map_err(|_| {
            AlienError::new(ErrorData::UnexpectedResponseFormat {
                provider: "gcp-agent-platform".to_string(),
                binding_name: JOB_START.to_string(),
                field: "jobId".to_string(),
                response_json: truncated(&body),
            })
            .context(ErrorData::SandboxOutcomeUnknown {
                operation: JOB_START.to_string(),
                reason: "the job started and its id could not be read, so it cannot be polled"
                    .to_string(),
            })
        })?;

        Ok(JobStart {
            job_id: started.job_id,
        })
    }

    async fn poll_job(
        &self,
        sandbox_id: &str,
        job_id: &str,
        since_seq: Option<u64>,
    ) -> Result<JobPoll> {
        Self::checked_sandbox_id(JOB_POLL, sandbox_id)?;
        let reply = poll_once(
            self.client.as_ref(),
            &self.engine,
            sandbox_id,
            job_id,
            since_seq,
        )
        .await?;

        Ok(JobPoll {
            running: reply.running,
            frames: reply
                .frames
                .into_iter()
                .map(WireFrame::into_output)
                .collect::<Result<Vec<_>>>()?,
            exit: reply.exit_code.map(|code| JobExit {
                code,
                truncated: reply.truncated.unwrap_or(false),
            }),
            error: reply.error.map(|error| JobError {
                code: error.code,
                message: error.message,
            }),
        })
    }

    async fn cancel_job(&self, sandbox_id: &str, job_id: &str) -> Result<()> {
        Self::checked_sandbox_id(JOB_CANCEL, sandbox_id)?;
        let body = self
            .client
            .execute(&self.engine, sandbox_id, &cancel_body(job_id))
            .await
            .map_err(|error| unanswered_job(JOB_CANCEL, error))?;

        if !cancel_confirmed(&body) {
            return Err(AlienError::new(ErrorData::SandboxCommandFailed {
                failure: "agentRefused".to_string(),
                reason: format!("{JOB_CANCEL}: {}", truncated(&body)),
            }));
        }

        Ok(())
    }

    async fn read_file(&self, sandbox_id: &str, path: &str) -> Result<Vec<u8>> {
        Self::checked_sandbox_id("sandbox.readFile", sandbox_id)?;
        let body = self
            .execute_op(
                sandbox_id,
                "sandbox.readFile",
                json!({ "v": AGENT_PROTOCOL_VERSION, "op": "readFile", "path": path }),
            )
            .await?;

        #[derive(Deserialize)]
        #[serde(rename_all = "camelCase")]
        struct ReadFile {
            contents_base64: String,
        }
        let read: ReadFile = serde_json::from_slice(&body).map_err(|_| {
            AlienError::new(ErrorData::SandboxCommandFailed {
                failure: "agentRefused".to_string(),
                reason: format!("sandbox.readFile was refused: {}", truncated(&body)),
            })
        })?;

        BASE64
            .decode(read.contents_base64.as_bytes())
            .map_err(|error| {
                AlienError::new(ErrorData::UnexpectedResponseFormat {
                    provider: "gcp-agent-platform".to_string(),
                    binding_name: "sandbox.readFile".to_string(),
                    field: "contentsBase64".to_string(),
                    response_json: format!("the agent returned data that is not base64: {error}"),
                })
            })
    }

    async fn write_files(&self, sandbox_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()> {
        Self::checked_sandbox_id("sandbox.writeFiles", sandbox_id)?;
        // One request per path, stopping at the first failure — the partial application every
        // backend performs, so a caller sees one contract rather than several. The agent's field
        // is `contentsBase64`; `contents` is dropped silently.
        for (path, contents) in files {
            let body = self
                .execute_op(
                    sandbox_id,
                    "sandbox.writeFiles",
                    json!({
                        "v": AGENT_PROTOCOL_VERSION,
                        "op": "writeFile",
                        "path": path,
                        "contentsBase64": BASE64.encode(&contents),
                    }),
                )
                .await?;
            confirm_empty_ok("sandbox.writeFiles", &body)?;
        }
        Ok(())
    }

    async fn preview(&self, _sandbox_id: &str, _port: u16) -> Result<PreviewCapability> {
        Err(self.unsupported(
            "preview",
            "Agent Platform mints no port-scoped ingress capability; the only ingress is :execute",
        ))
    }

    async fn pause(&self, sandbox_id: &str) -> Result<()> {
        Self::checked_sandbox_id("sandbox.pause", sandbox_id)?;
        let started = self.client.pause(&self.engine, sandbox_id).await.context(
            ErrorData::SandboxCommandFailed {
                failure: "pauseFailed".to_string(),
                reason: format!("sandbox.pause: sandbox '{sandbox_id}' could not be paused"),
            },
        )?;
        self.await_operation("sandbox.pause", started).await?;
        Ok(())
    }

    async fn resume(&self, sandbox_id: &str) -> Result<()> {
        Self::checked_sandbox_id("sandbox.resume", sandbox_id)?;
        let started = self.client.resume(&self.engine, sandbox_id).await.context(
            ErrorData::SandboxCommandFailed {
                failure: "resumeFailed".to_string(),
                reason: format!("sandbox.resume: sandbox '{sandbox_id}' could not be resumed"),
            },
        )?;
        self.await_operation("sandbox.resume", started).await?;
        Ok(())
    }

    async fn snapshot(&self, sandbox_id: &str) -> Result<String> {
        Self::checked_sandbox_id("sandbox.snapshot", sandbox_id)?;
        // A generated display name, because the API takes one and the caller does not supply it.
        // The trait has no restore verb, so the returned name is not yet consumable through it —
        // restore is `create` from a snapshot, which this backend can do but the trait cannot ask.
        let display_name = format!("snap-{}", uuid::Uuid::new_v4().simple());
        let started = self
            .client
            .snapshot(&self.engine, sandbox_id, &display_name)
            .await
            .context(ErrorData::SandboxCommandFailed {
                failure: "snapshotFailed".to_string(),
                reason: format!("sandbox.snapshot: sandbox '{sandbox_id}' could not be captured"),
            })?;

        let snapshot: SandboxSnapshot =
            serde_json::from_value(self.await_operation("sandbox.snapshot", started).await?)
                .map_err(|error| {
                    AlienError::new(ErrorData::UnexpectedResponseFormat {
                        provider: "gcp-agent-platform".to_string(),
                        binding_name: "sandbox.snapshot".to_string(),
                        field: "response".to_string(),
                        response_json: format!(
                            "the snapshot operation resolved to a non-snapshot: {error}"
                        ),
                    })
                })?;

        snapshot.name.ok_or_else(|| {
            AlienError::new(ErrorData::UnexpectedResponseFormat {
                provider: "gcp-agent-platform".to_string(),
                binding_name: "sandbox.snapshot".to_string(),
                field: "name".to_string(),
                response_json: "the snapshot completed without a resource name".to_string(),
            })
        })
    }

    async fn terminate(&self, sandbox_id: &str) -> Result<()> {
        Self::checked_sandbox_id(TERMINATE, sandbox_id)?;
        // Accepted, not completed: the client returns before the sandbox is gone. Returning here
        // would report containment while the code may still run, which is the whole point of
        // terminate — so the delete is confirmed by polling to not-found.
        // A sandbox that is already gone is the state terminate exists to reach, so not-found is
        // success. Narrowed to exactly that: mapping any failure to `Ok` would report containment
        // for a sandbox another deployment owns and this one was refused.
        if let Err(error) = self.client.delete_sandbox(&self.engine, sandbox_id).await {
            if !is_not_found(&error) {
                return Err(error.context(ErrorData::SandboxUnreachable {
                    operation: TERMINATE.to_string(),
                    reason: format!("the delete of sandbox '{sandbox_id}' was not accepted"),
                }));
            }
            return Ok(());
        }

        for _ in 0..TERMINATE_POLL_ATTEMPTS {
            match self.client.get_sandbox(&self.engine, sandbox_id).await {
                Err(error) if is_not_found(&error) => return Ok(()),
                // A read that fails is not a sandbox that is gone, and one throttled response must
                // not end the poll: the attempt budget decides.
                Err(error) => {
                    warn!(sandbox = %sandbox_id, %error, "could not confirm a sandbox is gone")
                }
                Ok(_) => {}
            }
            tokio::time::sleep(TERMINATE_POLL_INTERVAL).await;
        }

        Err(AlienError::new(ErrorData::SandboxUnreachable {
            operation: TERMINATE.to_string(),
            reason: format!(
                "deletion of '{sandbox_id}' was accepted but the sandbox was still present after \
                 {}s; it may still be running",
                TERMINATE_POLL_ATTEMPTS as u64 * TERMINATE_POLL_INTERVAL.as_secs()
            ),
        }))
    }
}

/// One step of a detached job's poll loop, yielding output frames as they arrive and a terminal
/// item once the job ends.
async fn job_poll_step(mut state: JobPollState) -> Option<(Result<CommandOutput>, JobPollState)> {
    loop {
        if let Some(item) = state.pending.pop_front() {
            return Some((item, state));
        }
        if state.finished {
            return None;
        }

        if tokio::time::Instant::now() >= state.deadline_at {
            // The deadline is this client's decision, so it only names an outcome once the cancel
            // that makes it true has landed. A cancel that fails leaves the job running, and the
            // caller has to be told that rather than that the command was stopped.
            let cancelled = state
                .client
                .execute(
                    &state.engine,
                    &state.sandbox_id,
                    &cancel_body(&state.job_id),
                )
                .await;
            let confirmed = cancelled.as_ref().is_ok_and(|body| cancel_confirmed(body));
            state.pending.push_back(Err(match cancelled {
                Ok(_) if confirmed => AlienError::new(ErrorData::SandboxCommandFailed {
                    failure: "timeoutExceeded".to_string(),
                    reason: "the command's deadline elapsed before its job reported an outcome"
                        .to_string(),
                }),
                Ok(_) => AlienError::new(ErrorData::SandboxOutcomeUnknown {
                    operation: RUN_COMMAND.to_string(),
                    reason: "the command's deadline elapsed and its job did not confirm the cancel"
                        .to_string(),
                }),
                Err(error) => error.context(ErrorData::SandboxOutcomeUnknown {
                    operation: RUN_COMMAND.to_string(),
                    reason: "the command's deadline elapsed and its job could not be cancelled"
                        .to_string(),
                }),
            }));
            state.finished = true;
            // Only a landed cancel stops the job; the other two arms report an outcome nobody
            // established, so the drop is left armed to try again.
            state.stopped = confirmed;
            continue;
        }

        let poll = match poll_once(
            state.client.as_ref(),
            &state.engine,
            &state.sandbox_id,
            &state.job_id,
            state.since_seq,
        )
        .await
        {
            Ok(poll) => poll,
            // A standalone poll is repeatable, but this one watches a running command, and giving
            // up on it leaves that command's outcome unestablished.
            Err(error) => {
                state
                    .pending
                    .push_back(Err(error.context(ErrorData::SandboxOutcomeUnknown {
                        operation: RUN_COMMAND.to_string(),
                        reason: "the job is no longer watched".to_string(),
                    })));
                state.finished = true;
                continue;
            }
        };

        for frame in poll.frames {
            // A seq gap is truncated output, not a frame still to come, so the cursor takes the
            // highest seq seen and the loop never waits for a "missing" one; `max` rather than the
            // last frame's seq so an out-of-order frame cannot walk the cursor backwards.
            state.since_seq = state.since_seq.max(frame.seq());
            let output = frame.into_output();
            // As in the synchronous path: a frame that will not convert ends the poll rather than
            // being queued ahead of a terminal result that would contradict it.
            let failed = output.is_err();
            state.pending.push_back(output);
            if failed {
                state.finished = true;
                break;
            }
        }

        if !poll.running {
            // The terminal outcome is the envelope's, not a frame's: a clean exit carries a code,
            // and a deadline, spawn failure or cancel carries an error object with no code.
            let terminal = match poll.error {
                Some(error) => Err(AlienError::new(ErrorData::SandboxCommandFailed {
                    failure: error.code,
                    reason: error.message,
                })),
                // A job that finished without an exit code never established its outcome; an
                // invented code is indistinguishable from one the command really exited with.
                None => match poll.exit_code {
                    Some(code) => Ok(CommandOutput::Exit {
                        code,
                        truncated: poll.truncated.unwrap_or(false),
                    }),
                    None => Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
                        operation: RUN_COMMAND.to_string(),
                        reason: "the job finished without reporting an exit code".to_string(),
                    })),
                },
            };
            state.pending.push_back(terminal);
            state.finished = true;
            state.stopped = true;
            continue;
        }

        if state.pending.is_empty() {
            tokio::time::sleep(JOB_POLL_INTERVAL).await;
        }
    }
}

/// One `jobPoll` against a sandbox, classified as a standalone poll: nothing about the job
/// changes, so a call that fails is worth repeating. `run_command`'s loop re-contexts it.
async fn poll_once(
    client: &dyn AgentPlatformApi,
    engine: &str,
    sandbox_id: &str,
    job_id: &str,
    since_seq: Option<u64>,
) -> Result<JobPollReply> {
    let body = client
        .execute(engine, sandbox_id, &poll_body(job_id, since_seq))
        .await
        .map_err(|error| unanswered_job(JOB_POLL, error))?;

    serde_json::from_slice(&body).map_err(|_| {
        AlienError::new(ErrorData::UnexpectedResponseFormat {
            provider: "gcp-agent-platform".to_string(),
            binding_name: JOB_POLL.to_string(),
            field: "jobPoll".to_string(),
            response_json: truncated(&body),
        })
    })
}

/// A poll or cancel that did not complete. Both leave the job exactly as it was, so unlike a
/// command they carry the retry signal; a sandbox that is gone is an answer rather than a failure.
fn unanswered_job(
    operation: &str,
    error: AlienError<AgentPlatformErrorData>,
) -> AlienError<ErrorData> {
    if is_not_found(&error) {
        return error.context(ErrorData::SandboxCommandFailed {
            failure: "sandboxGone".to_string(),
            reason: format!("{operation}: the sandbox does not exist"),
        });
    }
    error.context(ErrorData::SandboxUnreachable {
        operation: operation.to_string(),
        reason: "the sandbox did not complete the call".to_string(),
    })
}

/// Whether a `jobCancel` reply is the cancel landing.
///
/// A reply arriving is not the cancel succeeding: a cancelled job answers with the empty object and
/// a refusal — `JobNotFound`, say — with the agent's error text, both through a successful
/// `:execute`. A success body that grew a field then reads as unknown, never a refusal as landed.
fn cancel_confirmed(body: &[u8]) -> bool {
    serde_json::from_slice::<serde_json::Value>(body)
        .is_ok_and(|value| value.as_object().is_some_and(serde_json::Map::is_empty))
}

/// The id a started job answers to, as the agent's `jobStart` returns it.
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobStartReply {
    job_id: String,
}

/// The bookkeeping a detached job's poll loop carries between steps.
struct JobPollState {
    client: Arc<dyn AgentPlatformApi>,
    engine: String,
    sandbox_id: String,
    job_id: String,
    since_seq: Option<u64>,
    pending: VecDeque<Result<CommandOutput>>,
    finished: bool,
    /// Set only where the job itself ended, so a stream that ends any other way still cancels.
    stopped: bool,
    deadline_at: tokio::time::Instant,
}

/// Kills the command a dropped stream stops reading.
///
/// Where a command's frames ride a transport the caller can close, that close is the kill. A
/// detached job has no such transport, so the cancel is sent here or the job runs to its full
/// timeout, billing, with nobody able to reach it.
impl Drop for JobPollState {
    fn drop(&mut self) {
        if self.stopped {
            return;
        }
        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
            warn!(
                sandbox = %self.sandbox_id,
                job = %self.job_id,
                "a job's stream was dropped outside a runtime, so the job runs to its timeout"
            );
            return;
        };

        let client = self.client.clone();
        let engine = self.engine.clone();
        let sandbox_id = self.sandbox_id.clone();
        let job_id = self.job_id.clone();
        runtime.spawn(async move {
            // The reader is gone, so a failure has nobody to be returned to; what it leaves
            // running is named instead.
            match tokio::time::timeout(
                AGENT_PROBE_BUDGET,
                client.execute(&engine, &sandbox_id, &cancel_body(&job_id)),
            )
            .await
            {
                Ok(Ok(body)) if cancel_confirmed(&body) => {}
                Ok(Ok(body)) => warn!(
                    sandbox = %sandbox_id,
                    job = %job_id,
                    reply = %truncated(&body),
                    "a dropped job's cancel was refused, so the job runs to its timeout"
                ),
                Ok(Err(error)) => warn!(
                    sandbox = %sandbox_id,
                    job = %job_id,
                    %error,
                    "a dropped job's cancel did not land, so the job may run to its timeout"
                ),
                Err(_) => warn!(
                    sandbox = %sandbox_id,
                    job = %job_id,
                    budget_secs = AGENT_PROBE_BUDGET.as_secs(),
                    "a dropped job's cancel went unanswered, so the job may run to its timeout"
                ),
            }
        });
    }
}

/// A job's output so far, and how it ended once it has. Mirrors the agent's `jobPoll` reply.
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobPollReply {
    running: bool,
    #[serde(default)]
    frames: Vec<WireFrame>,
    #[serde(default)]
    exit_code: Option<i32>,
    #[serde(default)]
    truncated: Option<bool>,
    #[serde(default)]
    error: Option<JobErrorReply>,
}

#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JobErrorReply {
    code: String,
    message: String,
}

/// A frame as the agent writes it, shared by the synchronous NDJSON body and the job frames.
#[derive(Deserialize)]
#[serde(rename_all = "camelCase", tag = "t")]
enum WireFrame {
    Stdout {
        seq: u64,
        data: String,
    },
    Stderr {
        seq: u64,
        data: String,
    },
    Exit {
        code: i32,
        #[serde(default)]
        truncated: bool,
    },
    Error {
        code: String,
        message: String,
    },
}

impl WireFrame {
    fn is_terminal(&self) -> bool {
        matches!(self, Self::Exit { .. } | Self::Error { .. })
    }

    fn seq(&self) -> Option<u64> {
        match self {
            Self::Stdout { seq, .. } | Self::Stderr { seq, .. } => Some(*seq),
            _ => None,
        }
    }

    fn into_output(self) -> Result<CommandOutput> {
        match self {
            Self::Stdout { seq, data } => Ok(CommandOutput::Stdout {
                seq,
                data: decode_frame_data(&data)?,
            }),
            Self::Stderr { seq, data } => Ok(CommandOutput::Stderr {
                seq,
                data: decode_frame_data(&data)?,
            }),
            Self::Exit { code, truncated } => Ok(CommandOutput::Exit { code, truncated }),
            // An error frame is the command's outcome, so it surfaces as an error rather than a
            // stream that simply stopped.
            Self::Error { code, message } => {
                Err(AlienError::new(ErrorData::SandboxCommandFailed {
                    failure: code,
                    reason: message,
                }))
            }
        }
    }
}

/// A frame that arrived is proof the command ran, so a payload that will not decode leaves the
/// outcome unestablished rather than merely malformed.
fn decode_frame_data(data: &str) -> Result<Vec<u8>> {
    BASE64.decode(data).map_err(|error| {
        AlienError::new(ErrorData::UnexpectedResponseFormat {
            provider: "gcp-agent-platform".to_string(),
            binding_name: RUN_COMMAND.to_string(),
            field: "data".to_string(),
            response_json: format!("an output frame's data is not base64: {error}"),
        })
        .context(ErrorData::SandboxOutcomeUnknown {
            operation: RUN_COMMAND.to_string(),
            reason: "an output frame did not decode".to_string(),
        })
    })
}

/// Turns the agent's buffered NDJSON body into output frames.
///
/// A body that is not frames at all is the agent's error, reported as a refusal. A body that ends
/// without a terminal frame is a transport failure, not a command that finished: the command had
/// started, so the trailing item says the outcome is unknown rather than letting a truncated
/// stream read as success.
fn parse_exec_frames(body: &[u8]) -> Result<Vec<Result<CommandOutput>>> {
    let mut frames = Vec::new();
    let mut saw_any = false;
    let mut saw_terminal = false;

    for line in body.split(|byte| *byte == b'\n') {
        if line.is_empty() {
            continue;
        }
        match serde_json::from_slice::<WireFrame>(line) {
            Ok(frame) => {
                saw_any = true;
                saw_terminal |= frame.is_terminal();
                let output = frame.into_output();
                // A frame that will not convert ends the body: letting a later exit follow would
                // answer the question this item just reported as unanswerable. `saw_terminal`
                // stops the trailing item below from saying the same thing twice.
                let failed = output.is_err();
                frames.push(output);
                if failed {
                    saw_terminal = true;
                    break;
                }
            }
            Err(error) => {
                if !saw_any {
                    return Err(AlienError::new(ErrorData::SandboxCommandFailed {
                        failure: "agentRefused".to_string(),
                        reason: format!("run_command was refused: {}", truncated(body)),
                    }));
                }
                // Frames already arrived, so the command ran and this leaves its end unknown.
                // `saw_terminal` stops the trailing item below from saying the same thing twice.
                frames.push(Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
                    provider: "gcp-agent-platform".to_string(),
                    binding_name: RUN_COMMAND.to_string(),
                    field: "frame".to_string(),
                    response_json: format!("an output frame did not parse: {error}"),
                })
                .context(ErrorData::SandboxOutcomeUnknown {
                    operation: RUN_COMMAND.to_string(),
                    reason: "an output frame did not parse".to_string(),
                })));
                saw_terminal = true;
                break;
            }
        }
    }

    if !saw_any {
        return Err(AlienError::new(ErrorData::SandboxCommandFailed {
            failure: "agentRefused".to_string(),
            reason: "run_command returned an empty body".to_string(),
        }));
    }
    if !saw_terminal {
        frames.push(Err(AlienError::new(ErrorData::SandboxOutcomeUnknown {
            operation: RUN_COMMAND.to_string(),
            reason: "the command's output ended without a terminal frame".to_string(),
        })));
    }
    Ok(frames)
}

/// The envelope for `exec` or `jobStart`. `timeoutMs` is the field the agent reads; both ops take
/// the identical body.
fn exec_envelope(op: &str, _sandbox_id: &str, request: &RunCommandRequest) -> serde_json::Value {
    json!({
        "v": AGENT_PROTOCOL_VERSION,
        "op": op,
        "command": request.argv(),
        "timeoutMs": timeout_millis(request.timeout),
        "cwd": request.cwd,
        "env": request.env,
    })
}

fn poll_body(job_id: &str, since_seq: Option<u64>) -> Vec<u8> {
    serde_json::to_vec(&json!({
        "v": AGENT_PROTOCOL_VERSION,
        "op": "jobPoll",
        "jobId": job_id,
        "sinceSeq": since_seq,
    }))
    .unwrap_or_default()
}

fn cancel_body(job_id: &str) -> Vec<u8> {
    serde_json::to_vec(&json!({
        "v": AGENT_PROTOCOL_VERSION,
        "op": "jobCancel",
        "jobId": job_id,
    }))
    .unwrap_or_default()
}

/// Milliseconds, saturated: a timeout long enough to overflow `u64` ms is not one anyone meant,
/// and wrapping it would turn "effectively forever" into "immediately".
fn timeout_millis(timeout: Duration) -> u64 {
    u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX)
}

/// Reads a `writeFile` reply, which succeeds with an empty body.
///
/// A non-empty body from the op is the agent's error text, not a success shape, so it is
/// surfaced as a refusal rather than ignored.
fn confirm_empty_ok(operation: &str, body: &[u8]) -> Result<()> {
    if body.iter().all(|byte| byte.is_ascii_whitespace()) {
        return Ok(());
    }
    Err(AlienError::new(ErrorData::SandboxCommandFailed {
        failure: "agentRefused".to_string(),
        reason: format!("{operation} was refused: {}", truncated(body)),
    }))
}

/// The last path segment, if it is a usable id. Used for both minted names and listed ones.
fn sandbox_segment(name: &str) -> Option<&str> {
    let segment = name.rsplit('/').next()?;
    is_addressable_id(segment).then_some(segment)
}

fn is_addressable_id(id: &str) -> bool {
    !id.is_empty()
        && id.len() <= MAX_SANDBOX_ID
        && id
            .chars()
            .all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_')
}

/// Maps a container boot id to a numeric generation deterministically.
///
/// A caller may compare generations across processes, so this is an explicit FNV-1a rather than a
/// `Hash` impl — the same boot id must yield the same number in any build, and std's hashers
/// promise no cross-release stability. `| 1` keeps the result clear of `NO_GENERATION`.
fn generation_from_boot_id(boot_id: &str) -> u64 {
    const FNV_OFFSET_BASIS: u64 = 0xcbf2_9ce4_8422_2325;
    const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
    let mut hash = FNV_OFFSET_BASIS;
    for byte in boot_id.as_bytes() {
        hash ^= u64::from(*byte);
        hash = hash.wrapping_mul(FNV_PRIME);
    }
    hash | 1
}

/// The API's runtime states, in ours. An unrecognised one is an error rather than a default,
/// because every default here is a lie a caller acts on.
fn sandbox_state(operation: &str, state: Option<&str>) -> Result<SandboxState> {
    match state {
        Some("STATE_RUNNING") => Ok(SandboxState::Running),
        Some("STATE_CREATING" | "STATE_PENDING" | "STATE_RESUMING") => Ok(SandboxState::Starting),
        Some("STATE_PAUSED" | "STATE_PAUSING" | "STATE_SUSPENDED") => Ok(SandboxState::Paused),
        Some("STATE_STOPPED" | "STATE_FAILED" | "STATE_DELETING" | "STATE_DELETED") => {
            Ok(SandboxState::Terminated)
        }
        other => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
            provider: "gcp-agent-platform".to_string(),
            binding_name: operation.to_string(),
            field: "state".to_string(),
            response_json: other
                .map_or_else(|| "absent".to_string(), |state| format!("\"{state}\"")),
        })),
    }
}

/// Turns a completed operation into its response payload, or the error it reported.
fn finish_operation(operation: &str, name: &str, op: Operation) -> Result<serde_json::Value> {
    match op.result {
        Some(OperationResult::Response { response }) => Ok(response),
        Some(OperationResult::Error { error }) => {
            Err(AlienError::new(ErrorData::SandboxCommandFailed {
                failure: "operationFailed".to_string(),
                reason: format!(
                    "{operation}: operation '{name}' failed (grpc {}): {}",
                    error.code, error.message
                ),
            }))
        }
        None => Err(AlienError::new(ErrorData::UnexpectedResponseFormat {
            provider: "gcp-agent-platform".to_string(),
            binding_name: operation.to_string(),
            field: "response".to_string(),
            response_json: format!("operation '{name}' reported done without a result"),
        })),
    }
}

/// Whether a client error means the sandbox is already gone.
///
/// The client wraps a 404 as `RequestFailed` and leaves the `RemoteResourceNotFound` on the
/// source chain, so the classification is read by walking that chain rather than off the outer
/// variant — a path or trace id mentioning 404 in a message never reaches this.
fn is_not_found(error: &AlienError<AgentPlatformErrorData>) -> bool {
    const NOT_FOUND: &str = "REMOTE_RESOURCE_NOT_FOUND";
    if error.code == NOT_FOUND {
        return true;
    }
    let mut node = error.source.as_deref();
    while let Some(current) = node {
        if current.code == NOT_FOUND {
            return true;
        }
        node = current.source.as_deref();
    }
    false
}

/// A body short enough to sit in an error message without carrying a whole response into it.
fn truncated(body: &[u8]) -> String {
    const LIMIT: usize = 200;
    let text = String::from_utf8_lossy(body);
    let text = text.trim();
    if text.len() <= LIMIT {
        return text.to_string();
    }
    let end = (0..=LIMIT)
        .rev()
        .find(|at| text.is_char_boundary(*at))
        .unwrap_or(0);
    format!("{}…", &text[..end])
}

#[cfg(test)]
#[path = "gcp_agent_platform_tests.rs"]
mod tests;