Coverage Report

Created: 2026-10-01 05:28

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/build/source/nativelink-util/src/action_messages.rs
Line
Count
Source
1
// Copyright 2024 The NativeLink Authors. All rights reserved.
2
//
3
// Licensed under the Functional Source License, Version 1.1, Apache 2.0 Future License (the "License");
4
// you may not use this file except in compliance with the License.
5
// You may obtain a copy of the License at
6
//
7
//    See LICENSE file for details
8
//
9
// Unless required by applicable law or agreed to in writing, software
10
// distributed under the License is distributed on an "AS IS" BASIS,
11
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12
// See the License for the specific language governing permissions and
13
// limitations under the License.
14
15
use core::cmp::Ordering;
16
use core::convert::Into;
17
use core::fmt::Display;
18
use core::hash::Hash;
19
use core::time::Duration;
20
use std::collections::HashMap;
21
use std::time::SystemTime;
22
23
use humantime::format_duration;
24
use nativelink_error::{Error, ErrorContext, ResultExt, error_if, make_input_err};
25
use nativelink_metric::{
26
    MetricFieldData, MetricKind, MetricPublishKnownKindData, MetricsComponent, publish,
27
};
28
use nativelink_proto::build::bazel::remote::execution::v2::{
29
    Action, ActionResult as ProtoActionResult, ExecuteOperationMetadata, ExecuteRequest,
30
    ExecuteResponse, ExecutedActionMetadata, FileNode, LogFile, OutputDirectory, OutputFile,
31
    OutputSymlink, SymlinkNode, execution_stage,
32
};
33
use nativelink_proto::google::longrunning::Operation;
34
use nativelink_proto::google::longrunning::operation::Result as LongRunningResult;
35
use nativelink_proto::google::rpc::{PreconditionFailure, Status, precondition_failure};
36
use prost::Message;
37
use prost::bytes::Bytes;
38
use prost_types::Any;
39
use serde::ser::Error as SerdeError;
40
use serde::{Deserialize, Serialize};
41
use tonic::Code;
42
use uuid::Uuid;
43
44
use crate::common::{self, DigestInfo, HashMapExt, VecExt};
45
use crate::digest_hasher::DigestHasherFunc;
46
47
/// Default priority remote execution jobs will get when not provided.
48
pub const DEFAULT_EXECUTION_PRIORITY: i32 = 0;
49
50
/// Exit code sent if there is an internal error.
51
pub const INTERNAL_ERROR_EXIT_CODE: i32 = -178;
52
53
/// Holds an id that is unique to the client for a requested operation.
54
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
55
pub enum OperationId {
56
    Uuid(Uuid),
57
    String(String),
58
}
59
60
impl OperationId {
61
5
    pub fn into_string(self) -> String {
62
5
        match self {
63
1
            Self::Uuid(uuid) => uuid.to_string(),
64
4
            Self::String(name) => name,
65
        }
66
5
    }
67
}
68
69
impl Default for OperationId {
70
521
    fn default() -> Self {
71
521
        Self::Uuid(Uuid::new_v4())
72
521
    }
73
}
74
75
impl Display for OperationId {
76
1.88k
    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
77
1.88k
        match self {
78
1.28k
            Self::Uuid(uuid) => uuid.fmt(f),
79
598
            Self::String(name) => f.write_str(name),
80
        }
81
1.88k
    }
82
}
83
84
impl MetricsComponent for OperationId {
85
0
    fn publish(
86
0
        &self,
87
0
        _kind: MetricKind,
88
0
        _field_metadata: MetricFieldData,
89
0
    ) -> Result<MetricPublishKnownKindData, nativelink_metric::Error> {
90
0
        Ok(MetricPublishKnownKindData::String(self.to_string()))
91
0
    }
92
}
93
94
impl From<&str> for OperationId {
95
362
    fn from(value: &str) -> Self {
96
362
        Uuid::parse_str(value).map_or_else(|_| Self::String(
value280
.
to_string280
()), Self::Uuid)
97
362
    }
98
}
99
100
impl From<String> for OperationId {
101
41
    fn from(value: String) -> Self {
102
41
        Uuid::parse_str(&value).map_or(Self::String(value), Self::Uuid)
103
41
    }
104
}
105
106
impl TryFrom<Bytes> for OperationId {
107
    type Error = Error;
108
109
0
    fn try_from(value: Bytes) -> Result<Self, Self::Error> {
110
        // This is an optimized path to attempt to do the conversion in-place
111
        // to avoid an extra allocation/copy.
112
0
        match value.try_into_mut() {
113
            // We are the only reference to the Bytes, so we can convert it into a Vec<u8>
114
            // for free then convert the Vec<u8> to a String for free too.
115
0
            Ok(value) => {
116
0
                let value = String::from_utf8(value.into()).map_err(|e| {
117
0
                    Error::from_std_err(Code::InvalidArgument, &e).append(
118
                        "Failed to convert bytes to string in try_from<Bytes> for OperationId",
119
                    )
120
0
                })?;
121
0
                Ok(Self::from(value))
122
            }
123
            // We could not take ownership of the Bytes, so we may need to copy our data.
124
0
            Err(value) => {
125
0
                let value = core::str::from_utf8(&value).map_err(|e| {
126
0
                    Error::from_std_err(Code::InvalidArgument, &e).append(
127
                        "Failed to convert bytes to string in try_from<Bytes> for OperationId",
128
                    )
129
0
                })?;
130
0
                Ok(Self::from(value))
131
            }
132
        }
133
0
    }
134
}
135
136
/// Unique id of worker.
137
#[derive(Default, Eq, PartialEq, Hash, Clone, Serialize, Deserialize)]
138
pub struct WorkerId(pub String);
139
140
impl MetricsComponent for WorkerId {
141
0
    fn publish(
142
0
        &self,
143
0
        _kind: MetricKind,
144
0
        _field_metadata: MetricFieldData,
145
0
    ) -> Result<MetricPublishKnownKindData, nativelink_metric::Error> {
146
0
        Ok(MetricPublishKnownKindData::String(self.0.clone()))
147
0
    }
148
}
149
150
impl Display for WorkerId {
151
831
    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
152
831
        f.write_fmt(format_args!("{}", self.0))
153
831
    }
154
}
155
156
impl core::fmt::Debug for WorkerId {
157
524
    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
158
524
        Display::fmt(&self, f)
159
524
    }
