Skip to main content

ferrum_interfaces/vnext/event/
identity.rs

1use serde::{Deserialize, Serialize, Serializer};
2use std::collections::BTreeSet;
3use std::sync::Arc;
4
5use super::{
6    invalid_event, validate_sha256, ContractVersion, DeviceId, ExecutionFrameId, NodeId,
7    NodeInvocationId, OperationId, PlanHash, PlanId, ProviderId, RequestIdentity, ResourceId,
8    ResourcePoolId, RunId, SpanId, TransactionId, VNextError, EXECUTION_IDENTITY_VERSION,
9};
10
11#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
12pub struct ExecutionIdentityParts {
13    pub version: ContractVersion,
14    pub run_id: RunId,
15    pub request_id: RequestIdentity,
16    pub sequence: u64,
17    pub plan_id: Option<PlanId>,
18    pub plan_hash: Option<PlanHash>,
19    pub frame_id: Option<ExecutionFrameId>,
20    pub node_invocation_id: Option<NodeInvocationId>,
21    pub node_id: Option<NodeId>,
22    pub operation_id: Option<OperationId>,
23    pub provider_id: Option<ProviderId>,
24    pub device_id: Option<DeviceId>,
25    pub resource_pool_id: Option<ResourcePoolId>,
26    pub resource_pool_identity_fingerprint: Option<String>,
27    pub provisioning_run_id: Option<RunId>,
28    pub provisioning_request_id: Option<RequestIdentity>,
29    pub transaction_id: Option<TransactionId>,
30    pub active_sequence_slot: Option<u32>,
31    pub admission_generation: Option<u64>,
32    pub activation_epoch: Option<u64>,
33    pub runtime_implementation_fingerprint: Option<String>,
34    pub active_sequence_fingerprint: Option<String>,
35    pub completed_sequence_fingerprint: Option<String>,
36    pub aborted_sequence_fingerprint: Option<String>,
37    pub resource_id: Option<ResourceId>,
38    pub resource_generation: Option<u64>,
39    pub resource_batch_fingerprint: Option<String>,
40    pub span_id: SpanId,
41    pub parent_span_id: Option<SpanId>,
42    pub async_links: Vec<SpanId>,
43}
44
45#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
46#[serde(deny_unknown_fields)]
47pub struct UnvalidatedExecutionIdentityParts {
48    pub version: ContractVersion,
49    pub run_id: RunId,
50    pub request_id: RequestIdentity,
51    pub sequence: u64,
52    pub plan_id: Option<PlanId>,
53    pub plan_hash: Option<PlanHash>,
54    pub frame_id: Option<ExecutionFrameId>,
55    pub node_invocation_id: Option<NodeInvocationId>,
56    pub node_id: Option<NodeId>,
57    pub operation_id: Option<OperationId>,
58    pub provider_id: Option<ProviderId>,
59    pub device_id: Option<DeviceId>,
60    pub resource_pool_id: Option<ResourcePoolId>,
61    pub resource_pool_identity_fingerprint: Option<String>,
62    pub provisioning_run_id: Option<RunId>,
63    pub provisioning_request_id: Option<RequestIdentity>,
64    pub transaction_id: Option<TransactionId>,
65    pub active_sequence_slot: Option<u32>,
66    pub admission_generation: Option<u64>,
67    pub activation_epoch: Option<u64>,
68    pub runtime_implementation_fingerprint: Option<String>,
69    pub active_sequence_fingerprint: Option<String>,
70    pub completed_sequence_fingerprint: Option<String>,
71    pub aborted_sequence_fingerprint: Option<String>,
72    pub resource_id: Option<ResourceId>,
73    pub resource_generation: Option<u64>,
74    pub resource_batch_fingerprint: Option<String>,
75    pub span_id: SpanId,
76    pub parent_span_id: Option<SpanId>,
77    pub async_links: Vec<SpanId>,
78}
79
80impl From<UnvalidatedExecutionIdentityParts> for ExecutionIdentityParts {
81    fn from(parts: UnvalidatedExecutionIdentityParts) -> Self {
82        Self {
83            version: parts.version,
84            run_id: parts.run_id,
85            request_id: parts.request_id,
86            sequence: parts.sequence,
87            plan_id: parts.plan_id,
88            plan_hash: parts.plan_hash,
89            frame_id: parts.frame_id,
90            node_invocation_id: parts.node_invocation_id,
91            node_id: parts.node_id,
92            operation_id: parts.operation_id,
93            provider_id: parts.provider_id,
94            device_id: parts.device_id,
95            resource_pool_id: parts.resource_pool_id,
96            resource_pool_identity_fingerprint: parts.resource_pool_identity_fingerprint,
97            provisioning_run_id: parts.provisioning_run_id,
98            provisioning_request_id: parts.provisioning_request_id,
99            transaction_id: parts.transaction_id,
100            active_sequence_slot: parts.active_sequence_slot,
101            admission_generation: parts.admission_generation,
102            activation_epoch: parts.activation_epoch,
103            runtime_implementation_fingerprint: parts.runtime_implementation_fingerprint,
104            active_sequence_fingerprint: parts.active_sequence_fingerprint,
105            completed_sequence_fingerprint: parts.completed_sequence_fingerprint,
106            aborted_sequence_fingerprint: parts.aborted_sequence_fingerprint,
107            resource_id: parts.resource_id,
108            resource_generation: parts.resource_generation,
109            resource_batch_fingerprint: parts.resource_batch_fingerprint,
110            span_id: parts.span_id,
111            parent_span_id: parts.parent_span_id,
112            async_links: parts.async_links,
113        }
114    }
115}
116
117#[derive(Debug, Clone, PartialEq, Eq)]
118pub struct ExecutionIdentityEnvelope {
119    parts: Arc<ExecutionIdentityParts>,
120}
121
122impl Serialize for ExecutionIdentityEnvelope {
123    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
124    where
125        S: Serializer,
126    {
127        self.parts.as_ref().serialize(serializer)
128    }
129}
130
131impl ExecutionIdentityEnvelope {
132    pub fn new(parts: ExecutionIdentityParts) -> Result<Self, VNextError> {
133        if parts.version != EXECUTION_IDENTITY_VERSION || parts.sequence == 0 {
134            return Err(invalid_event(
135                "execution identity version or sequence is invalid",
136            ));
137        }
138        if parts.plan_id.is_some() != parts.plan_hash.is_some()
139            || parts.node_invocation_id.is_some() != parts.node_id.is_some()
140            || parts.node_id.is_some() && parts.frame_id.is_none()
141            || parts.operation_id.is_some() != parts.provider_id.is_some()
142            || parts.operation_id.is_some()
143                && (parts.node_id.is_none() || parts.device_id.is_none())
144            || parts.device_id.is_some() != parts.runtime_implementation_fingerprint.is_some()
145        {
146            return Err(invalid_event(
147                "plan, frame, node invocation, operation, and provider identity shape is invalid",
148            ));
149        }
150
151        let pool_present = parts.resource_pool_id.is_some();
152        let pool_fields = [
153            parts.resource_pool_identity_fingerprint.is_some(),
154            parts.provisioning_run_id.is_some(),
155            parts.provisioning_request_id.is_some(),
156            parts.transaction_id.is_some(),
157        ];
158        if pool_fields.iter().any(|present| *present != pool_present) {
159            return Err(invalid_event(
160                "pool identity requires fingerprint and exact provisioning transaction",
161            ));
162        }
163        let active_present = parts.active_sequence_slot.is_some();
164        let active_fields = [
165            parts.admission_generation.is_some(),
166            parts.activation_epoch.is_some(),
167            parts.active_sequence_fingerprint.is_some(),
168        ];
169        if active_fields
170            .iter()
171            .any(|present| *present != active_present)
172            || active_present && parts.device_id.is_none()
173        {
174            return Err(invalid_event(
175                "active identity requires slot, admission, epoch, runtime, and binding fingerprint",
176            ));
177        }
178        if parts.admission_generation == Some(0) || parts.activation_epoch == Some(0) {
179            return Err(invalid_event(
180                "active admission generation and activation epoch must be non-zero",
181            ));
182        }
183        if (parts.completed_sequence_fingerprint.is_some()
184            || parts.aborted_sequence_fingerprint.is_some())
185            && !active_present
186            || parts.completed_sequence_fingerprint.is_some()
187                && parts.aborted_sequence_fingerprint.is_some()
188        {
189            return Err(invalid_event(
190                "sequence disposition requires one full active binding and cannot be both completed and aborted",
191            ));
192        }
193        for (value, label) in [
194            (
195                parts.resource_pool_identity_fingerprint.as_deref(),
196                "resource pool identity fingerprint",
197            ),
198            (
199                parts.runtime_implementation_fingerprint.as_deref(),
200                "runtime implementation fingerprint",
201            ),
202            (
203                parts.active_sequence_fingerprint.as_deref(),
204                "active sequence fingerprint",
205            ),
206            (
207                parts.completed_sequence_fingerprint.as_deref(),
208                "completed sequence fingerprint",
209            ),
210            (
211                parts.aborted_sequence_fingerprint.as_deref(),
212                "aborted sequence fingerprint",
213            ),
214            (
215                parts.resource_batch_fingerprint.as_deref(),
216                "resource batch fingerprint",
217            ),
218        ] {
219            if let Some(value) = value {
220                validate_sha256(value, label)?;
221            }
222        }
223        if parts.resource_id.is_some() != parts.resource_generation.is_some()
224            || parts.resource_generation == Some(0)
225            || parts.resource_id.is_some() && parts.resource_batch_fingerprint.is_some()
226            || (parts.resource_id.is_some() || parts.resource_batch_fingerprint.is_some())
227                && !pool_present
228        {
229            return Err(invalid_event(
230                "resource item/batch identity is incomplete, ambiguous, or lacks a pool",
231            ));
232        }
233        if parts.parent_span_id.as_ref() == Some(&parts.span_id) {
234            return Err(invalid_event("an execution span cannot parent itself"));
235        }
236        let mut links = BTreeSet::new();
237        if parts.async_links.iter().any(|link| {
238            link == &parts.span_id
239                || parts.parent_span_id.as_ref() == Some(link)
240                || !links.insert(link.clone())
241        }) {
242            return Err(invalid_event(
243                "async links must be unique and distinct from span and parent",
244            ));
245        }
246        Ok(Self {
247            parts: Arc::new(parts),
248        })
249    }
250
251    pub fn parts(&self) -> &ExecutionIdentityParts {
252        self.parts.as_ref()
253    }
254}