lenso-kernel 0.3.11

Portable Lenso execution Kernel with bounded lifecycle settlement.
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
use std::time::Duration;
use std::{collections::BTreeMap, fmt, rc::Rc};

use super::{RequestId, lifecycle::CancellationToken};

/// An opaque extension supplied by a caller Plugin.
#[derive(Clone, Eq, PartialEq)]
pub struct InvocationExtension {
    key: String,
    value: Vec<u8>,
}

impl fmt::Debug for InvocationExtension {
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        formatter
            .debug_struct("InvocationExtension")
            .field("key", &self.key)
            .field("value", &"<redacted>")
            .finish()
    }
}

impl InvocationExtension {
    /// Creates one ordinary, caller-supplied extension value.
    pub fn new(key: impl Into<String>, value: Vec<u8>) -> Self {
        Self {
            key: key.into(),
            value,
        }
    }

    /// Returns the stable extension key.
    pub fn key(&self) -> &str {
        &self.key
    }

    /// Returns the opaque extension bytes.
    pub fn value(&self) -> &[u8] {
        &self.value
    }
}

/// An opaque extension whose issuer and audience must survive Adapter hops.
#[derive(Clone, Eq, PartialEq)]
pub struct SealedInvocationExtension {
    key: String,
    issuer: String,
    audience: Vec<String>,
    value: Vec<u8>,
    proof: String,
}

impl SealedInvocationExtension {
    /// Carries one domain-signed extension without granting it validity.
    ///
    /// Domain provider bindings must validate `proof` before projecting the
    /// payload. The Kernel preserves the signed fields and prevents replacement.
    pub fn signed(
        key: impl Into<String>,
        issuer: impl Into<String>,
        audience: impl IntoIterator<Item = impl Into<String>>,
        value: Vec<u8>,
        proof: impl Into<String>,
    ) -> Self {
        Self {
            key: key.into(),
            issuer: issuer.into(),
            audience: audience.into_iter().map(Into::into).collect(),
            value,
            proof: proof.into(),
        }
    }

    /// Returns the stable extension key.
    pub fn key(&self) -> &str {
        &self.key
    }

    /// Returns the issuer provenance without interpreting its domain.
    pub fn issuer(&self) -> &str {
        &self.issuer
    }

    /// Returns the intended Capability/Operation audience.
    pub fn audience(&self) -> &[String] {
        &self.audience
    }

    /// Returns the opaque extension bytes.
    pub fn value(&self) -> &[u8] {
        &self.value
    }

    /// Returns the domain proof covering issuer, audience, and payload.
    pub fn proof(&self) -> &str {
        &self.proof
    }

    /// Returns whether the signed audience covers one exact target Operation.
    pub fn covers(&self, capability_id: &str, operation: &str) -> bool {
        let target = format!("{capability_id}:{operation}");
        self.audience.iter().any(|audience| audience == &target)
    }
}

impl fmt::Debug for SealedInvocationExtension {
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        formatter
            .debug_struct("SealedInvocationExtension")
            .field("key", &self.key)
            .field("issuer", &self.issuer)
            .field("audience", &self.audience)
            .field("value", &"<redacted>")
            .field("proof", &"<redacted>")
            .finish()
    }
}

/// Failure returned when an Invocation Context extension cannot be attached.
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum InvocationContextError {
    /// An extension key cannot be empty.
    EmptyExtensionKey,
    /// An ordinary extension already occupies the requested key.
    ExtensionAlreadySet { key: String },
    /// A sealed extension cannot be replaced by another extension value.
    SealedExtensionAlreadySet { key: String },
    /// Sealed provenance must name an issuer and at least one audience entry.
    InvalidSealedExtension { key: String },
}

impl std::fmt::Display for InvocationContextError {
    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            Self::EmptyExtensionKey => {
                formatter.write_str("Invocation Context extension key is empty")
            }
            Self::ExtensionAlreadySet { key } => {
                write!(
                    formatter,
                    "Invocation Context extension `{key}` is already set"
                )
            }
            Self::SealedExtensionAlreadySet { key } => {
                write!(
                    formatter,
                    "sealed Invocation Context extension `{key}` is already set"
                )
            }
            Self::InvalidSealedExtension { key } => {
                write!(
                    formatter,
                    "sealed Invocation Context extension `{key}` has invalid provenance"
                )
            }
        }
    }
}