160
}
161
162
impl From<WorkerId> for String {
163
314
    fn from(val: WorkerId) -> Self {
164
314
        val.0
165
314
    }
166
}
167
168
impl From<String> for WorkerId {
169
40
    fn from(s: String) -> Self {
170
40
        Self(s)
171
40
    }
172
}
173
174
/// Holds the information needed to uniquely identify an action
175
/// and if it is cacheable or not.
176
#[derive(Debug, Clone, Hash, PartialEq, Eq, Serialize, Deserialize)]
177
pub enum ActionUniqueQualifier {
178
    /// The action is cacheable.
179
    #[serde(alias = "Cachable")] // Pre 0.7.0 spelling
180
    Cacheable(ActionUniqueKey),
181
    /// The action is uncacheable.
182
    #[serde(alias = "Uncachable")] // Pre 0.7.0 spelling
183
    Uncacheable(ActionUniqueKey),
184
}
185
186
impl MetricsComponent for ActionUniqueQualifier {
187
0
    fn publish(
188
0
        &self,
189
0
        _kind: MetricKind,
190
0
        field_metadata: MetricFieldData,
191
0
    ) -> Result<MetricPublishKnownKindData, nativelink_metric::Error> {
192
0
        let (cacheable, action) = match self {
193
0
            Self::Cacheable(action) => (true, action),
194
0
            Self::Uncacheable(action) => (false, action),
195
        };
196
0
        publish!(
197
0
            cacheable,
198
0
            &cacheable,
199
0
            MetricKind::Default,
200
0
            "If the action is cacheable.",
201
0
            ""
202
        );
203
0
        action.publish(MetricKind::Component, field_metadata)?;
204
0
        Ok(MetricPublishKnownKindData::Component)
205
0
    }
206
}
207
208
impl ActionUniqueQualifier {
209
    /// Get the `instance_name` of the action.
210
0
    pub const fn instance_name(&self) -> &String {
211
0
        match self {
212
0
            Self::Cacheable(action) | Self::Uncacheable(action) => &action.instance_name,
213
        }
214
0
    }
215
216
    /// Get the digest function of the action.
217
0
    pub const fn digest_function(&self) -> DigestHasherFunc {
218
0
        match self {
219
0
            Self::Cacheable(action) | Self::Uncacheable(action) => action.digest_function,
220
        }
221
0
    }
222
223
    /// Get the digest of the action.
224
846
    pub const fn digest(&self) -> DigestInfo {
225
846
        match self {
226
846
            Self::Cacheable(
action715
) | Self::Uncacheable(
action131
) => action.digest,
227
        }
228
846
    }
229
}
230
231
impl Display for ActionUniqueQualifier {
232
69
    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
233
69
        let (cacheable, unique_key) = match self {
234
69
            Self::Cacheable(action) => (true, action),
235
0
            Self::Uncacheable(action) => (false, action),
236
        };
237
69
        f.write_fmt(format_args!(
238
            // Note: We use underscores because it makes escaping easier
239
            // for redis.
240
            "{}_{}_{}_{}_{}",
241
            unique_key.instance_name,
242
            unique_key.digest_function,
243
69
            unique_key.digest.packed_hash(),
244
69
            unique_key.digest.size_bytes(),
245
69
            if cacheable { 'c' } else { 
'u'0
},
246
0
        ))?;
247
69
        if let Some(
scope15
) = &unique_key.execution_scope {
248
15
            write!(f, "_scope_{scope}")
?0
;
249
54
        }
250
69
        Ok(())
251
69
    }
252
}
253
254
/// This is a utility struct used to make it easier to match `ActionInfos` in a
255
/// `HashMap` without needing to construct an entire `ActionInfo`.
256
#[derive(Debug, Clone, Eq, PartialEq, Hash, Serialize, Deserialize, MetricsComponent)]
257
pub struct ActionUniqueKey {
258
    /// Optional in-flight execution identity. This never changes the action
259
    /// digest, instance, or completed action-cache lookup key.
260
    #[serde(default, skip_serializing_if = "Option::is_none")]
261
    pub execution_scope: Option<String>,
262
    /// Name of instance group this action belongs to.
263
    #[metric(help = "Name of instance group this action belongs to.")]
264
    pub instance_name: String,
265
    /// The digest function this action expects.
266
    #[metric(help = "The digest function this action expects.")]
267
    pub digest_function: DigestHasherFunc,
268
    /// Digest of the underlying `Action`.
269
    #[metric(help = "Digest of the underlying Action.")]
270
    pub digest: DigestInfo,
271
}
272
273
impl Display for ActionUniqueKey {
274
0
    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
275
0
        f.write_fmt(format_args!(
276
            "{}/{}/{}",
277
            self.instance_name, self.digest_function, self.digest,
278
        ))
279
0
    }
