use std::sync::Arc;
use std::time::Duration;
use tuitbot_core::automation::circuit_breaker::CircuitBreaker;
use tuitbot_core::automation::{
run_approval_poster, run_posting_queue_with_approval, run_token_refresh_loop,
scheduler_from_config, status_reporter::run_status_reporter, AnalyticsLoop, ContentLoop,
DiscoveryLoop, MentionsLoop, PostExecutor, Runtime, TargetLoop, ThreadLoop,
};
use tuitbot_core::config::{Config, OperatingMode};
use tuitbot_core::startup::format_startup_banner;
use crate::deps::RuntimeDeps;
pub async fn execute(config: &Config, status_interval: u64) -> anyhow::Result<()> {
let mut deps = RuntimeDeps::init(config, false).await?;
let effective_interval = if status_interval > 0 {
status_interval
} else {
config.logging.status_interval_seconds
};
let banner = format_startup_banner(deps.tier, &deps.capabilities, effective_interval);
eprintln!("{banner}");
let mut runtime = Runtime::new();
let min_delay = Duration::from_secs(config.limits.min_action_delay_seconds);
let max_delay = Duration::from_secs(config.limits.max_action_delay_seconds);
let circuit_breaker = CircuitBreaker::new(
config.circuit_breaker.error_threshold,
Duration::from_secs(config.circuit_breaker.window_seconds),
Duration::from_secs(config.circuit_breaker.cooldown_seconds),
);
let cancel = runtime.cancel_token();
let post_rx = deps.post_rx.take().expect("post_rx not yet consumed");
runtime.spawn("posting-queue", {
let executor = deps.post_executor.clone() as Arc<dyn PostExecutor>;
let approval_queue = deps.approval_queue.clone();
let cb = circuit_breaker.clone();
async move {
run_posting_queue_with_approval(
post_rx,
executor,
approval_queue,
min_delay,
max_delay,
Some(cb),
cancel,
)
.await;
}
});
if let (Some(tm), Some(xc)) = (&deps.token_manager, &deps.x_client) {
let cancel = runtime.cancel_token();
let tm = tm.clone();
let xc = xc.clone();
runtime.spawn("token-refresh", run_token_refresh_loop(tm, xc, cancel));
}
{
let cancel = runtime.cancel_token();
let pool = deps.pool.clone();
let xc = deps.dyn_client.clone();
let account_id = tuitbot_core::storage::accounts::DEFAULT_ACCOUNT_ID.to_string();
runtime.spawn(
"approval-poster",
run_approval_poster(pool, xc, account_id, min_delay, max_delay, cancel),
);
}
let is_composer = config.mode == OperatingMode::Composer;
if !is_composer {
{
let content_loop = ContentLoop::new(
deps.tweet_gen.clone(),
deps.content_safety.clone(),
deps.content_storage.clone(),
config.business.effective_industry_topics().to_vec(),
config.intervals.content_post_window_seconds,
false,
)
.with_topic_scorer(deps.topic_scorer.clone())
.with_thread_poster(deps.thread_poster.clone());
let cancel = runtime.cancel_token();
let scheduler = scheduler_from_config(
config.intervals.content_post_window_seconds,
config.limits.min_action_delay_seconds,
config.limits.max_action_delay_seconds,
);
let schedule = deps.active_schedule.clone();
runtime.spawn("content-loop", async move {
content_loop.run(cancel, scheduler, schedule).await;
});
}
{
let thread_loop = ThreadLoop::new(
deps.thread_gen.clone(),
deps.content_safety.clone(),
deps.content_storage.clone(),
deps.thread_poster.clone(),
config.business.effective_industry_topics().to_vec(),
config.intervals.thread_interval_seconds,
false,
);
let cancel = runtime.cancel_token();
let scheduler = scheduler_from_config(
config.intervals.thread_interval_seconds,
config.limits.min_action_delay_seconds,
config.limits.max_action_delay_seconds,
);
let schedule = deps.active_schedule.clone();
runtime.spawn("thread-loop", async move {
thread_loop.run(cancel, scheduler, schedule).await;
});
}
}
if deps.capabilities.discovery {
let discovery_loop = DiscoveryLoop::new(
deps.searcher.clone(),
deps.scorer.clone(),
deps.reply_gen.clone(),
deps.safety.clone(),
deps.loop_storage.clone(),
deps.post_sender.clone(),
deps.keywords.clone(),
config.scoring.threshold as f32,
is_composer, );
let cancel = runtime.cancel_token();
let scheduler = scheduler_from_config(
config.intervals.discovery_search_seconds,
config.limits.min_action_delay_seconds,
config.limits.max_action_delay_seconds,
);
let schedule = deps.active_schedule.clone();
runtime.spawn("discovery-loop", async move {
discovery_loop.run(cancel, scheduler, schedule).await;
});
}
if deps.capabilities.mentions && !is_composer {
let mentions_loop = MentionsLoop::new(
deps.mentions_fetcher.clone(),
deps.reply_gen.clone(),
deps.safety.clone(),
deps.post_sender.clone(),
false,
);
let cancel = runtime.cancel_token();
let scheduler = scheduler_from_config(
config.intervals.mentions_check_seconds,
config.limits.min_action_delay_seconds,
config.limits.max_action_delay_seconds,
);
let schedule = deps.active_schedule.clone();
let storage_clone = deps.loop_storage.clone();
runtime.spawn("mentions-loop", async move {
mentions_loop
.run(cancel, scheduler, schedule, storage_clone)
.await;
});
let target_loop = TargetLoop::new(
deps.target_adapter.clone(),
deps.target_adapter.clone(),
deps.reply_gen.clone(),
deps.safety.clone(),
deps.target_storage.clone(),
deps.post_sender.clone(),
deps.target_loop_config.clone(),
);
let cancel = runtime.cancel_token();
let scheduler = scheduler_from_config(
config.intervals.mentions_check_seconds,
config.limits.min_action_delay_seconds,
config.limits.max_action_delay_seconds,
);
let schedule = deps.active_schedule.clone();
runtime.spawn("target-loop", async move {
target_loop.run(cancel, scheduler, schedule).await;
});
}
if deps.capabilities.mentions {
let analytics_loop = AnalyticsLoop::new(
deps.profile_adapter.clone(),
deps.profile_adapter.clone(),
deps.analytics_storage.clone(),
);
let cancel = runtime.cancel_token();
let scheduler = scheduler_from_config(3600, 0, 0);
runtime.spawn("analytics-loop", async move {
analytics_loop.run(cancel, scheduler).await;
});
}
if effective_interval > 0 {
let scheduler = scheduler_from_config(effective_interval, 0, 0);
let cancel = runtime.cancel_token();
let status_querier = deps.status_querier.clone();
runtime.spawn("status-reporter", async move {
run_status_reporter(status_querier, scheduler, cancel).await;
});
}
tracing::info!(
tasks = runtime.task_count(),
"All automation loops spawned, running until shutdown"
);
runtime.run_until_shutdown().await;
tracing::info!("Shutdown complete.");
Ok(())
}