/// Kernel-owned context propagated across one native request invocation.
#[derive(Clone, Debug)]
pub struct InvocationContext {
    pub(crate) execution: Option<super::settlement::ExecutionScope>,
    pub(super) caller_instance: Option<Rc<str>>,
    pub(super) request_id: RequestId,
    pub(super) deadline: Option<Duration>,
    pub(super) remaining_budget: Option<Duration>,
    pub(super) shutdown_dependency_call: bool,
    pub(super) cancellation: CancellationToken,
    pub(super) extensions: BTreeMap<String, InvocationExtension>,
    pub(super) sealed_extensions: BTreeMap<String, SealedInvocationExtension>,
}

impl InvocationContext {
    /// Retains execution capacity for Adapter-managed work beyond its reply.
    /// The returned lease must be settled on observed termination, not cancellation acknowledgement.
    pub fn retain_execution(&self) -> Result<super::ExecutionLease, super::RuntimeFailure> {
        self.execution
            .as_ref()
            .ok_or(super::RuntimeFailure::AdmissionClosed)?
            .retain()
    }
    /// Creates an invocation context with an absolute Driver-monotonic deadline.
    pub fn new(
        request_id: RequestId,
        deadline: Option<Duration>,
        cancellation: CancellationToken,
    ) -> Self {
        Self {
            execution: None,
            caller_instance: None,
            request_id,
            deadline,
            remaining_budget: None,
            shutdown_dependency_call: false,
            cancellation,
            extensions: BTreeMap::new(),
            sealed_extensions: BTreeMap::new(),
        }
    }

    /// Attaches the resolved Caller Plugin Instance to this context.
    #[must_use]
    pub fn with_caller_instance(mut self, caller_instance: impl Into<String>) -> Self {
        self.caller_instance = Some(Rc::from(caller_instance.into()));
        self
    }

    pub(crate) fn with_shared_caller_instance(mut self, caller_instance: Rc<str>) -> Self {
        self.caller_instance = Some(caller_instance);
        self
    }

    pub(super) fn for_caller(mut self, caller_instance: &str) -> Self {
        if self.caller_instance.as_deref() != Some(caller_instance) {
            self.caller_instance = Some(Rc::from(caller_instance));
        }
        self
    }

    pub(super) fn for_shutdown_dependency_call(mut self) -> Self {
        self.shutdown_dependency_call = true;
        self
    }

    pub(super) const fn is_shutdown_dependency_call(&self) -> bool {
        self.shutdown_dependency_call
    }

    /// Returns the Caller Plugin Instance, when the App attached one.
    pub fn caller_instance(&self) -> Option<&str> {
        self.caller_instance.as_deref()
    }

    /// Returns the Kernel Request ID used for correlation and cancellation.
    pub const fn request_id(&self) -> RequestId {
        self.request_id
    }

    /// Returns the absolute Driver-monotonic deadline, when one was supplied.
    pub const fn deadline(&self) -> Option<Duration> {
        self.deadline
    }

    /// Returns the deadline budget captured immediately before provider dispatch.
    ///
    /// A missing value means the invocation has no deadline. Adapter code can
    /// forward this relative duration across an execution boundary without
    /// learning or reproducing the Driver's monotonic clock.
    pub const fn remaining_budget(&self) -> Option<Duration> {
        self.remaining_budget
    }

    /// Returns the caller-owned cooperative cancellation signal.
    pub fn cancellation(&self) -> CancellationToken {
        self.cancellation.clone()
    }

    pub(crate) fn for_child_request(mut self, request_id: RequestId) -> Self {
        self.request_id = request_id;
        self.cancellation = self.cancellation.child();
        self
    }

    /// Adds one ordinary opaque extension without replacing an existing value.
    pub fn with_extension(
        mut self,
        key: impl Into<String>,
        value: Vec<u8>,
    ) -> Result<Self, InvocationContextError> {
        let extension = InvocationExtension::new(key, value);
        if extension.key().is_empty() {
            return Err(InvocationContextError::EmptyExtensionKey);
        }
        if self.sealed_extensions.contains_key(extension.key()) {
            return Err(InvocationContextError::SealedExtensionAlreadySet {
                key: extension.key().to_owned(),
            });
        }
        if self.extensions.contains_key(extension.key()) {
            return Err(InvocationContextError::ExtensionAlreadySet {
                key: extension.key().to_owned(),
            });
        }
        self.extensions
            .insert(extension.key().to_owned(), extension);
        Ok(self)
    }