280
}
281
282
/// Information needed to execute an action. This struct is used over bazel's proto `Action`
283
/// for simplicity and offers a `salt`, which is useful to ensure during hashing (for dicts)
284
/// to ensure we never match against another `ActionInfo` (when a task should never be cached).
285
/// This struct must be 100% compatible with `ExecuteRequest` struct in `remote_execution.proto`
286
/// except for the salt field.
287
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize, MetricsComponent)]
288
pub struct ActionInfo {
289
    /// Digest of the underlying `Command`.
290
    #[metric(help = "Digest of the underlying Command.")]
291
    pub command_digest: DigestInfo,
292
    /// Digest of the underlying `Directory`.
293
    #[metric(help = "Digest of the underlying Directory.")]
294
    pub input_root_digest: DigestInfo,
295
    /// Timeout of the action.
296
    #[metric(help = "Timeout of the action.")]
297
    pub timeout: Duration,
298
    /// The properties rules that must be applied when finding a worker that can run this action.
299
    #[metric(group = "platform_properties")]
300
    pub platform_properties: HashMap<String, String>,
301
    /// The priority of the action. Higher value means it should execute faster.
302
    #[metric(help = "The priority of the action. Higher value means it should execute faster.")]
303
    pub priority: i32,
304
    /// When this action started to be loaded from the CAS.
305
    #[metric(help = "When this action started to be loaded from the CAS.")]
306
    pub load_timestamp: SystemTime,
307
    /// When this action was created.
308
    #[metric(help = "When this action was created.")]
309
    pub insert_timestamp: SystemTime,
310
    /// Info used to uniquely identify this `ActionInfo` and if it is cacheable.
311
    /// This is primarily used to join actions/operations together using this key.
312
    #[metric(help = "Info used to uniquely identify this ActionInfo and if it is cacheable.")]
313
    pub unique_qualifier: ActionUniqueQualifier,
314
}
315
316
impl ActionInfo {
317
    #[inline]
318
0
    pub const fn instance_name(&self) -> &String {
319
0
        self.unique_qualifier.instance_name()
320
0
    }
321
322
    /// Returns the underlying digest of the `Action`.
323
    #[inline]
324
589
    pub const fn digest(&self) -> DigestInfo {
325
589
        self.unique_qualifier.digest()
326
589
    }
327
328
256
    pub fn try_from_action_and_execute_request(
329
256
        execute_request: ExecuteRequest,
330
256
        action: Action,
331
256
        load_timestamp: SystemTime,
332
256
        queued_timestamp: SystemTime,
333
256
    ) -> Result<Self, Error> {
334
256
        let unique_key = ActionUniqueKey {
335
256
            execution_scope: None,
336
256
            instance_name: execute_request.instance_name,
337
256
            digest_function: DigestHasherFunc::try_from(execute_request.digest_function)
338
256
                .err_tip(|| 
format!0
("Could not find digest_function in try_from_action_and_execute_request {:?}", execute_request.digest_function))
?0
,
339
256
            digest: execute_request
340
256
                .action_digest
341
256
                .err_tip(|| "Expected action_digest to exist on ExecuteRequest")
?0
342
256
                .try_into()
?0
,
343
        };
344
256
        let unique_qualifier = if execute_request.skip_cache_lookup {
345
0
            ActionUniqueQualifier::Uncacheable(unique_key)
346
        } else {
347
256
            ActionUniqueQualifier::Cacheable(unique_key)
348
        };
349
350
256
        let proto_properties = action.platform.unwrap_or_default();
351
256
        let mut platform_properties = HashMap::with_capacity(proto_properties.properties.len());
352
256
        for 
property5
in proto_properties.properties {
353
5
            platform_properties.insert(property.name, property.value);
354
5
        }
355
356
        Ok(Self {
357
256
            command_digest: action
358
256
                .command_digest
359
256
                .err_tip(|| "Expected command_digest to exist on Action")
?0
360
256
                .try_into()
?0
,
361
256
            input_root_digest: action
362
256
                .input_root_digest
363
256
                .err_tip(|| "Expected input_root_digest to exist on Action")
?0
364
256
                .try_into()
?0
,
365
256
            timeout: action
366
256
                .timeout
367
256
                .unwrap_or_default()
368
256
                .try_into()
369
256
                .map_err(|err| 
{0
370
0
                    Error::from_std_err(Code::InvalidArgument, &err)
371
0
                        .append("Failed convert proto duration to system duration")
372
0
                })?,
373
256
            platform_properties,
374
256
            priority: execute_request
375
256
                .execution_policy
376
256
                .unwrap_or_default()
377
256
                .priority,
378
256
            load_timestamp,
379
256
            insert_timestamp: queued_timestamp,
380
256
            unique_qualifier,
381
        })
382
256
    }
383
}
384
385
impl From<&ActionInfo> for ExecuteRequest {
386
157
    fn from(val: &ActionInfo) -> Self {
387
157
        let digest = val.digest().into();
388
157
        let (skip_cache_lookup, unique_qualifier) = match &val.unique_qualifier {
389
115
            ActionUniqueQualifier::Cacheable(unique_qualifier) => (false, unique_qualifier),
390
42
            ActionUniqueQualifier::Uncacheable(unique_qualifier) => (true, unique_qualifier),
391
        };
392
157
        Self {
393
157
            instance_name: unique_qualifier.instance_name.clone(),
394
157
            action_digest: Some(digest),
395
157
            skip_cache_lookup,
396
157
            execution_policy: None,     // Not used in the worker.
397
157
            results_cache_policy: None, // Not used in the worker.
398
157
            digest_function: unique_qualifier.digest_function.proto_digest_func().into(),
399
157
        }
400
157
    }
401
}
402
403
/// Simple utility struct to determine if a string is representing a full path or
404
/// just the name of the file.
405
/// This is in order to be able to reuse the same struct instead of building different
406
/// structs when converting `FileInfo` -> {`OutputFile`, `FileNode`} and other similar
407
/// structs.
408
#[derive(Eq, PartialEq, Debug, Clone, Serialize, Deserialize)]
409
pub enum NameOrPath {
410
    Name(String),
411
    Path(String),
412
}
413
414
impl PartialOrd for NameOrPath {
415
0
    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
416
0
        Some(self.cmp(other))
417
0
    }
