1use super::{
2 invalid_plan, Deserialize, Deserializer, Digest, Range, Serialize, Sha256, VNextError,
3};
4
5pub const MAX_PROVIDER_WORKSPACE_SHAPE_BUCKETS: usize = 64;
6
7#[derive(Debug, Clone, Copy, PartialEq, Eq)]
10pub(crate) struct DynamicResourceShape {
11 pub(super) sequences: u32,
12 pub(super) tokens: u64,
13 pub(super) pages: u64,
14}
15
16impl DynamicResourceShape {
17 pub(crate) fn new(sequences: u32, tokens: u64, pages: u64) -> Result<Self, VNextError> {
18 if sequences == 0 || tokens == 0 || pages == 0 {
19 return Err(invalid_plan(
20 "dynamic resource shape dimensions must be non-zero",
21 ));
22 }
23 Ok(Self {
24 sequences,
25 tokens,
26 pages,
27 })
28 }
29
30 pub(crate) const fn sequences(self) -> u32 {
31 self.sequences
32 }
33
34 pub(crate) const fn tokens(self) -> u64 {
35 self.tokens
36 }
37
38 pub(crate) const fn pages(self) -> u64 {
39 self.pages
40 }
41
42 pub(crate) const fn from_validated(sequences: u32, tokens: u64, pages: u64) -> Self {
43 Self {
44 sequences,
45 tokens,
46 pages,
47 }
48 }
49}
50
51#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
55pub struct TokenSpanWork {
56 pub(super) immediate_tokens: u64,
57 pub(super) full_input_tokens: u64,
58 pub(super) fit_input_tokens: u64,
59 pub(super) immediate_start_token: u64,
60 pub(super) immediate_end_token: u64,
61 pub(super) fingerprint: String,
62}
63
64impl TokenSpanWork {
65 pub fn from_token_ids(
66 full_input: &[u32],
67 immediate_range: Range<usize>,
68 ) -> Result<Self, VNextError> {
69 Self::from_token_ids_with_fit(full_input, immediate_range, full_input.len())
70 }
71
72 pub fn from_token_ids_with_fit(
76 full_input: &[u32],
77 immediate_range: Range<usize>,
78 fit_input_tokens: usize,
79 ) -> Result<Self, VNextError> {
80 if full_input.is_empty()
81 || immediate_range.start >= immediate_range.end
82 || immediate_range.end > full_input.len()
83 || fit_input_tokens < full_input.len()
84 {
85 return Err(invalid_plan(
86 "token work requires a non-empty in-bounds immediate span and a non-regressing fit ceiling",
87 ));
88 }
89 let immediate_tokens = u64::try_from(immediate_range.len())
90 .map_err(|_| invalid_plan("immediate token span exceeds u64"))?;
91 let full_input_tokens = u64::try_from(full_input.len())
92 .map_err(|_| invalid_plan("full token input exceeds u64"))?;
93 let fit_input_tokens = u64::try_from(fit_input_tokens)
94 .map_err(|_| invalid_plan("fit token input exceeds u64"))?;
95 let immediate_start_token = u64::try_from(immediate_range.start)
96 .map_err(|_| invalid_plan("token span start exceeds u64"))?;
97 let immediate_end_token = u64::try_from(immediate_range.end)
98 .map_err(|_| invalid_plan("token span end exceeds u64"))?;
99 let mut digest = Sha256::new();
100 digest.update(b"ferrum.runtime-vnext.token-span-work.v3\0");
101 digest.update(full_input_tokens.to_le_bytes());
102 digest.update(fit_input_tokens.to_le_bytes());
103 digest.update(immediate_start_token.to_le_bytes());
104 digest.update(immediate_end_token.to_le_bytes());
105 for token in full_input {
106 digest.update(token.to_le_bytes());
107 }
108 Ok(Self {
109 immediate_tokens,
110 full_input_tokens,
111 fit_input_tokens,
112 immediate_start_token,
113 immediate_end_token,
114 fingerprint: format!("{:x}", digest.finalize()),
115 })
116 }
117
118 pub const fn immediate_tokens(&self) -> u64 {
119 self.immediate_tokens
120 }
121
122 pub const fn full_input_tokens(&self) -> u64 {
123 self.full_input_tokens
124 }
125
126 pub const fn fit_input_tokens(&self) -> u64 {
127 self.fit_input_tokens
128 }
129
130 pub fn immediate_token_range(&self) -> Range<u64> {
131 self.immediate_start_token..self.immediate_end_token
132 }
133
134 pub fn fingerprint(&self) -> &str {
135 &self.fingerprint
136 }
137}
138
139#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
140pub(crate) struct CommittedPageWork {
141 pub(super) pages: u64,
142 pub(super) fingerprint: String,
143}
144
145impl CommittedPageWork {
146 pub(crate) fn new(pages: u64, fingerprint: String) -> Result<Self, VNextError> {
147 if pages == 0
148 || fingerprint.len() != 64
149 || !fingerprint
150 .bytes()
151 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
152 {
153 return Err(invalid_plan("committed page work evidence is invalid"));
154 }
155 Ok(Self { pages, fingerprint })
156 }
157
158 pub(crate) const fn pages(&self) -> u64 {
159 self.pages
160 }
161}
162
163#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
166pub struct ResourceWorkShape {
167 pub(super) token_spans: Vec<TokenSpanWork>,
168 pub(super) committed_pages: Vec<CommittedPageWork>,
169 pub(super) immediate_sequences: u32,
170 pub(super) immediate_tokens: u64,
171 pub(super) immediate_pages: u64,
172 pub(super) fit_sequences: u32,
173 pub(super) fit_tokens: u64,
174 pub(super) fit_pages: u64,
175 pub(super) fingerprint: String,
176}
177
178impl ResourceWorkShape {
179 pub fn from_token_spans(token_spans: Vec<TokenSpanWork>) -> Result<Self, VNextError> {
180 Self::from_sources(token_spans, Vec::new())
181 }
182
183 pub fn single(token_span: TokenSpanWork) -> Result<Self, VNextError> {
184 Self::from_token_spans(vec![token_span])
185 }
186
187 pub(crate) fn from_sources(
188 token_spans: Vec<TokenSpanWork>,
189 committed_pages: Vec<CommittedPageWork>,
190 ) -> Result<Self, VNextError> {
191 if token_spans.is_empty() {
192 return Err(invalid_plan("resource work requires token evidence"));
193 }
194 let immediate_sequences = u32::try_from(token_spans.len())
195 .map_err(|_| invalid_plan("resource work sequence count exceeds u32"))?;
196 let immediate_tokens = token_spans.iter().try_fold(0_u64, |total, span| {
197 total
198 .checked_add(span.immediate_tokens())
199 .ok_or_else(|| invalid_plan("resource work immediate tokens overflow u64"))
200 })?;
201 let fit_tokens = token_spans.iter().try_fold(0_u64, |total, span| {
202 total
203 .checked_add(span.fit_input_tokens())
204 .ok_or_else(|| invalid_plan("resource work fit-input tokens overflow u64"))
205 })?;
206 let pages = committed_pages.iter().try_fold(0_u64, |total, page_work| {
207 total
208 .checked_add(page_work.pages())
209 .ok_or_else(|| invalid_plan("resource work pages overflow u64"))
210 })?;
211 #[derive(Serialize)]
212 struct FingerprintInput<'a> {
213 domain: &'static str,
214 token_spans: &'a [TokenSpanWork],
215 committed_pages: &'a [CommittedPageWork],
216 }
217 let bytes = serde_json::to_vec(&FingerprintInput {
218 domain: "ferrum.runtime-vnext.resource-work-shape.v1",
219 token_spans: &token_spans,
220 committed_pages: &committed_pages,
221 })
222 .map_err(|error| invalid_plan(format!("resource work encode failed: {error}")))?;
223 Ok(Self {
224 token_spans,
225 committed_pages,
226 immediate_sequences,
227 immediate_tokens,
228 immediate_pages: pages,
229 fit_sequences: immediate_sequences,
230 fit_tokens,
231 fit_pages: pages,
232 fingerprint: format!("{:x}", Sha256::digest(bytes)),
233 })
234 }
235
236 pub const fn immediate_sequences(&self) -> u32 {
237 self.immediate_sequences
238 }
239
240 pub const fn immediate_tokens(&self) -> u64 {
241 self.immediate_tokens
242 }
243
244 pub const fn immediate_pages(&self) -> u64 {
245 self.immediate_pages
246 }
247
248 pub const fn fit_sequences(&self) -> u32 {
249 self.fit_sequences
250 }
251
252 pub const fn fit_tokens(&self) -> u64 {
253 self.fit_tokens
254 }
255
256 pub const fn fit_pages(&self) -> u64 {
257 self.fit_pages
258 }
259
260 pub fn fingerprint(&self) -> &str {
261 &self.fingerprint
262 }
263
264 pub(crate) const fn immediate_shape(&self) -> DynamicResourceShape {
265 DynamicResourceShape::from_validated(
266 self.immediate_sequences,
267 self.immediate_tokens,
268 self.immediate_pages,
269 )
270 }
271
272 pub(crate) const fn fit_shape(&self) -> DynamicResourceShape {
273 DynamicResourceShape::from_validated(self.fit_sequences, self.fit_tokens, self.fit_pages)
274 }
275}
276
277#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
278#[serde(deny_unknown_fields)]
279pub struct DynamicResourceShapeBucket {
280 pub(super) maximum_sequences: u32,
281 pub(super) maximum_tokens: u64,
282 pub(super) maximum_pages: u64,
283 pub(super) bytes: u64,
284}
285
286#[derive(Deserialize)]
287#[serde(deny_unknown_fields)]
288pub(super) struct DynamicResourceShapeBucketWire {
289 pub(super) maximum_sequences: u32,
290 pub(super) maximum_tokens: u64,
291 pub(super) maximum_pages: u64,
292 pub(super) bytes: u64,
293}
294
295impl DynamicResourceShapeBucket {
296 pub fn new(
297 maximum_sequences: u32,
298 maximum_tokens: u64,
299 maximum_pages: u64,
300 bytes: u64,
301 ) -> Result<Self, VNextError> {
302 if maximum_sequences == 0 || maximum_tokens == 0 || maximum_pages == 0 || bytes == 0 {
303 return Err(invalid_plan(
304 "workspace shape bucket bounds and bytes must be non-zero",
305 ));
306 }
307 Ok(Self {
308 maximum_sequences,
309 maximum_tokens,
310 maximum_pages,
311 bytes,
312 })
313 }
314
315 pub(super) fn covers(&self, shape: DynamicResourceShape) -> bool {
316 shape.sequences <= self.maximum_sequences
317 && shape.tokens <= self.maximum_tokens
318 && shape.pages <= self.maximum_pages
319 }
320
321 pub const fn maximum_sequences(&self) -> u32 {
322 self.maximum_sequences
323 }
324
325 pub const fn maximum_tokens(&self) -> u64 {
326 self.maximum_tokens
327 }
328
329 pub const fn maximum_pages(&self) -> u64 {
330 self.maximum_pages
331 }
332
333 pub const fn bytes(&self) -> u64 {
334 self.bytes
335 }
336}
337
338impl<'de> Deserialize<'de> for DynamicResourceShapeBucket {
339 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
340 where
341 D: Deserializer<'de>,
342 {
343 let wire = DynamicResourceShapeBucketWire::deserialize(deserializer)?;
344 Self::new(
345 wire.maximum_sequences,
346 wire.maximum_tokens,
347 wire.maximum_pages,
348 wire.bytes,
349 )
350 .map_err(serde::de::Error::custom)
351 }
352}
353
354#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
358#[serde(rename_all = "snake_case")]
359pub enum DynamicResourceDemand {
360 Fixed {
361 bytes: u64,
362 },
363 ActualSequences {
364 bytes_per_sequence: u64,
365 maximum_sequences: u32,
366 },
367 Tokens {
368 bytes_per_token: u64,
369 maximum_tokens: u64,
370 },
371 Affine {
372 fixed_bytes: u64,
373 bytes_per_sequence: u64,
374 maximum_sequences: u32,
375 bytes_per_token: u64,
376 maximum_tokens: u64,
377 },
378 Pages {
379 bytes_per_page: u64,
380 maximum_pages: u64,
381 },
382 BoundedShapeBuckets {
383 buckets: Vec<DynamicResourceShapeBucket>,
384 },
385}
386
387#[derive(Deserialize)]
388#[serde(rename_all = "snake_case", deny_unknown_fields)]
389pub(super) enum DynamicResourceDemandWire {
390 Fixed {
391 bytes: u64,
392 },
393 ActualSequences {
394 bytes_per_sequence: u64,
395 maximum_sequences: u32,
396 },
397 Tokens {
398 bytes_per_token: u64,
399 maximum_tokens: u64,
400 },
401 Affine {
402 fixed_bytes: u64,
403 bytes_per_sequence: u64,
404 maximum_sequences: u32,
405 bytes_per_token: u64,
406 maximum_tokens: u64,
407 },
408 Pages {
409 bytes_per_page: u64,
410 maximum_pages: u64,
411 },
412 BoundedShapeBuckets {
413 buckets: Vec<DynamicResourceShapeBucket>,
414 },
415}
416
417impl DynamicResourceDemand {
418 pub fn fixed(bytes: u64) -> Result<Self, VNextError> {
419 Self::validated(Self::Fixed { bytes })
420 }
421
422 pub fn actual_sequences(
423 bytes_per_sequence: u64,
424 maximum_sequences: u32,
425 ) -> Result<Self, VNextError> {
426 Self::validated(Self::ActualSequences {
427 bytes_per_sequence,
428 maximum_sequences,
429 })
430 }
431
432 pub fn tokens(bytes_per_token: u64, maximum_tokens: u64) -> Result<Self, VNextError> {
433 Self::validated(Self::Tokens {
434 bytes_per_token,
435 maximum_tokens,
436 })
437 }
438
439 pub fn affine(
440 fixed_bytes: u64,
441 bytes_per_sequence: u64,
442 maximum_sequences: u32,
443 bytes_per_token: u64,
444 maximum_tokens: u64,
445 ) -> Result<Self, VNextError> {
446 Self::validated(Self::Affine {
447 fixed_bytes,
448 bytes_per_sequence,
449 maximum_sequences,
450 bytes_per_token,
451 maximum_tokens,
452 })
453 }
454
455 pub fn pages(bytes_per_page: u64, maximum_pages: u64) -> Result<Self, VNextError> {
456 Self::validated(Self::Pages {
457 bytes_per_page,
458 maximum_pages,
459 })
460 }
461
462 pub fn bounded_shape_buckets(
463 buckets: Vec<DynamicResourceShapeBucket>,
464 ) -> Result<Self, VNextError> {
465 Self::validated(Self::BoundedShapeBuckets { buckets })
466 }
467
468 pub(super) fn validated(demand: Self) -> Result<Self, VNextError> {
469 demand.validate()?;
470 Ok(demand)
471 }
472
473 pub(super) fn validate(&self) -> Result<(), VNextError> {
474 let valid = match self {
475 Self::Fixed { bytes } => *bytes > 0,
476 Self::ActualSequences {
477 bytes_per_sequence,
478 maximum_sequences,
479 } => {
480 *bytes_per_sequence > 0
481 && *maximum_sequences > 0
482 && bytes_per_sequence
483 .checked_mul(u64::from(*maximum_sequences))
484 .is_some()
485 }
486 Self::Tokens {
487 bytes_per_token,
488 maximum_tokens,
489 } => {
490 *bytes_per_token > 0
491 && *maximum_tokens > 0
492 && bytes_per_token.checked_mul(*maximum_tokens).is_some()
493 }
494 Self::Affine {
495 fixed_bytes,
496 bytes_per_sequence,
497 maximum_sequences,
498 bytes_per_token,
499 maximum_tokens,
500 } => {
501 (*bytes_per_sequence > 0 || *bytes_per_token > 0)
502 && *maximum_sequences > 0
503 && *maximum_tokens > 0
504 && bytes_per_sequence
505 .checked_mul(u64::from(*maximum_sequences))
506 .and_then(|sequence_bytes| fixed_bytes.checked_add(sequence_bytes))
507 .and_then(|bytes| {
508 bytes_per_token
509 .checked_mul(*maximum_tokens)
510 .and_then(|token_bytes| bytes.checked_add(token_bytes))
511 })
512 .is_some()
513 }
514 Self::Pages {
515 bytes_per_page,
516 maximum_pages,
517 } => {
518 *bytes_per_page > 0
519 && *maximum_pages > 0
520 && bytes_per_page.checked_mul(*maximum_pages).is_some()
521 }
522 Self::BoundedShapeBuckets { buckets } => {
523 !buckets.is_empty()
524 && buckets.len() <= MAX_PROVIDER_WORKSPACE_SHAPE_BUCKETS
525 && buckets.windows(2).all(|pair| {
526 let previous = &pair[0];
527 let next = &pair[1];
528 next.maximum_sequences >= previous.maximum_sequences
529 && next.maximum_tokens >= previous.maximum_tokens
530 && next.maximum_pages >= previous.maximum_pages
531 && (next.maximum_sequences > previous.maximum_sequences
532 || next.maximum_tokens > previous.maximum_tokens
533 || next.maximum_pages > previous.maximum_pages)
534 && next.bytes >= previous.bytes
535 })
536 }
537 };
538 if !valid {
539 return Err(invalid_plan(
540 "dynamic resource formula is zero, overflowing, or non-canonical",
541 ));
542 }
543 Ok(())
544 }
545
546 pub fn evaluate_bytes(&self, work: &ResourceWorkShape) -> Result<u64, VNextError> {
547 self.evaluate_shape_bytes(work.immediate_shape())
548 }
549
550 pub fn evaluate_fit_bytes(&self, work: &ResourceWorkShape) -> Result<u64, VNextError> {
551 self.evaluate_shape_bytes(work.fit_shape())
552 }
553
554 pub(crate) fn evaluate_shape_bytes(
555 &self,
556 shape: DynamicResourceShape,
557 ) -> Result<u64, VNextError> {
558 self.validate()?;
559 let bytes = match self {
560 Self::Fixed { bytes } => *bytes,
561 Self::ActualSequences {
562 bytes_per_sequence,
563 maximum_sequences,
564 } if shape.sequences <= *maximum_sequences => bytes_per_sequence
565 .checked_mul(u64::from(shape.sequences))
566 .ok_or_else(|| invalid_plan("sequence-scaled resource request overflows u64"))?,
567 Self::Tokens {
568 bytes_per_token,
569 maximum_tokens,
570 } if shape.tokens <= *maximum_tokens => bytes_per_token
571 .checked_mul(shape.tokens)
572 .ok_or_else(|| invalid_plan("token-scaled resource request overflows u64"))?,
573 Self::Affine {
574 fixed_bytes,
575 bytes_per_sequence,
576 maximum_sequences,
577 bytes_per_token,
578 maximum_tokens,
579 } if shape.sequences <= *maximum_sequences && shape.tokens <= *maximum_tokens => {
580 fixed_bytes
581 .checked_add(
582 bytes_per_sequence
583 .checked_mul(u64::from(shape.sequences))
584 .ok_or_else(|| {
585 invalid_plan("affine sequence resource request overflows u64")
586 })?,
587 )
588 .and_then(|bytes| {
589 bytes_per_token
590 .checked_mul(shape.tokens)
591 .and_then(|token_bytes| bytes.checked_add(token_bytes))
592 })
593 .ok_or_else(|| invalid_plan("affine token resource request overflows u64"))?
594 }
595 Self::Pages {
596 bytes_per_page,
597 maximum_pages,
598 } if shape.pages <= *maximum_pages => bytes_per_page
599 .checked_mul(shape.pages)
600 .ok_or_else(|| invalid_plan("page-scaled resource request overflows u64"))?,
601 Self::BoundedShapeBuckets { buckets } => buckets
602 .iter()
603 .find(|bucket| bucket.covers(shape))
604 .map(DynamicResourceShapeBucket::bytes)
605 .ok_or_else(|| invalid_plan("actual invocation shape exceeds workspace buckets"))?,
606 _ => {
607 return Err(invalid_plan(
608 "actual invocation shape exceeds its bounded resource formula",
609 ))
610 }
611 };
612 if bytes == 0 {
613 return Err(invalid_plan(
614 "dynamic resource request evaluates to zero bytes",
615 ));
616 }
617 Ok(bytes)
618 }
619
620 pub(crate) fn minimum_shape(&self) -> DynamicResourceShape {
621 DynamicResourceShape {
622 sequences: 1,
623 tokens: 1,
624 pages: 1,
625 }
626 }
627
628 pub(crate) fn theoretical_maximum_shape(&self) -> DynamicResourceShape {
629 match self {
630 Self::Fixed { .. } => self.minimum_shape(),
631 Self::ActualSequences {
632 maximum_sequences, ..
633 } => DynamicResourceShape {
634 sequences: *maximum_sequences,
635 tokens: 1,
636 pages: 1,
637 },
638 Self::Tokens { maximum_tokens, .. } => DynamicResourceShape {
639 sequences: 1,
640 tokens: *maximum_tokens,
641 pages: 1,
642 },
643 Self::Affine {
644 maximum_sequences,
645 maximum_tokens,
646 ..
647 } => DynamicResourceShape {
648 sequences: *maximum_sequences,
649 tokens: *maximum_tokens,
650 pages: 1,
651 },
652 Self::Pages { maximum_pages, .. } => DynamicResourceShape {
653 sequences: 1,
654 tokens: 1,
655 pages: *maximum_pages,
656 },
657 Self::BoundedShapeBuckets { buckets } => {
658 let bucket = buckets
659 .last()
660 .expect("validated bounded formula has at least one bucket");
661 DynamicResourceShape {
662 sequences: bucket.maximum_sequences,
663 tokens: bucket.maximum_tokens,
664 pages: bucket.maximum_pages,
665 }
666 }
667 }
668 }
669
670 pub(super) fn is_fixed(&self) -> bool {
671 matches!(self, Self::Fixed { .. })
672 }
673
674 pub(super) fn is_valid_for_sequence_scope(&self) -> bool {
675 match self {
676 Self::Fixed { .. } | Self::Tokens { .. } | Self::Pages { .. } => true,
677 Self::Affine {
678 bytes_per_sequence, ..
679 } => *bytes_per_sequence == 0,
680 Self::BoundedShapeBuckets { buckets } => {
681 buckets.iter().all(|bucket| bucket.maximum_sequences == 1)
682 }
683 Self::ActualSequences { .. } => false,
684 }
685 }
686}
687
688impl<'de> Deserialize<'de> for DynamicResourceDemand {
689 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
690 where
691 D: Deserializer<'de>,
692 {
693 let demand = match DynamicResourceDemandWire::deserialize(deserializer)? {
694 DynamicResourceDemandWire::Fixed { bytes } => Self::Fixed { bytes },
695 DynamicResourceDemandWire::ActualSequences {
696 bytes_per_sequence,
697 maximum_sequences,
698 } => Self::ActualSequences {
699 bytes_per_sequence,
700 maximum_sequences,
701 },
702 DynamicResourceDemandWire::Tokens {
703 bytes_per_token,
704 maximum_tokens,
705 } => Self::Tokens {
706 bytes_per_token,
707 maximum_tokens,
708 },
709 DynamicResourceDemandWire::Affine {
710 fixed_bytes,
711 bytes_per_sequence,
712 maximum_sequences,
713 bytes_per_token,
714 maximum_tokens,
715 } => Self::Affine {
716 fixed_bytes,
717 bytes_per_sequence,
718 maximum_sequences,
719 bytes_per_token,
720 maximum_tokens,
721 },
722 DynamicResourceDemandWire::Pages {
723 bytes_per_page,
724 maximum_pages,
725 } => Self::Pages {
726 bytes_per_page,
727 maximum_pages,
728 },
729 DynamicResourceDemandWire::BoundedShapeBuckets { buckets } => {
730 Self::BoundedShapeBuckets { buckets }
731 }
732 };
733 Self::validated(demand).map_err(serde::de::Error::custom)
734 }
735}