Skip to main content

beam_daemon/
lib.rs

1use std::collections::{BTreeMap, HashMap, HashSet};
2use std::path::PathBuf;
3use std::process::Stdio;
4use std::sync::{Arc, Mutex as StdMutex, OnceLock};
5use std::time::{Duration, Instant};
6mod api_token;
7mod ask;
8mod backend;
9mod card_i18n;
10mod connector_runtime;
11mod connector_store;
12mod daemon_types;
13mod dashboard_support;
14mod debug_simulate;
15mod dir_select;
16mod external_host_watcher;
17mod final_output;
18mod grant;
19mod herdr_adopt;
20mod herdr_lifecycle;
21mod herdr_probe;
22mod ip_resolver;
23mod lark_api_helpers;
24mod lark_card_builders;
25mod lark_delivery;
26mod lark_dispatch;
27mod lark_history;
28mod lark_identity;
29mod lark_ingress;
30mod lark_parse;
31mod lark_replies;
32mod lark_security;
33mod lark_session_cards;
34mod opencode_adopt_resolver;
35mod persistence;
36mod prompt;
37mod route_handlers;
38mod session_cards;
39mod session_creation;
40mod terminal_auth;
41mod terminal_proxy;
42pub mod test_hooks;
43mod trigger_log;
44mod utils;
45mod webhook_key;
46mod webhook_lifecycle;
47mod worker_health;
48mod worker_lifecycle;
49mod workflow_approval_cards;
50mod workflow_cancellation;
51mod workflow_catalog;
52mod workflow_commands;
53mod workflow_event_fanout;
54mod workflow_execution;
55mod workflow_host_executors;
56mod workflow_progress_card;
57mod workflow_reconcilers;
58mod workflow_resume;
59mod workflow_runtime_driver;
60mod zellij_adopt;
61mod zellij_web;
62
63// Re-export daemon types (used across modules and externally)
64pub use daemon_types::RunOptions;
65pub(crate) use daemon_types::*;
66
67// Re-export workflow catalog items for backward compatibility (used by route handlers and tests)
68pub(crate) use workflow_catalog::*;
69// Re-export workflow execution items for backward compatibility (used by route handlers and tests)
70pub(crate) use workflow_execution::*;
71// Re-export workflow resume items for backward compatibility (used by route handlers and tests)
72pub(crate) use api_token::*;
73pub(crate) use connector_runtime::*;
74pub(crate) use external_host_watcher::*;
75pub(crate) use final_output::*;
76pub(crate) use ip_resolver::*;
77pub(crate) use lark_api_helpers::*;
78pub(crate) use lark_card_builders::*;
79pub(crate) use lark_delivery::*;
80pub(crate) use lark_dispatch::*;
81#[allow(unused_imports)]
82pub(crate) use lark_history::*;
83pub(crate) use lark_identity::*;
84pub(crate) use lark_ingress::*;
85pub(crate) use lark_parse::*;
86pub(crate) use lark_replies::*;
87pub(crate) use lark_security::*;
88pub(crate) use lark_session_cards::*;
89pub(crate) use opencode_adopt_resolver::*;
90pub(crate) use persistence::*;
91pub(crate) use route_handlers::*;
92pub(crate) use session_cards::*;
93pub(crate) use session_creation::*;
94pub(crate) use utils::*;
95pub(crate) use worker_lifecycle::*;
96pub(crate) use workflow_approval_cards::*;
97pub(crate) use workflow_resume::*;
98pub(crate) use zellij_adopt::*;
99
100#[cfg(test)]
101use dashboard_support::{
102    WebhookTriggerRecord, extract_dashboard_token, write_webhook_trigger_records,
103};
104use dashboard_support::{
105    dashboard_gate, dashboard_token_is_valid, load_observed_bot_open_ids_for_app,
106    load_observed_bots_for_chat, mint_dashboard_token, read_webhook_trigger_records,
107    record_observed_bots,
108};
109
110use anyhow::{Context, Result};
111use axum::{
112    Json, Router,
113    body::Bytes,
114    extract::{Path as AxumPath, Query, State},
115    http::{HeaderMap, StatusCode},
116    middleware,
117    response::{IntoResponse, Redirect},
118    routing::{get, get_service, post, put},
119};
120use base64::Engine;
121use beam_core::{
122    AdoptedFrom, AgentAttention, ApiHealth, AttemptResumeRequest, AttentionRequest, BackendKind,
123    BeamPaths, BotConfig, BotSummary, ChatMode, CliUsageLimitState, ColdWorkflowRun, Config,
124    CreateSessionRequest, CustomTrigger, DaemonOverview, DaemonRuntimeState, DaemonToWorker,
125    DisplayMode, EventDraft, EventLog, EventWindowOpts, FinalOutputKind, FinalOutputRequest,
126    InitConfig, PendingResponseCardState, RestartSessionRequest, ResumeSessionRequest,
127    RunChatBinding, RunStatus, ScheduleChatType, ScreenStatus, Session, SessionGroup,
128    SessionInputRequest, SessionLocateInfo, SessionScope, SessionStatus, SessionSummary,
129    TalkEvaluation, TermActionKey, TranscriptChoice, TuiPromptOption, WaitResolution,
130    WorkerToDaemon, WorkflowActor, WorkflowOutputRef, can_operate, evaluate_talk,
131    parse_workflow_definition, read_event_window, read_run_events_pure, read_run_snapshot,
132    resolve_custom_trigger, resolve_trigger_message, scan_cold_workflow_runs,
133};
134use chrono::Utc;
135use connector_store::{
136    ConnectorDefinition, ConnectorLifecycleExtractors, ConnectorLoggingPolicy,
137    ConnectorPromptEnvelope, ConnectorRateLimit, ConnectorTarget, ConnectorVerify,
138    delete_connector, get_connector, list_connectors, upsert_connector,
139};
140use feishu_sdk::{
141    card::CardAction,
142    core as feishu_core,
143    event::{
144        Event, EventDispatcher, EventDispatcherConfig, EventHandler, EventHandlerResult, EventResp,
145    },
146    ws::{StreamClient, StreamConfig},
147};
148use hmac::{Hmac, KeyInit, Mac};
149use reqwest::Client;
150use serde_json::Value;
151use sha2::{Digest, Sha256};
152use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
153use tokio::net::TcpListener;
154use tokio::process::{Child, ChildStdin, Command};
155use tokio::sync::{Mutex, RwLock};
156use tower_http::services::ServeDir;
157use tracing::{debug, error, info, warn};
158use trigger_log::{
159    TriggerLogStats, list_trigger_logs, new_trigger_id as new_trigger_log_id, prune_trigger_logs,
160    summarize_trigger_logs,
161};
162use uuid::Uuid;
163use webhook_key::{
164    create_webhook_secret, delete_webhook_secret, generate_webhook_secret_plaintext,
165    get_webhook_secret, list_webhook_secret_refs, set_webhook_secret,
166};
167use webhook_lifecycle::{begin_webhook_lifecycle_firing, resolve_webhook_lifecycle_group};
168
169pub async fn run(paths: BeamPaths, options: RunOptions) -> Result<()> {
170    let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
171    tokio::fs::create_dir_all(paths.run_dir()).await?;
172    tokio::fs::create_dir_all(paths.logs_dir()).await?;
173    tokio::fs::create_dir_all(paths.sessions_dir()).await?;
174
175    let config = load_config(&paths)?;
176    let bots = load_bot_configs(&paths)?;
177    herdr_probe::probe_herdr_at_startup(&config, &bots).await?;
178    let mut sessions = load_sessions(&paths).await?;
179    for session in sessions.values_mut() {
180        let marker = read_pending_response_patch_marker(&paths, &session.session_id).await?;
181        if should_treat_pending_card_as_patched_by_marker(
182            session.pending_response_card_id.as_deref(),
183            marker.as_ref(),
184        ) {
185            mark_pending_response_card_patched(session);
186            let _ = clear_pending_response_patch_marker(&paths, &session.session_id).await;
187        }
188    }
189    let listener = TcpListener::bind("127.0.0.1:7893").await?;
190    let addr = listener.local_addr()?;
191    let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
192    let external_host = resolve_external_host(&config.web.host);
193    let started_at = Utc::now();
194    let runtime = DaemonRuntimeState {
195        pid: std::process::id(),
196        api_addr: addr.to_string(),
197        started_at,
198        log_path: paths.daemon_log().display().to_string(),
199    };
200    persist_runtime_state(&paths, &runtime).await?;
201
202    let workflow_progress_cards =
203        workflow_runtime_driver::load_workflow_progress_cards(&paths).await;
204    info!(
205        "loaded {} persisted workflow progress cards",
206        workflow_progress_cards.len()
207    );
208
209    let ask_pending_map = ask::load_ask_pending(&paths).await;
210    info!(
211        "loaded {} persisted ask pending entries",
212        ask_pending_map.len()
213    );
214
215    let grant_pending_map = grant::load_grant_pending(&paths);
216    info!(
217        "loaded {} persisted grant pending entries",
218        grant_pending_map.len()
219    );
220
221    let pending_creates_map = dir_select::load_pending_creates(&paths).await;
222    info!(
223        "loaded {} persisted pending creates",
224        pending_creates_map.len()
225    );
226
227    // Load persisted recent Lark events (dedupe state).
228    let recent_lark_events = load_recent_lark_events(&paths).await;
229    info!(
230        "loaded {} persisted recent lark events",
231        recent_lark_events.len()
232    );
233
234    // Load or rotate the local api token (daily rotation, 1h grace for the
235    // previous token) before any route can require it.
236    let api_token_state = load_or_create_api_token(&paths).await?;
237
238    let state = AppState {
239        paths: paths.clone(),
240        started_at,
241        sessions: Arc::new(Mutex::new(sessions)),
242        workers: Arc::new(Mutex::new(HashMap::new())),
243        worker_health: Arc::new(Mutex::new(HashMap::new())),
244        attempt_resumes: Arc::new(Mutex::new(HashMap::new())),
245        shutdown: Arc::new(Mutex::new(Some(shutdown_tx))),
246        options,
247        http: Client::new(),
248        config,
249        bots: Arc::new(bots),
250        lark_tokens: Arc::new(Mutex::new(HashMap::new())),
251        chat_mode_cache: Arc::new(Mutex::new(HashMap::new())),
252        recent_lark_events: Arc::new(Mutex::new(recent_lark_events)),
253        inflight_final_output_turns: Arc::new(Mutex::new(HashSet::new())),
254        workflow_progress_cards: Arc::new(Mutex::new(workflow_progress_cards)),
255        ask_pending: Arc::new(Mutex::new(ask_pending_map)),
256        grant_pending: Arc::new(Mutex::new(grant_pending_map)),
257        pending_creates: Arc::new(Mutex::new(pending_creates_map)),
258        dashboard_token: Arc::new(Mutex::new(None)),
259        api_token: Arc::new(RwLock::new(api_token_state)),
260        external_host: std::sync::Arc::new(tokio::sync::RwLock::new(external_host)),
261    };
262
263    refresh_external_host(&state, true).await?;
264    spawn_external_host_watcher(state.clone());
265    spawn_api_token_rotator(state.clone());
266
267    spawn_lark_ws_clients(&state);
268
269    // Load replay nonces and rate buckets from disk into static stores.
270    {
271        let path = paths.replay_nonces_json();
272        if let Ok(Some(map)) = beam_core::persist::read_json::<HashMap<String, u64>>(&path) {
273            let now = now_ms();
274            let mut guard = replay_nonce_store()
275                .lock()
276                .unwrap_or_else(|p| p.into_inner());
277            for (key, expiry) in map {
278                if expiry > now {
279                    guard.insert(key, expiry);
280                }
281            }
282            info!("loaded {} persisted replay nonces", guard.len());
283        }
284        let path = paths.rate_buckets_json();
285        if let Ok(Some(map)) = beam_core::persist::read_json::<HashMap<String, (u64, u64)>>(&path) {
286            let mut guard = rate_bucket_store()
287                .lock()
288                .unwrap_or_else(|p| p.into_inner());
289            for (key, val) in map {
290                guard.insert(key, val);
291            }
292            info!("loaded {} persisted rate buckets", guard.len());
293        }
294    }
295
296    // Probe bot open_id / app_name from Lark API (best-effort).
297    for bot in state.bots.values() {
298        let paths = state.paths.clone();
299        let bot = bot.clone();
300        tokio::spawn(async move {
301            probe_and_persist_bot_info(&paths, &bot).await;
302        });
303    }
304
305    let restore_candidates = {
306        let mut sessions = state.sessions.lock().await;
307        let restore_candidates = reconcile_restored_sessions_with(
308            &mut sessions,
309            state.config.daemon.quiet_restart,
310            crate::herdr_lifecycle::mux_target_alive_sync,
311        );
312        let snapshot = sessions.clone();
313        drop(sessions);
314        persist_sessions(&state.paths, &snapshot).await?;
315        restore_candidates
316    };
317    for session in restore_candidates {
318        match build_init_from_session(&session, &state.config, &state.bots) {
319            Ok(init) => {
320                if let Err(err) = spawn_worker(state.clone(), session.clone(), init).await {
321                    warn!("failed to restore session {}: {}", session.session_id, err);
322                }
323            }
324            Err(err) => warn!(
325                "failed to rebuild init for session {}: {}",
326                session.session_id, err
327            ),
328        }
329    }
330    {
331        let sessions = state.sessions.lock().await;
332        for session in sessions.values() {
333            if let Some(usage_limit) = session.usage_limit.clone() {
334                arm_usage_limit_retry_timer(state.clone(), session.session_id.clone(), usage_limit);
335            }
336        }
337    }
338    // Recover pending final output retries from before restart.
339    {
340        let markers = load_final_output_retry_markers(&state.paths);
341        if !markers.is_empty() {
342            let active_sessions: HashSet<String> = {
343                let sessions = state.sessions.lock().await;
344                sessions
345                    .values()
346                    .filter(|s| s.status == SessionStatus::Active)
347                    .map(|s| s.session_id.clone())
348                    .collect()
349            };
350            let mut recovered = 0usize;
351            let mut skipped = 0usize;
352            for marker in &markers {
353                if !active_sessions.contains(&marker.session_id) {
354                    skipped += 1;
355                    continue; // Session was closed during restart
356                }
357                // Check idempotency: resume only if this turn has NOT yet been delivered
358                let should_resume = {
359                    let sessions = state.sessions.lock().await;
360                    sessions
361                        .get(&marker.session_id)
362                        .map(|s| {
363                            marker.turn_id.as_deref().is_none_or(|tid| {
364                                // Resume if last_final_output_turn_id doesn't match this turn
365                                s.last_final_output_turn_id.as_deref() != Some(tid)
366                            })
367                        })
368                        .unwrap_or(false) // session not found → skip
369                };
370                if should_resume {
371                    info!(
372                        "final output retry: resuming delivery for session {} turn {:?} attempt {}",
373                        marker.session_id, marker.turn_id, marker.attempt
374                    );
375                    schedule_final_output_delivery(
376                        state.clone(),
377                        marker.session_id.clone(),
378                        marker.content.clone(),
379                        marker.turn_id.clone(),
380                        marker.kind,
381                        marker.user_text.clone(),
382                        marker.attempt,
383                    );
384                    recovered += 1;
385                } else {
386                    skipped += 1;
387                }
388            }
389            info!(
390                "final output retry: {} markers recovered, {} skipped (closed/duplicate)",
391                recovered, skipped
392            );
393        }
394    }
395
396    let cold_scan_bots: Vec<String> = state.bots.keys().cloned().collect();
397    for lark_app_id in &cold_scan_bots {
398        match scan_cold_workflow_runs(&state.paths, lark_app_id).await {
399            Ok((runs, stats)) => {
400                if stats.discovered > 0 {
401                    info!(
402                        "cold-scan: discovered {} non-terminal workflow runs for bot {}",
403                        stats.discovered, lark_app_id
404                    );
405                }
406                for skipped in &stats.skipped {
407                    warn!("cold-scan skipped: {}", skipped);
408                }
409                for run in runs {
410                    let run_id = run.run_id.clone();
411                    info!("cold-attaching workflow run {}", run_id);
412                    let s = state.clone();
413                    tokio::spawn(async move {
414                        if let Err(err) = drive_workflow_run_after_cold_attach(s, run).await {
415                            warn!("cold-attach workflow run {} failed: {}", run_id, err);
416                        }
417                    });
418                }
419            }
420            Err(err) => {
421                warn!("cold-scan failed for bot {}: {}", lark_app_id, err);
422            }
423        }
424    }
425
426    async fn drive_workflow_run_after_cold_attach(
427        state: AppState,
428        run: ColdWorkflowRun,
429    ) -> Result<()> {
430        let workflow_json =
431            serde_json::to_string(&run.def).context("failed to serialize workflow definition")?;
432        workflow_runtime_driver::run(&state, &run.run_id, &workflow_json).await;
433        Ok(())
434    }
435
436    async fn create_schedule(
437        State(state): State<AppState>,
438        Json(body): Json<Value>,
439    ) -> Json<Value> {
440        let content = body.get("content").and_then(Value::as_str).unwrap_or("");
441        let schedule_id = uuid::Uuid::new_v4().to_string();
442        let task = serde_json::json!({
443            "scheduleId": schedule_id,
444            "content": content,
445            "createdAt": chrono::Utc::now().to_rfc3339(),
446            "status": "active",
447        });
448        let schedules_path = state.paths.schedules_json();
449        let mut schedules: Vec<Value> = tokio::fs::read_to_string(&schedules_path)
450            .await
451            .ok()
452            .and_then(|raw| serde_json::from_str(&raw).ok())
453            .unwrap_or_default();
454        schedules.push(task.clone());
455        let _ = tokio::fs::write(
456            &schedules_path,
457            serde_json::to_string_pretty(&schedules).unwrap_or_default(),
458        )
459        .await;
460        Json(task)
461    }
462
463    async fn report_session(
464        State(state): State<AppState>,
465        AxumPath(session_id): AxumPath<String>,
466        Json(body): Json<Value>,
467    ) -> Result<Json<Value>, (StatusCode, String)> {
468        let content = body
469            .get("content")
470            .and_then(Value::as_str)
471            .unwrap_or("")
472            .trim()
473            .to_string();
474        if content.is_empty() {
475            return Err((
476                StatusCode::BAD_REQUEST,
477                "content must not be empty".to_string(),
478            ));
479        }
480        let session = {
481            let sessions = state.sessions.lock().await;
482            sessions.get(&session_id).cloned()
483        }
484        .ok_or_else(|| (StatusCode::NOT_FOUND, "session not found".to_string()))?;
485        if session.lark_app_id == "local" {
486            return Ok(Json(serde_json::json!({
487                "ok": true,
488                "sessionId": session_id,
489                "local": true,
490            })));
491        }
492        let Some(bot) = state.bots.get(&session.lark_app_id) else {
493            return Err((StatusCode::NOT_FOUND, "bot not registered".to_string()));
494        };
495        let post = build_report_post_content(&session, &content);
496        let target_message_id = session
497            .quote_target_id
498            .as_deref()
499            .filter(|value| !value.trim().is_empty());
500        let message_id = if let Some(target_message_id) = target_message_id {
501            match lark_reply_post_message(&state, bot, target_message_id, &post).await {
502                Ok(message_id) => message_id,
503                Err(err) => return Err((StatusCode::BAD_GATEWAY, err.to_string())),
504            }
505        } else {
506            match lark_send_post_message(&state, bot, &session.chat_id, &post).await {
507                Ok(message_id) => message_id,
508                Err(err) => return Err((StatusCode::BAD_GATEWAY, err.to_string())),
509            }
510        };
511        Ok(Json(serde_json::json!({
512            "ok": true,
513            "sessionId": session_id,
514            "messageId": message_id,
515            "targetMessageId": target_message_id,
516        })))
517    }
518
519    async fn list_bots(State(state): State<AppState>) -> Json<Vec<BotSummary>> {
520        let sessions = state.sessions.lock().await;
521        Json(
522            state
523                .bots
524                .iter()
525                .map(|(app_id, bot)| {
526                    let active = sessions
527                        .values()
528                        .filter(|s| s.lark_app_id == *app_id && s.status == SessionStatus::Active)
529                        .count();
530                    BotSummary {
531                        lark_app_id: app_id.clone(),
532                        name: bot.name.clone(),
533                        cli_id: bot.cli_id.clone(),
534                        model: bot.model.clone(),
535                        allowed_users: bot.allowed_users.clone(),
536                        allowed_chat_groups: bot.allowed_chat_groups.clone(),
537                        oncall_chats: bot
538                            .oncall_chats
539                            .iter()
540                            .map(|oc| oc.chat_id.clone())
541                            .collect(),
542                        private_card: bot.private_card,
543                        active_sessions: active,
544                    }
545                })
546                .collect(),
547        )
548    }
549
550    async fn get_bot(
551        State(state): State<AppState>,
552        AxumPath(app_id): AxumPath<String>,
553    ) -> Result<Json<BotSummary>, (StatusCode, String)> {
554        let sessions = state.sessions.lock().await;
555        let bot = state
556            .bots
557            .get(&app_id)
558            .ok_or_else(|| (StatusCode::NOT_FOUND, format!("bot {} not found", app_id)))?;
559        let active = sessions
560            .values()
561            .filter(|s| s.lark_app_id == app_id && s.status == SessionStatus::Active)
562            .count();
563        Ok(Json(BotSummary {
564            lark_app_id: app_id,
565            name: bot.name.clone(),
566            cli_id: bot.cli_id.clone(),
567            model: bot.model.clone(),
568            allowed_users: bot.allowed_users.clone(),
569            allowed_chat_groups: bot.allowed_chat_groups.clone(),
570            oncall_chats: bot
571                .oncall_chats
572                .iter()
573                .map(|oc| oc.chat_id.clone())
574                .collect(),
575            private_card: bot.private_card,
576            active_sessions: active,
577        }))
578    }
579
580    async fn list_session_groups(State(state): State<AppState>) -> Json<Vec<SessionGroup>> {
581        let sessions = state.sessions.lock().await;
582        let mut groups: HashMap<String, SessionGroup> = HashMap::new();
583        for session in sessions.values() {
584            let key = session.chat_id.clone();
585            let summary = SessionSummary::from(session);
586            groups
587                .entry(key)
588                .and_modify(|g| g.sessions.push(summary.clone()))
589                .or_insert_with(|| SessionGroup {
590                    chat_id: session.chat_id.clone(),
591                    title: Some(session.title.clone()),
592                    sessions: vec![summary],
593                });
594        }
595        Json(groups.into_values().collect())
596    }
597
598    async fn locate_session(
599        State(state): State<AppState>,
600        AxumPath(session_id): AxumPath<String>,
601    ) -> Result<Json<SessionLocateInfo>, (StatusCode, String)> {
602        let sessions = state.sessions.lock().await;
603        let session = sessions.get(&session_id).ok_or_else(|| {
604            (
605                StatusCode::NOT_FOUND,
606                format!("session {} not found", session_id),
607            )
608        })?;
609        Ok(Json(SessionLocateInfo {
610            session_id: session.session_id.clone(),
611            terminal_url: session.terminal_url.clone(),
612            worker_pid: session.worker_pid,
613        }))
614    }
615
616    async fn overview(State(state): State<AppState>) -> Json<DaemonOverview> {
617        let sessions = state.sessions.lock().await;
618        let active = sessions
619            .values()
620            .filter(|s| s.status == SessionStatus::Active)
621            .count();
622        let closed = sessions
623            .values()
624            .filter(|s| s.status == SessionStatus::Closed)
625            .count();
626        Json(DaemonOverview {
627            pid: std::process::id(),
628            started_at: state.started_at,
629            session_count: sessions.len(),
630            active_session_count: active,
631            closed_session_count: closed,
632            bot_count: state.bots.len(),
633            worker_count: state.workers.lock().await.len(),
634            config_path: state.paths.config_toml().display().to_string(),
635            data_dir: state.paths.root().display().to_string(),
636        })
637    }
638
639    async fn preferences(State(state): State<AppState>) -> Json<Value> {
640        Json(serde_json::json!({
641            "web": state.config.web,
642            "daemon": state.config.daemon,
643            "lark": state.config.lark,
644            "screenAnalyzer": state.config.screen_analyzer,
645        }))
646    }
647
648    async fn auth(State(state): State<AppState>) -> Json<Value> {
649        let token = mint_dashboard_token();
650        let expires_at = Instant::now() + Duration::from_secs(24 * 60 * 60);
651        {
652            let mut guard = state.dashboard_token.lock().await;
653            *guard = Some(DashboardAuthToken {
654                token: token.clone(),
655                expires_at,
656            });
657        }
658        Json(serde_json::json!({
659            "authenticated": true,
660            "token": token,
661            "loginPath": format!("/dashboard/login?token={}", token),
662            "dashboardPath": "/dashboard/",
663            "expiresInSeconds": expires_at
664                .checked_duration_since(Instant::now())
665                .map(|d| d.as_secs())
666                .unwrap_or(0),
667            "mode": "ws",
668            "botCount": state.bots.len(),
669            "daemonPid": std::process::id(),
670            "dashboard": {
671                "host": state.config.web.host,
672                "proxyBasePort": state.config.web.proxy_base_port,
673            },
674        }))
675    }
676
677    async fn dashboard_login(
678        State(state): State<AppState>,
679        Query(query): Query<HashMap<String, String>>,
680    ) -> Result<impl IntoResponse, (StatusCode, String)> {
681        let Some(token) = query
682            .get("token")
683            .map(|s| s.trim())
684            .filter(|s| !s.is_empty())
685        else {
686            return Err((
687                StatusCode::BAD_REQUEST,
688                "missing dashboard token".to_string(),
689            ));
690        };
691        if !dashboard_token_is_valid(&state, token).await {
692            return Err((
693                StatusCode::UNAUTHORIZED,
694                "dashboard token expired".to_string(),
695            ));
696        }
697        let mut response = Redirect::temporary("/dashboard/").into_response();
698        response.headers_mut().insert(
699            axum::http::header::SET_COOKIE,
700            axum::http::HeaderValue::from_str(&format!(
701                "beam-dashboard-token={}; Path=/; HttpOnly; SameSite=Lax; Max-Age=86400",
702                token
703            ))
704            .map_err(internal_error)?,
705        );
706        Ok(response)
707    }
708
709    let protected_dashboard = Router::new()
710        .route(
711            "/api/workflows/definitions",
712            get(list_workflow_definitions_api),
713        )
714        .route(
715            "/api/workflows/definitions/{workflow_id}",
716            get(get_workflow_definition_api),
717        )
718        .route(
719            "/api/workflows/definitions/{workflow_id}/run",
720            post(trigger_workflow_definition_run_api),
721        )
722        .route("/api/workflows/runs", get(list_workflow_runs_api))
723        .route(
724            "/api/workflows/runs/{run_id}/snapshot",
725            get(get_workflow_run_snapshot_api),
726        )
727        .route(
728            "/api/workflows/runs/{run_id}/events",
729            get(get_workflow_run_events),
730        )
731        .route(
732            "/api/workflows/runs/{run_id}/approve",
733            post(approve_workflow_run),
734        )
735        .route(
736            "/api/workflows/runs/{run_id}/reject",
737            post(reject_workflow_run),
738        )
739        .route(
740            "/api/workflows/runs/{run_id}/attempts/{activity_id}/{attempt_id}/resume",
741            post(start_workflow_attempt_resume),
742        )
743        .route(
744            "/api/workflows/runs/{run_id}/attempts/{activity_id}/{attempt_id}/resume/end",
745            post(end_workflow_attempt_resume),
746        )
747        .route(
748            "/api/workflows/runs/{run_id}/cancel",
749            post(cancel_workflow_run),
750        )
751        .route(
752            "/api/workflows/runs/{run_id}/resume",
753            post(resume_workflow_run),
754        )
755        .route("/sessions", post(create_session))
756        .route("/sessions/{session_id}", get(get_session))
757        .route("/sessions/{session_id}/input", post(send_input))
758        .route("/sessions/{session_id}/report", post(report_session))
759        .route("/sessions/{session_id}/refresh", post(refresh_session))
760        .route("/sessions/{session_id}/restart", post(restart_session))
761        .route("/sessions/{session_id}/resume", post(resume_session))
762        .route("/sessions/{session_id}/close", post(close_session))
763        .route(
764            "/api/workflows/{workflow_id}/run",
765            post(trigger_workflow_run),
766        )
767        .route("/api/workflows/{run_id}", get(get_workflow_run))
768        .route("/api/trigger", post(api_trigger))
769        .route(
770            "/adopt/zellij",
771            get(list_zellij_adopt_candidates).post(adopt_zellij_session),
772        )
773        .route("/api/bots", get(list_bots))
774        .route("/api/bots/{app_id}", get(get_bot))
775        .route("/api/preferences", get(preferences))
776        .route("/api/connectors", get(connectors).post(create_connector))
777        .route("/api/connectors/stats", get(connector_stats))
778        .route(
779            "/api/connectors/{id}",
780            get(get_connector_api)
781                .put(update_connector_api)
782                .patch(patch_connector_api)
783                .delete(delete_connector_api),
784        )
785        .route(
786            "/api/webhook-secrets",
787            get(list_webhook_secrets_api).post(create_webhook_secret_api),
788        )
789        .route(
790            "/api/webhook-secrets/{ref}",
791            put(update_webhook_secret_api).delete(delete_webhook_secret_api),
792        )
793        .route("/api/trigger-logs", get(trigger_logs_api))
794        .route("/api/trigger-logs/prune", post(prune_trigger_logs_api))
795        .route("/api/connectors/webhooks", get(list_webhook_triggers))
796        .route("/api/sessions/groups", get(list_session_groups))
797        .route("/api/sessions/{session_id}/locate", get(locate_session))
798        .route("/api/overview", get(overview))
799        .nest_service(
800            "/dashboard",
801            get_service(ServeDir::new("src/dashboard/web")),
802        )
803        .route_layer(middleware::from_fn_with_state(
804            state.clone(),
805            dashboard_gate,
806        ));
807
808    let open_routes = Router::new()
809        .route("/health", get(health))
810        .route("/shutdown", post(shutdown))
811        .route("/sessions", get(list_sessions))
812        .route("/api/auth", get(auth))
813        .route("/dashboard/login", get(dashboard_login))
814        .route("/api/schedules", post(create_schedule))
815        .route("/webhook/{workflow_id}", post(handle_webhook_trigger))
816        .route(
817            "/sessions/{session_id}/history",
818            get(lark_history::session_history),
819        )
820        .route(
821            "/sessions/{session_id}/quoted/{message_id}",
822            get(lark_history::quoted_message),
823        )
824        .route("/sessions/{session_id}/final-output", post(final_output))
825        .route("/api/asks", post(ask::create_ask))
826        .route("/api/attention", post(set_attention_route))
827        .route(
828            "/debug/simulate/lark-message",
829            post(debug_simulate::simulate_lark_message_handler),
830        );
831
832    // Start zellij web server and ensure tokens
833    let zellij_web_port = state.config.web.proxy_base_port + 1;
834    // PR5 gate: a herdr-only daemon may skip zellij web entirely. Default is
835    // true (existing behavior); only tested herdr-only deployments flip it.
836    let zellij_tokens = zellij_web::start_zellij_web_if_enabled(
837        state.config.web.zellij_web,
838        zellij_web_port,
839        &state.paths.zellij_web_tokens_json(),
840    )?;
841
842    // Start terminal proxy with auth bridge
843    let proxy_host = state.config.web.host.clone();
844    let proxy_port = state.config.web.proxy_base_port;
845    let proxy_sessions = state.sessions.clone();
846    let auth_state = terminal_auth::TerminalAuthState::new();
847    // Load persisted used tickets for terminal auth anti-replay.
848    auth_state
849        .load_used_tickets(&paths.used_tickets_json())
850        .await;
851    terminal_proxy::start_proxy(
852        &proxy_host,
853        proxy_port,
854        zellij_web_port,
855        proxy_sessions,
856        zellij_tokens,
857        auth_state.clone(),
858        terminal_proxy::herdr_ws::HerdrWebLimits {
859            enabled: state.config.web.herdr_terminal,
860            max_observers_per_session: state.config.web.herdr_terminal_max_observers_per_session,
861            max_observers_global: state.config.web.herdr_terminal_max_observers_global,
862        },
863    )
864    .await
865    .with_context(|| format!("failed to start terminal proxy on {proxy_host}:{proxy_port}"))?;
866
867    let periodic_paths = paths.clone();
868    let periodic_auth = auth_state.clone();
869    let periodic_recent_events = state.recent_lark_events.clone();
870    tokio::spawn(async move {
871        loop {
872            tokio::time::sleep(std::time::Duration::from_secs(30)).await;
873            // Save replay nonces
874            let replay_map: HashMap<String, u64> = {
875                let guard = replay_nonce_store()
876                    .lock()
877                    .unwrap_or_else(|p| p.into_inner());
878                guard.clone()
879            };
880            if !replay_map.is_empty() {
881                let path = periodic_paths.replay_nonces_json();
882                let _ = tokio::task::spawn_blocking(move || {
883                    let _ = beam_core::persist::atomic_write_json(&path, &replay_map);
884                })
885                .await;
886            } else {
887                let _ = tokio::fs::remove_file(periodic_paths.replay_nonces_json()).await;
888            }
889            // Save rate buckets
890            let rate_map: HashMap<String, (u64, u64)> = {
891                let guard = rate_bucket_store()
892                    .lock()
893                    .unwrap_or_else(|p| p.into_inner());
894                guard.clone()
895            };
896            if !rate_map.is_empty() {
897                let path = periodic_paths.rate_buckets_json();
898                let _ = tokio::task::spawn_blocking(move || {
899                    let _ = beam_core::persist::atomic_write_json(&path, &rate_map);
900                })
901                .await;
902            } else {
903                let _ = tokio::fs::remove_file(periodic_paths.rate_buckets_json()).await;
904            }
905            // Save used tickets
906            periodic_auth
907                .save_used_tickets(&periodic_paths.used_tickets_json())
908                .await;
909            // Save recent Lark events dedupe state
910            {
911                let events = periodic_recent_events.lock().await;
912                save_recent_lark_events(&periodic_paths, &events).await;
913            }
914        }
915    });
916
917    worker_health::spawn_worker_health_watchdog(state.clone());
918
919    // Schedule loop: periodically check schedules and trigger due tasks.
920    let schedule_paths = paths.clone();
921    let schedule_state = state.clone();
922    tokio::spawn(async move {
923        loop {
924            tokio::time::sleep(std::time::Duration::from_secs(30)).await;
925            let tasks = match beam_core::list_tasks(&schedule_paths) {
926                Ok(tasks) => tasks,
927                Err(_) => continue,
928            };
929            let now = chrono::Utc::now();
930            for task in &tasks {
931                if !task.enabled {
932                    continue;
933                }
934                let Some(next_run_at_str) = task.next_run_at.as_deref() else {
935                    continue;
936                };
937                let Ok(next_run_at) = chrono::DateTime::parse_from_rfc3339(next_run_at_str) else {
938                    continue;
939                };
940                let next_run_at = next_run_at.with_timezone(&chrono::Utc);
941                if next_run_at > now {
942                    continue;
943                }
944                // Task is due: execute it.
945                info!("schedule: triggering task {} ({})", task.id, task.name);
946                let state = schedule_state.clone();
947                let task_id = task.id.clone();
948                let task_prompt = task.prompt.clone();
949                let task_working_dir = task.working_dir.clone();
950                let task_chat_id = task.chat_id.clone();
951                let task_lark_app_id = task.lark_app_id.clone();
952                let task_root_message_id = task.root_message_id.clone();
953                let task_scope = task.scope.clone();
954                let task_chat_type = task.chat_type.clone();
955                let task_name = task.name.clone();
956                let paths = schedule_paths.clone();
957                tokio::spawn(async move {
958                    let result = execute_schedule_task(
959                        &state,
960                        &task_id,
961                        &task_prompt,
962                        &task_working_dir,
963                        &task_chat_id,
964                        task_lark_app_id.as_deref(),
965                        task_root_message_id.as_deref(),
966                        task_scope.as_deref(),
967                        task_chat_type.as_ref(),
968                        &task_name,
969                    )
970                    .await;
971                    let success = result.is_ok();
972                    let error = result.as_ref().err().map(|e| e.to_string());
973                    let _ = beam_core::mark_run(&paths, &task_id, success, error.as_deref(), None);
974                    if let Err(err) = result {
975                        warn!("schedule task {} failed: {}", task_id, err);
976                    }
977                });
978            }
979        }
980    });
981
982    let app = Router::new()
983        .merge(protected_dashboard)
984        .merge(open_routes)
985        .with_state(state);
986
987    info!("beam daemon listening on {}", addr);
988    axum::serve(listener, app)
989        .with_graceful_shutdown(async move {
990            let _ = shutdown_rx.await;
991        })
992        .await?;
993
994    let _ = tokio::fs::remove_file(paths.runtime_state_json()).await;
995    Ok(())
996}
997
998#[cfg(test)]
999#[path = "tests/lib_integration/mod.rs"]
1000mod tests;