418
}
419
420
impl Ord for NameOrPath {
421
0
    fn cmp(&self, other: &Self) -> Ordering {
422
0
        let self_lexical_name = match self {
423
0
            Self::Name(name) => name,
424
0
            Self::Path(path) => path,
425
        };
426
0
        let other_lexical_name = match other {
427
0
            Self::Name(name) => name,
428
0
            Self::Path(path) => path,
429
        };
430
0
        self_lexical_name.cmp(other_lexical_name)
431
0
    }
432
}
433
434
/// Represents an individual file and associated metadata.
435
/// This struct must be 100% compatible with `OutputFile` and `FileNode` structs
436
/// in `remote_execution.proto`.
437
#[derive(Eq, PartialEq, Debug, Clone, Serialize, Deserialize)]
438
pub struct FileInfo {
439
    pub name_or_path: NameOrPath,
440
    pub digest: DigestInfo,
441
    pub is_executable: bool,
442
}
443
444
impl TryFrom<FileInfo> for FileNode {
445
    type Error = Error;
446
447
263
    fn try_from(val: FileInfo) -> Result<Self, Error> {
448
263
        match val.name_or_path {
449
0
            NameOrPath::Path(_) => Err(make_input_err!(
450
0
                "Cannot return a FileInfo that uses a NameOrPath::Path(), it must be a NameOrPath::Name()"
451
0
            )),
452
263
            NameOrPath::Name(name) => Ok(Self {
453
263
                name,
454
263
                digest: Some((&val.digest).into()),
455
263
                is_executable: val.is_executable,
456
263
                node_properties: None, // Not supported.
457
263
            }),
458
        }
459
263
    }
460
}
461
462
impl TryFrom<OutputFile> for FileInfo {
463
    type Error = Error;
464
465
2
    fn try_from(output_file: OutputFile) -> Result<Self, Error> {
466
        Ok(Self {
467
2
            name_or_path: NameOrPath::Path(output_file.path),
468
2
            digest: output_file
469
2
                .digest
470
2
                .err_tip(|| "Expected digest to exist on OutputFile")
?0
471
2
                .try_into()
?0
,
472
2
            is_executable: output_file.is_executable,
473
        })
474
2
    }
475
}
476
477
impl TryFrom<FileInfo> for OutputFile {
478
    type Error = Error;
479
480
7
    fn try_from(val: FileInfo) -> Result<Self, Error> {
481
7
        match val.name_or_path {
482
0
            NameOrPath::Name(_) => Err(make_input_err!(
483
0
                "Cannot return a FileInfo that uses a NameOrPath::Name(), it must be a NameOrPath::Path()"
484
0
            )),
485
7
            NameOrPath::Path(path) => Ok(Self {
486
7
                path,
487
7
                digest: Some((&val.digest).into()),
488
7
                is_executable: val.is_executable,
489
7
                contents: Bytes::default(),
490
7
                node_properties: None, // Not supported.
491
7
            }),
492
        }
493
7
    }
494
}
495
496
/// Represents an individual symlink file and associated metadata.
497
/// This struct must be 100% compatible with `SymlinkNode` and `OutputSymlink`.
498
#[derive(Eq, PartialEq, Debug, Clone, Serialize, Deserialize)]
499
pub struct SymlinkInfo {
500
    pub name_or_path: NameOrPath,
501
    pub target: String,
502
}
503
504
impl TryFrom<SymlinkNode> for SymlinkInfo {
505
    type Error = Error;
506
507
0
    fn try_from(symlink_node: SymlinkNode) -> Result<Self, Error> {
508
0
        Ok(Self {
509
0
            name_or_path: NameOrPath::Name(symlink_node.name),
510
0
            target: symlink_node.target,
511
0
        })
512
0
    }
513
}
514
515
impl TryFrom<SymlinkInfo> for SymlinkNode {
516
    type Error = Error;
517
518
1
    fn try_from(val: SymlinkInfo) -> Result<Self, Error> {
519
1
        match val.name_or_path {
520
0
            NameOrPath::Path(_) => Err(make_input_err!(
521
0
                "Cannot return a SymlinkInfo that uses a NameOrPath::Path(), it must be a NameOrPath::Name()"
522
0
            )),
523
1
            NameOrPath::Name(name) => Ok(Self {
524
1
                name,
525
1
                target: val.target,
526
1
                node_properties: None, // Not supported.
527
1
            }),
528
        }
529
1
    }
530
}
531
532
impl TryFrom<OutputSymlink> for SymlinkInfo {
533
    type Error = Error;
534
535
2
    fn try_from(output_symlink: OutputSymlink) -> Result<Self, Error> {
536
2
        Ok(Self {
537
2
            name_or_path: NameOrPath::Path(output_symlink.path),
538
2
            target: output_symlink.target,
539
2
        })
540
2
    }
541
}
542
543
impl TryFrom<SymlinkInfo> for OutputSymlink {
544
    type Error = Error;
545
546
2
    fn try_from(val: SymlinkInfo) -> Result<Self, Error> {
547
2
        match val.name_or_path {
548
2
            NameOrPath::Path(path) => {
549
2
                Ok(Self {
550
2
                    path,
551
2
                    target: val.target,
552
2
                    node_properties: None, // Not supported.
553
2
                })
554
            }
555
0
            NameOrPath::Name(_) => Err(make_input_err!(
556
0
                "Cannot return a SymlinkInfo that uses a NameOrPath::Name(), it must be a NameOrPath::Path()"
557
0
            )),
558
        }
559
2
    }
560
}
561
562
/// Represents an individual directory file and associated metadata.
563
/// This struct must be 100% compatible with `SymlinkNode` and `OutputSymlink`.
564
#[derive(Eq, PartialEq, Debug, Clone, Serialize, Deserialize)]
565
pub struct DirectoryInfo {
566
    pub path: String,
567
    pub tree_digest: DigestInfo,
568
}
569
570
impl TryFrom<OutputDirectory> for DirectoryInfo {
571
    type Error = Error;
572
573
2
    fn try_from(output_directory: OutputDirectory) -> Result<Self, Error> {
574
        Ok(Self {
575
2
            path: output_directory.path,
576
2
            tree_digest: output_directory
577
2
                .tree_digest
578
2
                .err_tip(|| "Expected tree_digest to exist in OutputDirectory")
?0
579
2
                .try_into()
?0
,
580
        })
581
2
    }
582
}
583
584
impl From<DirectoryInfo> for OutputDirectory {
585
1
    fn from(val: DirectoryInfo) -> Self {
586
1
        Self {
587
1
            path: val.path,
588
1
            tree_digest: Some(val.tree_digest.into()),
589
1
            is_topologically_sorted: false,
590
1
        }
591
1
    }
592
}
593
594
/// Represents the metadata associated with the execution result.
595
/// This struct must be 100% compatible with `ExecutedActionMetadata`.
596
#[derive(Eq, PartialEq, Debug, Clone, Serialize, Deserialize)]
597
pub struct ExecutionMetadata {
598
    pub worker: String,
599
    pub queued_timestamp: SystemTime,
600
    pub worker_start_timestamp: SystemTime,
601
    pub worker_completed_timestamp: SystemTime,
602
    pub input_fetch_start_timestamp: SystemTime,
603
    pub input_fetch_completed_timestamp: SystemTime,
604
    pub execution_start_timestamp: SystemTime,
605
    pub execution_completed_timestamp: SystemTime,
606
    pub output_upload_start_timestamp: SystemTime,
607
    pub output_upload_completed_timestamp: SystemTime,
608
}
609
610
impl Default for ExecutionMetadata {
611
5
    fn default() -> Self {
612
5
        Self {
613
5
            worker: String::new(),
614
5
            queued_timestamp: SystemTime::UNIX_EPOCH,
615
5
            worker_start_timestamp: SystemTime::UNIX_EPOCH,
616
5
            worker_completed_timestamp: SystemTime::UNIX_EPOCH,
617
5
            input_fetch_start_timestamp: SystemTime::UNIX_EPOCH,
618
5
            input_fetch_completed_timestamp: SystemTime::UNIX_EPOCH,
619
5
            execution_start_timestamp: SystemTime::UNIX_EPOCH,
620
5
            execution_completed_timestamp: SystemTime::UNIX_EPOCH,
621
5
            output_upload_start_timestamp: SystemTime::UNIX_EPOCH,
622
5
            output_upload_completed_timestamp: SystemTime::UNIX_EPOCH,
623
5
        }
624
5
    }
625
}
626
627
impl From<ExecutionMetadata> for ExecutedActionMetadata {
628
52
    fn from(val: ExecutionMetadata) -> Self {
629
        Self {
630
52
            worker: val.worker,
631
52
            queued_timestamp: Some(val.queued_timestamp.into()),
632
52
            worker_start_timestamp: Some(val.worker_start_timestamp.into()),
633
52
            worker_completed_timestamp: Some(val.worker_completed_timestamp.into()),
634
52
            input_fetch_start_timestamp: Some(val.input_fetch_start_timestamp.into()),
635
52
            input_fetch_completed_timestamp: Some(val.input_fetch_completed_timestamp.into()),
636
52
            execution_start_timestamp: Some(val.execution_start_timestamp.into()),
637
52
            execution_completed_timestamp: Some(val.execution_completed_timestamp.into()),
638
52
            output_upload_start_timestamp: Some(val.output_upload_start_timestamp.into()),
639
52
            output_upload_completed_timestamp: Some(val.output_upload_completed_timestamp.into()),
640
52
            virtual_execution_duration: val
641
52
                .execution_completed_timestamp
642
52
                .duration_since(val.execution_start_timestamp)
643
52
                .ok()
644
52
                .and_then(|duration| prost_types::Duration::try_from(duration).ok()),
645
52
            auxiliary_metadata: Vec::default(),
646
        }
647
52
    }
648
}
649
650
impl TryFrom<ExecutedActionMetadata> for ExecutionMetadata {
651
    type Error = Error;
652
653
3
    fn try_from(eam: ExecutedActionMetadata) -> Result<Self, Error> {
654
        Ok(Self {
655
3
            worker: eam.worker,
656
3
            queued_timestamp: eam
657
3
                .queued_timestamp
658
3
                .err_tip(|| "Expected queued_timestamp to exist in ExecutedActionMetadata")
?0
659
3
                .try_into()
?0
,
660
3
            worker_start_timestamp: eam
661
3
                .worker_start_timestamp
662
3
                .err_tip(|| "Expected worker_start_timestamp to exist in ExecutedActionMetadata")
?0
663
3
                .try_into()
?0
,
664
3
            worker_completed_timestamp: eam
665
3
                .worker_completed_timestamp
666
3
                .err_tip(|| 
{0
667
0
                    "Expected worker_completed_timestamp to exist in ExecutedActionMetadata"
668
0
                })?
669
3
                .try_into()
?0
,
670
3
            input_fetch_start_timestamp: eam
671
3
                .input_fetch_start_timestamp
672
3
                .err_tip(|| 
{0
673
0
                    "Expected input_fetch_start_timestamp to exist in ExecutedActionMetadata"
674
0
                })?
675
3
                .try_into()
?0
,
676
3
            input_fetch_completed_timestamp: eam
677
3
                .input_fetch_completed_timestamp
678
3
                .err_tip(|| 
{0
679
0
                    "Expected input_fetch_completed_timestamp to exist in ExecutedActionMetadata"
680
0
                })?
681
3
                .try_into()
?0
,
682
3
            execution_start_timestamp: eam
683
3
                .execution_start_timestamp
684
3
                .err_tip(|| 
{0
685
0
                    "Expected execution_start_timestamp to exist in ExecutedActionMetadata"
686
0
                })?
687
3
                .try_into()
?0
,
688
3
            execution_completed_timestamp: eam
689
3
                .execution_completed_timestamp
690
3
                .err_tip(|| 
{0
691
0
                    "Expected execution_completed_timestamp to exist in ExecutedActionMetadata"
692
0
                })?
693
3
                .try_into()
?0
,
694
3
            output_upload_start_timestamp: eam
695
3
                .output_upload_start_timestamp
696
3
                .err_tip(|| 
{0
697
0
                    "Expected output_upload_start_timestamp to exist in ExecutedActionMetadata"
698
0
                })?
699
3
                .try_into()
?0
,
700
3
            output_upload_completed_timestamp: eam
701
3
                .output_upload_completed_timestamp
702
3
                .err_tip(|| 
{0
703
0
                    "Expected output_upload_completed_timestamp to exist in ExecutedActionMetadata"
704
0
                })?
705
3
                .try_into()
?0
,
706
        })
707
3
    }
708
}
709
710
/// Represents the results of an execution.
711
/// This struct must be 100% compatible with `ActionResult` in `remote_execution.proto`.
712
#[derive(Eq, PartialEq, Debug, Clone, Serialize, Deserialize)]
713
pub struct ActionResult {
714
    pub output_files: Vec<FileInfo>,
715
    pub output_folders: Vec<DirectoryInfo>,
716
    pub output_directory_symlinks: Vec<SymlinkInfo>,
717
    pub output_file_symlinks: Vec<SymlinkInfo>,
718
    pub exit_code: i32,
719
    pub stdout_digest: DigestInfo,
720
    pub stderr_digest: DigestInfo,
721
    pub execution_metadata: ExecutionMetadata,
722
    pub server_logs: HashMap<String, DigestInfo>,
723
    pub error: Option<Error>,
724
    pub message: String,
725
}
726
727
impl Default for ActionResult {
728
82
    fn default() -> Self {
729
82
        Self {
730
82
            output_files: Vec::default(),
731
82
            output_folders: Vec::default(),
732
82
            output_directory_symlinks: Vec::default(),
733
82
            output_file_symlinks: Vec::default(),
734
82
            exit_code: INTERNAL_ERROR_EXIT_CODE,
735
82
            stdout_digest: DigestInfo::new([0u8; 32], 0),
736
82
            stderr_digest: DigestInfo::new([0u8; 32], 0),
737
82
            execution_metadata: ExecutionMetadata {
738
82
                worker: String::new(),
739
82
                queued_timestamp: SystemTime::UNIX_EPOCH,
740
82
                worker_start_timestamp: SystemTime::UNIX_EPOCH,
741
82
                worker_completed_timestamp: SystemTime::UNIX_EPOCH,
742
82
                input_fetch_start_timestamp: SystemTime::UNIX_EPOCH,
743
82
                input_fetch_completed_timestamp: SystemTime::UNIX_EPOCH,
744
82
                execution_start_timestamp: SystemTime::UNIX_EPOCH,
745
82
                execution_completed_timestamp: SystemTime::UNIX_EPOCH,
746
82
                output_upload_start_timestamp: SystemTime::UNIX_EPOCH,
747
82
                output_upload_completed_timestamp: SystemTime::UNIX_EPOCH,
748
82
            },
749
82
            server_logs: HashMap::default(),
750
82
            error: None,
751
82
            message: String::new(),
752
82
        }
753
82
    }
754
}
755
756
/// The execution status/stage. This should match `ExecutionStage::Value` in `remote_execution.proto`.
757
#[derive(PartialEq, Debug, Clone, Serialize, Deserialize)]
758
#[allow(
759
    clippy::large_enum_variant,
