1use std::convert::Infallible;
4use std::future::Future;
5use std::net::{IpAddr, Ipv6Addr, SocketAddr};
6use std::time::Duration as StdDuration;
7
8use async_stream::stream;
9use axum::extract::{FromRequestParts, Path, Query, State};
10use axum::http::{HeaderMap, HeaderValue, StatusCode, header, request::Parts};
11use axum::response::sse::{Event, KeepAlive, Sse};
12use axum::response::{IntoResponse, Response};
13use axum::routing::{any, get};
14use axum::{Json, Router};
15use bhtune_cli::commands::tune::{drive, prepare_owned};
16use bhtune_cli::config::{
17 DEMO_COOKIE_NAME, DEMO_CYCLES_COUNT_DEFAULT, DEMO_CYCLES_COUNT_MAX, DEMO_CYCLES_COUNT_MIN,
18 DEMO_CYCLES_SKIP_DEFAULT, DEMO_CYCLES_SKIP_MAX, DEMO_CYCLES_SKIP_MIN,
19 DEMO_NOISE_PROTECTION_SECS_DEFAULT, DEMO_NOISE_PROTECTION_SECS_MAX,
20 DEMO_NOISE_PROTECTION_SECS_MIN, DEMO_RANGE_ENDPOINT_MAX, DEMO_RANGE_ENDPOINT_MIN,
21 DEMO_RANGE_HIGH, DEMO_RANGE_LOW, DEMO_RANGE_SPAN_MAX, DEMO_RANGE_SPAN_MIN,
22 DEMO_RELAY_AMP_DEFAULT, DEMO_RELAY_AMP_MAX, DEMO_RELAY_AMP_MIN, DEMO_SIM_DEAD_TIME_DEFAULT,
23 DEMO_SIM_DEAD_TIME_MAX, DEMO_SIM_DEAD_TIME_MIN, DEMO_SIM_GAIN_ABS_MIN, DEMO_SIM_GAIN_DEFAULT,
24 DEMO_SIM_GAIN_MAX, DEMO_SIM_INITIAL_VALUE_DEFAULT, DEMO_SIM_NOISE_DEFAULT,
25 DEMO_SIM_NOISE_MAX_PV_SPAN_FRACTION, DEMO_SIM_SEED_DEFAULT, DEMO_SIM_SEED_MAX,
26 DEMO_SIM_TAU_DEFAULT, DEMO_SIM_TAU_MAX, DEMO_SIM_TAU_MIN, DEMO_TAG_NAME, DemoPolicy,
27 ServerMode,
28};
29use bhtune_core::{ControllerDirection, template::built_in_templates};
30use bhtune_db::models::{
31 DemoSessionRow, Pagination, TemplateOrigin, TuneDriver, TuneMvActuationRow, TuneOutcome,
32 TuneResultRow, TuneRunRow, TuneSampleRow, TuneWriteRow,
33};
34use chrono::{Duration, Utc};
35use rand::random;
36use serde::Serialize;
37use sha2::{Digest, Sha256};
38use tower_http::limit::RequestBodyLimitLayer;
39use tower_http::timeout::TimeoutLayer;
40
41use crate::AppState;
42use crate::error::ApiError;
43use crate::routes::history::{
44 InitialReadingsResponse, MvActuationResponse, PidConstantTagsResponse,
45 PidParameterLabelsResponse, ResultResponse, RunDetailResponse, RunExportFormat, RunExportQuery,
46 RunListQuery, RunListResponse, RunSummaryResponse, SampleResponse, WriteResponse,
47 filter_from_query, parse_stored_request,
48};
49use crate::routes::runs::StartRunRequest;
50use crate::routes::templates::TemplateResponse;
51use crate::state::DemoQuotaExceeded;
52
53const SSE_POLL_INTERVAL: StdDuration = StdDuration::from_millis(300);
54const FORWARDED_CLIENT_IP_HEADER: &str = "X-BHTune-Client-IP";
55
56pub(crate) struct PeerAddress(SocketAddr);
57
58impl<S> FromRequestParts<S> for PeerAddress
59where
60 S: Send + Sync,
61{
62 type Rejection = ApiError;
63
64 fn from_request_parts(
65 parts: &mut Parts,
66 _state: &S,
67 ) -> impl Future<Output = Result<Self, Self::Rejection>> + Send {
68 std::future::ready(
69 parts
70 .extensions
71 .get::<axum::extract::ConnectInfo<SocketAddr>>()
72 .map(|value| Self(value.0))
73 .ok_or_else(|| {
74 ApiError::Internal(anyhow::anyhow!("Demo request is missing peer ConnectInfo"))
75 }),
76 )
77 }
78}
79
80fn ensure_demo(state: &AppState) -> Result<(), ApiError> {
81 (state.mode == ServerMode::Demo)
82 .then_some(())
83 .ok_or_else(|| ApiError::NotFound("demo mode is not enabled".into()))
84}
85
86fn bad_request(message: impl Into<String>) -> ApiError {
87 ApiError::BadRequest(message.into())
88}
89
90fn too_many(message: impl Into<String>, retry_after_secs: u64) -> ApiError {
91 ApiError::TooManyRequests {
92 message: message.into(),
93 retry_after_secs,
94 }
95}
96
97fn global_capacity(message: impl Into<String>, retry_after_secs: u64) -> ApiError {
98 ApiError::GlobalCapacity {
99 message: message.into(),
100 retry_after_secs,
101 }
102}
103
104pub(crate) fn ordinary_request_permit(
105 state: &AppState,
106) -> Result<tokio::sync::OwnedSemaphorePermit, ApiError> {
107 state
108 .demo_runtime
109 .try_acquire_ordinary_request()
110 .map_err(|_| {
111 global_capacity(
112 "demo request capacity is temporarily exhausted",
113 state.demo_policy.ordinary_request_timeout_secs,
114 )
115 })
116}
117
118fn quota_error(error: DemoQuotaExceeded) -> ApiError {
119 too_many(
120 "demo accepted-start quota exceeded; wait before starting another tune",
121 error.retry_after_secs,
122 )
123}
124
125fn validate_demo_request(request: &StartRunRequest) -> Result<(), ApiError> {
126 validate_demo_identity(request)?;
127 validate_demo_scalar_ranges(request)?;
128 validate_demo_gain_and_direction(request)?;
129 validate_demo_discrete_fields(request)?;
130 validate_demo_process_ranges(request)?;
131 Ok(())
132}
133
134fn validate_demo_identity(request: &StartRunRequest) -> Result<(), ApiError> {
135 if request.driver != TuneDriver::Simulator {
136 return Err(bad_request("demo mode only supports the simulator driver"));
137 }
138 if !built_in_templates()
139 .iter()
140 .any(|template| template.name == request.template)
141 {
142 return Err(bad_request(
143 "demo mode only supports built-in DCS/PLC templates",
144 ));
145 }
146 if request.tagname != DEMO_TAG_NAME {
147 return Err(bad_request(format!(
148 "demo requests must use the fixed tag label '{DEMO_TAG_NAME}'"
149 )));
150 }
151 if !request.controller_type.is_allowed_for(request.process_type) {
152 return Err(bad_request(
153 "the selected controller type is not valid for the selected process type",
154 ));
155 }
156 Ok(())
157}
158
159fn validate_demo_scalar_ranges(request: &StartRunRequest) -> Result<(), ApiError> {
160 let ranges = [
161 (
162 "relay_amp",
163 request.relay_amp,
164 DEMO_RELAY_AMP_MIN,
165 DEMO_RELAY_AMP_MAX,
166 ),
167 (
168 "sim_tau",
169 request.sim_tau,
170 DEMO_SIM_TAU_MIN,
171 DEMO_SIM_TAU_MAX,
172 ),
173 (
174 "sim_dead_time",
175 request.sim_dead_time,
176 DEMO_SIM_DEAD_TIME_MIN,
177 DEMO_SIM_DEAD_TIME_MAX,
178 ),
179 ];
180 for (field, value, min, max) in ranges {
181 if !value.is_finite() || !(min..=max).contains(&value) {
182 return Err(bad_request(format!(
183 "demo field '{field}' must be finite and between {min} and {max}"
184 )));
185 }
186 }
187 Ok(())
188}
189
190fn validate_demo_gain_and_direction(request: &StartRunRequest) -> Result<(), ApiError> {
191 if !request.sim_gain.is_finite()
192 || !(DEMO_SIM_GAIN_ABS_MIN..=DEMO_SIM_GAIN_MAX).contains(&request.sim_gain.abs())
193 {
194 return Err(bad_request(format!(
195 "demo field 'sim_gain' must be finite and between -{DEMO_SIM_GAIN_MAX} and \
196 -{DEMO_SIM_GAIN_ABS_MIN} or between {DEMO_SIM_GAIN_ABS_MIN} and {DEMO_SIM_GAIN_MAX}"
197 )));
198 }
199 let expected_direction = if request.sim_gain.is_sign_positive() {
200 ControllerDirection::Reverse
201 } else {
202 ControllerDirection::Direct
203 };
204 if request.direction != Some(expected_direction) {
205 return Err(bad_request(format!(
206 "demo direction must be '{}' for a {} process gain so the simulated loop uses negative feedback",
207 match expected_direction {
208 ControllerDirection::Direct => "direct",
209 ControllerDirection::Reverse => "reverse",
210 },
211 if request.sim_gain.is_sign_positive() {
212 "positive"
213 } else {
214 "negative"
215 }
216 )));
217 }
218 Ok(())
219}
220
221fn validate_demo_discrete_fields(request: &StartRunRequest) -> Result<(), ApiError> {
222 if request.sim_seed > DEMO_SIM_SEED_MAX {
223 return Err(bad_request(
224 "demo field 'sim_seed' exceeds the supported maximum",
225 ));
226 }
227 if !request
228 .cycles_skip
229 .is_some_and(|value| (DEMO_CYCLES_SKIP_MIN..=DEMO_CYCLES_SKIP_MAX).contains(&value))
230 {
231 return Err(bad_request(format!(
232 "demo field 'cycles_skip' must be between {DEMO_CYCLES_SKIP_MIN} and {DEMO_CYCLES_SKIP_MAX}"
233 )));
234 }
235 if !request
236 .cycles_count
237 .is_some_and(|value| (DEMO_CYCLES_COUNT_MIN..=DEMO_CYCLES_COUNT_MAX).contains(&value))
238 {
239 return Err(bad_request(format!(
240 "demo field 'cycles_count' must be between {DEMO_CYCLES_COUNT_MIN} and {DEMO_CYCLES_COUNT_MAX}"
241 )));
242 }
243 if !request.noise_protection_secs.is_some_and(|value| {
244 (DEMO_NOISE_PROTECTION_SECS_MIN..=DEMO_NOISE_PROTECTION_SECS_MAX).contains(&value)
245 }) {
246 return Err(bad_request(format!(
247 "demo field 'noise_protection_secs' must be between {DEMO_NOISE_PROTECTION_SECS_MIN} and {DEMO_NOISE_PROTECTION_SECS_MAX}"
248 )));
249 }
250 Ok(())
251}
252
253fn validate_demo_process_ranges(request: &StartRunRequest) -> Result<(), ApiError> {
254 let pv_low = validate_range_endpoint("pv_range_low", request.pv_range_low)?;
255 let pv_high = validate_range_endpoint("pv_range_high", request.pv_range_high)?;
256 let mv_low = validate_range_endpoint("mv_range_low", request.mv_range_low)?;
257 let mv_high = validate_range_endpoint("mv_range_high", request.mv_range_high)?;
258 let pv_span = validate_range_span("PV", pv_low, pv_high)?;
259 validate_range_span("MV", mv_low, mv_high)?;
260 validate_initial_value("sim_initial_pv", request.sim_initial_pv, pv_low, pv_high)?;
261 validate_initial_value("sim_initial_mv", request.sim_initial_mv, mv_low, mv_high)?;
262 let max_noise = pv_span * DEMO_SIM_NOISE_MAX_PV_SPAN_FRACTION;
263 if !request.sim_noise.is_finite() || !(0.0..=max_noise).contains(&request.sim_noise) {
264 return Err(bad_request(format!(
265 "demo field 'sim_noise' must be finite and between 0 and {max_noise} (5% of the PV span)"
266 )));
267 }
268 Ok(())
269}
270
271fn validate_range_endpoint(field: &str, value: Option<f32>) -> Result<f32, ApiError> {
272 let value = value.ok_or_else(|| bad_request(format!("demo field '{field}' is required")))?;
273 if value.is_finite() && (DEMO_RANGE_ENDPOINT_MIN..=DEMO_RANGE_ENDPOINT_MAX).contains(&value) {
274 Ok(value)
275 } else {
276 Err(bad_request(format!(
277 "demo field '{field}' must be finite and between {DEMO_RANGE_ENDPOINT_MIN} and {DEMO_RANGE_ENDPOINT_MAX}"
278 )))
279 }
280}
281
282fn validate_range_span(label: &str, low: f32, high: f32) -> Result<f32, ApiError> {
283 let span = high - low;
284 if (DEMO_RANGE_SPAN_MIN..=DEMO_RANGE_SPAN_MAX).contains(&span) {
285 Ok(span)
286 } else {
287 Err(bad_request(format!(
288 "demo {label} range must have an ordered span between {DEMO_RANGE_SPAN_MIN} and {DEMO_RANGE_SPAN_MAX}"
289 )))
290 }
291}
292
293fn validate_initial_value(field: &str, value: f32, low: f32, high: f32) -> Result<(), ApiError> {
294 if value.is_finite() && (low..=high).contains(&value) {
295 Ok(())
296 } else {
297 Err(bad_request(format!(
298 "demo field '{field}' must be finite and within its configured range"
299 )))
300 }
301}
302
303fn apply_demo_defaults(object: &mut serde_json::Map<String, serde_json::Value>) {
304 let defaults = [
305 ("relay_amp", serde_json::json!(DEMO_RELAY_AMP_DEFAULT)),
306 ("cycles_skip", serde_json::json!(DEMO_CYCLES_SKIP_DEFAULT)),
307 ("cycles_count", serde_json::json!(DEMO_CYCLES_COUNT_DEFAULT)),
308 (
309 "noise_protection_secs",
310 serde_json::json!(DEMO_NOISE_PROTECTION_SECS_DEFAULT),
311 ),
312 ("sim_gain", serde_json::json!(DEMO_SIM_GAIN_DEFAULT)),
313 ("sim_tau", serde_json::json!(DEMO_SIM_TAU_DEFAULT)),
314 (
315 "sim_dead_time",
316 serde_json::json!(DEMO_SIM_DEAD_TIME_DEFAULT),
317 ),
318 ("sim_noise", serde_json::json!(DEMO_SIM_NOISE_DEFAULT)),
319 ("sim_seed", serde_json::json!(DEMO_SIM_SEED_DEFAULT)),
320 (
321 "sim_initial_pv",
322 serde_json::json!(DEMO_SIM_INITIAL_VALUE_DEFAULT),
323 ),
324 (
325 "sim_initial_mv",
326 serde_json::json!(DEMO_SIM_INITIAL_VALUE_DEFAULT),
327 ),
328 ("pv_range_high", serde_json::json!(DEMO_RANGE_HIGH)),
329 ("pv_range_low", serde_json::json!(DEMO_RANGE_LOW)),
330 ("mv_range_high", serde_json::json!(DEMO_RANGE_HIGH)),
331 ("mv_range_low", serde_json::json!(DEMO_RANGE_LOW)),
332 ("direction", serde_json::json!(ControllerDirection::Reverse)),
333 ];
334 for (field, value) in defaults {
335 object.entry(field.to_owned()).or_insert(value);
336 }
337}
338
339fn parse_demo_request(mut value: serde_json::Value) -> Result<StartRunRequest, ApiError> {
340 let object = value
341 .as_object_mut()
342 .ok_or_else(|| bad_request("demo run request body must be a JSON object"))?;
343 const ALLOWED_FIELDS: &[&str] = &[
344 "tagname",
345 "template",
346 "process_type",
347 "controller_type",
348 "relay_amp",
349 "cycles_skip",
350 "cycles_count",
351 "noise_protection_secs",
352 "driver",
353 "sim_gain",
354 "sim_tau",
355 "sim_dead_time",
356 "sim_noise",
357 "sim_seed",
358 "sim_initial_pv",
359 "sim_initial_mv",
360 "pv_range_high",
361 "pv_range_low",
362 "mv_range_high",
363 "mv_range_low",
364 "direction",
365 ];
366 if let Some(field) = object
367 .keys()
368 .find(|field| !ALLOWED_FIELDS.contains(&field.as_str()))
369 {
370 return Err(bad_request(format!("demo requests may not set '{field}'")));
371 }
372 apply_demo_defaults(object);
373 let request: StartRunRequest = serde_json::from_value(value)
374 .map_err(|error| bad_request(format!("invalid demo run request: {error}")))?;
375 validate_demo_request(&request)?;
376 Ok(request)
377}
378
379fn token_hash(token: &str) -> String {
380 let mut hasher = Sha256::new();
381 hasher.update(token.as_bytes());
382 hex_encode(&hasher.finalize())
383}
384
385fn hex_encode(bytes: &[u8]) -> String {
386 bytes.iter().map(|byte| format!("{byte:02x}")).collect()
387}
388
389fn trusted_proxy(peer: SocketAddr, configured: Option<&str>) -> bool {
390 let Some(configured) = configured else {
391 return false;
392 };
393 let configured = configured.trim();
394 if let Ok(address) = configured.parse::<IpAddr>() {
395 return address == peer.ip();
396 }
397 let Some((network, bits)) = configured.split_once('/') else {
398 return false;
399 };
400 let (Ok(network), Ok(bits)) = (network.parse::<IpAddr>(), bits.parse::<u32>()) else {
401 return false;
402 };
403 match (peer.ip(), network) {
404 (IpAddr::V4(peer), IpAddr::V4(network)) if bits <= 32 => {
405 let mask = if bits == 0 {
406 0
407 } else {
408 u32::MAX << (32 - bits)
409 };
410 u32::from(peer) & mask == u32::from(network) & mask
411 }
412 (IpAddr::V6(peer), IpAddr::V6(network)) if bits <= 128 => {
413 let mask = if bits == 0 {
414 0
415 } else {
416 u128::MAX << (128 - bits)
417 };
418 u128::from(peer) & mask == u128::from(network) & mask
419 }
420 _ => false,
421 }
422}
423
424fn quota_ip(address: IpAddr) -> String {
425 match address {
426 IpAddr::V4(address) => address.to_string(),
427 IpAddr::V6(address) => {
428 let network = Ipv6Addr::from(u128::from(address) & (u128::MAX << 64));
429 format!("{network}/64")
430 }
431 }
432}
433
434fn client_ip(headers: &HeaderMap, peer: PeerAddress, configured_proxy: Option<&str>) -> String {
435 if trusted_proxy(peer.0, configured_proxy) {
436 let mut values = headers.get_all(FORWARDED_CLIENT_IP_HEADER).iter();
437 if let (Some(value), None) = (values.next(), values.next())
438 && let Ok(value) = value.to_str()
439 && value.trim() == value
440 && !value.contains(',')
441 && let Ok(address) = value.parse::<IpAddr>()
442 {
443 return quota_ip(address);
444 }
445 }
446 quota_ip(peer.0.ip())
447}
448
449fn valid_token(token: &str) -> bool {
450 token.len() == 64
451 && token
452 .bytes()
453 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
454}
455
456fn cookie_token(headers: &HeaderMap) -> Result<Option<String>, ()> {
457 let mut found = None;
458 for value in headers.get_all(header::COOKIE) {
459 let value = value.to_str().map_err(|_| ())?;
460 for part in value.split(';') {
461 let Some((name, value)) = part.trim().split_once('=') else {
462 continue;
463 };
464 if name != DEMO_COOKIE_NAME {
465 continue;
466 }
467 if found.is_some() || !valid_token(value) {
468 return Err(());
469 }
470 found = Some(value.to_owned());
471 }
472 }
473 Ok(found)
474}
475
476fn parsed_token_hash(headers: &HeaderMap) -> Result<String, ApiError> {
477 let token = cookie_token(headers)
478 .map_err(|()| ApiError::Unauthorized("invalid demo session cookie".into()))?
479 .ok_or_else(|| ApiError::Unauthorized("a demo session cookie is required".into()))?;
480 Ok(token_hash(&token))
481}
482
483pub(crate) fn session_cookie_header(
484 headers: &HeaderMap,
485 policy: DemoPolicy,
486) -> Result<Option<HeaderValue>, ApiError> {
487 if matches!(cookie_token(headers), Ok(Some(_))) {
488 return Ok(None);
489 }
490 let token: [u8; 32] = random();
491 let cookie = format!(
492 "{DEMO_COOKIE_NAME}={}; Path=/; Max-Age={}; HttpOnly; SameSite=Strict; Secure",
493 hex_encode(&token),
494 policy.session_ttl_secs
495 );
496 HeaderValue::from_str(&cookie)
497 .map(Some)
498 .map_err(|error| ApiError::Internal(error.into()))
499}
500
501struct DemoIdentity {
502 token_hash: String,
503 persisted: Option<DemoSessionRow>,
504}
505
506async fn identify(state: &AppState, headers: &HeaderMap) -> Result<DemoIdentity, ApiError> {
507 let token_hash = parsed_token_hash(headers)?;
508 let persisted = DemoSessionRow::get_by_token_hash(&state.pool, &token_hash, Utc::now()).await?;
509 Ok(DemoIdentity {
510 token_hash,
511 persisted,
512 })
513}
514
515fn identify_without_lookup(headers: &HeaderMap) -> Result<DemoIdentity, ApiError> {
516 Ok(DemoIdentity {
517 token_hash: parsed_token_hash(headers)?,
518 persisted: None,
519 })
520}
521
522fn owner_id(identity: &DemoIdentity, run_id: i64) -> Result<i64, ApiError> {
523 identity
524 .persisted
525 .as_ref()
526 .map(|session| session.id)
527 .ok_or_else(|| ApiError::NotFound(format!("no demo run with id {run_id}")))
528}
529
530async fn build_owned_run_detail(
531 state: &AppState,
532 run_id: i64,
533 owner_id: i64,
534) -> Result<Option<RunDetailResponse>, ApiError> {
535 let Some(run) = TuneRunRow::get_for_demo_session(&state.pool, run_id, owner_id).await? else {
536 return Ok(None);
537 };
538 let samples = TuneSampleRow::list_for_run(&state.pool, run_id).await?;
539 let results = TuneResultRow::list_for_run(&state.pool, run_id).await?;
540 let writes = TuneWriteRow::list_for_run(&state.pool, run_id).await?;
541 let mv_actuations = TuneMvActuationRow::list_for_run(&state.pool, run_id).await?;
542 let pid_constant_tags = match (
543 &run.tags.proportional_constant,
544 &run.tags.integral_constant,
545 &run.tags.derivative_constant,
546 ) {
547 (Some(proportional), Some(integral), Some(derivative)) => Some(PidConstantTagsResponse {
548 proportional: proportional.clone(),
549 integral: integral.clone(),
550 derivative: derivative.clone(),
551 }),
552 _ => None,
553 };
554 let original_request = serde_json::from_str(&run.request_json).ok();
555 let pid_parameter_labels = PidParameterLabelsResponse::from(&run.template);
556 Ok(Some(RunDetailResponse {
557 id: run.id,
558 tag_name: run.loop_name,
559 notes: run.notes,
560 driver: run.driver,
561 outcome: run.outcome,
562 failure_reason: run.failure_reason,
563 started_at: run.started_at,
564 completed_at: run.completed_at,
565 template_name: run.template.name,
566 template_origin: run.template_origin,
567 allow_uncertain_quality: run.allow_uncertain_quality,
568 config: run.config,
569 effective_tuning: run.effective_tuning,
570 opc_server: run.opc_server,
571 bridge_host: run.bridge_host,
572 pid_constant_tags,
573 pid_parameter_labels,
574 initial_readings: run.initial_readings.map(InitialReadingsResponse::from),
575 timing_metrics: run.timing_metrics,
576 samples: samples.iter().map(SampleResponse::from).collect(),
577 results: results.iter().map(ResultResponse::from).collect(),
578 writes: writes.iter().map(WriteResponse::from).collect(),
579 mv_actuations: mv_actuations
580 .iter()
581 .map(MvActuationResponse::from)
582 .collect(),
583 restore_status: run.restore_status,
584 restore_detail: run.restore_detail,
585 original_request,
586 }))
587}
588
589pub(crate) async fn list_runs(
590 State(state): State<AppState>,
591 headers: HeaderMap,
592 Query(query): Query<RunListQuery>,
593) -> Result<Json<RunListResponse>, ApiError> {
594 ensure_demo(&state)?;
595 let _request_permit = ordinary_request_permit(&state)?;
596 let identity = identify(&state, &headers).await?;
597 let Some(owner_id) = identity.persisted.map(|session| session.id) else {
598 return Ok(Json(RunListResponse {
599 runs: Vec::new(),
600 returned: 0,
601 total: 0,
602 }));
603 };
604 if query.offset.is_some_and(|offset| offset < 0) || query.limit.is_some_and(|limit| limit < 1) {
605 return Err(bad_request(
606 "demo run pagination requires limit >= 1 and offset >= 0",
607 ));
608 }
609 let max_page = i64::from(
610 state.demo_policy.retained_runs_per_visitor + state.demo_policy.max_active_runs_per_visitor,
611 );
612 let pagination = Pagination::new(
613 query.limit.unwrap_or(max_page).min(max_page),
614 query.offset.unwrap_or(0),
615 );
616 let filter = filter_from_query(&query).with_demo_session_id(owner_id);
617 let runs = TuneRunRow::list(&state.pool, &filter, pagination).await?;
618 let total = TuneRunRow::count(&state.pool, &filter).await?;
619 Ok(Json(RunListResponse {
620 returned: runs.len(),
621 runs: runs.iter().map(RunSummaryResponse::from).collect(),
622 total,
623 }))
624}
625
626pub(crate) async fn last_request(
627 State(state): State<AppState>,
628 headers: HeaderMap,
629) -> Result<Json<Option<StartRunRequest>>, ApiError> {
630 ensure_demo(&state)?;
631 let _request_permit = ordinary_request_permit(&state)?;
632 let identity = identify(&state, &headers).await?;
633 let Some(owner_id) = identity.persisted.map(|session| session.id) else {
634 return Ok(Json(None));
635 };
636 let request = TuneRunRow::newest_for_demo_session(&state.pool, owner_id)
637 .await?
638 .and_then(|run| parse_stored_request(run.id, &run.request_json));
639 Ok(Json(request))
640}
641
642pub(crate) async fn get_run(
643 State(state): State<AppState>,
644 headers: HeaderMap,
645 Path(run_id): Path<i64>,
646) -> Result<Json<RunDetailResponse>, ApiError> {
647 ensure_demo(&state)?;
648 let _request_permit = ordinary_request_permit(&state)?;
649 let identity = identify(&state, &headers).await?;
650 let owner_id = owner_id(&identity, run_id)?;
651 build_owned_run_detail(&state, run_id, owner_id)
652 .await?
653 .map(Json)
654 .ok_or_else(|| ApiError::NotFound(format!("no demo run with id {run_id}")))
655}
656
657fn ok_event(event: Event) -> Result<Event, Infallible> {
658 Ok(event)
659}
660
661#[derive(Serialize)]
662struct DemoStreamDone {
663 outcome: TuneOutcome,
664}
665
666pub(crate) async fn stream_run(
667 State(state): State<AppState>,
668 headers: HeaderMap,
669 Path(run_id): Path<i64>,
670) -> Result<impl IntoResponse, ApiError> {
671 ensure_demo(&state)?;
672 let identity = identify(&state, &headers).await?;
673 let owner_id = owner_id(&identity, run_id)?;
674 TuneRunRow::get_for_demo_session(&state.pool, run_id, owner_id)
675 .await?
676 .ok_or_else(|| ApiError::NotFound(format!("no demo run with id {run_id}")))?;
677
678 let global_permit = state.demo_runtime.try_acquire_global_sse().map_err(|_| {
679 global_capacity(
680 "demo SSE capacity is temporarily exhausted",
681 state.demo_policy.sse_lifetime_secs,
682 )
683 })?;
684 let visitor_permit = state
685 .demo_runtime
686 .try_acquire_visitor_sse(&identity.token_hash, state.demo_policy.max_sse_per_visitor)
687 .await
688 .map_err(|_| {
689 too_many(
690 "demo SSE connection limit exceeded for this visitor",
691 state.demo_policy.sse_lifetime_secs,
692 )
693 })?;
694
695 let pool = state.pool.clone();
696 let lifetime = StdDuration::from_secs(state.demo_policy.sse_lifetime_secs);
697 let events = stream! {
698 let _global_permit = global_permit;
699 let _visitor_permit = visitor_permit;
700 if lifetime.is_zero() {
701 yield ok_event(Event::default().event("error").data("demo stream lifetime exceeded"));
702 } else {
703 let deadline = tokio::time::Instant::now() + lifetime;
704 let mut last_tick = -1;
705 let mut sent_initial = false;
706 loop {
707 let run = match tokio::time::timeout_at(
708 deadline,
709 TuneRunRow::get_for_demo_session(&pool, run_id, owner_id),
710 ).await {
711 Ok(Ok(Some(run))) => run,
712 Ok(Ok(None)) => {
713 yield ok_event(Event::default().event("error").data("demo run is no longer available"));
714 break;
715 }
716 Ok(Err(error)) => {
717 tracing::error!(run_id, owner_id, %error, "failed to poll owned demo run");
718 yield ok_event(Event::default().event("error").data("demo stream unavailable"));
719 break;
720 }
721 Err(_) => {
722 yield ok_event(Event::default().event("error").data("demo stream lifetime exceeded"));
723 break;
724 }
725 };
726 if !sent_initial
727 && let Some(initial) = run.initial_readings.clone()
728 {
729 match Event::default()
730 .event("initial")
731 .json_data(InitialReadingsResponse::from(initial))
732 {
733 Ok(event) => {
734 sent_initial = true;
735 yield ok_event(event);
736 }
737 Err(error) => {
738 tracing::error!(run_id, owner_id, %error, "failed to encode owned demo initial readings");
739 yield ok_event(Event::default().event("error").data("demo stream unavailable"));
740 break;
741 }
742 }
743 }
744 match tokio::time::timeout_at(
745 deadline,
746 TuneSampleRow::list_for_run_since(&pool, run_id, last_tick),
747 ).await {
748 Ok(Ok(samples)) => {
749 let mut encoding_failed = false;
750 for sample in &samples {
751 last_tick = sample.tick_index;
752 match Event::default()
753 .event("sample")
754 .json_data(SampleResponse::from(sample))
755 {
756 Ok(event) => yield ok_event(event),
757 Err(error) => {
758 tracing::error!(run_id, owner_id, %error, "failed to encode owned demo sample");
759 yield ok_event(Event::default().event("error").data("demo stream unavailable"));
760 encoding_failed = true;
761 break;
762 }
763 }
764 }
765 if encoding_failed {
766 break;
767 }
768 }
769 Ok(Err(error)) => {
770 tracing::error!(run_id, owner_id, %error, "failed to poll owned demo samples");
771 yield ok_event(Event::default().event("error").data("demo stream unavailable"));
772 break;
773 }
774 Err(_) => {
775 yield ok_event(Event::default().event("error").data("demo stream lifetime exceeded"));
776 break;
777 }
778 }
779 if run.outcome != TuneOutcome::Running {
780 match Event::default()
781 .event("done")
782 .json_data(DemoStreamDone { outcome: run.outcome })
783 {
784 Ok(event) => yield ok_event(event),
785 Err(error) => {
786 tracing::error!(run_id, owner_id, %error, "failed to encode owned demo terminal event");
787 yield ok_event(Event::default().event("error").data("demo stream unavailable"));
788 }
789 }
790 break;
791 }
792 if tokio::time::timeout_at(deadline, tokio::time::sleep(SSE_POLL_INTERVAL))
793 .await
794 .is_err()
795 {
796 yield ok_event(Event::default().event("error").data("demo stream lifetime exceeded"));
797 break;
798 }
799 }
800 }
801 };
802 Ok(Sse::new(events).keep_alive(KeepAlive::default()))
803}
804
805pub(crate) async fn cancel_run(
806 State(state): State<AppState>,
807 headers: HeaderMap,
808 Path(run_id): Path<i64>,
809) -> Result<StatusCode, ApiError> {
810 ensure_demo(&state)?;
811 let _request_permit = ordinary_request_permit(&state)?;
812 let identity = identify(&state, &headers).await?;
813 let owner_id = owner_id(&identity, run_id)?;
814 let run = TuneRunRow::get_for_demo_session(&state.pool, run_id, owner_id)
815 .await?
816 .ok_or_else(|| ApiError::NotFound(format!("no demo run with id {run_id}")))?;
817 let _cancelled_or_already_terminal =
818 run.outcome != TuneOutcome::Running || state.active_run.cancel(run_id).await;
819 Ok(StatusCode::NO_CONTENT)
820}
821
822pub(crate) async fn delete_run(
823 State(state): State<AppState>,
824 headers: HeaderMap,
825 Path(run_id): Path<i64>,
826) -> Result<StatusCode, ApiError> {
827 ensure_demo(&state)?;
828 let _request_permit = ordinary_request_permit(&state)?;
829 let identity = identify(&state, &headers).await?;
830 let owner_id = owner_id(&identity, run_id)?;
831 let run = TuneRunRow::get_for_demo_session(&state.pool, run_id, owner_id)
832 .await?
833 .ok_or_else(|| ApiError::NotFound(format!("no demo run with id {run_id}")))?;
834 if run.outcome == TuneOutcome::Running {
835 return Err(ApiError::Conflict(
836 "cancel the demo run before deleting it".into(),
837 ));
838 }
839 delete_existing_demo_run(&state, run_id, owner_id).await?;
840 Ok(StatusCode::NO_CONTENT)
841}
842
843async fn delete_existing_demo_run(
844 state: &AppState,
845 run_id: i64,
846 owner_id: i64,
847) -> Result<(), ApiError> {
848 if !TuneRunRow::delete_for_demo_session(&state.pool, run_id, owner_id).await? {
849 return Err(ApiError::NotFound(format!("no demo run with id {run_id}")));
850 }
851 Ok(())
852}
853
854pub(crate) async fn export_run(
855 State(state): State<AppState>,
856 headers: HeaderMap,
857 Path(run_id): Path<i64>,
858 Query(query): Query<RunExportQuery>,
859) -> Result<Response, ApiError> {
860 ensure_demo(&state)?;
861 let _request_permit = ordinary_request_permit(&state)?;
862 let identity = identify(&state, &headers).await?;
863 let owner_id = owner_id(&identity, run_id)?;
864 TuneRunRow::get_for_demo_session(&state.pool, run_id, owner_id)
865 .await?
866 .ok_or_else(|| ApiError::NotFound(format!("no demo run with id {run_id}")))?;
867 let samples = TuneSampleRow::list_for_run(&state.pool, run_id).await?;
868 if samples.is_empty() {
869 return Err(ApiError::NotFound(format!(
870 "demo run {run_id} has no samples"
871 )));
872 }
873 let format = query.format.unwrap_or(RunExportFormat::Csv);
874 let bytes = bhtune_cli::commands::export::samples_to_bytes(&samples, format.into())?;
875 let (content_type, extension) = match format {
876 RunExportFormat::Csv => ("text/csv", "csv"),
877 RunExportFormat::Json => ("application/json", "json"),
878 };
879 let mut response = bytes.into_response();
880 response
881 .headers_mut()
882 .insert(header::CONTENT_TYPE, HeaderValue::from_static(content_type));
883 response.headers_mut().insert(
884 header::CONTENT_DISPOSITION,
885 HeaderValue::from_str(&format!(
886 "attachment; filename=\"demo-run-{run_id}.{extension}\""
887 ))
888 .map_err(|error| ApiError::Internal(error.into()))?,
889 );
890 Ok(response)
891}
892
893pub(crate) async fn list_templates(
894 State(state): State<AppState>,
895) -> Result<Json<Vec<TemplateResponse>>, ApiError> {
896 ensure_demo(&state)?;
897 let _request_permit = ordinary_request_permit(&state)?;
898 let rows = bhtune_db::models::DcsTemplateRow::list(&state.pool).await?;
899 Ok(Json(
900 rows.into_iter()
901 .filter(|row| row.origin == TemplateOrigin::Builtin)
902 .map(TemplateResponse::from)
903 .collect(),
904 ))
905}
906
907pub(crate) async fn get_template(
908 State(state): State<AppState>,
909 Path(name): Path<String>,
910) -> Result<Json<TemplateResponse>, ApiError> {
911 ensure_demo(&state)?;
912 let _request_permit = ordinary_request_permit(&state)?;
913 let row = bhtune_db::models::DcsTemplateRow::get_by_name(&state.pool, &name)
914 .await?
915 .filter(|row| row.origin == TemplateOrigin::Builtin)
916 .ok_or_else(|| ApiError::NotFound(format!("no built-in template named '{name}'")))?;
917 Ok(Json(row.into()))
918}
919
920async fn trim_owned_history(state: &AppState, owner_id: i64) {
921 if let Err(error) = TuneRunRow::prune_terminal_for_demo_session(
922 &state.pool,
923 owner_id,
924 state.demo_policy.retained_runs_per_visitor,
925 )
926 .await
927 {
928 tracing::warn!(owner_id, %error, "failed to prune excess demo history");
929 }
930}
931
932async fn ensure_global_run_capacity(state: &AppState) -> Result<(), ApiError> {
933 let current = TuneRunRow::count_demo_owned(&state.pool).await?;
934 if current >= i64::from(state.demo_policy.max_tune_run_rows_global) {
935 return Err(ApiError::GlobalCapacity {
936 message: "demo history capacity is temporarily exhausted".into(),
937 retry_after_secs: state.demo_policy.cleanup_interval_secs,
938 });
939 }
940 Ok(())
941}
942
943fn prepare_error(error: anyhow::Error, state: &AppState) -> ApiError {
944 if error.chain().any(|cause| {
945 cause
946 .to_string()
947 .contains("demo tune run row limit reached")
948 }) {
949 ApiError::GlobalCapacity {
950 message: "demo history capacity is temporarily exhausted".into(),
951 retry_after_secs: state.demo_policy.cleanup_interval_secs,
952 }
953 } else {
954 ApiError::Internal(error)
955 }
956}
957
958pub(crate) async fn start_run(
959 State(state): State<AppState>,
960 headers: HeaderMap,
961 peer: PeerAddress,
962 Json(value): Json<serde_json::Value>,
963) -> Result<(StatusCode, Json<RunDetailResponse>), ApiError> {
964 ensure_demo(&state)?;
965 let request = parse_demo_request(value)?;
970 let _request_permit = ordinary_request_permit(&state)?;
971
972 let identity = identify_without_lookup(&headers)?;
973 let client_ip = client_ip(&headers, peer, state.trusted_proxy.as_deref());
974 let _start_admission = state.demo_runtime.lock_start_admission().await;
975 let global_permit = state.demo_runtime.try_acquire_global_run().map_err(|_| {
976 global_capacity(
977 "demo run capacity is temporarily exhausted",
978 state.demo_policy.run_timeout_secs,
979 )
980 })?;
981 let visitor_permit = state
982 .demo_runtime
983 .try_acquire_visitor_run(
984 &identity.token_hash,
985 state.demo_policy.max_active_runs_per_visitor,
986 )
987 .await
988 .map_err(|_| {
989 too_many(
990 "a demo tune is already active for this visitor",
991 state.demo_policy.run_timeout_secs,
992 )
993 })?;
994
995 let now = Utc::now();
996 let accepted_start = state
997 .demo_runtime
998 .reserve_accepted_start(&identity.token_hash, &client_ip, now, state.demo_policy)
999 .await
1000 .map_err(quota_error)?;
1001 let preparation =
1002 async {
1003 DemoSessionRow::cleanup_expired(&state.pool, now).await?;
1006 bhtune_db::models::TuneRunRow::prune_terminal_demo_owned(
1007 &state.pool,
1008 state.demo_policy.retained_runs_per_visitor,
1009 )
1010 .await?;
1011 ensure_global_run_capacity(&state).await?;
1012 let session =
1013 match DemoSessionRow::get_by_token_hash(&state.pool, &identity.token_hash, now)
1014 .await?
1015 {
1016 Some(session) => session,
1017 None => {
1018 DemoSessionRow::get_or_create(
1019 &state.pool,
1020 &identity.token_hash,
1021 now,
1022 now + Duration::seconds(state.demo_policy.session_ttl_secs as i64),
1023 )
1024 .await?
1025 }
1026 };
1027 trim_owned_history(&state, session.id).await;
1028 let accepted_runs = TuneRunRow::count_for_demo_session(&state.pool, session.id).await?;
1029 if accepted_runs >= i64::from(state.demo_policy.max_runs_per_session) {
1030 return Err(too_many(
1031 "maximum demo runs for this session has been reached",
1032 state.demo_policy.accepted_start_window_secs,
1033 ));
1034 }
1035
1036 let args = request.into_tune_args()?;
1037 let mut config = state.config_snapshot()?;
1038 config.tuning.mrft_delay_secs = Some(0);
1039 config.tuning.poll_interval_ms = Some(state.demo_policy.poll_interval_ms);
1040 config.tuning.timeout_secs = Some(state.demo_policy.run_timeout_secs);
1041 let prepared = prepare_owned(&state.pool, args, &config, session.id)
1042 .await
1043 .map_err(|error| prepare_error(error, &state))?;
1044 Ok::<_, ApiError>((session.id, prepared))
1045 }
1046 .await;
1047 let (session_id, prepared) = match preparation {
1048 Ok(prepared) => prepared,
1049 Err(error) => {
1050 state
1051 .demo_runtime
1052 .release_accepted_start(accepted_start)
1053 .await;
1054 return Err(error);
1055 }
1056 };
1057 let run_id = prepared.run_id();
1058
1059 let (mut ctrl_c, cancel_handle) = bhtune_cli::cancel::CtrlC::manual();
1060 let pool = state.pool.clone();
1061 let state_for_task = state.clone();
1062 let task = async move {
1063 let _global_permit = global_permit;
1064 let _visitor_permit = visitor_permit;
1065 let _ = drive(&pool, prepared, &mut ctrl_c).await;
1066 trim_owned_history(&state_for_task, session_id).await;
1067 };
1068 if state
1069 .active_run
1070 .start(run_id, cancel_handle, task)
1071 .await
1072 .is_err()
1073 {
1074 state
1075 .demo_runtime
1076 .release_accepted_start(accepted_start)
1077 .await;
1078 discard_unscheduled_demo_run(&state, run_id, session_id).await;
1079 return Err(global_capacity(
1080 "demo run could not be scheduled; retry shortly",
1081 state.demo_policy.run_timeout_secs,
1082 ));
1083 }
1084
1085 let detail = build_owned_run_detail(&state, run_id, session_id)
1086 .await?
1087 .expect("prepare_owned inserted this owner-scoped demo run");
1088 Ok((StatusCode::CREATED, Json(detail)))
1089}
1090
1091async fn discard_unscheduled_demo_run(state: &AppState, run_id: i64, session_id: i64) {
1092 if let Err(error) = TuneRunRow::delete_for_demo_session(&state.pool, run_id, session_id).await {
1093 tracing::error!(run_id, session_id, %error, "failed to discard unscheduled demo run");
1094 }
1095}
1096
1097async fn api_not_found() -> ApiError {
1098 ApiError::NotFound("API route is not available in Demo mode".into())
1099}
1100
1101pub fn router(policy: DemoPolicy) -> Router<AppState> {
1102 let ordinary = Router::new()
1103 .route("/api/templates", get(list_templates))
1104 .route("/api/templates/{name}", get(get_template))
1105 .route("/api/runs", get(list_runs).post(start_run))
1106 .route("/api/runs/last-request", get(last_request))
1107 .route("/api/runs/{id}", get(get_run).delete(delete_run))
1108 .route("/api/runs/{id}/cancel", axum::routing::post(cancel_run))
1109 .route("/api/runs/{id}/export", get(export_run))
1110 .route("/api", any(api_not_found))
1111 .route("/api/{*path}", any(api_not_found))
1112 .layer(TimeoutLayer::with_status_code(
1113 StatusCode::REQUEST_TIMEOUT,
1114 StdDuration::from_secs(policy.ordinary_request_timeout_secs),
1115 ));
1116 let streaming = Router::new().route("/api/runs/{id}/stream", get(stream_run));
1117 ordinary.merge(streaming).layer(RequestBodyLimitLayer::new(
1118 policy.max_json_body_bytes as usize,
1119 ))
1120}
1121
1122#[cfg(test)]
1123mod tests {
1124 use super::*;
1125 use axum::body::{Body, to_bytes};
1126 use axum::extract::ConnectInfo;
1127 use axum::http::Request;
1128 use bhtune_cli::config::DEMO_TEMPLATE_NAME;
1129 use bhtune_core::{ControllerType, LoopConfig, LoopTags, ProcessType, Tick};
1130 use bhtune_db::models::{DcsTemplateRow, SampleQuality, TemplateOrigin};
1131 use tower::ServiceExt;
1132
1133 const DEMO_ORIGIN: &str = "https://demo.test";
1134
1135 fn raw_token(seed: &str) -> String {
1136 seed.repeat(64 / seed.len())
1137 }
1138
1139 fn cookie(seed: &str) -> String {
1140 format!("{DEMO_COOKIE_NAME}={}", raw_token(seed))
1141 }
1142
1143 fn with_peer(mut request: Request<Body>, peer: &str) -> Request<Body> {
1144 request
1145 .extensions_mut()
1146 .insert(ConnectInfo::<SocketAddr>(peer.parse().unwrap()));
1147 request
1148 }
1149
1150 async fn body_text(response: axum::response::Response) -> String {
1151 String::from_utf8(
1152 to_bytes(response.into_body(), usize::MAX)
1153 .await
1154 .unwrap()
1155 .to_vec(),
1156 )
1157 .unwrap()
1158 }
1159
1160 async fn body_json(response: axum::response::Response) -> serde_json::Value {
1161 serde_json::from_str(&body_text(response).await).unwrap()
1162 }
1163
1164 fn valid_request(
1165 process_type: ProcessType,
1166 controller_type: ControllerType,
1167 ) -> serde_json::Value {
1168 serde_json::json!({
1169 "tagname": DEMO_TAG_NAME,
1170 "template": DEMO_TEMPLATE_NAME,
1171 "process_type": process_type,
1172 "controller_type": controller_type,
1173 "relay_amp": DEMO_RELAY_AMP_DEFAULT,
1174 "cycles_skip": DEMO_CYCLES_SKIP_DEFAULT,
1175 "cycles_count": DEMO_CYCLES_COUNT_DEFAULT,
1176 "noise_protection_secs": DEMO_NOISE_PROTECTION_SECS_DEFAULT,
1177 "driver": "simulator",
1178 "sim_gain": DEMO_SIM_GAIN_DEFAULT,
1179 "sim_tau": DEMO_SIM_TAU_DEFAULT,
1180 "sim_dead_time": DEMO_SIM_DEAD_TIME_DEFAULT,
1181 "sim_noise": DEMO_SIM_NOISE_DEFAULT,
1182 "sim_seed": DEMO_SIM_SEED_DEFAULT,
1183 "sim_initial_pv": DEMO_SIM_INITIAL_VALUE_DEFAULT,
1184 "sim_initial_mv": DEMO_SIM_INITIAL_VALUE_DEFAULT,
1185 "pv_range_high": DEMO_RANGE_HIGH,
1186 "pv_range_low": DEMO_RANGE_LOW,
1187 "mv_range_high": DEMO_RANGE_HIGH,
1188 "mv_range_low": DEMO_RANGE_LOW,
1189 "direction": ControllerDirection::Reverse,
1190 })
1191 }
1192
1193 async fn seed_owned_run(
1194 state: &AppState,
1195 raw_token: &str,
1196 now: chrono::DateTime<Utc>,
1197 ) -> (i64, i64) {
1198 let session = DemoSessionRow::get_or_create(
1199 &state.pool,
1200 &token_hash(raw_token),
1201 now,
1202 now + Duration::hours(1),
1203 )
1204 .await
1205 .unwrap();
1206 let row = DcsTemplateRow::get_by_name(&state.pool, DEMO_TEMPLATE_NAME)
1207 .await
1208 .unwrap()
1209 .unwrap();
1210 let config = LoopConfig {
1211 process_type: ProcessType::Flow,
1212 controller_type: ControllerType::Pi,
1213 relay_amp_percent: 5.0,
1214 num_cycles_skip: 1,
1215 num_cycles_count: 2,
1216 noise_protection_secs: 0,
1217 mrft_delay_secs: 0,
1218 };
1219 let tags = LoopTags::derive_from_pv_tag(DEMO_TAG_NAME, &row.template);
1220 let run = TuneRunRow::start_owned(
1221 &state.pool,
1222 session.id,
1223 DEMO_TAG_NAME,
1224 config,
1225 row.origin,
1226 &row.template,
1227 &tags,
1228 now,
1229 )
1230 .await
1231 .unwrap();
1232 (session.id, run.id)
1233 }
1234
1235 async fn insert_one_sample(state: &AppState, run_id: i64, now: chrono::DateTime<Utc>) {
1236 TuneSampleRow::insert(
1237 &state.pool,
1238 run_id,
1239 0,
1240 Tick {
1241 time: now,
1242 pv: 50.0,
1243 },
1244 bhtune_core::MrftState {
1245 hysteresis: 1.0,
1246 mv_value_current: 45.0,
1247 mv_sign_next_step: 1,
1248 counter_all_switches: 0,
1249 cycles_completed: 0,
1250 cycles_remaining: 1,
1251 },
1252 SampleQuality::Good,
1253 )
1254 .await
1255 .unwrap();
1256 }
1257
1258 async fn wait_until_inactive(state: &AppState, run_id: i64) {
1259 tokio::time::timeout(StdDuration::from_secs(2), async {
1260 while state.active_run.active_run_ids().await.contains(&run_id) {
1261 tokio::task::yield_now().await;
1262 }
1263 })
1264 .await
1265 .expect("demo tune task should finish or honour cancellation");
1266 }
1267
1268 #[test]
1269 fn trusted_proxy_and_dedicated_client_ip_rules_fail_closed_and_bucket_ipv6() {
1270 let peer = "10.0.0.4:443".parse().unwrap();
1271 assert!(trusted_proxy(peer, Some("10.0.0.4")));
1272 assert!(!trusted_proxy(peer, Some("10.0.0.5")));
1273 assert!(trusted_proxy(peer, Some("10.0.0.0/24")));
1274 assert!(trusted_proxy(peer, Some("0.0.0.0/0")));
1275 assert!(!trusted_proxy(peer, Some("10.0.1.0/24")));
1276 assert!(!trusted_proxy(peer, Some("not-a-proxy")));
1277 assert!(!trusted_proxy(peer, Some("10.0.0.0/not-bits")));
1278 assert!(!trusted_proxy(peer, Some("2001:db8::/32")));
1279 assert!(trusted_proxy(
1280 "[2001:db8::4]:443".parse().unwrap(),
1281 Some("2001:db8::/32")
1282 ));
1283 assert!(trusted_proxy(
1284 "[2001:db8::4]:443".parse().unwrap(),
1285 Some("::/0")
1286 ));
1287 let mut headers = HeaderMap::new();
1288 headers.insert(FORWARDED_CLIENT_IP_HEADER, "198.51.100.7".parse().unwrap());
1289 assert_eq!(
1290 client_ip(&headers, PeerAddress(peer), Some("10.0.0.0/24")),
1291 "198.51.100.7"
1292 );
1293 assert_eq!(client_ip(&headers, PeerAddress(peer), None), "10.0.0.4");
1294 headers.insert(
1295 FORWARDED_CLIENT_IP_HEADER,
1296 "198.51.100.7, 203.0.113.9".parse().unwrap(),
1297 );
1298 assert_eq!(
1299 client_ip(&headers, PeerAddress(peer), Some("10.0.0.0/24")),
1300 "10.0.0.4"
1301 );
1302 headers.insert(FORWARDED_CLIENT_IP_HEADER, "not-an-ip".parse().unwrap());
1303 assert_eq!(
1304 client_ip(&headers, PeerAddress(peer), Some("10.0.0.0/24")),
1305 "10.0.0.4"
1306 );
1307 headers.clear();
1308 headers.append(FORWARDED_CLIENT_IP_HEADER, "198.51.100.7".parse().unwrap());
1309 headers.append(FORWARDED_CLIENT_IP_HEADER, "203.0.113.9".parse().unwrap());
1310 assert_eq!(
1311 client_ip(&headers, PeerAddress(peer), Some("10.0.0.0/24")),
1312 "10.0.0.4"
1313 );
1314 headers.clear();
1315 headers.insert("x-forwarded-for", "203.0.113.9".parse().unwrap());
1316 assert_eq!(
1317 client_ip(&headers, PeerAddress(peer), Some("10.0.0.0/24")),
1318 "10.0.0.4"
1319 );
1320
1321 let ipv6_peer = PeerAddress("[2001:db8:1:2:1234::1]:443".parse().unwrap());
1322 assert_eq!(
1323 client_ip(&HeaderMap::new(), ipv6_peer, None),
1324 "2001:db8:1:2::/64"
1325 );
1326 let mut forwarded_ipv6 = HeaderMap::new();
1327 forwarded_ipv6.insert(
1328 FORWARDED_CLIENT_IP_HEADER,
1329 "2001:db8:abcd:12:ffff::1".parse().unwrap(),
1330 );
1331 assert_eq!(
1332 client_ip(&forwarded_ipv6, PeerAddress(peer), Some("10.0.0.0/24")),
1333 "2001:db8:abcd:12::/64"
1334 );
1335 }
1336
1337 #[test]
1338 fn cookie_parser_rejects_ambiguous_or_malformed_values() {
1339 let mut headers = HeaderMap::new();
1340 headers.append(
1341 header::COOKIE,
1342 format!("ignored; {DEMO_COOKIE_NAME}={}", raw_token("ab"))
1343 .parse()
1344 .unwrap(),
1345 );
1346 assert_eq!(cookie_token(&headers).unwrap().unwrap(), raw_token("ab"));
1347 headers.append(
1348 header::COOKIE,
1349 format!("{DEMO_COOKIE_NAME}={}", raw_token("cd"))
1350 .parse()
1351 .unwrap(),
1352 );
1353 assert!(cookie_token(&headers).is_err());
1354
1355 let mut malformed = HeaderMap::new();
1356 malformed.insert(
1357 header::COOKIE,
1358 format!("{DEMO_COOKIE_NAME}=UPPERCASE").parse().unwrap(),
1359 );
1360 assert!(cookie_token(&malformed).is_err());
1361 malformed.insert(header::COOKIE, HeaderValue::from_bytes(&[0xff]).unwrap());
1362 assert!(cookie_token(&malformed).is_err());
1363 }
1364
1365 #[test]
1366 fn session_cookie_header_reuses_valid_cookies_and_replaces_missing_or_invalid_ones() {
1367 let policy = DemoPolicy::default();
1368 let mut headers = HeaderMap::new();
1369 let generated = session_cookie_header(&headers, policy).unwrap().unwrap();
1370 let generated = generated.to_str().unwrap();
1371 assert!(generated.starts_with(&format!("{DEMO_COOKIE_NAME}=")));
1372 assert!(generated.contains("HttpOnly"));
1373 assert!(generated.contains("SameSite=Strict"));
1374
1375 headers.insert(header::COOKIE, cookie("ab").parse().unwrap());
1376 assert!(session_cookie_header(&headers, policy).unwrap().is_none());
1377
1378 headers.insert(
1379 header::COOKIE,
1380 format!("{DEMO_COOKIE_NAME}=bad").parse().unwrap(),
1381 );
1382 assert!(session_cookie_header(&headers, policy).unwrap().is_some());
1383 }
1384
1385 #[tokio::test]
1386 async fn demo_helpers_map_full_mode_and_capacity_errors() {
1387 let state = crate::test_support::in_memory_state().await;
1388 assert!(matches!(ensure_demo(&state), Err(ApiError::NotFound(_))));
1389
1390 let base = crate::test_support::in_memory_state().await;
1391 let state = AppState::for_mode(
1392 base.pool,
1393 base.config_store,
1394 ServerMode::Demo,
1395 DemoPolicy {
1396 ordinary_request_concurrency: 1,
1397 ..DemoPolicy::default()
1398 },
1399 );
1400 let ordinary = ordinary_request_permit(&state).unwrap();
1401 let second = ordinary_request_permit(&state);
1402 assert!(matches!(second, Err(ApiError::GlobalCapacity { .. })));
1403 drop(ordinary);
1404
1405 assert!(matches!(
1406 prepare_error(
1407 anyhow::anyhow!("outer").context("demo tune run row limit reached"),
1408 &state
1409 ),
1410 ApiError::GlobalCapacity { .. }
1411 ));
1412 assert!(matches!(
1413 prepare_error(anyhow::anyhow!("ordinary failure"), &state),
1414 ApiError::Internal(_)
1415 ));
1416
1417 assert!(matches!(
1418 delete_existing_demo_run(&state, 12345, 67890).await,
1419 Err(ApiError::NotFound(_))
1420 ));
1421
1422 let closed = crate::test_support::in_memory_state().await;
1423 closed.pool.close().await;
1424 discard_unscheduled_demo_run(&closed, 1, 1).await;
1425 }
1426
1427 #[tokio::test]
1428 async fn shared_routes_require_a_cookie() {
1429 let mut state = crate::test_support::in_memory_state().await;
1430 state.mode = ServerMode::Demo;
1431 let app = crate::build_router(state.clone());
1432 let shared = app
1433 .oneshot(Request::get("/api/runs").body(Body::empty()).unwrap())
1434 .await
1435 .unwrap();
1436 assert_eq!(shared.status(), StatusCode::UNAUTHORIZED);
1437 }
1438
1439 #[tokio::test]
1440 async fn demo_mode_replaces_unavailable_full_api_routes_with_404() {
1441 let mut state = crate::test_support::in_memory_state().await;
1442 state.mode = ServerMode::Demo;
1443 let response = crate::build_router(state)
1444 .oneshot(
1445 Request::get("/api/openapi.json")
1446 .body(Body::empty())
1447 .unwrap(),
1448 )
1449 .await
1450 .unwrap();
1451 assert_eq!(response.status(), StatusCode::NOT_FOUND);
1452 assert!(body_text(response).await.contains("Demo mode"));
1453 }
1454
1455 #[tokio::test]
1456 async fn start_requires_peer_connect_info() {
1457 let mut state = crate::test_support::in_memory_state().await;
1458 state.mode = ServerMode::Demo;
1459 state.allowed_origin = Some(DEMO_ORIGIN.into());
1460 let response = crate::build_router(state)
1461 .oneshot(
1462 Request::post("/api/runs")
1463 .header(header::ORIGIN, DEMO_ORIGIN)
1464 .header(header::COOKIE, cookie("ab"))
1465 .header(header::CONTENT_TYPE, "application/json")
1466 .body(Body::from(
1467 valid_request(ProcessType::Flow, ControllerType::Pi).to_string(),
1468 ))
1469 .unwrap(),
1470 )
1471 .await
1472 .unwrap();
1473 assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR);
1474 assert!(body_text(response).await.contains("internal server error"));
1475 }
1476
1477 #[tokio::test]
1478 async fn forbidden_fields_are_rejected_by_presence_before_database_io() {
1479 let mut state = crate::test_support::in_memory_state().await;
1480 state.mode = ServerMode::Demo;
1481 state.allowed_origin = Some(DEMO_ORIGIN.into());
1482 state.pool.close().await;
1483 let mut body = valid_request(ProcessType::Flow, ControllerType::Pi);
1484 body["notes"] = serde_json::Value::Null;
1485 let request = with_peer(
1486 Request::post("/api/runs")
1487 .header(header::ORIGIN, DEMO_ORIGIN)
1488 .header(header::COOKIE, cookie("ab"))
1489 .header(header::CONTENT_TYPE, "application/json")
1490 .body(Body::from(body.to_string()))
1491 .unwrap(),
1492 "127.0.0.1:12345",
1493 );
1494 let response = crate::build_router(state).oneshot(request).await.unwrap();
1495 assert_eq!(response.status(), StatusCode::BAD_REQUEST);
1496 let body = to_bytes(response.into_body(), usize::MAX).await.unwrap();
1497 assert!(String::from_utf8_lossy(&body).contains("'notes'"));
1498 }
1499
1500 #[test]
1501 fn every_process_type_uses_authoritative_controller_compatibility() {
1502 for process_type in ProcessType::ALL {
1503 for controller_type in ControllerType::ALL {
1504 assert_eq!(
1505 parse_demo_request(valid_request(process_type, controller_type)).is_ok(),
1506 controller_type.is_allowed_for(process_type),
1507 "{process_type:?}/{controller_type:?}"
1508 );
1509 }
1510 }
1511 }
1512
1513 #[test]
1514 fn malformed_and_non_demo_start_requests_are_rejected_before_defaults() {
1515 assert!(parse_demo_request(serde_json::Value::Null).is_err());
1516
1517 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1518 request.as_object_mut().unwrap().remove("tagname");
1519 assert!(parse_demo_request(request).is_err());
1520
1521 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1522 request["relay_amp"] = serde_json::json!("large");
1523 assert!(parse_demo_request(request).is_err());
1524
1525 for (field, value) in [
1526 ("driver", serde_json::json!("opcda")),
1527 ("template", serde_json::json!("Private custom template")),
1528 ("tagname", serde_json::json!("Other tag")),
1529 ("sim_seed", serde_json::json!(DEMO_SIM_SEED_MAX + 1)),
1530 ("cycles_skip", serde_json::Value::Null),
1531 ("cycles_count", serde_json::Value::Null),
1532 ("noise_protection_secs", serde_json::Value::Null),
1533 ("pv_range_high", serde_json::Value::Null),
1534 ("mv_range_low", serde_json::Value::Null),
1535 ] {
1536 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1537 request[field] = value;
1538 assert!(parse_demo_request(request).is_err(), "{field}");
1539 }
1540
1541 for field in ["sim_tau", "sim_dead_time", "sim_initial_mv"] {
1542 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1543 request[field] = serde_json::json!(f32::NAN);
1544 assert!(parse_demo_request(request).is_err(), "{field}");
1545 }
1546
1547 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1548 request["mv_range_low"] = serde_json::json!(10.0);
1549 request["mv_range_high"] = serde_json::json!(9.0);
1550 assert!(parse_demo_request(request).is_err());
1551 }
1552
1553 #[test]
1554 fn omitted_demo_values_use_the_approved_defaults() {
1555 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1556 for field in [
1557 "relay_amp",
1558 "cycles_skip",
1559 "cycles_count",
1560 "noise_protection_secs",
1561 "sim_gain",
1562 "sim_tau",
1563 "sim_dead_time",
1564 "sim_noise",
1565 "sim_seed",
1566 "sim_initial_pv",
1567 "sim_initial_mv",
1568 "pv_range_high",
1569 "pv_range_low",
1570 "mv_range_high",
1571 "mv_range_low",
1572 "direction",
1573 ] {
1574 request.as_object_mut().unwrap().remove(field);
1575 }
1576
1577 let parsed = parse_demo_request(request).unwrap();
1578 assert_eq!(parsed.relay_amp, 10.0);
1579 assert_eq!(parsed.cycles_skip, Some(1));
1580 assert_eq!(parsed.cycles_count, Some(2));
1581 assert_eq!(parsed.noise_protection_secs, Some(0));
1582 assert_eq!(parsed.sim_gain, 1.0);
1583 assert_eq!(parsed.sim_tau, 0.5);
1584 assert_eq!(parsed.sim_dead_time, 1.0);
1585 assert_eq!(parsed.sim_noise, 0.0);
1586 assert_eq!(parsed.sim_seed, 0);
1587 assert_eq!(parsed.sim_initial_pv, 50.0);
1588 assert_eq!(parsed.sim_initial_mv, 50.0);
1589 assert_eq!(parsed.pv_range_low, Some(0.0));
1590 assert_eq!(parsed.pv_range_high, Some(100.0));
1591 assert_eq!(parsed.mv_range_low, Some(0.0));
1592 assert_eq!(parsed.mv_range_high, Some(100.0));
1593 assert_eq!(parsed.direction, Some(ControllerDirection::Reverse));
1594 }
1595
1596 #[test]
1597 fn demo_bounds_cover_test_parameters_ranges_initials_and_noise() {
1598 let mut lower_bounds = valid_request(ProcessType::Flow, ControllerType::Pi);
1599 lower_bounds["relay_amp"] = serde_json::json!(DEMO_RELAY_AMP_MIN);
1600 lower_bounds["cycles_skip"] = serde_json::json!(DEMO_CYCLES_SKIP_MIN);
1601 lower_bounds["cycles_count"] = serde_json::json!(DEMO_CYCLES_COUNT_MIN);
1602 lower_bounds["noise_protection_secs"] = serde_json::json!(DEMO_NOISE_PROTECTION_SECS_MIN);
1603 lower_bounds["pv_range_low"] = serde_json::json!(-1_000.0);
1604 lower_bounds["pv_range_high"] = serde_json::json!(-999.0);
1605 lower_bounds["mv_range_low"] = serde_json::json!(999.0);
1606 lower_bounds["mv_range_high"] = serde_json::json!(1_000.0);
1607 lower_bounds["sim_initial_pv"] = serde_json::json!(-999.5);
1608 lower_bounds["sim_initial_mv"] = serde_json::json!(999.5);
1609 lower_bounds["sim_noise"] = serde_json::json!(0.05);
1610 assert!(parse_demo_request(lower_bounds).is_ok());
1611
1612 let mut upper_bounds = valid_request(ProcessType::Flow, ControllerType::Pi);
1613 upper_bounds["relay_amp"] = serde_json::json!(DEMO_RELAY_AMP_MAX);
1614 upper_bounds["cycles_skip"] = serde_json::json!(DEMO_CYCLES_SKIP_MAX);
1615 upper_bounds["cycles_count"] = serde_json::json!(DEMO_CYCLES_COUNT_MAX);
1616 upper_bounds["noise_protection_secs"] = serde_json::json!(DEMO_NOISE_PROTECTION_SECS_MAX);
1617 upper_bounds["pv_range_low"] = serde_json::json!(-500.0);
1618 upper_bounds["pv_range_high"] = serde_json::json!(500.0);
1619 upper_bounds["sim_initial_pv"] = serde_json::json!(500.0);
1620 upper_bounds["sim_noise"] = serde_json::json!(50.0);
1621 assert!(parse_demo_request(upper_bounds).is_ok());
1622
1623 for (field, value) in [
1624 ("relay_amp", serde_json::json!(20.1)),
1625 ("cycles_skip", serde_json::json!(3)),
1626 ("cycles_count", serde_json::json!(4)),
1627 ("noise_protection_secs", serde_json::json!(4)),
1628 ("pv_range_low", serde_json::json!(-1000.1)),
1629 ("sim_initial_pv", serde_json::json!(100.1)),
1630 ("sim_noise", serde_json::json!(5.1)),
1631 ] {
1632 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1633 request[field] = value;
1634 assert!(parse_demo_request(request).is_err(), "{field}");
1635 }
1636
1637 for (low, high) in [(0.0, 0.5), (0.0, 1_000.1), (10.0, 9.0)] {
1638 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1639 request["pv_range_low"] = serde_json::json!(low);
1640 request["pv_range_high"] = serde_json::json!(high);
1641 request["sim_initial_pv"] = serde_json::json!(low);
1642 assert!(parse_demo_request(request).is_err(), "{low}..{high}");
1643 }
1644 }
1645
1646 #[test]
1647 fn gain_magnitude_and_direction_must_form_negative_feedback() {
1648 for (gain, direction) in [
1649 (1.0, ControllerDirection::Reverse),
1650 (-1.0, ControllerDirection::Direct),
1651 (DEMO_SIM_GAIN_MAX, ControllerDirection::Reverse),
1652 (-DEMO_SIM_GAIN_MAX, ControllerDirection::Direct),
1653 ] {
1654 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1655 request["sim_gain"] = serde_json::json!(gain);
1656 request["direction"] = serde_json::json!(direction);
1657 assert!(parse_demo_request(request).is_ok());
1658 }
1659 for (gain, direction) in [
1660 (0.0, ControllerDirection::Reverse),
1661 (DEMO_SIM_GAIN_ABS_MIN / 2.0, ControllerDirection::Reverse),
1662 (1.0, ControllerDirection::Direct),
1663 (-1.0, ControllerDirection::Reverse),
1664 ] {
1665 let mut request = valid_request(ProcessType::Flow, ControllerType::Pi);
1666 request["sim_gain"] = serde_json::json!(gain);
1667 request["direction"] = serde_json::json!(direction);
1668 assert!(parse_demo_request(request).is_err());
1669 }
1670 }
1671
1672 #[tokio::test]
1673 async fn start_rejects_when_session_history_reaches_the_configured_limit() {
1674 let mut state = crate::test_support::in_memory_state().await;
1675 state.mode = ServerMode::Demo;
1676 state.demo_policy.max_runs_per_session = 1;
1677 let token = raw_token("ab");
1678 let now = Utc::now();
1679 let (owner_id, run_id) = seed_owned_run(&state, &token, now).await;
1680 TuneRunRow::complete(&state.pool, run_id, now)
1681 .await
1682 .unwrap();
1683
1684 let mut headers = HeaderMap::new();
1685 headers.insert(header::COOKIE, cookie("ab").parse().unwrap());
1686 let result = start_run(
1687 State(state.clone()),
1688 headers,
1689 PeerAddress("127.0.0.1:12345".parse().unwrap()),
1690 Json(valid_request(ProcessType::Flow, ControllerType::Pi)),
1691 )
1692 .await;
1693
1694 assert!(matches!(
1695 result,
1696 Err(ApiError::TooManyRequests {
1697 message,
1698 retry_after_secs
1699 }) if message == "maximum demo runs for this session has been reached"
1700 && retry_after_secs == state.demo_policy.accepted_start_window_secs
1701 ));
1702 assert_eq!(
1703 TuneRunRow::count_for_demo_session(&state.pool, owner_id)
1704 .await
1705 .unwrap(),
1706 1
1707 );
1708 }
1709
1710 #[tokio::test]
1711 async fn built_in_templates_are_read_only_and_user_templates_are_hidden() {
1712 let mut state = crate::test_support::in_memory_state().await;
1713 state.mode = ServerMode::Demo;
1714 state.allowed_origin = Some(DEMO_ORIGIN.into());
1715 let mut custom = DcsTemplateRow::get_by_name(&state.pool, DEMO_TEMPLATE_NAME)
1716 .await
1717 .unwrap()
1718 .unwrap()
1719 .template;
1720 custom.name = "Private custom template".into();
1721 DcsTemplateRow::insert(&state.pool, &custom, TemplateOrigin::User, Utc::now())
1722 .await
1723 .unwrap();
1724 let app = crate::build_router(state.clone());
1725 let list = app
1726 .clone()
1727 .oneshot(Request::get("/api/templates").body(Body::empty()).unwrap())
1728 .await
1729 .unwrap();
1730 let templates: serde_json::Value =
1731 serde_json::from_slice(&to_bytes(list.into_body(), usize::MAX).await.unwrap()).unwrap();
1732 assert!(
1733 templates
1734 .as_array()
1735 .unwrap()
1736 .iter()
1737 .all(|row| row["origin"] == "builtin")
1738 );
1739 assert!(
1740 !templates
1741 .as_array()
1742 .unwrap()
1743 .iter()
1744 .any(|row| row["name"] == custom.name)
1745 );
1746
1747 let hidden = app
1748 .clone()
1749 .oneshot(
1750 Request::get("/api/templates/Private%20custom%20template")
1751 .body(Body::empty())
1752 .unwrap(),
1753 )
1754 .await
1755 .unwrap();
1756 assert_eq!(hidden.status(), StatusCode::NOT_FOUND);
1757 let mutation = app
1758 .clone()
1759 .oneshot(
1760 Request::post("/api/templates")
1761 .header(header::ORIGIN, DEMO_ORIGIN)
1762 .header(header::CONTENT_TYPE, "application/json")
1763 .body(Body::from("{}"))
1764 .unwrap(),
1765 )
1766 .await
1767 .unwrap();
1768 assert_eq!(mutation.status(), StatusCode::METHOD_NOT_ALLOWED);
1769
1770 let shown = app
1771 .oneshot(
1772 Request::get("/api/templates/Yokogawa%20CentumVP")
1773 .body(Body::empty())
1774 .unwrap(),
1775 )
1776 .await
1777 .unwrap();
1778 assert_eq!(shown.status(), StatusCode::OK);
1779 assert_eq!(body_json(shown).await["origin"], "builtin");
1780 }
1781
1782 #[tokio::test]
1783 async fn shared_history_is_owner_scoped_and_uses_the_normal_list_shape() {
1784 let mut state = crate::test_support::in_memory_state().await;
1785 state.mode = ServerMode::Demo;
1786 let token_a = raw_token("ab");
1787 let token_b = raw_token("cd");
1788 let now = Utc::now();
1789 let (_, run_a) = seed_owned_run(&state, &token_a, now).await;
1790 let (_, run_b) = seed_owned_run(&state, &token_b, now + Duration::seconds(1)).await;
1791 let app = crate::build_router(state);
1792
1793 let list = app
1794 .clone()
1795 .oneshot(
1796 Request::get("/api/runs?limit=50")
1797 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token_a}"))
1798 .body(Body::empty())
1799 .unwrap(),
1800 )
1801 .await
1802 .unwrap();
1803 let body: serde_json::Value =
1804 serde_json::from_slice(&to_bytes(list.into_body(), usize::MAX).await.unwrap()).unwrap();
1805 assert_eq!(body["returned"], 1);
1806 assert_eq!(body["total"], 1);
1807 assert_eq!(body["runs"][0]["id"], run_a);
1808
1809 let foreign = app
1810 .oneshot(
1811 Request::get(format!("/api/runs/{run_b}"))
1812 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token_a}"))
1813 .body(Body::empty())
1814 .unwrap(),
1815 )
1816 .await
1817 .unwrap();
1818 assert_eq!(foreign.status(), StatusCode::NOT_FOUND);
1819 }
1820
1821 #[tokio::test]
1822 async fn unpersisted_history_is_empty_and_pagination_is_validated_and_capped() {
1823 let mut state = crate::test_support::in_memory_state().await;
1824 state.mode = ServerMode::Demo;
1825 state.demo_policy.retained_runs_per_visitor = 1;
1826 state.demo_policy.max_active_runs_per_visitor = 1;
1827 let token = raw_token("ab");
1828 let now = Utc::now();
1829 let (_, first) = seed_owned_run(&state, &token, now).await;
1830 let (_, second) = seed_owned_run(&state, &token, now + Duration::seconds(1)).await;
1831 let app = crate::build_router(state);
1832
1833 let empty = app
1834 .clone()
1835 .oneshot(
1836 Request::get("/api/runs")
1837 .header(header::COOKIE, cookie("cd"))
1838 .body(Body::empty())
1839 .unwrap(),
1840 )
1841 .await
1842 .unwrap();
1843 assert_eq!(empty.status(), StatusCode::OK);
1844 assert_eq!(body_json(empty).await["returned"], 0);
1845
1846 for uri in ["/api/runs?limit=0", "/api/runs?offset=-1"] {
1847 let response = app
1848 .clone()
1849 .oneshot(
1850 Request::get(uri)
1851 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
1852 .body(Body::empty())
1853 .unwrap(),
1854 )
1855 .await
1856 .unwrap();
1857 assert_eq!(response.status(), StatusCode::BAD_REQUEST);
1858 }
1859
1860 let capped = app
1861 .oneshot(
1862 Request::get("/api/runs?limit=50")
1863 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
1864 .body(Body::empty())
1865 .unwrap(),
1866 )
1867 .await
1868 .unwrap();
1869 let capped = body_json(capped).await;
1870 assert_eq!(capped["returned"], 2);
1871 assert_eq!(capped["total"], 2);
1872 assert_eq!(capped["runs"][0]["id"], second);
1873 assert_eq!(capped["runs"][1]["id"], first);
1874 }
1875
1876 #[tokio::test]
1877 async fn owned_run_detail_and_last_request_include_the_stored_request_when_valid() {
1878 let mut state = crate::test_support::in_memory_state().await;
1879 state.mode = ServerMode::Demo;
1880 let token = raw_token("ab");
1881 let now = Utc::now();
1882 let (_, run_id) = seed_owned_run(&state, &token, now).await;
1883 let request =
1884 parse_demo_request(valid_request(ProcessType::Flow, ControllerType::Pi)).unwrap();
1885 TuneRunRow::record_connection(
1886 &state.pool,
1887 run_id,
1888 None,
1889 None,
1890 &serde_json::to_string(&request).unwrap(),
1891 )
1892 .await
1893 .unwrap();
1894 let app = crate::build_router(state);
1895
1896 let detail = app
1897 .clone()
1898 .oneshot(
1899 Request::get(format!("/api/runs/{run_id}"))
1900 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
1901 .body(Body::empty())
1902 .unwrap(),
1903 )
1904 .await
1905 .unwrap();
1906 assert_eq!(detail.status(), StatusCode::OK);
1907 let detail = body_json(detail).await;
1908 assert_eq!(detail["id"], run_id);
1909 assert_eq!(detail["pid_constant_tags"]["proportional"], "P");
1910 assert_eq!(detail["original_request"]["tagname"], DEMO_TAG_NAME);
1911
1912 let last = app
1913 .oneshot(
1914 Request::get("/api/runs/last-request")
1915 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
1916 .body(Body::empty())
1917 .unwrap(),
1918 )
1919 .await
1920 .unwrap();
1921 assert_eq!(last.status(), StatusCode::OK);
1922 assert_eq!(body_json(last).await["tagname"], DEMO_TAG_NAME);
1923 }
1924
1925 #[tokio::test]
1926 async fn every_run_resource_returns_404_for_a_different_owner() {
1927 let mut state = crate::test_support::in_memory_state().await;
1928 state.mode = ServerMode::Demo;
1929 state.allowed_origin = Some(DEMO_ORIGIN.into());
1930 let owner_token = raw_token("ab");
1931 let other_token = raw_token("cd");
1932 let now = Utc::now();
1933 let (_, run_id) = seed_owned_run(&state, &owner_token, now).await;
1934 TuneRunRow::complete(&state.pool, run_id, now)
1935 .await
1936 .unwrap();
1937 let app = crate::build_router(state);
1938 let requests = [
1939 Request::get(format!("/api/runs/{run_id}"))
1940 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={other_token}"))
1941 .body(Body::empty())
1942 .unwrap(),
1943 Request::get(format!("/api/runs/{run_id}/stream"))
1944 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={other_token}"))
1945 .body(Body::empty())
1946 .unwrap(),
1947 Request::get(format!("/api/runs/{run_id}/export"))
1948 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={other_token}"))
1949 .body(Body::empty())
1950 .unwrap(),
1951 Request::post(format!("/api/runs/{run_id}/cancel"))
1952 .header(header::ORIGIN, DEMO_ORIGIN)
1953 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={other_token}"))
1954 .body(Body::empty())
1955 .unwrap(),
1956 Request::delete(format!("/api/runs/{run_id}"))
1957 .header(header::ORIGIN, DEMO_ORIGIN)
1958 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={other_token}"))
1959 .body(Body::empty())
1960 .unwrap(),
1961 ];
1962 for request in requests {
1963 let response = app.clone().oneshot(request).await.unwrap();
1964 assert_eq!(response.status(), StatusCode::NOT_FOUND);
1965 }
1966 }
1967
1968 #[tokio::test]
1969 async fn last_request_is_null_until_this_owner_has_started_a_run() {
1970 let mut state = crate::test_support::in_memory_state().await;
1971 state.mode = ServerMode::Demo;
1972 let response = crate::build_router(state)
1973 .oneshot(
1974 Request::get("/api/runs/last-request")
1975 .header(header::COOKIE, cookie("ab"))
1976 .body(Body::empty())
1977 .unwrap(),
1978 )
1979 .await
1980 .unwrap();
1981 assert_eq!(response.status(), StatusCode::OK);
1982 assert_eq!(
1983 to_bytes(response.into_body(), usize::MAX)
1984 .await
1985 .unwrap()
1986 .as_ref(),
1987 b"null"
1988 );
1989 }
1990
1991 #[tokio::test]
1992 async fn valid_shared_start_is_lazily_persisted_and_can_be_cancelled() {
1993 let mut state = crate::test_support::in_memory_state().await;
1994 state.mode = ServerMode::Demo;
1995 state.allowed_origin = Some(DEMO_ORIGIN.into());
1996 let app = crate::build_router(state.clone());
1997 let start = app
1998 .clone()
1999 .oneshot(with_peer(
2000 Request::post("/api/runs")
2001 .header(header::ORIGIN, DEMO_ORIGIN)
2002 .header(header::COOKIE, cookie("ab"))
2003 .header(header::CONTENT_TYPE, "application/json")
2004 .body(Body::from(
2005 valid_request(ProcessType::Flow, ControllerType::Pi).to_string(),
2006 ))
2007 .unwrap(),
2008 "127.0.0.1:12345",
2009 ))
2010 .await
2011 .unwrap();
2012 assert_eq!(start.status(), StatusCode::CREATED);
2013 let body: serde_json::Value =
2014 serde_json::from_slice(&to_bytes(start.into_body(), usize::MAX).await.unwrap())
2015 .unwrap();
2016 let run_id = body["id"].as_i64().unwrap();
2017 assert_eq!(
2018 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM demo_sessions")
2019 .fetch_one(&state.pool)
2020 .await
2021 .unwrap(),
2022 1
2023 );
2024 let cancel = app
2025 .clone()
2026 .oneshot(
2027 Request::post(format!("/api/runs/{run_id}/cancel"))
2028 .header(header::ORIGIN, DEMO_ORIGIN)
2029 .header(header::COOKIE, cookie("ab"))
2030 .body(Body::empty())
2031 .unwrap(),
2032 )
2033 .await
2034 .unwrap();
2035 assert_eq!(cancel.status(), StatusCode::NO_CONTENT);
2036 wait_until_inactive(&state, run_id).await;
2037
2038 let second_start = app
2039 .clone()
2040 .oneshot(with_peer(
2041 Request::post("/api/runs")
2042 .header(header::ORIGIN, DEMO_ORIGIN)
2043 .header(header::COOKIE, cookie("ab"))
2044 .header(header::CONTENT_TYPE, "application/json")
2045 .body(Body::from(
2046 valid_request(ProcessType::Flow, ControllerType::Pi).to_string(),
2047 ))
2048 .unwrap(),
2049 "127.0.0.1:12345",
2050 ))
2051 .await
2052 .unwrap();
2053 assert_eq!(second_start.status(), StatusCode::CREATED);
2054 let second_run_id = body_json(second_start).await["id"].as_i64().unwrap();
2055 state.active_run.cancel(second_run_id).await;
2056 wait_until_inactive(&state, second_run_id).await;
2057 }
2058
2059 #[tokio::test]
2060 async fn start_rejects_global_visitor_and_accepted_start_quota_exhaustion() {
2061 for (policy, expected) in [
2062 (
2063 DemoPolicy {
2064 max_active_runs_global: 0,
2065 ..DemoPolicy::default()
2066 },
2067 StatusCode::SERVICE_UNAVAILABLE,
2068 ),
2069 (
2070 DemoPolicy {
2071 max_active_runs_per_visitor: 0,
2072 ..DemoPolicy::default()
2073 },
2074 StatusCode::TOO_MANY_REQUESTS,
2075 ),
2076 ] {
2077 let base = crate::test_support::in_memory_state().await;
2078 let mut state =
2079 AppState::for_mode(base.pool, base.config_store, ServerMode::Demo, policy);
2080 state.allowed_origin = Some(DEMO_ORIGIN.into());
2081 let response = crate::build_router(state)
2082 .oneshot(with_peer(
2083 Request::post("/api/runs")
2084 .header(header::ORIGIN, DEMO_ORIGIN)
2085 .header(header::COOKIE, cookie("ab"))
2086 .header(header::CONTENT_TYPE, "application/json")
2087 .body(Body::from(
2088 valid_request(ProcessType::Flow, ControllerType::Pi).to_string(),
2089 ))
2090 .unwrap(),
2091 "127.0.0.1:12345",
2092 ))
2093 .await
2094 .unwrap();
2095 assert_eq!(response.status(), expected);
2096 }
2097
2098 let mut state = crate::test_support::in_memory_state().await;
2099 state.mode = ServerMode::Demo;
2100 state.allowed_origin = Some(DEMO_ORIGIN.into());
2101 state.demo_policy.accepted_starts_per_token = 1;
2102 state
2103 .demo_runtime
2104 .reserve_accepted_start(
2105 &token_hash(&raw_token("ab")),
2106 "127.0.0.1",
2107 Utc::now(),
2108 state.demo_policy,
2109 )
2110 .await
2111 .unwrap();
2112 let response = crate::build_router(state)
2113 .oneshot(with_peer(
2114 Request::post("/api/runs")
2115 .header(header::ORIGIN, DEMO_ORIGIN)
2116 .header(header::COOKIE, cookie("ab"))
2117 .header(header::CONTENT_TYPE, "application/json")
2118 .body(Body::from(
2119 valid_request(ProcessType::Flow, ControllerType::Pi).to_string(),
2120 ))
2121 .unwrap(),
2122 "127.0.0.1:12345",
2123 ))
2124 .await
2125 .unwrap();
2126 assert_eq!(response.status(), StatusCode::TOO_MANY_REQUESTS);
2127 }
2128
2129 #[tokio::test]
2130 async fn preparation_failure_releases_accepted_start_quota() {
2131 let mut state = crate::test_support::in_memory_state().await;
2132 state.mode = ServerMode::Demo;
2133 state.allowed_origin = Some(DEMO_ORIGIN.into());
2134 state.pool.close().await;
2135 let response = crate::build_router(state)
2136 .oneshot(with_peer(
2137 Request::post("/api/runs")
2138 .header(header::ORIGIN, DEMO_ORIGIN)
2139 .header(header::COOKIE, cookie("ab"))
2140 .header(header::CONTENT_TYPE, "application/json")
2141 .body(Body::from(
2142 valid_request(ProcessType::Flow, ControllerType::Pi).to_string(),
2143 ))
2144 .unwrap(),
2145 "127.0.0.1:12345",
2146 ))
2147 .await
2148 .unwrap();
2149 assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR);
2150 }
2151
2152 #[tokio::test]
2153 async fn scheduling_conflict_discards_the_prepared_run() {
2154 let mut state = crate::test_support::in_memory_state().await;
2155 state.mode = ServerMode::Demo;
2156 state.allowed_origin = Some(DEMO_ORIGIN.into());
2157 state.active_run.reserve(999).await.unwrap();
2158 let response = crate::build_router(state.clone())
2159 .oneshot(with_peer(
2160 Request::post("/api/runs")
2161 .header(header::ORIGIN, DEMO_ORIGIN)
2162 .header(header::COOKIE, cookie("ab"))
2163 .header(header::CONTENT_TYPE, "application/json")
2164 .body(Body::from(
2165 valid_request(ProcessType::Flow, ControllerType::Pi).to_string(),
2166 ))
2167 .unwrap(),
2168 "127.0.0.1:12345",
2169 ))
2170 .await
2171 .unwrap();
2172 assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
2173 assert_eq!(TuneRunRow::count_demo_owned(&state.pool).await.unwrap(), 0);
2174 state.active_run.release(999).await;
2175 }
2176
2177 #[tokio::test]
2178 async fn owned_completed_stream_ends_with_exactly_one_done_event() {
2179 let mut state = crate::test_support::in_memory_state().await;
2180 state.mode = ServerMode::Demo;
2181 let token = raw_token("ab");
2182 let now = Utc::now();
2183 let (_, run_id) = seed_owned_run(&state, &token, now).await;
2184 TuneRunRow::complete(&state.pool, run_id, now)
2185 .await
2186 .unwrap();
2187 let response = crate::build_router(state)
2188 .oneshot(
2189 Request::get(format!("/api/runs/{run_id}/stream"))
2190 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2191 .body(Body::empty())
2192 .unwrap(),
2193 )
2194 .await
2195 .unwrap();
2196 let body = String::from_utf8(
2197 to_bytes(response.into_body(), usize::MAX)
2198 .await
2199 .unwrap()
2200 .to_vec(),
2201 )
2202 .unwrap();
2203 assert_eq!(body.matches("event: done").count(), 1);
2204 assert!(body.contains("\"outcome\":\"completed\""));
2205 assert!(!body.contains("event: error"));
2206 }
2207
2208 #[tokio::test]
2209 async fn streams_report_missing_deleted_or_unavailable_runs_and_permit_limits() {
2210 let base = crate::test_support::in_memory_state().await;
2211 let policy = DemoPolicy {
2212 max_sse_global: 1,
2213 max_sse_per_visitor: 1,
2214 sse_lifetime_secs: 5,
2215 ..DemoPolicy::default()
2216 };
2217 let state = AppState::for_mode(base.pool, base.config_store, ServerMode::Demo, policy);
2218 let token = raw_token("ab");
2219 let other = raw_token("cd");
2220 let now = Utc::now();
2221 let (_, run_id) = seed_owned_run(&state, &token, now).await;
2222 let (_, other_run_id) = seed_owned_run(&state, &other, now).await;
2223 let app = crate::build_router(state.clone());
2224
2225 let missing = app
2226 .clone()
2227 .oneshot(
2228 Request::get("/api/runs/999999/stream")
2229 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2230 .body(Body::empty())
2231 .unwrap(),
2232 )
2233 .await
2234 .unwrap();
2235 assert_eq!(missing.status(), StatusCode::NOT_FOUND);
2236
2237 let held = app
2238 .clone()
2239 .oneshot(
2240 Request::get(format!("/api/runs/{run_id}/stream"))
2241 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2242 .body(Body::empty())
2243 .unwrap(),
2244 )
2245 .await
2246 .unwrap();
2247 assert_eq!(held.status(), StatusCode::OK);
2248
2249 let global_limited = app
2250 .clone()
2251 .oneshot(
2252 Request::get(format!("/api/runs/{other_run_id}/stream"))
2253 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={other}"))
2254 .body(Body::empty())
2255 .unwrap(),
2256 )
2257 .await
2258 .unwrap();
2259 assert_eq!(global_limited.status(), StatusCode::SERVICE_UNAVAILABLE);
2260 drop(held);
2261
2262 let base = crate::test_support::in_memory_state().await;
2263 let visitor_policy = DemoPolicy {
2264 max_sse_global: 2,
2265 max_sse_per_visitor: 1,
2266 sse_lifetime_secs: 5,
2267 ..DemoPolicy::default()
2268 };
2269 let visitor_state = AppState::for_mode(
2270 base.pool,
2271 base.config_store,
2272 ServerMode::Demo,
2273 visitor_policy,
2274 );
2275 let (_, visitor_run_id) = seed_owned_run(&visitor_state, &token, now).await;
2276 let visitor_app = crate::build_router(visitor_state);
2277 let held = visitor_app
2278 .clone()
2279 .oneshot(
2280 Request::get(format!("/api/runs/{visitor_run_id}/stream"))
2281 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2282 .body(Body::empty())
2283 .unwrap(),
2284 )
2285 .await
2286 .unwrap();
2287 let visitor_limited = visitor_app
2288 .clone()
2289 .oneshot(
2290 Request::get(format!("/api/runs/{visitor_run_id}/stream"))
2291 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2292 .body(Body::empty())
2293 .unwrap(),
2294 )
2295 .await
2296 .unwrap();
2297 assert_eq!(visitor_limited.status(), StatusCode::TOO_MANY_REQUESTS);
2298 drop(held);
2299
2300 let deleted_before_body = app
2301 .clone()
2302 .oneshot(
2303 Request::get(format!("/api/runs/{run_id}/stream"))
2304 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2305 .body(Body::empty())
2306 .unwrap(),
2307 )
2308 .await
2309 .unwrap();
2310 TuneRunRow::delete_for_demo_session(&state.pool, run_id, 1)
2311 .await
2312 .unwrap();
2313 let body = body_text(deleted_before_body).await;
2314 assert!(body.contains("demo run is no longer available"));
2315
2316 let unavailable = app
2317 .oneshot(
2318 Request::get(format!("/api/runs/{other_run_id}/stream"))
2319 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={other}"))
2320 .body(Body::empty())
2321 .unwrap(),
2322 )
2323 .await
2324 .unwrap();
2325 state.pool.close().await;
2326 let body = body_text(unavailable).await;
2327 assert!(body.contains("demo stream unavailable"));
2328 }
2329
2330 #[tokio::test]
2331 async fn stream_lifetime_emits_a_generic_error_and_releases_its_permits() {
2332 let mut state = crate::test_support::in_memory_state().await;
2333 state.mode = ServerMode::Demo;
2334 state.demo_policy.sse_lifetime_secs = 0;
2335 let token = raw_token("ab");
2336 let (_, run_id) = seed_owned_run(&state, &token, Utc::now()).await;
2337 let app = crate::build_router(state.clone());
2338 for _ in 0..2 {
2339 let response = app
2340 .clone()
2341 .oneshot(
2342 Request::get(format!("/api/runs/{run_id}/stream"))
2343 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2344 .body(Body::empty())
2345 .unwrap(),
2346 )
2347 .await
2348 .unwrap();
2349 let body = String::from_utf8(
2350 to_bytes(response.into_body(), usize::MAX)
2351 .await
2352 .unwrap()
2353 .to_vec(),
2354 )
2355 .unwrap();
2356 assert_eq!(body.matches("event: error").count(), 1);
2357 assert!(body.contains("demo stream lifetime exceeded"));
2358 assert!(!body.contains("event: done"));
2359 }
2360 }
2361
2362 #[tokio::test]
2363 async fn global_demo_run_capacity_uses_the_global_policy_limit() {
2364 let mut state = crate::test_support::in_memory_state().await;
2365 state.mode = ServerMode::Demo;
2366 state.demo_policy.max_tune_run_rows_global = 1;
2367 assert!(ensure_global_run_capacity(&state).await.is_ok());
2368 seed_owned_run(&state, &raw_token("ab"), Utc::now()).await;
2369 assert!(matches!(
2370 ensure_global_run_capacity(&state).await.unwrap_err(),
2371 ApiError::GlobalCapacity { .. }
2372 ));
2373 }
2374
2375 #[tokio::test]
2376 async fn owned_history_retention_keeps_running_rows_and_newest_terminal_rows() {
2377 let mut state = crate::test_support::in_memory_state().await;
2378 state.demo_policy.retained_runs_per_visitor = 1;
2379 let token = raw_token("ab");
2380 let now = Utc::now();
2381 let (owner_id, first) = seed_owned_run(&state, &token, now).await;
2382 TuneRunRow::fail(&state.pool, first, now, "finished")
2383 .await
2384 .unwrap();
2385 let (_, second) = seed_owned_run(&state, &token, now + Duration::seconds(1)).await;
2386 TuneRunRow::fail(&state.pool, second, now + Duration::seconds(1), "finished")
2387 .await
2388 .unwrap();
2389 let (_, running) = seed_owned_run(&state, &token, now + Duration::seconds(2)).await;
2390 trim_owned_history(&state, owner_id).await;
2391 assert!(TuneRunRow::get(&state.pool, first).await.unwrap().is_none());
2392 assert!(
2393 TuneRunRow::get(&state.pool, second)
2394 .await
2395 .unwrap()
2396 .is_some()
2397 );
2398 assert!(
2399 TuneRunRow::get(&state.pool, running)
2400 .await
2401 .unwrap()
2402 .is_some()
2403 );
2404 }
2405
2406 #[tokio::test]
2407 async fn capacity_and_retention_database_failures_are_safe() {
2408 let state = crate::test_support::in_memory_state().await;
2409 state.pool.close().await;
2410 assert!(matches!(
2411 ensure_global_run_capacity(&state).await.unwrap_err(),
2412 ApiError::Internal(_)
2413 ));
2414 trim_owned_history(&state, 1).await;
2415 }
2416
2417 #[tokio::test]
2418 async fn cancel_delete_and_export_cover_owned_terminal_and_running_variants() {
2419 let mut state = crate::test_support::in_memory_state().await;
2420 state.mode = ServerMode::Demo;
2421 state.allowed_origin = Some(DEMO_ORIGIN.into());
2422 let token = raw_token("ab");
2423 let now = Utc::now();
2424 let (_, running) = seed_owned_run(&state, &token, now).await;
2425 let (_ctrl_c, cancel_handle) = bhtune_cli::cancel::CtrlC::manual();
2426 let (finish, finished) = tokio::sync::oneshot::channel::<()>();
2427 state
2428 .active_run
2429 .start(running, cancel_handle, async move {
2430 let _ = finished.await;
2431 })
2432 .await
2433 .unwrap();
2434 let (_, empty_terminal) = seed_owned_run(&state, &token, now + Duration::seconds(1)).await;
2435 TuneRunRow::complete(&state.pool, empty_terminal, now)
2436 .await
2437 .unwrap();
2438 let (_, sampled_terminal) =
2439 seed_owned_run(&state, &token, now + Duration::seconds(2)).await;
2440 insert_one_sample(&state, sampled_terminal, now).await;
2441 TuneRunRow::complete(&state.pool, sampled_terminal, now)
2442 .await
2443 .unwrap();
2444 let app = crate::build_router(state.clone());
2445
2446 let cancel = app
2447 .clone()
2448 .oneshot(
2449 Request::post(format!("/api/runs/{running}/cancel"))
2450 .header(header::ORIGIN, DEMO_ORIGIN)
2451 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2452 .body(Body::empty())
2453 .unwrap(),
2454 )
2455 .await
2456 .unwrap();
2457 assert_eq!(cancel.status(), StatusCode::NO_CONTENT);
2458 finish.send(()).unwrap();
2459 wait_until_inactive(&state, running).await;
2460
2461 let conflict = app
2462 .clone()
2463 .oneshot(
2464 Request::delete(format!("/api/runs/{running}"))
2465 .header(header::ORIGIN, DEMO_ORIGIN)
2466 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2467 .body(Body::empty())
2468 .unwrap(),
2469 )
2470 .await
2471 .unwrap();
2472 assert_eq!(conflict.status(), StatusCode::CONFLICT);
2473
2474 let empty_export = app
2475 .clone()
2476 .oneshot(
2477 Request::get(format!("/api/runs/{empty_terminal}/export"))
2478 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2479 .body(Body::empty())
2480 .unwrap(),
2481 )
2482 .await
2483 .unwrap();
2484 assert_eq!(empty_export.status(), StatusCode::NOT_FOUND);
2485
2486 for (format, content_type, extension) in [
2487 ("csv", "text/csv", "csv"),
2488 ("json", "application/json", "json"),
2489 ] {
2490 let export = app
2491 .clone()
2492 .oneshot(
2493 Request::get(format!(
2494 "/api/runs/{sampled_terminal}/export?format={format}"
2495 ))
2496 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2497 .body(Body::empty())
2498 .unwrap(),
2499 )
2500 .await
2501 .unwrap();
2502 assert_eq!(export.status(), StatusCode::OK);
2503 assert_eq!(
2504 export.headers().get(header::CONTENT_TYPE).unwrap(),
2505 content_type
2506 );
2507 assert_eq!(
2508 export.headers().get(header::CONTENT_DISPOSITION).unwrap(),
2509 &format!("attachment; filename=\"demo-run-{sampled_terminal}.{extension}\"")
2510 );
2511 assert!(!body_text(export).await.is_empty());
2512 }
2513
2514 let deleted = app
2515 .oneshot(
2516 Request::delete(format!("/api/runs/{sampled_terminal}"))
2517 .header(header::ORIGIN, DEMO_ORIGIN)
2518 .header(header::COOKIE, format!("{DEMO_COOKIE_NAME}={token}"))
2519 .body(Body::empty())
2520 .unwrap(),
2521 )
2522 .await
2523 .unwrap();
2524 assert_eq!(deleted.status(), StatusCode::NO_CONTENT);
2525 }
2526}