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
63pub use daemon_types::RunOptions;
65pub(crate) use daemon_types::*;
66
67pub(crate) use workflow_catalog::*;
69pub(crate) use workflow_execution::*;
71pub(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 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 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 {
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 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 {
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; }
357 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 s.last_final_output_turn_id.as_deref() != Some(tid)
366 })
367 })
368 .unwrap_or(false) };
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 let zellij_web_port = state.config.web.proxy_base_port + 1;
834 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 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 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 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 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 periodic_auth
907 .save_used_tickets(&periodic_paths.used_tickets_json())
908 .await;
909 {
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 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 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;