    /// Adds one sealed extension while preserving issuer, audience, and key ownership.
    pub fn with_sealed_extension(
        mut self,
        extension: SealedInvocationExtension,
    ) -> Result<Self, InvocationContextError> {
        if extension.key().is_empty() {
            return Err(InvocationContextError::EmptyExtensionKey);
        }
        if extension.issuer().is_empty()
            || extension.audience().is_empty()
            || extension.proof().is_empty()
            || extension.audience().iter().any(String::is_empty)
        {
            return Err(InvocationContextError::InvalidSealedExtension {
                key: extension.key().to_owned(),
            });
        }
        if self.sealed_extensions.contains_key(extension.key())
            || self.extensions.contains_key(extension.key())
        {
            return Err(InvocationContextError::SealedExtensionAlreadySet {
                key: extension.key().to_owned(),
            });
        }
        self.sealed_extensions
            .insert(extension.key().to_owned(), extension);
        Ok(self)
    }

    /// Returns one ordinary extension's opaque bytes.
    pub fn extension(&self, key: &str) -> Option<&[u8]> {
        self.extensions.get(key).map(InvocationExtension::value)
    }

    /// Returns ordinary extensions in deterministic key order.
    pub fn extensions(&self) -> impl Iterator<Item = &InvocationExtension> {
        self.extensions.values()
    }

    /// Returns one sealed extension by key.
    pub fn sealed_extension(&self, key: &str) -> Option<&SealedInvocationExtension> {
        self.sealed_extensions.get(key)
    }

    /// Returns sealed extensions in deterministic key order.
    pub fn sealed_extensions(&self) -> impl Iterator<Item = &SealedInvocationExtension> {
        self.sealed_extensions.values()
    }

    /// Restricts sealed extensions to one exact Capability/Operation target.
    ///
    /// Ordinary baggage is preserved. A sealed extension whose audience does
    /// not cover the target is not disclosed to that provider.
    #[must_use]
    pub fn for_target(mut self, capability_id: &str, operation: &str) -> Self {
        self.sealed_extensions
            .retain(|_, extension| extension.covers(capability_id, operation));
        self
    }

    /// Returns whether the caller has already cancelled this invocation.
    pub fn is_cancelled(&self) -> bool {
        self.cancellation.is_cancelled()
    }

    /// Returns whether the context deadline has passed at a Driver instant.
    pub fn is_expired(&self, now: Duration) -> bool {
        self.deadline.is_some_and(|deadline| deadline <= now)
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn resolved_caller_reuses_matching_storage_and_overrides_spoofed_identity() {
        let context = InvocationContext::new(1, None, CancellationToken::new())
            .with_caller_instance("consumer".to_owned());
        let original = context.caller_instance().unwrap().as_ptr();
        let context = context.for_caller("consumer");
        assert_eq!(context.caller_instance().unwrap().as_ptr(), original);

        let context = context.for_caller("resolved-consumer");
        assert_eq!(context.caller_instance(), Some("resolved-consumer"));
    }

    #[test]
    fn cloning_context_reuses_caller_storage() {
        let context = InvocationContext::new(1, None, CancellationToken::new())
            .with_caller_instance("consumer".to_owned());
        let cloned = context.clone();

        assert_eq!(
            context.caller_instance().unwrap().as_ptr(),
            cloned.caller_instance().unwrap().as_ptr()
        );
    }

    #[test]
    fn child_request_has_fresh_identity_and_one_way_cancellation() {
        let parent =
            InvocationContext::new(7, Some(Duration::from_secs(2)), CancellationToken::new())
                .with_extension("trace", b"kept".to_vec())
                .unwrap();
        let child = parent.clone().for_child_request(8);
        assert_eq!(child.request_id(), 8);
        assert_eq!(child.deadline(), parent.deadline());
        assert_eq!(child.extension("trace"), Some(b"kept".as_slice()));

        child.cancellation().cancel();
        assert!(child.is_cancelled());
        assert!(!parent.is_cancelled());

        let second_child = parent.clone().for_child_request(9);
        parent.cancellation().cancel();
        assert!(second_child.is_cancelled());
    }
}