-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathbudget.rs
More file actions
1478 lines (1364 loc) · 60.6 KB
/
Copy pathbudget.rs
File metadata and controls
1478 lines (1364 loc) · 60.6 KB
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
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
//! A single [`ExecutionBudget`] shared across parsing, evaluation, pattern
//! matching, web handling, and module loading.
//!
//! # Why one object
//!
//! Before this, the runtime enforced a dozen unrelated ceilings from a dozen
//! unrelated places: the interpreter's wall-clock timeout (`max_duration` +
//! `op_count`), the pattern VM's `MAX_STEPS`, the web server's
//! `web_server_max_body_size` and `web_server_request_queue_bound`, an
//! `execute file` depth constant, and several things that were simply
//! *unbounded* (recursion depth, WebSocket queues and connection counts, HTTP
//! response size, source-file size). Each was a separate audit finding.
//!
//! `ExecutionBudget` replaces that scatter with one coherent object. It carries
//! the immutable [`BudgetLimits`] for a run plus the small amount of shared
//! mutable accounting (operations charged, live pending requests, live
//! WebSocket connections) needed to enforce them. Every dimension the task
//! enumerates lives here:
//!
//! * **Deadline and cancellation** — [`ExecutionBudget::charge_operation`] /
//! [`ExecutionBudget::check_deadline`] / [`ExecutionBudget::cancel`].
//! * **Remaining interpreter operations** — the operation counter and its
//! optional ceiling.
//! * **Recursion and import depth** — [`ExecutionBudget::check_call_depth`],
//! [`ExecutionBudget::check_import_depth`],
//! [`ExecutionBudget::check_execute_file_depth`].
//! * **Pattern transitions and active states** —
//! [`ExecutionBudget::check_pattern_steps`] /
//! [`ExecutionBudget::check_pattern_states`].
//! * **Source, file-read, body, and response bytes** —
//! [`ExecutionBudget::check_source_bytes`],
//! [`ExecutionBudget::check_file_read_bytes`],
//! [`ExecutionBudget::check_request_body_bytes`],
//! [`ExecutionBudget::check_response_bytes`].
//! * **Pending HTTP requests** — [`ExecutionBudget::max_pending_requests`].
//! * **WebSocket queue and connection limits** —
//! [`ExecutionBudget::ws_queue_bound`] /
//! [`ExecutionBudget::try_acquire_ws_connection`].
//!
//! # Thread-safety
//!
//! The interpreter core stays `!Send` (`Rc`/`RefCell`), but the budget must
//! also be readable from the multi-threaded web transport (warp accept tasks,
//! per-connection WebSocket tasks). It is therefore `Send + Sync`: every
//! mutable field is an atomic, so an `Arc<ExecutionBudget>` can be cloned into a
//! transport task without any `Rc`/`RefCell` crossing a thread boundary. Sharing
//! a small atomic-only object this way does not violate the "no `Rc`→`Arc`
//! rewrite of the interpreter" rule — it is exactly how `Arc<WflConfig>` is
//! already shared.
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::time::{Duration, Instant};
use crate::config::WflConfig;
/// How often (in charged operations) the wall-clock deadline and cancellation
/// flag are sampled on the interpreter hot path. Reading the clock on every
/// operation is a measurable cost in tight loops, so those checks run only on
/// this stride — preserving the interpreter's historic `op_count & 1023`
/// throttle. Must stay a power of two so `index & (STRIDE - 1)` is exact.
const CLOCK_SAMPLE_STRIDE: u64 = 1024;
/// Immutable per-run ceilings. Construct via [`BudgetLimits::from_config`] (maps
/// the existing `.wflcfg` keys) or [`BudgetLimits::default`].
///
/// The two "opt-in" fields ([`BudgetLimits::max_duration`] and
/// [`BudgetLimits::max_operations`]) are `Option`; `None` means "no limit".
/// Every other field is a concrete ceiling that is always enforced — its
/// default is chosen generously so existing programs never trip it while
/// runaway behaviour still gets a clean error instead of a crash or unbounded
/// growth.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BudgetLimits {
/// Wall-clock deadline for non-`main loop` execution. `None` disables it.
/// Mapped from `.wflcfg` `timeout_seconds`.
pub max_duration: Option<Duration>,
/// Hard ceiling on charged interpreter operations. `None` (the default)
/// disables it, matching the historic behaviour where the operation counter
/// only throttled clock reads. Mapped from `.wflcfg` `max_operations`
/// (`0` = unlimited).
pub max_operations: Option<u64>,
/// Maximum WFL call/recursion depth. Mapped from `.wflcfg` `max_call_depth`.
pub max_call_depth: usize,
/// Maximum nested `load module` / `include` depth. Mapped from `.wflcfg`
/// `max_import_depth`.
pub max_import_depth: usize,
/// Maximum `execute file` nesting depth. Kept small because each level
/// re-enters the whole interpreter pipeline. The CLI (and
/// `run_with_interpreter_stack`) give the run enough native stack for the
/// guard to fire as a clean error; bare embedders on a ~1 MB stack can
/// overflow first (#681). Mapped from `.wflcfg` `max_execute_file_depth`.
pub max_execute_file_depth: usize,
/// Maximum pattern-VM transitions per match attempt (ReDoS guard). Mapped
/// from `.wflcfg` `max_pattern_steps`.
pub max_pattern_steps: usize,
/// Maximum simultaneously-active pattern-VM states per match attempt.
/// Mapped from `.wflcfg` `max_pattern_states`.
pub max_pattern_states: usize,
/// Maximum WFL source-file size in bytes. Mapped from `.wflcfg`
/// `max_source_size`.
pub max_source_bytes: usize,
/// Maximum bytes buffered by one text or binary file read. Mapped from
/// `.wflcfg` `max_file_read_size`.
pub max_file_read_bytes: usize,
/// Maximum accepted HTTP request body size in bytes. Mapped from `.wflcfg`
/// `web_server_max_body_size`.
pub max_request_body_bytes: usize,
/// Maximum HTTP response body size in bytes, for both handler responses and
/// bodies read by outbound `open url` statements. Mapped from `.wflcfg`
/// `web_server_max_response_size`.
pub max_response_bytes: usize,
/// Maximum accepted-but-unhandled HTTP requests held in the transport
/// queue. Mapped from `.wflcfg` `web_server_request_queue_bound`.
pub max_pending_requests: usize,
/// Maximum wall-clock time the transport waits for a handler to answer an
/// accepted request before shedding it with 504 and releasing its in-flight
/// slot. `None` disables the timeout. Mapped from `.wflcfg`
/// `web_server_response_timeout_seconds` (`0` = disabled).
pub max_request_duration: Option<Duration>,
/// Maximum queued frames/events per WebSocket channel. Mapped from
/// `.wflcfg` `web_socket_queue_bound`.
pub max_ws_queue: usize,
/// Maximum simultaneous live WebSocket connections. Mapped from `.wflcfg`
/// `web_socket_max_connections`.
pub max_ws_connections: usize,
/// Maximum size in bytes of a single WebSocket text message (inbound or
/// outbound). Mapped from `.wflcfg` `web_socket_max_message_size`.
pub max_ws_message_bytes: usize,
/// Global ceiling in bytes on WebSocket payloads queued across every
/// connection. Mapped from `.wflcfg` `web_socket_max_queued_bytes`.
pub max_ws_queued_bytes: usize,
}
impl Default for BudgetLimits {
fn default() -> Self {
Self {
// Matches the historic interpreter default (`timeout_seconds: 60`).
max_duration: Some(Duration::from_secs(60)),
// Off by default: no operation ceiling existed before.
max_operations: None,
// No runtime recursion guard existed before (only a debug assert at
// 10_000). WFL runs on a dedicated 1 GiB stack (see main.rs), and
// 1_000 frames fit comfortably within it in both debug and release,
// so runaway recursion gets a clean error well before the stack
// overflows — while clearing any realistic program's depth.
max_call_depth: 1_000,
// Module/include nesting was unbounded before; 64 is far beyond any
// real dependency chain.
max_import_depth: 64,
// Preserves the previous `MAX_EXECUTE_FILE_DEPTH` constant exactly.
max_execute_file_depth: 4,
// Per-instruction charging (not per-wave), so this is far above the
// old per-wave `MAX_STEPS` (100_000) while still catching runaway
// (e.g. ReDoS) matches that blow past millions of transitions.
max_pattern_steps: 5_000_000,
// Active-state fan-out was unbounded before; 10_000 is generous for
// any non-pathological pattern.
max_pattern_states: 10_000,
// Source size was unchecked before; 64 MiB clears any real program.
max_source_bytes: 64 * 1024 * 1024,
// Preserve the file-I/O guide's existing 50 MiB per-read policy,
// now enforced for both text and binary reads while streaming.
max_file_read_bytes: 50 * 1024 * 1024,
// Preserves the previous `web_server_max_body_size` default (1 MiB).
max_request_body_bytes: 1_048_576,
// Response size was unchecked before; 64 MiB clears any real payload.
max_response_bytes: 64 * 1024 * 1024,
// Preserves the previous `web_server_request_queue_bound` default.
max_pending_requests: 256,
// Bound how long an accepted request may await its handler, so a
// dequeued-but-unanswered request cannot pin its in-flight slot
// forever. 300s is far longer than any serial handler needs.
max_request_duration: Some(Duration::from_secs(300)),
// WebSocket channels were unbounded before; 1_024 clears normal use.
max_ws_queue: 1_024,
// Connection count was uncapped before; 1_024 clears normal use.
max_ws_connections: 1_024,
// Per-message size was unbounded (only frame count was capped); 1 MiB
// clears normal chat/JSON traffic while bounding a single frame.
max_ws_message_bytes: 1_048_576,
// Global queued-byte ceiling across all WS channels; 16 MiB bounds
// total buffered payload regardless of connection/frame counts.
max_ws_queued_bytes: 16 * 1_048_576,
}
}
}
impl BudgetLimits {
/// Derive the limits from a loaded [`WflConfig`], mapping the existing
/// `.wflcfg` keys (`timeout_seconds`, `web_server_max_body_size`,
/// `web_server_request_queue_bound`) and the budget-specific keys onto their
/// budget fields. Any field the config does not carry keeps its default.
pub fn from_config(config: &WflConfig) -> Self {
Self {
max_duration: Some(Duration::from_secs(config.timeout_seconds)),
max_operations: config.max_operations,
max_call_depth: config.max_call_depth,
max_import_depth: config.max_import_depth,
max_execute_file_depth: config.max_execute_file_depth,
max_pattern_steps: config.max_pattern_steps,
max_pattern_states: config.max_pattern_states,
max_source_bytes: config.max_source_size,
max_file_read_bytes: config.max_file_read_size,
max_request_body_bytes: config.web_server_max_body_size,
max_response_bytes: config.web_server_max_response_size,
max_pending_requests: config.web_server_request_queue_bound.max(1),
max_request_duration: match config.web_server_response_timeout_seconds {
0 => None,
secs => Some(Duration::from_secs(secs)),
},
max_ws_queue: config.web_socket_queue_bound.max(1),
max_ws_connections: config.web_socket_max_connections.max(1),
max_ws_message_bytes: config.web_socket_max_message_size.max(1),
max_ws_queued_bytes: config.web_socket_max_queued_bytes.max(1),
}
}
/// Limits with every ceiling effectively disabled. Used by standalone
/// pattern helpers and tests that must not be constrained by a run budget.
/// Pattern limits keep the historic `MAX_STEPS`/state defaults so a
/// bare [`crate::pattern::PatternVM::new`] still resists ReDoS.
pub fn unlimited() -> Self {
Self {
max_duration: None,
max_operations: None,
max_call_depth: usize::MAX,
max_import_depth: usize::MAX,
max_execute_file_depth: usize::MAX,
max_pattern_steps: 5_000_000,
max_pattern_states: 10_000,
max_source_bytes: usize::MAX,
max_file_read_bytes: usize::MAX,
max_request_body_bytes: usize::MAX,
max_response_bytes: usize::MAX,
max_pending_requests: usize::MAX,
max_request_duration: None,
max_ws_queue: usize::MAX,
max_ws_connections: usize::MAX,
max_ws_message_bytes: usize::MAX,
max_ws_queued_bytes: usize::MAX,
}
}
}
/// The specific ceiling a run tripped. Callers map this onto their own error
/// type (the interpreter to `RuntimeError`, the pattern VM to `PatternError`).
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BudgetExceeded {
/// The wall-clock deadline elapsed.
Deadline { limit_secs: u64 },
/// The run was cancelled cooperatively via [`ExecutionBudget::cancel`].
Cancelled,
/// The interpreter-operation ceiling was reached.
Operations { limit: u64 },
/// Call/recursion depth would exceed the ceiling.
CallDepth { limit: usize },
/// `load module` / `include` nesting would exceed the ceiling.
ImportDepth { limit: usize },
/// `execute file` nesting would exceed the ceiling.
ExecuteFileDepth { limit: usize },
/// A pattern match exceeded its transition ceiling.
PatternSteps { limit: usize },
/// A pattern match exceeded its active-state ceiling.
PatternStates { limit: usize },
/// A source file exceeded the byte ceiling.
SourceBytes { limit: usize, actual: usize },
/// A text or binary file read exceeded its per-operation byte ceiling.
FileReadBytes { limit: usize, actual: usize },
/// An HTTP request body exceeded the byte ceiling.
RequestBodyBytes { limit: usize, actual: usize },
/// An HTTP response body exceeded the byte ceiling.
ResponseBytes { limit: usize, actual: usize },
/// The in-flight/pending HTTP request ceiling was reached.
PendingRequests { limit: usize },
/// The WebSocket connection ceiling was reached.
WsConnections { limit: usize },
}
impl BudgetExceeded {
/// A human-facing, Elm-style message describing the breach.
pub fn message(&self) -> String {
match self {
// Preserves the historic interpreter timeout wording verbatim so
// existing timeout diagnostics/tests keep matching.
BudgetExceeded::Deadline { limit_secs } => {
format!("Execution exceeded timeout ({limit_secs}s)")
}
BudgetExceeded::Cancelled => "Execution was cancelled".to_string(),
BudgetExceeded::Operations { limit } => {
format!("Execution exceeded the operation budget ({limit} operations)")
}
BudgetExceeded::CallDepth { limit } => {
format!("Maximum call depth ({limit}) exceeded - possible infinite recursion")
}
BudgetExceeded::ImportDepth { limit } => {
format!("Maximum import depth ({limit}) exceeded - possible circular imports")
}
BudgetExceeded::ExecuteFileDepth { limit } => format!(
"Maximum execute file nesting depth ({limit}) exceeded - possible circular execution"
),
BudgetExceeded::PatternSteps { limit } => {
format!("Pattern execution step limit exceeded ({limit} steps)")
}
BudgetExceeded::PatternStates { limit } => {
format!("Pattern active-state limit exceeded ({limit} states)")
}
BudgetExceeded::SourceBytes { limit, actual } => {
format!("Source file too large: {actual} bytes (limit: {limit} bytes)")
}
BudgetExceeded::FileReadBytes { limit, actual } => {
format!("File read too large: {actual} bytes (limit: {limit} bytes)")
}
BudgetExceeded::RequestBodyBytes { limit, actual } => {
format!("Request body too large: {actual} bytes (limit: {limit} bytes)")
}
BudgetExceeded::ResponseBytes { limit, actual } => {
format!("Response body too large: {actual} bytes (limit: {limit} bytes)")
}
BudgetExceeded::PendingRequests { limit } => {
format!("Pending request limit reached ({limit} in flight)")
}
BudgetExceeded::WsConnections { limit } => {
format!("WebSocket connection limit reached ({limit} connections)")
}
}
}
}
impl std::fmt::Display for BudgetExceeded {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.message())
}
}
impl std::error::Error for BudgetExceeded {}
/// The shared runtime budget. Cheap to clone as an `Arc`; every method takes
/// `&self` and mutates only atomics.
#[derive(Debug)]
pub struct ExecutionBudget {
limits: BudgetLimits,
/// Original finite operation ceiling for deadline-exempt main-loop HTTP.
/// An explicit CLI invocation override changes the run lifetime only.
main_loop_duration: Option<Duration>,
started: Instant,
cancelled: AtomicBool,
/// Total interpreter operations charged. Also drives the clock-sampling
/// stride, exactly as the old `op_count` field did.
operations: AtomicU64,
/// Accepted-but-unfinished HTTP requests currently in flight.
pending_requests: AtomicUsize,
/// Live WebSocket connections currently registered.
ws_connections: AtomicUsize,
/// WebSocket payload bytes currently queued across every connection's
/// inbound event and outbound frame channels. Bounded by
/// `limits.max_ws_queued_bytes`; each queued frame holds a [`WsBytePermit`]
/// that releases its bytes when the frame is consumed or shed.
ws_queued_bytes: AtomicUsize,
/// Number of `main loop`s currently active (a *depth*, not a flag). The
/// wall-clock deadline is exempt while this is `> 0` — a long-lived server
/// must not time out on its own uptime. A **depth counter** (rather than a
/// bool) makes the exemption nestable and correct under `execute file`: a
/// child interpreter that shares this budget inherits the parent's active
/// main loop, and an [`MainLoopGuard`] restores the depth on *every* exit,
/// including a caught error. Pattern matching reads this live so a match
/// launched inside a `main loop` gets the same exemption as ordinary
/// operations (see [`ExecutionBudget::charge_operation`]).
main_loop_depth: AtomicUsize,
}
// NOTE: per-match pattern accounting (transitions + active states) lives on a
// separate per-match [`PatternMeter`], *not* on the shared run budget. Two
// matches that share one `Arc<ExecutionBudget>` (e.g. concurrent web handlers)
// must never reset or share a single transition counter, or one could grant the
// other unbounded extra quota. The budget owns only the *limits* and the
// cross-cutting deadline/cancellation the meter samples.
impl ExecutionBudget {
/// Build a budget from explicit limits, starting the deadline clock now.
pub fn new(limits: BudgetLimits) -> Self {
Self {
main_loop_duration: limits.max_duration,
limits,
started: Instant::now(),
cancelled: AtomicBool::new(false),
operations: AtomicU64::new(0),
pending_requests: AtomicUsize::new(0),
ws_connections: AtomicUsize::new(0),
ws_queued_bytes: AtomicUsize::new(0),
main_loop_depth: AtomicUsize::new(0),
}
}
/// Override only the invocation lifetime, preserving the original finite
/// main-loop operation ceiling. Existing explicit-budget constructors keep
/// their historical shared lifetime/operation-duration behavior.
pub fn with_invocation_timeout(limits: BudgetLimits, duration: Duration) -> Self {
let mut budget = Self::new(limits);
budget.limits.max_duration = Some(duration);
budget
}
/// Original operation duration, unaffected by a CLI lifetime override.
pub fn main_loop_duration(&self) -> Option<Duration> {
self.main_loop_duration
}
/// Build a budget from a loaded configuration. See
/// [`BudgetLimits::from_config`].
pub fn from_config(config: &WflConfig) -> Self {
Self::new(BudgetLimits::from_config(config))
}
/// A budget with no effective ceilings (pattern ReDoS guards aside). For
/// standalone pattern helpers and tests.
pub fn unlimited() -> Self {
Self::new(BudgetLimits::unlimited())
}
/// The immutable limits backing this budget.
pub fn limits(&self) -> &BudgetLimits {
&self.limits
}
/// Time elapsed since the budget was created.
pub fn elapsed(&self) -> Duration {
self.started.elapsed()
}
// ----- Deadline & cancellation -----------------------------------------
/// Request cooperative cancellation. The next sampled checkpoint (and any
/// [`ExecutionBudget::check_cancelled`] call) observes it.
pub fn cancel(&self) {
self.cancelled.store(true, Ordering::Relaxed);
}
/// Whether cancellation has been requested.
pub fn is_cancelled(&self) -> bool {
self.cancelled.load(Ordering::Relaxed)
}
/// Fail if cancellation has been requested. Cheap; safe to call anywhere.
pub fn check_cancelled(&self) -> Result<(), BudgetExceeded> {
if self.is_cancelled() {
Err(BudgetExceeded::Cancelled)
} else {
Ok(())
}
}
/// Enter a `main loop`: bump the main-loop depth and return an
/// [`MainLoopGuard`] that restores it on drop — on *every* exit path,
/// including a caught error or a nested loop, so the wall-clock exemption is
/// never leaked or cleared early. While the depth is `> 0` the deadline is
/// exempt (a long-lived server must not time out on its own uptime).
pub fn enter_main_loop(self: &Arc<Self>) -> MainLoopGuard {
self.main_loop_depth.fetch_add(1, Ordering::AcqRel);
MainLoopGuard {
budget: Arc::clone(self),
}
}
/// Whether the wall-clock deadline is currently exempt — i.e. at least one
/// `main loop` is active on this (shared) budget. Read live, so a match or
/// operation launched inside a `main loop` is exempt and one launched after
/// it exits is not.
pub fn is_deadline_exempt(&self) -> bool {
self.main_loop_depth.load(Ordering::Acquire) > 0
}
/// The number of `main loop`s currently active on this budget.
pub fn main_loop_depth(&self) -> usize {
self.main_loop_depth.load(Ordering::Acquire)
}
/// Fail if the wall-clock deadline has elapsed. Reads the clock every call;
/// prefer [`ExecutionBudget::charge_operation`] on hot paths, which samples.
pub fn check_deadline(&self) -> Result<(), BudgetExceeded> {
if let Some(limit) = self.limits.max_duration
&& self.started.elapsed() > limit
{
return Err(BudgetExceeded::Deadline {
limit_secs: limit.as_secs(),
});
}
Ok(())
}
// ----- Interpreter operations ------------------------------------------
/// Charge one interpreter operation.
///
/// Always: increments the operation counter and, on a throttled stride,
/// honours cancellation. When `enforce_limits` is true it also enforces the
/// operation ceiling (every call) and the wall-clock deadline (on the
/// stride). `enforce_limits` is set to `false` while inside a `main loop`,
/// preserving the historic rule that a long-lived server loop is exempt from
/// the timeout — cancellation still applies so a server can be stopped.
pub fn charge_operation(&self, enforce_limits: bool) -> Result<(), BudgetExceeded> {
// `main loop` exemption: do not consume the operation budget or read the
// clock (a long-lived server would otherwise exhaust the ceiling), but
// still honour cooperative cancellation so the loop can be stopped. The
// operation counter is left untouched so exempt work cannot later push a
// post-loop `Operations` breach.
if !enforce_limits {
if self.cancelled.load(Ordering::Relaxed) {
return Err(BudgetExceeded::Cancelled);
}
return Ok(());
}
// `fetch_add` returns the previous value; use it as this op's index so
// the very first op (index 0) is a sample point, matching the old code.
let index = self.operations.fetch_add(1, Ordering::Relaxed);
let sample = index & (CLOCK_SAMPLE_STRIDE - 1) == 0;
if let Some(limit) = self.limits.max_operations
&& index >= limit
{
return Err(BudgetExceeded::Operations { limit });
}
if sample {
if self.cancelled.load(Ordering::Relaxed) {
return Err(BudgetExceeded::Cancelled);
}
if let Some(limit) = self.limits.max_duration
&& self.started.elapsed() > limit
{
return Err(BudgetExceeded::Deadline {
limit_secs: limit.as_secs(),
});
}
}
Ok(())
}
/// Operations charged so far.
pub fn operations_charged(&self) -> u64 {
self.operations.load(Ordering::Relaxed)
}
// ----- Depth guards -----------------------------------------------------
/// Fail if entering another call frame would exceed the recursion ceiling.
/// `current_depth` is the number of frames already on the stack.
pub fn check_call_depth(&self, current_depth: usize) -> Result<(), BudgetExceeded> {
if current_depth >= self.limits.max_call_depth {
Err(BudgetExceeded::CallDepth {
limit: self.limits.max_call_depth,
})
} else {
Ok(())
}
}
/// Fail if entering another import would exceed the import ceiling.
/// `current_depth` is the number of modules already loading.
pub fn check_import_depth(&self, current_depth: usize) -> Result<(), BudgetExceeded> {
if current_depth >= self.limits.max_import_depth {
Err(BudgetExceeded::ImportDepth {
limit: self.limits.max_import_depth,
})
} else {
Ok(())
}
}
/// Fail if entering another `execute file` level would exceed the ceiling.
pub fn check_execute_file_depth(&self, current_depth: usize) -> Result<(), BudgetExceeded> {
if current_depth >= self.limits.max_execute_file_depth {
Err(BudgetExceeded::ExecuteFileDepth {
limit: self.limits.max_execute_file_depth,
})
} else {
Ok(())
}
}
// ----- Pattern matching -------------------------------------------------
/// The per-match transition ceiling for the pattern VM.
pub fn pattern_step_limit(&self) -> usize {
self.limits.max_pattern_steps
}
/// The per-match active-state ceiling for the pattern VM.
pub fn pattern_state_limit(&self) -> usize {
self.limits.max_pattern_states
}
/// Fail if a pattern match has taken more transitions than allowed.
pub fn check_pattern_steps(&self, steps: usize) -> Result<(), BudgetExceeded> {
if steps > self.limits.max_pattern_steps {
Err(BudgetExceeded::PatternSteps {
limit: self.limits.max_pattern_steps,
})
} else {
Ok(())
}
}
/// Fail if a pattern match holds more active states than allowed.
pub fn check_pattern_states(&self, states: usize) -> Result<(), BudgetExceeded> {
if states > self.limits.max_pattern_states {
Err(BudgetExceeded::PatternStates {
limit: self.limits.max_pattern_states,
})
} else {
Ok(())
}
}
// ----- Byte ceilings ----------------------------------------------------
/// The source-file byte ceiling. A bounded loader reads at most this many
/// bytes (plus one) so an oversized file is refused without allocating it.
pub fn max_source_bytes(&self) -> usize {
self.limits.max_source_bytes
}
/// Fail if a source file exceeds the byte ceiling. `len` is a raw file
/// length (`u64`); a value that does not fit in `usize` (huge file on a
/// 32-bit target) is treated as over the limit rather than truncated.
pub fn check_source_len(&self, len: u64) -> Result<(), BudgetExceeded> {
match usize::try_from(len) {
Ok(len) => self.check_source_bytes(len),
Err(_) => Err(BudgetExceeded::SourceBytes {
limit: self.limits.max_source_bytes,
actual: usize::MAX,
}),
}
}
/// Fail if a source file exceeds the byte ceiling.
pub fn check_source_bytes(&self, len: usize) -> Result<(), BudgetExceeded> {
if len > self.limits.max_source_bytes {
Err(BudgetExceeded::SourceBytes {
limit: self.limits.max_source_bytes,
actual: len,
})
} else {
Ok(())
}
}
/// The per-operation ceiling for buffered text and binary file reads.
pub fn max_file_read_bytes(&self) -> usize {
self.limits.max_file_read_bytes
}
/// Fail if a text or binary file read exceeds its byte ceiling.
pub fn check_file_read_bytes(&self, len: usize) -> Result<(), BudgetExceeded> {
if len > self.limits.max_file_read_bytes {
Err(BudgetExceeded::FileReadBytes {
limit: self.limits.max_file_read_bytes,
actual: len,
})
} else {
Ok(())
}
}
/// Fail if an HTTP request body exceeds the byte ceiling.
pub fn check_request_body_bytes(&self, len: usize) -> Result<(), BudgetExceeded> {
if len > self.limits.max_request_body_bytes {
Err(BudgetExceeded::RequestBodyBytes {
limit: self.limits.max_request_body_bytes,
actual: len,
})
} else {
Ok(())
}
}
/// Fail if an HTTP response body exceeds the byte ceiling.
pub fn check_response_bytes(&self, len: usize) -> Result<(), BudgetExceeded> {
if len > self.limits.max_response_bytes {
Err(BudgetExceeded::ResponseBytes {
limit: self.limits.max_response_bytes,
actual: len,
})
} else {
Ok(())
}
}
/// The accepted HTTP request body ceiling in bytes.
pub fn max_request_body_bytes(&self) -> usize {
self.limits.max_request_body_bytes
}
// ----- Pending HTTP requests -------------------------------------------
/// The pending/in-flight HTTP request ceiling. The web transport sizes its
/// bounded queue and admission semaphore from this value.
pub fn max_pending_requests(&self) -> usize {
self.limits.max_pending_requests
}
/// The maximum time the transport waits for a handler to answer an accepted
/// request before shedding it (504) and freeing its slot. `None` = no limit.
pub fn max_request_duration(&self) -> Option<Duration> {
self.limits.max_request_duration
}
/// Try to reserve a pending-request slot, returning an RAII guard that
/// releases it on drop. `None` when already at the ceiling. Provided for
/// callers that want to account pending requests directly; the web server's
/// bounded queue is the primary enforcement path.
pub fn try_acquire_request(self: &Arc<Self>) -> Option<RequestGuard> {
let limit = self.limits.max_pending_requests;
let mut current = self.pending_requests.load(Ordering::Acquire);
loop {
if current >= limit {
return None;
}
match self.pending_requests.compare_exchange_weak(
current,
current + 1,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
return Some(RequestGuard {
budget: Arc::clone(self),
});
}
Err(observed) => current = observed,
}
}
}
/// Pending HTTP requests currently reserved.
pub fn pending_requests(&self) -> usize {
self.pending_requests.load(Ordering::Relaxed)
}
// ----- WebSocket --------------------------------------------------------
/// The per-channel WebSocket queue bound. Transport tasks size their
/// bounded channels from this value and shed on `Full`.
pub fn ws_queue_bound(&self) -> usize {
self.limits.max_ws_queue
}
/// The maximum size in bytes of a single WebSocket text message; larger
/// frames are dropped rather than queued.
pub fn max_ws_message_bytes(&self) -> usize {
self.limits.max_ws_message_bytes
}
/// Try to reserve `bytes` of the global WebSocket queued-byte budget for one
/// frame, returning an RAII [`WsBytePermit`] that releases them when the
/// frame is consumed or shed. `None` when the frame alone exceeds
/// `max_ws_message_bytes`, or when reserving would exceed the global
/// `max_ws_queued_bytes` ceiling — in which case the transport sheds it.
pub fn try_reserve_ws_bytes(self: &Arc<Self>, bytes: usize) -> Option<WsBytePermit> {
if bytes > self.limits.max_ws_message_bytes {
return None;
}
let limit = self.limits.max_ws_queued_bytes;
let mut current = self.ws_queued_bytes.load(Ordering::Acquire);
loop {
if current.saturating_add(bytes) > limit {
return None;
}
match self.ws_queued_bytes.compare_exchange_weak(
current,
current + bytes,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
return Some(WsBytePermit {
budget: Arc::clone(self),
bytes,
});
}
Err(observed) => current = observed,
}
}
}
/// WebSocket payload bytes currently queued across every connection.
pub fn ws_queued_bytes(&self) -> usize {
self.ws_queued_bytes.load(Ordering::Relaxed)
}
/// Try to reserve a WebSocket connection slot, returning an RAII guard that
/// releases it when the connection ends. `None` when already at the
/// ceiling, in which case the transport should refuse the connection.
pub fn try_acquire_ws_connection(self: &Arc<Self>) -> Option<WsConnectionGuard> {
let limit = self.limits.max_ws_connections;
let mut current = self.ws_connections.load(Ordering::Acquire);
loop {
if current >= limit {
return None;
}
match self.ws_connections.compare_exchange_weak(
current,
current + 1,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
return Some(WsConnectionGuard {
budget: Arc::clone(self),
});
}
Err(observed) => current = observed,
}
}
}
/// Live WebSocket connections currently reserved.
pub fn ws_connections(&self) -> usize {
self.ws_connections.load(Ordering::Relaxed)
}
}
impl Default for ExecutionBudget {
fn default() -> Self {
Self::new(BudgetLimits::default())
}
}
tokio::task_local! {
/// The budget in effect for the current async **task** — an interpreter run
/// or a REPL command. Task-local, NOT thread-local, so two interpreter
/// futures interleaved on one thread (a library embedder that `join!`s or
/// `spawn_local`s two `Interpreter`s — both are re-exported from the crate
/// root and are `!Send`, so this is legal) never observe each other's budget
/// or restore stale state across an `.await`. Any async run establishes this
/// scope, so it always takes precedence over the synchronous fallback below.
static CURRENT_BUDGET_TASK: Arc<ExecutionBudget>;
}
thread_local! {
/// Synchronous fallback current budget, consulted only when no task-local
/// scope is active. It exists for code that runs to completion **without
/// awaiting** and cannot interleave — specifically the CLI front-end
/// (lex/parse/analyze/type-check) installed by `main`, which runs on a
/// single-future runtime. Because every async run (`interpret`, REPL
/// `process_line`) wraps itself in a [`ExecutionBudget::scope`] that shadows
/// this, the fallback can never cross-contaminate an interleaved run.
static CURRENT_BUDGET_THREAD: std::cell::RefCell<Option<Arc<ExecutionBudget>>> =
const { std::cell::RefCell::new(None) };
}
impl ExecutionBudget {
/// Run `future` with `budget` installed as the task-local current budget,
/// restoring the previous task-local (if any) when it completes. This is the
/// interleaving-safe way to scope a run; prefer it for every async run.
/// Nesting (e.g. an `execute file` child) is supported.
pub async fn scope<F>(budget: Arc<ExecutionBudget>, future: F) -> F::Output
where
F: std::future::Future,
{
CURRENT_BUDGET_TASK.scope(budget, future).await
}
/// Install `budget` as the synchronous thread-local fallback for the
/// lifetime of the returned guard (restoring the previous one on drop). Use
/// this ONLY for synchronous, non-interleaving contexts (the CLI front-end);
/// async runs must use [`ExecutionBudget::scope`], which takes precedence.
pub fn enter(budget: Arc<ExecutionBudget>) -> CurrentBudgetGuard {
let previous = CURRENT_BUDGET_THREAD.with(|c| c.borrow_mut().replace(budget));
CurrentBudgetGuard { previous }
}
/// The current budget: the task-local scope if one is active (an async run),
/// otherwise the synchronous thread-local fallback (the CLI front-end).
pub fn current() -> Option<Arc<ExecutionBudget>> {
CURRENT_BUDGET_TASK
.try_with(Arc::clone)
.ok()
.or_else(|| CURRENT_BUDGET_THREAD.with(|c| c.borrow().clone()))
}
/// The current budget, or a fresh unlimited one (which still carries the
/// pattern ReDoS ceilings) when no run is active — so a bare
/// [`crate::pattern::PatternVM::new`] is always bounded.
pub fn current_or_default() -> Arc<ExecutionBudget> {
Self::current().unwrap_or_else(|| Arc::new(Self::unlimited()))
}
}
/// Restores the previous synchronous thread-local fallback budget when dropped.
#[must_use]
pub struct CurrentBudgetGuard {
previous: Option<Arc<ExecutionBudget>>,
}
impl Drop for CurrentBudgetGuard {
fn drop(&mut self) {
CURRENT_BUDGET_THREAD.with(|c| *c.borrow_mut() = self.previous.take());
}
}
/// RAII slot for one in-flight HTTP request; releases on drop.
#[derive(Debug)]
pub struct RequestGuard {
budget: Arc<ExecutionBudget>,
}
impl Drop for RequestGuard {
fn drop(&mut self) {
self.budget.pending_requests.fetch_sub(1, Ordering::AcqRel);
}
}
/// RAII slot for one live WebSocket connection; releases on drop.
#[derive(Debug)]
pub struct WsConnectionGuard {
budget: Arc<ExecutionBudget>,
}
impl Drop for WsConnectionGuard {
fn drop(&mut self) {
self.budget.ws_connections.fetch_sub(1, Ordering::AcqRel);
}
}
/// RAII reservation of global WebSocket queued bytes for one queued frame;
/// releases the bytes when the frame is consumed (dequeued and dropped) or shed.
#[derive(Debug)]
pub struct WsBytePermit {
budget: Arc<ExecutionBudget>,
bytes: usize,
}
impl Drop for WsBytePermit {
fn drop(&mut self) {
self.budget
.ws_queued_bytes
.fetch_sub(self.bytes, Ordering::AcqRel);
}
}
/// RAII marker for one active `main loop`; decrements the shared main-loop depth
/// on drop. Because it restores on *every* exit — normal, early `return`, a
/// caught error unwinding through the loop, or a nested loop — the wall-clock
/// exemption can never leak past the loop or be cleared while an outer loop is
/// still active.
#[derive(Debug)]
#[must_use]
pub struct MainLoopGuard {
budget: Arc<ExecutionBudget>,
}
impl Drop for MainLoopGuard {
fn drop(&mut self) {
self.budget.main_loop_depth.fetch_sub(1, Ordering::AcqRel);
}
}
/// Per-top-level-match pattern metering.
///
/// A fresh `PatternMeter` is created for each top-level pattern operation
/// (`matches`/`find`/`find_all`) and cloned (as an `Arc`) **only** into nested
/// lookaround/lookbehind VMs, so their transitions and active states count
/// against the *same* per-match ceilings as the enclosing match. It borrows the
/// run's limits, wall-clock deadline, and cancellation flag from the shared
/// [`ExecutionBudget`], but keeps its own transition counter and active-state
/// accounting — so two matches sharing one run budget (e.g. concurrent web
/// handlers) never reset or share each other's meter.
///
/// Kept atomic (rather than `Cell`) so a [`crate::pattern::PatternVM`] stays
/// `Send`; a single match runs on one thread, so the atomics are uncontended.
#[derive(Debug)]
pub struct PatternMeter {
budget: Arc<ExecutionBudget>,
/// Transitions charged for this match (all frontiers, all nested VMs).
steps: AtomicU64,
/// State slots reserved live across every frontier (current + next
/// generation + any suspended nested lookaround/lookbehind frontiers).
active_states: AtomicUsize,
}
impl PatternMeter {
/// A fresh per-match meter bound to `budget`.
pub fn new(budget: Arc<ExecutionBudget>) -> Arc<Self> {
Arc::new(Self {
budget,
steps: AtomicU64::new(0),
active_states: AtomicUsize::new(0),
})
}
/// The shared run budget this meter borrows limits/deadline/cancellation from.
pub fn budget(&self) -> &Arc<ExecutionBudget> {
&self.budget
}
/// Reset the per-match counters. Called once at the start of each *direct*
/// top-level VM operation so reusing one VM does not accumulate transitions