133 lines
3.6 KiB
Rust
133 lines
3.6 KiB
Rust
use std::sync::Arc;
|
|
|
|
use tokio::net::TcpListener;
|
|
use tracing::info;
|
|
|
|
use playout::config::PlayoutConfig;
|
|
use playout::engine::{PlayoutEngine, SourceUriResolver};
|
|
use playout::http::{self, AppState};
|
|
use playout::segment_store::filesystem::FilesystemSegmentStore;
|
|
|
|
use domain::value_objects::SourceUri;
|
|
|
|
struct StubResolver;
|
|
|
|
impl SourceUriResolver for StubResolver {
|
|
fn resolve(&self, _provider_id: &str, _external_id: &str) -> Option<SourceUri> {
|
|
None
|
|
}
|
|
}
|
|
|
|
#[tokio::main]
|
|
async fn main() -> anyhow::Result<()> {
|
|
dotenvy::dotenv().ok();
|
|
|
|
tracing_subscriber::fmt()
|
|
.with_env_filter(
|
|
tracing_subscriber::EnvFilter::try_from_default_env()
|
|
.unwrap_or_else(|_| "info".into()),
|
|
)
|
|
.init();
|
|
|
|
let config = PlayoutConfig::from_env();
|
|
info!(addr = %config.listen_addr, "starting k-tv-playout");
|
|
|
|
let store = Arc::new(FilesystemSegmentStore::new(config.storage_path.clone()));
|
|
let resolver: Arc<dyn SourceUriResolver> = Arc::new(StubResolver);
|
|
|
|
// TODO: wire real ScheduleQuery from adapter-sqlite once DB URL is configured
|
|
let schedule_query: Arc<dyn domain::ports::ScheduleQuery> =
|
|
Arc::new(NoopScheduleQuery);
|
|
|
|
let engine = Arc::new(PlayoutEngine::new(
|
|
config.clone(),
|
|
schedule_query,
|
|
resolver,
|
|
store.clone(),
|
|
));
|
|
|
|
let tick_engine = engine.clone();
|
|
let tick_interval = config.tick_interval_ms;
|
|
tokio::spawn(async move {
|
|
let mut interval =
|
|
tokio::time::interval(tokio::time::Duration::from_millis(tick_interval));
|
|
loop {
|
|
interval.tick().await;
|
|
tick_engine.tick().await;
|
|
}
|
|
});
|
|
|
|
let state = AppState {
|
|
engine: engine.clone(),
|
|
store,
|
|
config: config.clone(),
|
|
};
|
|
|
|
let app = http::router(state);
|
|
let listener = TcpListener::bind(&config.listen_addr).await?;
|
|
info!("listening on {}", config.listen_addr);
|
|
|
|
axum::serve(listener, app)
|
|
.with_graceful_shutdown(shutdown_signal(engine))
|
|
.await?;
|
|
|
|
Ok(())
|
|
}
|
|
|
|
async fn shutdown_signal(engine: Arc<PlayoutEngine>) {
|
|
tokio::signal::ctrl_c()
|
|
.await
|
|
.expect("failed to install ctrl+c handler");
|
|
info!("shutting down");
|
|
engine.shutdown().await;
|
|
}
|
|
|
|
struct NoopScheduleQuery;
|
|
|
|
#[async_trait::async_trait]
|
|
impl domain::ports::ScheduleQuery for NoopScheduleQuery {
|
|
async fn find_active(
|
|
&self,
|
|
_channel_id: domain::value_objects::ChannelId,
|
|
_at: chrono::DateTime<chrono::Utc>,
|
|
) -> domain::DomainResult<Option<domain::GeneratedSchedule>> {
|
|
Ok(None)
|
|
}
|
|
|
|
async fn find_latest(
|
|
&self,
|
|
_channel_id: domain::value_objects::ChannelId,
|
|
) -> domain::DomainResult<Option<domain::GeneratedSchedule>> {
|
|
Ok(None)
|
|
}
|
|
|
|
async fn find_playback_history(
|
|
&self,
|
|
_channel_id: domain::value_objects::ChannelId,
|
|
) -> domain::DomainResult<Vec<domain::PlaybackRecord>> {
|
|
Ok(vec![])
|
|
}
|
|
|
|
async fn find_last_slot_per_block(
|
|
&self,
|
|
_channel_id: domain::value_objects::ChannelId,
|
|
) -> domain::DomainResult<std::collections::HashMap<domain::BlockId, domain::MediaItemId>> {
|
|
Ok(std::collections::HashMap::new())
|
|
}
|
|
|
|
async fn list_schedule_history(
|
|
&self,
|
|
_channel_id: domain::value_objects::ChannelId,
|
|
) -> domain::DomainResult<Vec<domain::GeneratedSchedule>> {
|
|
Ok(vec![])
|
|
}
|
|
|
|
async fn get_schedule_by_id(
|
|
&self,
|
|
_channel_id: domain::value_objects::ChannelId,
|
|
_schedule_id: domain::value_objects::ScheduleId,
|
|
) -> domain::DomainResult<Option<domain::GeneratedSchedule>> {
|
|
Ok(None)
|
|
}
|
|
}
|