760
    reason = "TODO box the two relevant variants in a breaking release. Unfulfilled on nightly"
761
)]
762
pub enum ActionStage {
763
    /// Stage is unknown.
764
    Unknown,
765
    /// Checking the cache to see if action exists.
766
    CacheCheck,
767
    /// Action has been accepted and waiting for worker to take it.
768
    Queued,
769
    // TODO(palfrey) We need a way to know if the job was sent to a worker, but hasn't begun
770
    // execution yet.
771
    /// Worker is executing the action.
772
    Executing,
773
    /// Worker completed the work with result.
774
    Completed(ActionResult),
775
    /// Result was found from cache, don't decode the proto just to re-encode it.
776
    #[serde(serialize_with = "serialize_proto_result", skip_deserializing)]
777
    // The serialization step decodes this to an ActionResult which is serializable.
778
    // Since it will always be serialized as an ActionResult, we do not need to support
779
    // deserialization on this type at all.
780
    // In theory, serializing this should never happen so performance shouldn't be affected.
781
    CompletedFromCache(ProtoActionResult),
782
}
783
784
0
fn serialize_proto_result<S>(v: &ProtoActionResult, serializer: S) -> Result<S::Ok, S::Error>
785
0
where
786
0
    S: serde::Serializer,
787
{
788
0
    let s = ActionResult::try_from(v.clone()).map_err(S::Error::custom)?;
789
0
    s.serialize(serializer)
790
0
}
791
792
impl ActionStage {
793
465
    pub const fn has_action_result(&self) -> bool {
794
465
        match self {
795
379
            Self::Unknown | Self::CacheCheck | Self::Queued | Self::Executing => false,
796
86
            Self::Completed(_) | Self::CompletedFromCache(_) => true,
797
        }
798
465
    }
799
800
    /// Returns true if the worker considers the action done and no longer needs to be tracked.
801
    // Note: This function is separate from `has_action_result()` to not mix the concept of
802
    //       "finished" with "has a result".
803
0
    pub const fn is_finished(&self) -> bool {
804
0
        self.has_action_result()
805
0
    }
806
807
    /// Returns if the stage enum is the same as the other stage enum, but
808
    /// does not compare the values of the enum.
809
0
    pub const fn is_same_stage(&self, other: &Self) -> bool {
810
0
        matches!(
811
0
            (self, other),
812
            (Self::Unknown, Self::Unknown)
813
                | (Self::CacheCheck, Self::CacheCheck)
814
                | (Self::Queued, Self::Queued)
815
                | (Self::Executing, Self::Executing)
816
                | (Self::Completed(_), Self::Completed(_))
817
                | (Self::CompletedFromCache(_), Self::CompletedFromCache(_))
818
        )
819
0
    }
820
821
0
    pub fn name(&self) -> String {
822
0
        match self {
823
0
            Self::Unknown => "Unknown".to_string(),
824
0
            Self::CacheCheck => "CacheCheck".to_string(),
825
0
            Self::Queued => "Queued".to_string(),
826
0
            Self::Executing => "Executing".to_string(),
827
0
            Self::Completed(_) => "Completed".to_string(),
828
0
            Self::CompletedFromCache(_) => "CompletedFromCache".to_string(),
829
        }
830
0
    }
831
}
832
833
impl MetricsComponent for ActionStage {
834
0
    fn publish(
835
0
        &self,
836
0
        _kind: MetricKind,
837
0
        _field_metadata: MetricFieldData,
838
0
    ) -> Result<MetricPublishKnownKindData, nativelink_metric::Error> {
839
0
        Ok(MetricPublishKnownKindData::String(self.name()))
840
0
    }
841
}
842
843
impl From<&ActionStage> for execution_stage::Value {
844
5
    fn from(val: &ActionStage) -> Self {
845
5
        match val {
846
0
            ActionStage::Unknown => Self::Unknown,
847
0
            ActionStage::CacheCheck => Self::CacheCheck,
848
3
            ActionStage::Queued => Self::Queued,
849
0
            ActionStage::Executing => Self::Executing,
850
2
            ActionStage::Completed(_) | ActionStage::CompletedFromCache(_) => Self::Completed,
851
        }
852
5
    }
853
}
854
855
/// Build a `google.rpc.Status` of code `FAILED_PRECONDITION` whose
856
/// details carry a `PreconditionFailure` naming the missing blob.
857
///
858
/// This is the worker-side counterpart to `execution_server`'s
859
/// `missing_blobs_failed_precondition` — both produce the `REv2`
860
/// subject format `blobs/{hash}/{size}` that Bazel auto-retries on.
861
4
fn missing_blob_failed_precondition_status(err: &Error, hash: &str, size: i64) -> Status {
862
4
    let pf = PreconditionFailure {
863
4
        violations: vec![precondition_failure::Violation {
864
4
            r#type: common::VIOLATION_TYPE_MISSING.to_string(),
865
4
            // REv2-mandated subject format for missing-blob violations.
866
4
            subject: format!("blobs/{hash}/{size}"),
867
4
            description: err.message_string(),
868
4
        }],
869
4
    };
870
4
    let mut buf: Vec<u8> = Vec::with_capacity(pf.encoded_len());
871
4
    pf.encode(&mut buf)
872
4
        .expect("encoding prost message into Vec<u8> cannot fail");
873
4
    let any = Any {
874
4
        type_url: PreconditionFailure::TYPE_URL.to_string(),
875
4
        value: buf,
876
4
    };
877
4
    Status {
878
4
        code: Code::FailedPrecondition as i32,
879
4
        message: err.message_string(),
880
4
        details: vec![any],
881
4
    }
882
4
}
883
884
45
pub fn to_execute_response(action_result: ActionResult) -> ExecuteResponse {
885
45
    fn logs_from(server_logs: HashMap<String, DigestInfo>) -> HashMap<String, LogFile> {
886
45
        let mut logs = HashMap::with_capacity(server_logs.len());
887
45
        for (
k1
,
v1
) in server_logs {
888
1
            logs.insert(
889
1
                k.clone(),
890
1
                LogFile {
891
1
                    digest: Some(v.into()),
892
1
                    human_readable: false,
893
1
                },
894
1
            );
895
1
        }
896
45
        logs
897
45
    }
898
899
    // If the action failed because a CAS blob is missing — most often a
900
    // `Directory` proto in the input tree (the Execute pre-check only
901
    // validates the top-level Action, command_digest, and
902
    // input_root_digest; nested Directories are fetched lazily by the
903
    // worker) — surface the failure as `FAILED_PRECONDITION` with a
904
    // `PreconditionFailure` detail naming the digest. Bazel sees the
905
    // detail, re-uploads the missing blob, and retries automatically;
906
    // without the detail it gives up and the build fails.
907
    //
908
    // The dispatch is on `Error::context` (typed metadata attached at
909
    // the production site in `fast_slow_store`), not the message text.
910
    // String-matching across crate boundaries silently regresses when
911
    // the producing crate reformats its error — see commit history.
912
45
    let status = Some(
913
45
        action_result
914
45
            .error
915
45
            .clone()
916
45
            .map(|err| match 
&err.context32
{
917
4
                ErrorContext::MissingDigest { hash, size } => {
918
4
                    let (hash, size) = (hash.clone(), *size);
919
4
                    missing_blob_failed_precondition_status(&err, &hash, size)
920
                }
921
28
                ErrorContext::None => err.into(),
922
32
            })
923
45
            .unwrap_or_default(),
924
    );
925
45
    let message = action_result.message.clone();
926
45
    ExecuteResponse {
927
45
        server_logs: logs_from(action_result.server_logs.clone()),
928
45
        result: action_result.try_into().ok(),
929
45
        cached_result: false,
930
45
        status,
931
45
        message,
932
45
    }
933
45
}
934
935
impl From<ActionStage> for ExecuteResponse {
936
11
    fn from(val: ActionStage) -> Self {
937
11
        match val {
938
            // We don't have an execute response if we don't have the results. It is defined
939
            // behavior to return an empty proto struct.
940
            ActionStage::Unknown
941
            | ActionStage::CacheCheck
942
            | ActionStage::Queued
943
0
            | ActionStage::Executing => Self::default(),
944
11
            ActionStage::Completed(action_result) => to_execute_response(action_result),
945
            // Handled separately as there are no server logs and the action
946
            // result is already in Proto format.
947
0
            ActionStage::CompletedFromCache(proto_action_result) => Self {
948
0
                server_logs: HashMap::new(),
949
0
                result: Some(proto_action_result),
950
0
                cached_result: true,
951
0
                status: Some(Status::default()),
952
0
                message: String::new(), // Will be populated later if applicable.
953
0
            },
954
        }
955
11
    }
956
}
957
958
impl TryFrom<ActionResult> for ProtoActionResult {
959
    type Error = Error;
960
961
52
    fn try_from(val: ActionResult) -> Result<Self, Error> {
962
52
        let mut output_symlinks = Vec::with_capacity(
963
52
            val.output_file_symlinks.len() + val.output_directory_symlinks.len(),
964
        );
965
52
        output_symlinks.extend_from_slice(val.output_file_symlinks.as_slice());
966
52
        output_symlinks.extend_from_slice(val.output_directory_symlinks.as_slice());
967
968
        Ok(Self {
969
52
            output_files: val
970
52
                .output_files
971
52
                .into_iter()
972
52
                .map(TryInto::try_into)
973
52
                .collect::<Result<_, _>>()
?0
,
974
52
            output_file_symlinks: val
975
52
                .output_file_symlinks
976
52
                .into_iter()
977
52
                .map(TryInto::try_into)
978
52
                .collect::<Result<_, _>>()
?0
,
979
52
            output_symlinks: output_symlinks
980
52
                .into_iter()
981
52
                .map(TryInto::try_into)
982
52
                .collect::<Result<_, _>>()
?0
,
983
52
            output_directories: val
984
52
                .output_folders
985
52
                .into_iter()
986
52
                .map(TryInto::try_into)
987
52
                .collect::<Result<_, _>>()
?0
,
988
52
            output_directory_symlinks: val
989
52
                .output_directory_symlinks
990
52
                .into_iter()
991
52
                .map(TryInto::try_into)
992
52
                .collect::<Result<_, _>>()
?0
,
993
52
            exit_code: val.exit_code,
994
52
            stdout_raw: Bytes::default(),
995
52
            stdout_digest: Some(val.stdout_digest.into()),
996
52
            stderr_raw: Bytes::default(),
997
52
            stderr_digest: Some(val.stderr_digest.into()),
998
52
            execution_metadata: Some(val.execution_metadata.into()),
999
        })
1000
52
    }
1001
}
1002
1003
impl TryFrom<ProtoActionResult> for ActionResult {
1004
    type Error = Error;
1005
1006
0
    fn try_from(val: ProtoActionResult) -> Result<Self, Error> {
1007
0
        let output_file_symlinks = val
1008
0
            .output_file_symlinks
1009
0
            .into_iter()
1010
0
            .map(|output_symlink| {
1011
0
                SymlinkInfo::try_from(output_symlink)
1012
0
                    .err_tip(|| "Output File Symlinks could not be converted to SymlinkInfo")
1013
0
            })
1014
0
            .collect::<Result<Vec<_>, _>>()?;
1015
1016
0
        let output_directory_symlinks = val
1017
0
            .output_directory_symlinks
1018
0
            .into_iter()
1019
0
            .map(|output_symlink| {
1020
0
                SymlinkInfo::try_from(output_symlink)
1021
0
                    .err_tip(|| "Output File Symlinks could not be converted to SymlinkInfo")
1022
0
            })
1023
0
            .collect::<Result<Vec<_>, _>>()?;
1024
1025
0
        let output_files = val
1026
0
            .output_files
1027
0
            .into_iter()
1028
0
            .map(|output_file| {
1029
0
                output_file
1030
0
                    .try_into()
1031
0
                    .err_tip(|| "Output File could not be converted")
1032
0
            })
1033
0
            .collect::<Result<Vec<_>, _>>()?;
1034
1035
0
        let output_folders = val
1036
0
            .output_directories
1037
0
            .into_iter()
1038
0
            .map(|output_directory| {
1039
0
                output_directory
1040
0
                    .try_into()
1041
0
                    .err_tip(|| "Output File could not be converted")
1042
0
            })
1043
0
            .collect::<Result<Vec<_>, _>>()?;
1044
1045
        Ok(Self {
1046
0
            output_files,
1047
0
            output_folders,
1048
0
            output_file_symlinks,
1049
0
            output_directory_symlinks,
1050
0
            exit_code: val.exit_code,
1051
0
            stdout_digest: val
1052
0
                .stdout_digest
1053
0
                .err_tip(|| "Expected stdout_digest to be set on ExecuteResponse msg")?
1054
0
                .try_into()?,
1055
0
            stderr_digest: val
1056
0
                .stderr_digest
1057
0
                .err_tip(|| "Expected stderr_digest to be set on ExecuteResponse msg")?
1058
0
                .try_into()?,
1059
0
            execution_metadata: val
1060
0
                .execution_metadata
1061
0
                .err_tip(|| "Expected execution_metadata to be set on ExecuteResponse msg")?
1062
0
                .try_into()?,
1063
0
            server_logs: HashMap::default(),
1064
0
            error: None,
1065
0
            message: String::new(),
1066
        })
1067
0
    }
1068
}
1069
1070
impl TryFrom<ExecuteResponse> for ActionStage {
1071
    type Error = Error;
1072
1073
3
    fn try_from(execute_response: ExecuteResponse) -> Result<Self, Error> {
1074
3
        let proto_action_result = execute_response
1075
3
            .result
1076
3
            .err_tip(|| "Expected result to be set on ExecuteResponse msg")
?0
;
1077
3
        let action_result = ActionResult {
1078
3
            output_files: proto_action_result
1079
3
                .output_files
1080
3
                .try_map(TryInto::try_into)
?0
,
1081
3
            output_directory_symlinks: proto_action_result
1082
3
                .output_directory_symlinks
1083
3
                .try_map(TryInto::try_into)
?0
,
1084
3
            output_file_symlinks: proto_action_result
1085
3
                .output_file_symlinks
1086
3
                .try_map(TryInto::try_into)
?0
,
1087
3
            output_folders: proto_action_result
1088
3
                .output_directories
1089
3
                .try_map(TryInto::try_into)
?0
,
1090
3
            exit_code: proto_action_result.exit_code,
1091
1092
3
            stdout_digest: proto_action_result
1093
3
                .stdout_digest
1094
3
                .err_tip(|| "Expected stdout_digest to be set on ExecuteResponse msg")
?0
1095
3
                .try_into()
?0
,
1096
3
            stderr_digest: proto_action_result
1097
3
                .stderr_digest
1098
3
                .err_tip(|| "Expected stderr_digest to be set on ExecuteResponse msg")
?0
1099
3
                .try_into()
?0
,
1100
3
            execution_metadata: proto_action_result
1101
3
                .execution_metadata
1102
3
                .err_tip(|| "Expected execution_metadata to be set on ExecuteResponse msg")
?0
1103
3
                .try_into()
?0
,
1104
3
            server_logs: execute_response.server_logs.try_map(|v| 
{2
1105
2
                v.digest
1106
2
                    .err_tip(|| "Expected digest to be set on LogFile msg")
?0
1107
2
                    .try_into()
1108
2
            })
?0
,
1109
3
            error: execute_response
1110
3
                .status
1111
3
                .clone()
1112
3
                .and_then(|v| if v.code == 0 { 
None1
} else {
Some(v.into())2
}),
1113
3
            message: execute_response.message,
1114
        };
1115
1116
3
        if execute_response.cached_result {
1117
0
            return Ok(Self::CompletedFromCache(action_result.try_into()?));
1118
3
        }
1119
3
        Ok(Self::Completed(action_result))
1120
3
    }
1121
}
1122
1123
// TODO: Should be able to remove this after tokio-rs/prost#299
1124
pub trait TypeUrl: Message {
1125
    const TYPE_URL: &'static str;
1126
}
1127
1128
impl TypeUrl for ExecuteResponse {
1129
    const TYPE_URL: &'static str =
1130
        "type.googleapis.com/build.bazel.remote.execution.v2.ExecuteResponse";
1131
}
1132
1133
impl TypeUrl for ExecuteOperationMetadata {
1134
    const TYPE_URL: &'static str =
1135
        "type.googleapis.com/build.bazel.remote.execution.v2.ExecuteOperationMetadata";
1136
}
1137
1138
impl TypeUrl for PreconditionFailure {
1139
    const TYPE_URL: &'static str = "type.googleapis.com/google.rpc.PreconditionFailure";
1140
}
1141
1142
2
fn from_any<T>(message: &Any) -> Result<T, Error>
1143
2
where
1144
2
    T: TypeUrl + Default,
