Skip to main content

bhtune_server/routes/
demo.rs

1//! Restricted simulator-only public Demo API.
2
3use 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    // The raw object allow-list is intentionally checked before acquiring a permit, touching
966    // session/history storage, reading configuration, or constructing a driver. Deserializing
967    // directly to `StartRunRequest` would lose the distinction between an omitted forbidden
968    // field and an explicitly supplied `null`/`false` field.
969    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            // On-demand cleanup prevents an expired visitor or over-retained history from causing a
1004            // friendly capacity rejection until the next periodic sweep.
1005            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}