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}