1145
{
1146
0
    error_if!(
1147
2
        message.type_url != T::TYPE_URL,
1148
        "Incorrect type when decoding Any. {} != {}",
1149
        message.type_url,
1150
0
        T::TYPE_URL.to_string()
1151
    );
1152
2
    Ok(T::decode(message.value.as_slice())
?0
)
1153
2
}
1154
1155
7
fn to_any<T>(message: &T) -> Any
1156
7
where
1157
7
    T: TypeUrl,
1158
{
1159
7
    Any {
1160
7
        type_url: T::TYPE_URL.to_string(),
1161
7
        value: message.encode_to_vec(),
1162
7
    }
1163
7
}
1164
1165
/// Current state of the action.
1166
/// This must be 100% compatible with `Operation` in `google/longrunning/operations.proto`.
1167
#[derive(Debug, Clone, Serialize, Deserialize, MetricsComponent)]
1168
pub struct ActionState {
1169
    #[metric(help = "The current stage of the action.")]
1170
    pub stage: ActionStage,
1171
    #[metric(help = "Last time this action changed stage")]
1172
    pub last_transition_timestamp: SystemTime,
1173
    #[metric(help = "The unique identifier of the action.")]
1174
    pub client_operation_id: OperationId,
1175
    #[metric(help = "The digest of the action.")]
1176
    pub action_digest: DigestInfo,
1177
}
1178
1179
impl Display for ActionState {
1180
0
    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
1181
0
        write!(
1182
0
            f,
1183
            "stage={} last_transition={} client_operation_id={} action_digest={}",
1184
0
            self.stage.name(),
1185
0
            self.last_transition_timestamp.elapsed().map_or_else(
1186
0
                |_| "<unknown duration>".to_string(),
1187
0
                |d| { format_duration(d).to_string() }
1188
            ),
1189
            self.client_operation_id,
1190
            self.action_digest
1191
        )
1192
0
    }
