1use core::mem;
30
31use sp_runtime::traits::Zero;
32mod benchmarking;
33pub mod migration;
34
35extern crate alloc;
36
37use crate::{configuration, paras};
38use alloc::{collections::BTreeSet, vec::Vec};
39use frame_support::{
40 pallet_prelude::*,
41 traits::{
42 defensive_prelude::*,
43 Currency,
44 ExistenceRequirement::{self, AllowDeath, KeepAlive},
45 WithdrawReasons,
46 },
47 PalletId,
48};
49use frame_system::{pallet_prelude::*, Pallet as System};
50use polkadot_primitives::{Id as ParaId, ON_DEMAND_MAX_QUEUE_MAX_SIZE};
51use sp_runtime::{
52 traits::{AccountIdConversion, One, SaturatedConversion},
53 FixedPointNumber, FixedPointOperand, FixedU128, Perbill, Saturating,
54};
55
56pub use pallet::*;
57
58mod mock_helpers;
59#[cfg(test)]
60mod tests;
61
62const LOG_TARGET: &str = "runtime::parachains::on-demand";
63
64pub trait WeightInfo {
65 fn place_order_allow_death() -> Weight;
66 fn place_order_keep_alive() -> Weight;
67 fn place_order_with_credits() -> Weight;
68}
69
70pub struct TestWeightInfo;
72
73impl WeightInfo for TestWeightInfo {
74 fn place_order_allow_death() -> Weight {
75 Weight::MAX
76 }
77
78 fn place_order_keep_alive() -> Weight {
79 Weight::MAX
80 }
81
82 fn place_order_with_credits() -> Weight {
83 Weight::MAX
84 }
85}
86
87#[derive(Encode, Decode, TypeInfo, Debug, PartialEq, Clone, Eq)]
89enum PaymentType {
90 Credits,
92 Balance,
94}
95
96pub type BalanceOf<T> =
98 <<T as Config>::Currency as Currency<<T as frame_system::Config>::AccountId>>::Balance;
99
100#[derive(Encode, Decode, TypeInfo)]
102pub struct OrderQueue<N> {
103 queue: BoundedVec<EnqueuedOrder<N>, ConstU32<ON_DEMAND_MAX_QUEUE_MAX_SIZE>>,
104}
105
106impl<N> OrderQueue<N> {
107 pub fn pop_assignment_for_cores<T: Config>(
109 &mut self,
110 now: N,
111 mut num_cores: u32,
112 ) -> impl Iterator<Item = ParaId>
113 where
114 N: Saturating + Ord + One + Copy,
115 {
116 let mut popped = BTreeSet::new();
117 let mut remaining_orders = Vec::with_capacity(self.queue.len());
118 for order in mem::take(&mut self.queue) {
119 let ready_at = order.ordered_at.saturating_plus_one().saturating_plus_one();
121 let is_ready = ready_at <= now;
122
123 if num_cores > 0 && is_ready && popped.insert(order.para_id) {
124 num_cores -= 1;
125 } else {
126 remaining_orders.push(order);
127 }
128 }
129 self.queue = BoundedVec::truncate_from(remaining_orders);
130 popped.into_iter()
131 }
132
133 fn new() -> Self {
134 OrderQueue { queue: BoundedVec::new() }
135 }
136
137 fn try_push(&mut self, now: N, para_id: ParaId) -> Result<(), ParaId> {
141 self.queue
142 .try_push(EnqueuedOrder { para_id, ordered_at: now })
143 .map_err(|o| o.para_id)
144 }
145
146 fn len(&self) -> usize {
147 self.queue.len()
148 }
149}
150
151#[derive(Encode, Decode, TypeInfo)]
153struct EnqueuedOrder<N> {
154 para_id: ParaId,
156 ordered_at: N,
158}
159
160#[derive(Encode, Decode, TypeInfo)]
162struct OrderStatus<N> {
163 traffic: FixedU128,
165
166 queue: OrderQueue<N>,
168}
169
170impl<N> Default for OrderStatus<N> {
171 fn default() -> OrderStatus<N> {
172 OrderStatus { traffic: FixedU128::default(), queue: OrderQueue::new() }
173 }
174}
175
176#[derive(PartialEq, Debug)]
178pub enum SpotTrafficCalculationErr {
179 QueueCapacityIsZero,
181 QueueSizeLargerThanCapacity,
183 Division,
185}
186
187#[frame_support::pallet]
188pub mod pallet {
189
190 use super::*;
191 use polkadot_primitives::Id as ParaId;
192
193 const STORAGE_VERSION: StorageVersion = StorageVersion::new(2);
194
195 #[pallet::pallet]
196 #[pallet::without_storage_info]
197 #[pallet::storage_version(STORAGE_VERSION)]
198 pub struct Pallet<T>(_);
199
200 #[pallet::config]
201 pub trait Config: frame_system::Config + configuration::Config + paras::Config {
202 #[allow(deprecated)]
204 type RuntimeEvent: From<Event<Self>> + IsType<<Self as frame_system::Config>::RuntimeEvent>;
205
206 type Currency: Currency<Self::AccountId>;
208
209 type WeightInfo: WeightInfo;
211
212 #[pallet::constant]
214 type TrafficDefaultValue: Get<FixedU128>;
215
216 #[pallet::constant]
219 type MaxHistoricalRevenue: Get<u32>;
220
221 #[pallet::constant]
223 type PalletId: Get<PalletId>;
224 }
225
226 #[pallet::storage]
228 pub(super) type OrderStatus<T: Config> =
229 StorageValue<_, super::OrderStatus<BlockNumberFor<T>>, ValueQuery>;
230
231 #[pallet::storage]
233 pub(super) type Revenue<T: Config> =
234 StorageValue<_, BoundedVec<BalanceOf<T>, T::MaxHistoricalRevenue>, ValueQuery>;
235
236 #[pallet::storage]
238 pub type Credits<T: Config> =
239 StorageMap<_, Blake2_128Concat, T::AccountId, BalanceOf<T>, ValueQuery>;
240
241 #[pallet::event]
242 #[pallet::generate_deposit(pub(super) fn deposit_event)]
243 pub enum Event<T: Config> {
244 OnDemandOrderPlaced { para_id: ParaId, spot_price: BalanceOf<T>, ordered_by: T::AccountId },
246 SpotPriceSet { spot_price: BalanceOf<T> },
248 AccountCredited { who: T::AccountId, amount: BalanceOf<T> },
250 UnexpectedQueueFull { dropped: u32 },
254 BatchQueued { batch: Vec<(ParaId, BlockNumberFor<T>)> },
256 }
257
258 #[pallet::error]
259 pub enum Error<T> {
260 QueueFull,
262 SpotPriceHigherThanMaxAmount,
265 InsufficientCredits,
267 }
268
269 #[pallet::hooks]
270 impl<T: Config> Hooks<BlockNumberFor<T>> for Pallet<T> {
271 fn on_initialize(_now: BlockNumberFor<T>) -> Weight {
272 Revenue::<T>::mutate(|revenue| {
274 if let Some(overdue) =
275 revenue.force_insert_keep_left(0, 0u32.into()).defensive_unwrap_or(None)
276 {
277 if let Some(last) = revenue.last_mut() {
280 *last = last.saturating_add(overdue);
281 }
282 }
283 });
284
285 let config = configuration::ActiveConfig::<T>::get();
286 OrderStatus::<T>::mutate(|order_status| {
289 Self::update_spot_traffic(&config, order_status);
290 });
291
292 T::DbWeight::get().reads_writes(3, 2)
295 }
296 }
297
298 #[pallet::call]
299 impl<T: Config> Pallet<T> {
300 #[pallet::call_index(0)]
316 #[pallet::weight(<T as Config>::WeightInfo::place_order_allow_death())]
317 #[allow(deprecated)]
318 #[deprecated(note = "This will be removed in favor of using `place_order_with_credits`")]
319 pub fn place_order_allow_death(
320 origin: OriginFor<T>,
321 max_amount: BalanceOf<T>,
322 para_id: ParaId,
323 ) -> DispatchResult {
324 let sender = ensure_signed(origin)?;
325 Pallet::<T>::do_place_order(
326 sender,
327 max_amount,
328 para_id,
329 AllowDeath,
330 PaymentType::Balance,
331 )
332 }
333
334 #[pallet::call_index(1)]
350 #[pallet::weight(<T as Config>::WeightInfo::place_order_keep_alive())]
351 #[allow(deprecated)]
352 #[deprecated(note = "This will be removed in favor of using `place_order_with_credits`")]
353 pub fn place_order_keep_alive(
354 origin: OriginFor<T>,
355 max_amount: BalanceOf<T>,
356 para_id: ParaId,
357 ) -> DispatchResult {
358 let sender = ensure_signed(origin)?;
359 Pallet::<T>::do_place_order(
360 sender,
361 max_amount,
362 para_id,
363 KeepAlive,
364 PaymentType::Balance,
365 )
366 }
367
368 #[pallet::call_index(2)]
386 #[pallet::weight(<T as Config>::WeightInfo::place_order_with_credits())]
387 pub fn place_order_with_credits(
388 origin: OriginFor<T>,
389 max_amount: BalanceOf<T>,
390 para_id: ParaId,
391 ) -> DispatchResult {
392 let sender = ensure_signed(origin)?;
393 Pallet::<T>::do_place_order(
394 sender,
395 max_amount,
396 para_id,
397 KeepAlive,
398 PaymentType::Credits,
399 )
400 }
401 }
402}
403
404impl<T: Config> Pallet<T>
406where
407 BalanceOf<T>: FixedPointOperand,
408{
409 pub fn pop_assignment_for_cores(
411 now: BlockNumberFor<T>,
412 num_cores: u32,
413 ) -> impl Iterator<Item = ParaId> {
414 pallet::OrderStatus::<T>::mutate(|order_status| {
415 order_status.queue.pop_assignment_for_cores::<T>(now, num_cores)
416 })
417 }
418
419 pub fn peek_order_queue() -> OrderQueue<BlockNumberFor<T>> {
429 pallet::OrderStatus::<T>::get().queue
430 }
431
432 pub fn push_back_order(para_id: ParaId) {
439 pallet::OrderStatus::<T>::mutate(|order_status| {
440 let now = <frame_system::Pallet<T>>::block_number();
441 if let Err(e) = order_status.queue.try_push(now, para_id) {
442 log::debug!(target: LOG_TARGET, "Pushing back order failed (queue too long): {:?}", e);
443 };
444 });
445 }
446
447 pub fn credit_account(who: T::AccountId, amount: BalanceOf<T>) {
453 Credits::<T>::mutate(who.clone(), |credits| {
454 *credits = credits.saturating_add(amount);
455 });
456 Pallet::<T>::deposit_event(Event::<T>::AccountCredited { who, amount });
457 }
458
459 fn do_place_order(
478 sender: <T as frame_system::Config>::AccountId,
479 max_amount: BalanceOf<T>,
480 para_id: ParaId,
481 existence_requirement: ExistenceRequirement,
482 payment_type: PaymentType,
483 ) -> DispatchResult {
484 let config = configuration::ActiveConfig::<T>::get();
485
486 pallet::OrderStatus::<T>::mutate(|order_status| {
487 Self::update_spot_traffic(&config, order_status);
488 let traffic = order_status.traffic;
489
490 let spot_price: BalanceOf<T> = traffic.saturating_mul_int(
492 config.scheduler_params.on_demand_base_fee.saturated_into::<BalanceOf<T>>(),
493 );
494
495 ensure!(spot_price.le(&max_amount), Error::<T>::SpotPriceHigherThanMaxAmount);
497
498 ensure!(
499 order_status.queue.len() <
500 config.scheduler_params.on_demand_queue_max_size as usize,
501 Error::<T>::QueueFull
502 );
503
504 match payment_type {
505 PaymentType::Balance => {
506 let amt = T::Currency::withdraw(
509 &sender,
510 spot_price,
511 WithdrawReasons::FEE,
512 existence_requirement,
513 )?;
514
515 let pot = Self::account_id();
518 if !System::<T>::account_exists(&pot) {
519 System::<T>::inc_providers(&pot);
520 }
521 T::Currency::resolve_creating(&pot, amt);
522 },
523 PaymentType::Credits => {
524 let credits = Credits::<T>::get(&sender);
525
526 let new_credits_value =
528 credits.checked_sub(&spot_price).ok_or(Error::<T>::InsufficientCredits)?;
529
530 if new_credits_value.is_zero() {
531 Credits::<T>::remove(&sender);
532 } else {
533 Credits::<T>::insert(&sender, new_credits_value);
534 }
535 },
536 }
537
538 Revenue::<T>::mutate(|bounded_revenue| {
540 if let Some(current_block) = bounded_revenue.get_mut(0) {
541 *current_block = current_block.saturating_add(spot_price);
542 } else {
543 bounded_revenue.try_push(spot_price).defensive_ok();
547 }
548 });
549
550 let now = <frame_system::Pallet<T>>::block_number();
551 order_status
552 .queue
553 .try_push(now, para_id)
554 .defensive_map_err(|_| Error::<T>::QueueFull)?;
555
556 Pallet::<T>::deposit_event(Event::<T>::OnDemandOrderPlaced {
557 para_id,
558 spot_price,
559 ordered_by: sender,
560 });
561
562 Ok(())
563 })
564 }
565
566 pub fn queue_order_batch(batch: &[(ParaId, BlockNumberFor<T>)]) {
568 pallet::OrderStatus::<T>::mutate(|order_status| {
569 for (queued, (para_id, ordered_at)) in batch.iter().enumerate() {
572 if let Err(err) = order_status.queue.try_push(*ordered_at, *para_id) {
573 log::debug!(
574 target: LOG_TARGET,
575 "Error trying to push an order to the queue: {:?}", err
576 );
577 Pallet::<T>::deposit_event(Event::<T>::UnexpectedQueueFull {
578 dropped: (batch.len() - queued) as u32,
579 });
580 return;
581 }
582 }
583 Pallet::<T>::deposit_event(Event::<T>::BatchQueued { batch: batch.to_vec() });
584 });
585 }
586
587 fn update_spot_traffic(
589 config: &configuration::HostConfiguration<BlockNumberFor<T>>,
590 order_status: &mut OrderStatus<BlockNumberFor<T>>,
591 ) {
592 let old_traffic = order_status.traffic;
593 match Self::calculate_spot_traffic(
594 old_traffic,
595 config.scheduler_params.on_demand_queue_max_size,
596 order_status.queue.len() as u32,
597 config.scheduler_params.on_demand_target_queue_utilization,
598 config.scheduler_params.on_demand_fee_variability,
599 ) {
600 Ok(new_traffic) => {
601 if new_traffic != old_traffic {
603 order_status.traffic = new_traffic;
604
605 let spot_price: BalanceOf<T> = new_traffic.saturating_mul_int(
607 config.scheduler_params.on_demand_base_fee.saturated_into::<BalanceOf<T>>(),
608 );
609
610 Pallet::<T>::deposit_event(Event::<T>::SpotPriceSet { spot_price });
612 }
613 },
614 Err(err) => {
615 log::debug!(
616 target: LOG_TARGET,
617 "Error calculating spot traffic: {:?}", err
618 );
619 },
620 };
621 }
622
623 fn calculate_spot_traffic(
646 traffic: FixedU128,
647 queue_capacity: u32,
648 queue_size: u32,
649 target_queue_utilisation: Perbill,
650 variability: Perbill,
651 ) -> Result<FixedU128, SpotTrafficCalculationErr> {
652 if queue_capacity == 0 {
654 return Err(SpotTrafficCalculationErr::QueueCapacityIsZero);
655 }
656
657 if queue_size > queue_capacity {
659 return Err(SpotTrafficCalculationErr::QueueSizeLargerThanCapacity);
660 }
661
662 let queue_util_ratio = FixedU128::from_rational(queue_size.into(), queue_capacity.into());
664 let positive = queue_util_ratio >= target_queue_utilisation.into();
665 let queue_util_diff = queue_util_ratio.max(target_queue_utilisation.into()) -
666 queue_util_ratio.min(target_queue_utilisation.into());
667
668 let var_times_qud = queue_util_diff.saturating_mul(variability.into());
670
671 let var_times_qud_pow = var_times_qud.saturating_mul(var_times_qud);
673
674 let div_by_two: FixedU128;
676 match var_times_qud_pow.const_checked_div(2.into()) {
677 Some(dbt) => div_by_two = dbt,
678 None => return Err(SpotTrafficCalculationErr::Division),
679 }
680
681 if positive {
683 let new_traffic = queue_util_diff
684 .saturating_add(div_by_two)
685 .saturating_add(One::one())
686 .saturating_mul(traffic);
687 Ok(new_traffic.max(<T as Config>::TrafficDefaultValue::get()))
688 } else {
689 let new_traffic = queue_util_diff.saturating_sub(div_by_two).saturating_mul(traffic);
690 Ok(new_traffic.max(<T as Config>::TrafficDefaultValue::get()))
691 }
692 }
693
694 pub fn claim_revenue_until(when: BlockNumberFor<T>) -> BalanceOf<T> {
696 let now = <frame_system::Pallet<T>>::block_number();
697 let mut amount: BalanceOf<T> = BalanceOf::<T>::zero();
698 Revenue::<T>::mutate(|revenue| {
699 while !revenue.is_empty() {
700 let index = (revenue.len() - 1) as u32;
701 if when > now.saturating_sub(index.into()) {
702 amount = amount.saturating_add(revenue.pop().defensive_unwrap_or(0u32.into()));
703 } else {
704 break;
705 }
706 }
707 });
708
709 amount
710 }
711
712 pub fn account_id() -> T::AccountId {
714 T::PalletId::get().into_account_truncating()
715 }
716
717 #[cfg(feature = "runtime-benchmarks")]
718 pub fn populate_queue(para_id: ParaId, num: u32) {
719 let now = <frame_system::Pallet<T>>::block_number();
720 pallet::OrderStatus::<T>::mutate(|order_status| {
721 for _ in 0..num {
722 order_status.queue.try_push(now, para_id).unwrap();
723 }
724 });
725 }
726
727 #[cfg(feature = "runtime-benchmarks")]
728 pub(crate) fn set_revenue(rev: BoundedVec<BalanceOf<T>, T::MaxHistoricalRevenue>) {
729 Revenue::<T>::put(rev);
730 }
731
732 #[cfg(test)]
733 fn set_order_status(new_status: OrderStatus<BlockNumberFor<T>>) {
734 pallet::OrderStatus::<T>::set(new_status);
735 }
736
737 #[cfg(test)]
738 fn get_order_status() -> OrderStatus<BlockNumberFor<T>> {
739 pallet::OrderStatus::<T>::get()
740 }
741
742 #[cfg(test)]
743 fn get_traffic_default_value() -> FixedU128 {
744 <T as Config>::TrafficDefaultValue::get()
745 }
746
747 #[cfg(test)]
748 fn get_revenue() -> Vec<BalanceOf<T>> {
749 Revenue::<T>::get().to_vec()
750 }
751}