1193
}
1194
1195
impl PartialOrd for ActionState {
1196
0
    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
1197
0
        Some(self.cmp(other))
1198
0
    }
1199
}
1200
1201
impl Ord for ActionState {
1202
0
    fn cmp(&self, other: &Self) -> Ordering {
1203
0
        self.last_transition_timestamp
1204
0
            .cmp(&other.last_transition_timestamp)
1205
0
    }
1206
}
1207
1208
impl PartialEq for ActionState {
1209
24
    fn eq(&self, other: &Self) -> bool {
1210
        // Ignore last_transition_timestamp as the actions can still be the same even if they happened at different times
1211
24
        self.stage == other.stage
1212
24
            && self.client_operation_id == other.client_operation_id
1213
24
            && self.action_digest == other.action_digest
1214
24
    }
1215
}
1216
1217
impl Eq for ActionState {}
1218
1219
impl ActionState {
1220
1
    pub fn try_from_operation(
1221
1
        operation: Operation,
1222
1
        client_operation_id: OperationId,
1223
1
    ) -> Result<Self, Error> {
1224
1
        let metadata = from_any::<ExecuteOperationMetadata>(
1225
1
            &operation
1226
1
                .metadata
1227
1
                .err_tip(|| "No metadata in upstream operation")
?0
,
1228
        )
1229
1
        .err_tip(|| "Could not decode metadata in upstream operation")
?0
;
1230
1231
1
        let stage = match execution_stage::Value::try_from(metadata.stage).err_tip(|| 
{0
1232
0
            format!(
1233
                "Could not convert {} to execution_stage::Value",
1234
                metadata.stage
1235
            )
1236
0
        })? {
1237
0
            execution_stage::Value::Unknown => ActionStage::Unknown,
1238
0
            execution_stage::Value::CacheCheck => ActionStage::CacheCheck,
1239
0
            execution_stage::Value::Queued => ActionStage::Queued,
1240
0
            execution_stage::Value::Executing => ActionStage::Executing,
1241
            execution_stage::Value::Completed => {
1242
1
                let execute_response = operation
1243
1
                    .result
1244
1
                    .err_tip(|| "No result data for completed upstream action")
?0
;
1245
1
                match execute_response {
1246
0
                    LongRunningResult::Error(error) => ActionStage::Completed(ActionResult {
1247
0
                        error: Some(error.into()),
1248
0
                        ..ActionResult::default()
1249
0
                    }),
1250
1
                    LongRunningResult::Response(response) => {
1251
                        // Could be Completed, CompletedFromCache or Error.
1252
1
                        from_any::<ExecuteResponse>(&response)
1253
1
                            .err_tip(|| 
{0
1254
0
                                "Could not decode result structure for completed upstream action"
1255
0
                            })?
1256
1
                            .try_into()
?0
1257
                    }
1258
                }
1259
            }
1260
        };
1261
1262
1
        let action_digest = metadata
1263
1
            .action_digest
1264
1
            .err_tip(|| "No action_digest in upstream operation")
?0
1265
1
            .try_into()
1266
1
            .err_tip(|| "Could not convert action_digest into DigestInfo")
?0
;
1267
1268
1
        Ok(Self {
1269
1
            stage,
1270
1
            client_operation_id,
1271
1
            action_digest,
1272
1
            last_transition_timestamp: SystemTime::now(),
1273
1
        })
1274
1
    }
1275
1276
5
    pub fn as_operation(&self, client_operation_id: OperationId) -> Operation {
1277
5
        let stage = Into::<execution_stage::Value>::into(&self.stage) as i32;
1278
5
        let name = client_operation_id.into_string();
1279
1280
5
        let result = if self.stage.has_action_result() {
1281
2
            let execute_response: ExecuteResponse = self.stage.clone().into();
1282
2
            Some(LongRunningResult::Response(to_any(&execute_response)))
1283
        } else {
1284
3
            None
1285
        };
1286
5
        let digest = Some(self.action_digest.into());
1287
1288
5
        let metadata = ExecuteOperationMetadata {
1289
5
            stage,
1290
5
            action_digest: digest,
1291
5
            // TODO(palfrey) We should support stderr/stdout streaming.
1292
5
            stdout_stream_name: String::default(),
1293
5
            stderr_stream_name: String::default(),
1294
5
            partial_execution_metadata: None,
1295
5
        };
1296
1297
5
        Operation {
1298
5
            name,
1299
5
            metadata: Some(to_any(&metadata)),
1300
5
            done: result.is_some(),
1301
5
            result,
1302
5
        }
1303
5
    }
1304
}